Добавляет Cursor/OffsetSequencer в :outbox-api и интегрирует PersistentOffsetSequencer через :outbox-ksqlite.

Введение монотонного offset'а как персистентного состояния агента:
offset'ы переживают рестарт standalone-агента, клиент продолжает
синхронизацию инкрементально, без полной re-sync с нуля.

outbox-api:
  - Cursor (offset: Long) — курсор в журнале событий агента.
  - OffsetSequencer — интерфейс резервирования уникального offset.
  - CursorStore — персистентное хранилище текущего offset'а.
  - PersistentOffsetSequencer — декоратор над любым OutboxStore,
    обновляет CursorStore на каждом append (atomic transaction).
  - OutboxGapException — клиент запросил after < earliestCursor() →
    сервер не может удовлетворить, клиент обязан делать full resync.
  - DurableEvent переименован из Event.kt → DurableEvent.kt (Event.kt
    был общим sealed-типом, теперь это термин из спеки).
  - MutableOutboxStore и OutboxStore теперь читают offset через
    CursorStore вместо in-memory counter'а.

outbox-inmemory:
  - InMemoryOffsetSequencer — для тестов и dev-режима.
  - InMemoryOutboxStore теперь принимает OffsetSequencer в конструкторе.

outbox-ksqlite (новый модуль):
  - KsqliteCursorStore — таблица outbox_cursor (agent_id TEXT PK,
    offset INTEGER NOT NULL DEFAULT 0, updated_at INTEGER NOT NULL).
  - KsqliteCursorStoreTest — 4 теста (set/get, monotonic, concurrent).

proto + server:
  - Snapshot.proto — server-state snapshot endpoint для клиентов,
    которым нужна полная материализация (использование TBD).
  - Routes.kt + SnapshotRouteTest — endpoint /agentik/snapshot (GET).

journal-ksqlite:
  - KsqliteJournalStore.listFlow/append — без изменений по API,
    нотации минимальные (codecs).

standalone:
  - DurableLog (бывший ChatAgent-orchestration) — атомарный commit
    события в OutboxStore + PersistentOffsetSequencer + materialization
    (через Reducer) одной транзакцией.
  - SqliteStores — добавляет KsqliteCursorStore в bundle, единая
    shared-connection для всех ksqlite-сторов standalone-агента.
  - ChatAgent / ConversationLoop / ConversationEvents / ReflectionScheduler /
    ToolDispatcher — переход на новые абстракции.
  - standalone/build.gradle.kts — implementation(project(':outbox-ksqlite'))
    включено (раньше было закомментировано — модуль только создавался).

client:
  - AgentikAgent / AgentClient / HttpEventStore / HttpJournalStore /
    ReconnectingOutbox — используют Cursor через transport API.
  - client/README.md — синхронизирован с новым поведением (468 строк
    diff — это в основном оформление и примеры).

kotlinx-io: 0.8.0 → 0.9.1 в libs.versions.toml (см. sync-core tests).

SYNC-SYSTEM.md (в корне) — спецификация, на которую ссылается и
:sync-core (эта сессия), и эта Cursor-абстракция в outbox-api.

Тесты: standalone 132, journal-ksqlite 25, outbox-inmemory 20,
outbox-ksqlite 4, client 10, sync-core 74 — все зелёные на jvm;
sync-core linuxX64 74 тоже зелёный.

sync2/ (заброшенный stub с одним build.gradle.kts) удалён.
This commit is contained in:
2026-10-02 01:16:15 +03:00
parent 92b76c4e3c
commit 5bdc517988
63 changed files with 2953 additions and 1017 deletions
@@ -18,7 +18,9 @@ import pw.binom.agentik.outbox.OnlineOutbox
import pw.binom.agentik.outbox.OutboxStore
import pw.binom.agentik.proto.Agent
import pw.binom.agentik.proto.AgentInfo
import pw.binom.agentik.proto.ChatSnapshot
import pw.binom.agentik.proto.Conversation
import pw.binom.agentik.proto.ConversationsSnapshot
import kotlin.time.Instant
/**
@@ -81,6 +83,22 @@ internal class AgentClient private constructor(
return rec.updatedAt
}
override suspend fun conversationsSnapshot(): ConversationsSnapshot {
val response = httpClient.get("$agentUrl/snapshot")
check(response.status == HttpStatusCode.OK) {
"snapshot: server returned ${response.status}"
}
return response.body()
}
override suspend fun chatSnapshot(conversationId: String): ChatSnapshot {
val response = httpClient.get("$agentUrl/conversations/$conversationId/snapshot")
check(response.status == HttpStatusCode.OK) {
"conversations/$conversationId/snapshot: server returned ${response.status}"
}
return response.body()
}
override fun close() {
httpClient.close()
}
@@ -1,12 +1,14 @@
package pw.binom.agentik.client
import io.ktor.client.engine.HttpClientEngineFactory
import kotlinx.coroutines.CancellationException
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.Job
import kotlinx.coroutines.SupervisorJob
import kotlinx.coroutines.cancel
import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.delay
import kotlinx.coroutines.isActive
import kotlinx.coroutines.launch
import kotlinx.coroutines.runBlocking
import pw.binom.agentik.journal.ConversationRecord
@@ -14,8 +16,9 @@ import pw.binom.agentik.journal.ConversationStore
import pw.binom.agentik.journal.MutableConversationStore
import pw.binom.agentik.journal.inmemory.InMemoryMutableConversationStore
import pw.binom.agentik.outbox.AgentEvent
import pw.binom.agentik.outbox.OutboxGapException
import pw.binom.agentik.proto.Agent
import kotlin.time.Instant
import kotlin.time.Duration.Companion.seconds
/**
* Создаёт [Agent], который ходит в HTTP-фасад `agentikAgent` (модуль `:server`).
@@ -34,8 +37,9 @@ import kotlin.time.Instant
* )
* val conv = agent.createConversation(temp = false)
* conv.send(listOf(Content.Text("hi")))
* agent.outbox.conversationEvents(Instant.DISTANT_PAST, conv.id)
* .map { it.event }
* // Курсор-протокол: сначала снапшот (state + cursor), потом подписка «после»:
* val snap = agent.chatSnapshot(conv.id)
* agent.outbox.conversationEvents(after = snap.cursor, conversationId = conv.id)
* .collect { ... }
* agent.close() // закрывает HttpClient + локальный кэш
* ```
@@ -76,10 +80,12 @@ import kotlin.time.Instant
*
* [conversationStore], который видит клиент — это **кэш**, не прямой HTTP.
* Внутри лежит [InMemoryMutableConversationStore], который:
* 1. На старте делает snapshot через `remote.listFlow(0)` → `local.upsert(...)`.
* 2. Подписывается на `outbox.agentEvents(after)` → для каждого
* 1. На старте берёт `conversationsSnapshot()` (полный список + курсор) и
* приводит к нему локальную копию.
* 2. Подписывается на `outbox.agentEvents(after = snapshot.cursor)` → для каждого
* [AgentEvent.Created] / `Deleted` / `Renamed` / `Touched` применяет
* соответствующий `upsert/delete/rename/touch` к локальной копии.
* 3. При `OutboxGapException` повторяет с шага 1 (полный resync).
*
* UI читает `agent.conversationStore.list(0, PAGE_SIZE)` — мгновенно,
* без HTTP, в т.ч. оффлайн. Команды (create/delete/rename) идут
@@ -103,17 +109,23 @@ fun AgentikAgent(
/**
* Оборачивает [Agent] так, что [Agent.conversationStore] становится
* локальным in-memory кэшем, синхронизированным с удалённым стором
* через outbox-события.
* по курсор-протоколу.
*
* - **Seed**: при создании делает один snapshot через
* `remote.listFlow(0)` и заливает в [InMemoryMutableConversationStore].
* - **Live**: подписка на `agent.outbox.agentEvents(after)` применяет
* `Created` / `Deleted` / `Renamed` / `Touched` к локальному кэшу.
* **Протокол синхронизации** (гарантирует актуальный список бесед):
* 1. `conversationsSnapshot()` — база (полный список) + курсор `C`.
* 2. `outbox.agentEvents(after = C)` — дельты, применяются поверх базы
* (`Created`/`Deleted`/`Renamed`/`Touched`, все абсолютные и идемпотентные).
* 3. [OutboxGapException] (курсор мёртв — retention / смена epoch) → повтор
* с шага 1 (полный resync: `reconcile` удаляет локальные беседы, которых
* нет в снапшоте, и upsert'ит все из снапшота).
* 4. Прочие ошибки (сеть) → пауза и повтор.
*
* Возвращает обёртку, у которой переопределён только [Agent.conversationStore]
* (на read-only projection локального [InMemoryMutableConversationStore]).
* Остальные методы [Agent] — delegated в [delegate].
*/
private val RESYNC_RETRY_DELAY = 2.seconds
private fun wrapWithLocalConversationCache(
delegate: Agent,
scopeClient: Agent,
@@ -124,38 +136,59 @@ private fun wrapWithLocalConversationCache(
private val syncJob: Job
init {
// Делаем cacheStore read-only view на localStore.
// (Через вложенный класс — см. ниже.)
// Запускаем seed + live-refresh параллельно.
syncJob = cacheScope.launch {
// 1. seed — snapshot всех текущих бесед с сервера
try {
delegate.conversationStore.listFlow(offset = 0, pageSize = ConversationStore.PAGE_SIZE)
.collect { rec -> localStore.upsert(rec) }
} catch (_: Throwable) {
// seed может упасть (offline / 5xx) — не критично,
// live-источник всё равно догонит при первом событии.
}
syncJob = cacheScope.launch { syncLoop() }
}
// 2. live — применяем outbox-события.
// Используем `first()` для knownId после Created — потом отписываемся,
// потому что Created нужно вытянуть полный record через `remote.get(id)`.
// Renamed/Touched меняют локальную копию без round-trip.
delegate.outbox.agentEvents(after = Instant.DISTANT_PAST).collect { ce ->
when (val ev = ce.event) {
is AgentEvent.Created -> {
// Created не несёт title/timestamps — нужно сходить в remote.
val rec = delegate.conversationStore.get(ev.conversationId)
if (rec != null) localStore.upsert(rec)
}
is AgentEvent.Deleted -> localStore.delete(ev.id)
is AgentEvent.Renamed -> localStore.rename(ev.id, ev.title)
is AgentEvent.Touched -> localStore.touch(ev.id, ev.updatedAt)
}
private suspend fun syncLoop() {
while (cacheScope.isActive) {
try {
val snap = delegate.conversationsSnapshot()
reconcile(snap.conversations)
delegate.outbox.agentEvents(after = snap.cursor).collect { ce -> apply(ce.event) }
// Штатное завершение потока (не должно) → переподключаемся.
} catch (e: CancellationException) {
throw e
} catch (_: OutboxGapException) {
// Курсор мёртв — немедленно новый снапшот.
} catch (_: Throwable) {
// Сеть/5xx — пауза и повтор (локальный кэш сохраняем).
delay(RESYNC_RETRY_DELAY)
}
}
}
/**
* Приводит локальный кэш к снапшоту: чего нет в снапшоте — удаляем,
* всё из снапшота — upsert. Делает полный resync корректным (в т.ч.
* «пропавшие» беседы = удалённые).
*/
private suspend fun reconcile(records: List<ConversationRecord>) {
val fresh = records.mapTo(HashSet()) { it.id }
val stale = ArrayList<String>()
var offset = 0
while (true) {
val page = localStore.list(offset, ConversationStore.PAGE_SIZE)
if (page.isEmpty()) break
page.forEach { if (it.id !in fresh) stale += it.id }
offset += page.size
}
stale.forEach { localStore.delete(it) }
records.forEach { localStore.upsert(it) }
}
private suspend fun apply(ev: AgentEvent) {
when (ev) {
is AgentEvent.Created -> {
// Created не несёт title/timestamps — нужно сходить в remote.
val rec = delegate.conversationStore.get(ev.conversationId)
if (rec != null) localStore.upsert(rec)
}
is AgentEvent.Deleted -> localStore.delete(ev.id)
is AgentEvent.Renamed -> localStore.rename(ev.id, ev.title)
is AgentEvent.Touched -> localStore.touch(ev.id, ev.updatedAt)
}
}
/**
* Read-only projection локального кэша — клиент через него только
* читает (`get` / `list` / `listFlow`).
@@ -1,45 +1,42 @@
package pw.binom.agentik.client
import io.ktor.client.HttpClient
import io.ktor.client.call.body
import io.ktor.client.request.get
import io.ktor.client.request.parameter
import io.ktor.client.request.prepareGet
import io.ktor.client.statement.HttpResponse
import io.ktor.client.statement.bodyAsChannel
import io.ktor.client.statement.bodyAsText
import io.ktor.http.HttpStatusCode
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.flow
import pw.binom.agentik.outbox.OutboxStore
import pw.binom.agentik.outbox.AgentEvent
import kotlinx.serialization.KSerializer
import kotlinx.serialization.Serializable
import pw.binom.agentik.outbox.CommonEvent
import pw.binom.agentik.outbox.Event
import kotlin.time.Clock
import kotlin.time.Instant
import pw.binom.agentik.outbox.Cursor
import pw.binom.agentik.outbox.OutboxGapException
import pw.binom.agentik.outbox.OutboxStore
/**
* HTTP-реализация [OutboxStore] (= [pw.binom.agentik.outbox.OutboxStore]),
* ходящая в `:server`-фасад.
* HTTP-реализация [OutboxStore], ходящая в `:server`-фасад.
*
* **Endpoint-раскладка** (новый дизайн — storage handles на [Agent]):
* - [events] → `GET {baseUrl}/outbox/events?after=` (полный поток
* [CommonEvent], bounded-tail + live SSE, см. [pw.binom.agentik.server.outboxRoutes])
* - [agentEvents] → `GET {baseUrl}/events?after=` (legacy proto-роут:
* сервер пробрасывает [pw.binom.agentik.outbox.agentEvents] и распаковывает
* `.event` для обратной совместимости с форматом AgentEvent)
* - [conversationEvents] с `conversationId != null` → `GET /conversations/{id}/events`
* **Endpoint-раскладка**:
* - [events] → `GET {baseUrl}/outbox/events?epoch=&offset=` (полный поток
* [CommonEvent], bounded-tail + live SSE). Без параметров — live-only.
* - [agentEvents] → `GET {baseUrl}/events?epoch=&offset=` (только
* `CommonEvent.Agent`).
* - [conversationEvents] с `conversationId != null` →
* `GET /conversations/{id}/events?epoch=&offset=`; с `null` — fallback на
* default [OutboxStore.conversationEvents] (общий `/outbox/events` + filter).
* - [currentCursor] / [oldestCursor] → `GET {baseUrl}/outbox/cursor`.
*
* Для [conversationEvents] с `conversationId == null` (события всех диалогов)
* fallback на default [OutboxStore.conversationEvents] — общий поток
* `/outbox/events` + filter. Это редкий кейс (admin-дашборды), и
* оптимизировать его отдельно нерационально.
* **Gap** (`410 Gone`): сервер отвечает `410` с [GapResponse] (oldest/current
* курсоры) — клиент конвертирует в [OutboxGapException]. Это сигнал сделать
* resync: `agent.conversationsSnapshot()` / `agent.chatSnapshot(id)`.
*
* [earliestEventDate] не имеет своего endpoint'а; возвращает `Clock.System.now()`
* (см. KDoc [OutboxStore.earliestEventDate] — для пустого буфера это и есть
* контрактное значение). Клиент, который полагался на gap detection через
* message store, продолжит работать — просто fallback никогда не сработает.
*
* **Импорты [CommonEvent]/[AgentEvent]/[Event] идут напрямую из
* `pw.binom.agentik.outbox`** — typealias'ы в `:proto.CommonEvent` и т.п.
* НЕ поддерживают nested-class access (`CommonEvent.Agent` через alias
* даёт "Unresolved qualified name"), поэтому приходится использовать
* конкретный пакет. Типы идентичны, alias только для удобства внешнего API.
* **Импорты [CommonEvent]/[Cursor] идут напрямую из `pw.binom.agentik.outbox`** —
* typealias'ы в `:proto` не поддерживают nested-class access.
*/
internal class HttpEventStore(
private val httpClient: HttpClient,
@@ -48,85 +45,80 @@ internal class HttpEventStore(
private val agentUrl: String = baseUrl.trimEnd('/')
override fun events(after: Instant?): Flow<CommonEvent> = flow {
val url = buildString {
append("$agentUrl/outbox/events")
if (after != null) append("?after=$after")
}
httpClient.prepareGet(url) { noReadTimeout() }
.execute { response ->
check(response.status == HttpStatusCode.OK) {
"events: server returned ${response.status}"
}
readSse(response.bodyAsChannel())
.collect { payload ->
emit(agentikJson.decodeFromString(CommonEvent.serializer(), payload))
}
}
}
override fun events(after: Cursor?): Flow<CommonEvent> =
sse("$agentUrl/outbox/events", after, CommonEvent.serializer())
/**
* Override: идём в `/events` напрямую — сервер фильтрует только lifecycle-события.
* Default из [EventStore.agentEvents] читал бы `/events/all` + `filterIsInstance`.
*/
override fun agentEvents(after: Instant?): Flow<CommonEvent.Agent> = flow {
val url = buildString {
append("$agentUrl/events")
if (after != null) append("?after=$after")
}
httpClient.prepareGet(url) { noReadTimeout() }
.execute { response ->
check(response.status == HttpStatusCode.OK) {
"agentEvents: server returned ${response.status}"
}
readSse(response.bodyAsChannel())
.collect { payload ->
val event = agentikJson.decodeFromString(AgentEvent.serializer(), payload)
emit(CommonEvent.Agent(date = event.date, event = event))
}
}
}
override fun agentEvents(after: Cursor?): Flow<CommonEvent.Agent> =
sse("$agentUrl/events", after, CommonEvent.Agent.serializer())
/**
* Override с `conversationId != null` — идём в `/conversations/{id}/events`.
* С `null` (события всех диалогов) — fallback на default impl из [EventStore]:
* общий `/events/all` + filter.
*/
override fun conversationEvents(
after: Instant?,
after: Cursor?,
conversationId: String?,
): Flow<CommonEvent.Conversation> {
if (conversationId == null) {
return super.conversationEvents(after, null)
if (conversationId == null) return super.conversationEvents(after, null)
return sse(
"$agentUrl/conversations/$conversationId/events",
after,
CommonEvent.Conversation.serializer(),
)
}
override suspend fun currentCursor(): Cursor = cursorResponse().current
override suspend fun oldestCursor(): Cursor = cursorResponse().oldest
private suspend fun cursorResponse(): CursorResponse {
val response = httpClient.get("$agentUrl/outbox/cursor")
check(response.status == HttpStatusCode.OK) {
"outbox.cursor: server returned ${response.status}"
}
return flow {
val url = buildString {
append("$agentUrl/conversations/$conversationId/events")
if (after != null) append("?after=$after")
return response.body()
}
private fun <T> sse(url: String, after: Cursor?, serializer: KSerializer<T>): Flow<T> = flow {
httpClient.prepareGet(url) {
noReadTimeout()
if (after != null) {
parameter("epoch", after.epoch)
parameter("offset", after.offset)
}
httpClient.prepareGet(url) { noReadTimeout() }
.execute { response ->
check(response.status == HttpStatusCode.OK) {
"conversationEvents: server returned ${response.status}"
}
readSse(response.bodyAsChannel())
.collect { payload ->
val event = agentikJson.decodeFromString(Event.serializer(), payload)
emit(CommonEvent.Conversation(date = event.date, conversationId = conversationId, event = event))
}
}.execute { response ->
if (response.status == HttpStatusCode.Gone) {
throw response.toGapException(after)
}
check(response.status == HttpStatusCode.OK) {
"$url: server returned ${response.status}"
}
readSse(response.bodyAsChannel())
.collect { payload ->
emit(agentikJson.decodeFromString(serializer, payload))
}
}
}
/**
* У HTTP-варианта нет своего endpoint'а для earliest-event-date.
* Контракт [EventStore.earliestEventDate] для пустого буфера говорит
* "сейчас" — для HTTP-клиента буфер на нашей стороне всегда "пуст"
* (мы не держим своё состояние), поэтому возвращаем `Clock.System.now()`.
*/
override suspend fun earliestEventDate(): Instant = Clock.System.now()
override fun close() {
// HttpClient закрывает владелец (AgentClient / AgentikAgent).
}
}
/** Тело `GET {baseUrl}/outbox/cursor`. */
@Serializable
internal data class CursorResponse(val current: Cursor, val oldest: Cursor)
/** Тело `410 Gone` (см. [pw.binom.agentik.server.OutboxGapResponse]). */
@Serializable
internal data class GapResponse(
val requested: Cursor? = null,
val oldest: Cursor,
val current: Cursor,
)
private suspend fun HttpResponse.toGapException(requested: Cursor?): OutboxGapException {
val dto = runCatching { agentikJson.decodeFromString(GapResponse.serializer(), bodyAsText()) }.getOrNull()
val fallback = requested ?: Cursor(epoch = "", offset = -1L)
return OutboxGapException(
requested = requested,
oldest = dto?.oldest ?: fallback,
current = dto?.current ?: fallback,
)
}
@@ -58,6 +58,23 @@ internal class HttpJournalStore(
return response.body<List<MessageRecord>>()
}
override suspend fun list(
conversationId: String,
afterSeq: Long,
upToSeq: Long,
limit: Int,
): List<MessageRecord> {
val response = httpClient.get("$agentUrl/journal/conversations/$conversationId/messages") {
parameter("afterSeq", afterSeq)
parameter("upToSeq", upToSeq)
parameter("limit", limit)
}
check(response.status == HttpStatusCode.OK) {
"journal.list(seq): server returned ${response.status}"
}
return response.body<List<MessageRecord>>()
}
override suspend fun count(conversationId: String): Long {
val response = httpClient.get("$agentUrl/journal/conversations/$conversationId/count")
check(response.status == HttpStatusCode.OK) {
@@ -76,6 +93,16 @@ internal class HttpJournalStore(
return response.body<CountResponse>().count
}
override suspend fun count(conversationId: String, afterSeq: Long): Long {
val response = httpClient.get("$agentUrl/journal/conversations/$conversationId/count") {
parameter("afterSeq", afterSeq)
}
check(response.status == HttpStatusCode.OK) {
"journal.count(afterSeq): server returned ${response.status}"
}
return response.body<CountResponse>().count
}
override fun close() {
// HttpClient закрывает владелец (AgentClient / AgentikAgent).
}
@@ -11,6 +11,8 @@ import kotlinx.coroutines.flow.MutableSharedFlow
import kotlinx.coroutines.isActive
import kotlinx.coroutines.launch
import pw.binom.agentik.outbox.CommonEvent
import pw.binom.agentik.outbox.Cursor
import pw.binom.agentik.outbox.OutboxGapException
import pw.binom.agentik.outbox.OutboxStore
import kotlin.concurrent.atomics.AtomicBoolean
import kotlin.concurrent.atomics.AtomicReference
@@ -60,6 +62,17 @@ sealed interface ConnectionStatus {
* outbox и т.п.
*/
data class Failed(val cause: Throwable) : ConnectionStatus
/**
* Курсор мёртв ([pw.binom.agentik.outbox.OutboxGapException]): клиент был
* оффлайн дольше retention'а outbox'а или эпоха сменилась. **Не** retry'ится
* (ретрай никогда не пройдёт). Создатель обязан сделать полный resync:
* взять снапшот (`Agent.conversationsSnapshot()` / `Agent.chatSnapshot(id)`),
* применить его и создать новый [ReconnectingOutbox] с курсором снапшота.
*
* Поток [events] закрывается после этого, background-loop останавливается.
*/
data class Gap(val cause: OutboxGapException) : ConnectionStatus
}
/**
@@ -116,8 +129,10 @@ data class BackoffPolicy(
* ```
* val outbox = ReconnectingOutbox(httpEventStore, scope)
*
* // Курсор берётся из снапшота: state + cursor, затем подписка «после него».
* val snap = agent.chatSnapshot(conversationId)
* scope.launch {
* outbox.events(after = Instant.DISTANT_PAST).collect { e -> handle(e) }
* outbox.events(after = snap.cursor).collect { e -> handle(e) }
* }
* scope.launch {
* outbox.connectionStatus().collect { s -> ui.showStatus(s) }
@@ -154,19 +169,27 @@ class ReconnectingOutbox(
private var job: Job? = null
@OptIn(ExperimentalAtomicApi::class)
private val lastSeen: AtomicReference<Instant?> = AtomicReference(null)
private val lastSeen: AtomicReference<Cursor?> = AtomicReference(null)
/**
* Live-события из [outbox] с авто-reconnect. [after] — начальный курсор;
* учитывается только при первом вызове (любом из [events] /
* [connectionStatus]). После reconnect курсор берётся из `date`
* последнего виденного события.
* [connectionStatus]). После reconnect курсор берётся из `offset`
* последнего виденного события (та же `epoch`, что и у подписки).
*
* Если [after] == null, при старте background-loop берётся
* [OutboxStore.currentCursor] — это live-only семантика (событий строго
* после текущего) плюс известная `epoch` для будущих reconnect.
*
* При мёртвом курсоре (retention / смена epoch) loop **не** ретраит, а
* эмитит [ConnectionStatus.Gap] и останавливается — клиент обязан сделать
* resync (снапшот + новый [ReconnectingOutbox] с курсором снапшота).
*
* Коллекторы независимы — каждый получает свою копию потока (shared).
* Медленный коллектор может пропускать события при переполнении буфера
* (`DROP_OLDEST`).
*/
fun events(after: Instant? = null): Flow<CommonEvent> {
fun events(after: Cursor? = null): Flow<CommonEvent> {
ensureStarted(after)
return _events
}
@@ -183,7 +206,7 @@ class ReconnectingOutbox(
}
@OptIn(ExperimentalAtomicApi::class)
private fun ensureStarted(initialCursor: Instant?) {
private fun ensureStarted(initialCursor: Cursor?) {
if (!started.compareAndSet(false, true)) return
lastSeen.store(initialCursor)
job = scope.launch { runLoop() }
@@ -196,9 +219,13 @@ class ReconnectingOutbox(
while (currentCoroutineContext().isActive) {
attempt++
_status.emit(ConnectionStatus.Connecting(attempt))
// Курсор подписки: сохранённый lastSeen, либо (при live-only)
// currentCursor() — чтобы знать epoch и не терять позицию.
val cursor: Cursor? = lastSeen.load() ?: runCatching { outbox.currentCursor() }.getOrNull()
var gap: OutboxGapException? = null
val error: Throwable? = try {
outbox.events(after = lastSeen.load()).collect { event ->
lastSeen.store(event.date)
outbox.events(after = cursor).collect { event ->
lastSeen.store(Cursor(epoch = cursor?.epoch ?: "", offset = event.offset))
_events.emit(event)
if (!connected) {
connected = true
@@ -208,10 +235,18 @@ class ReconnectingOutbox(
null
} catch (t: CancellationException) {
throw t
} catch (t: OutboxGapException) {
gap = t
null
} catch (t: Throwable) {
t
}
connected = false
if (gap != null) {
// Ретраить бессмысленно: курсор мёртв. Отдаём сигнал наружу.
_status.emit(ConnectionStatus.Gap(gap))
return
}
if (attempt >= policy.maxAttempts) {
_status.emit(
ConnectionStatus.Failed(error ?: RuntimeException("outbox flow ended normally"))
@@ -11,10 +11,11 @@ import kotlinx.coroutines.launch
import kotlinx.coroutines.test.advanceTimeBy
import kotlinx.coroutines.test.runCurrent
import kotlinx.coroutines.test.runTest
import pw.binom.agentik.outbox.AgentEvent
import pw.binom.agentik.outbox.CommonEvent
import pw.binom.agentik.outbox.Cursor
import pw.binom.agentik.outbox.OutboxStore
import pw.binom.agentik.outbox.Event
import pw.binom.agentik.outbox.OutboxGapException
import pw.binom.agentik.outbox.DurableEvent
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertNotNull
@@ -41,7 +42,7 @@ internal class FakeOutbox : OutboxStore {
private val channel = Channel<Msg>(Channel.UNLIMITED)
override fun events(after: Instant?): Flow<CommonEvent> = flow {
override fun events(after: Cursor?): Flow<CommonEvent> = flow {
for (msg in channel) {
when (msg) {
is Msg.Err -> throw msg.throwable
@@ -53,20 +54,24 @@ internal class FakeOutbox : OutboxStore {
fun push(event: CommonEvent) { channel.trySend(Msg.Ev(event)) }
fun throwAtNextEvent(t: Throwable) { channel.trySend(Msg.Err(t)) }
override fun agentEvents(after: Instant?): Flow<CommonEvent.Agent> = emptyFlow()
override fun agentEvents(after: Cursor?): Flow<CommonEvent.Agent> = emptyFlow()
override fun conversationEvents(
after: Instant?,
after: Cursor?,
conversationId: String?,
): Flow<CommonEvent.Conversation> = emptyFlow()
override suspend fun earliestEventDate(): Instant = Instant.DISTANT_PAST
override suspend fun currentCursor(): Cursor = Cursor(epoch = "test", offset = -1L)
override suspend fun oldestCursor(): Cursor = Cursor(epoch = "test", offset = -1L)
override fun close() { channel.close() }
}
private const val TEST_EPOCH = "test"
private fun testEvent(dateMs: Long): CommonEvent =
CommonEvent.Conversation(
date = Instant.fromEpochMilliseconds(dateMs),
offset = dateMs,
conversationId = "test",
event = Event.Interrupted(date = Instant.fromEpochMilliseconds(dateMs)),
event = DurableEvent.Interrupted(date = Instant.fromEpochMilliseconds(dateMs)),
)
@OptIn(ExperimentalCoroutinesApi::class)
@@ -141,6 +146,32 @@ class ReconnectingOutboxTest {
assertEquals(0, ctx.eventsLog.size)
}
@Test
fun `gap is not retried and emits Gap status`() = runConnectionTest(attempts = 5) { ctx ->
val fake = ctx.fake
val gap = OutboxGapException(
requested = Cursor("test", -1L),
oldest = Cursor("test", 10L),
current = Cursor("test", 20L),
)
fake.throwAtNextEvent(gap)
ctx.advanceAndDrain(50)
val gaps = ctx.statusLog.filterIsInstance<ConnectionStatus.Gap>()
assertEquals(1, gaps.size, "status=${ctx.statusLog}")
assertEquals(gap, gaps[0].cause)
// Ретрая быть не должно: курсор мёртв, следующая попытка ничего не изменит.
assertTrue(ctx.statusLog.none { it is ConnectionStatus.Disconnected }, "status=${ctx.statusLog}")
assertTrue(
ctx.statusLog.none { it is ConnectionStatus.Connecting && it.attempt == 2 },
"status=${ctx.statusLog}",
)
// Поток событий закрыт — новые эмиссии не доходят.
fake.push(testEvent(2000))
ctx.advanceAndDrain(50)
assertEquals(0, ctx.eventsLog.size)
}
@Test
fun `close cancels background loop`() = runConnectionTest(
attempts = 5,