From c0a933d2510b71abc4fc2417b6bc373b9d741aad Mon Sep 17 00:00:00 2001 From: subochev Date: Tue, 22 Sep 2026 02:55:19 +0300 Subject: [PATCH] feat(client): add ReconnectingOutbox with parallel connectionStatus flow MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Android-client review item 9: every client reimplements SSE reconnect with cursor preservation, exponential backoff, and connection-status UI signals. ReconnectingOutbox extracts that into the lib. Design: - wraps any OutboxStore (HttpEventStore or local InMemoryJournalStore) - two INDEPENDENT parallel flows — never mixed: - events(after): Flow with auto-reconnect, cursor (lastSeen) preserved across retries, so client never loses events - connectionStatus(): Flow = Connecting(attempt) / Connected(since) / Disconnected(reason, willRetryIn) / Failed(cause) — for UI banner / spinner; NOT emitted into CommonEvent stream - BackoffPolicy.Default: initial=1s, max=30s, multiplier=2.0, jitter=0.2 (±20% spread), maxAttempts=∞ - BackoffPolicy.Fixed(delay, attempts) for tests - After maxAttempts exhaustion: Failed + flow closes - recon.close() cancels background job, both flows terminate Tests (4 cases, all green): - first event → Connecting(1) + Connected + event delivered - disconnect mid-stream → Disconnected → Connecting(2) → resume from lastSeen cursor (no duplicate) - exhausted attempts → Failed + 0 events - close() → background loop cancelled, no further emissions client/README.md: new 'Auto-reconnect для живого outbox' section with usage example (two parallel scope.launch blocks) + parameter table. jvmTest green (95 tasks, includes 4 new ReconnectingOutboxTest cases). --- client/README.md | 73 +++++- .../agentik/client/ReconnectingOutbox.kt | 239 ++++++++++++++++++ .../agentik/client/ReconnectingOutboxTest.kt | 211 ++++++++++++++++ 3 files changed, 522 insertions(+), 1 deletion(-) create mode 100644 client/src/commonMain/kotlin/pw/binom/agentik/client/ReconnectingOutbox.kt create mode 100644 client/src/commonTest/kotlin/pw/binom/agentik/client/ReconnectingOutboxTest.kt diff --git a/client/README.md b/client/README.md index cd0f517..4ffc871 100644 --- a/client/README.md +++ b/client/README.md @@ -16,6 +16,11 @@ - `HttpJournalStore` — `list(convId, after, offset, limit)` → `List` со всеми типами записей (User/Assistant/ToolCall/ToolResult/Error + tokens). - `HttpEventStore` — `events` / `agentEvents` / `conversationEvents` (SSE). +- `ReconnectingOutbox(outbox, scope, policy)` — обёртка над `OutboxStore` с + авто-reconnect при обрыве стрима (exponential backoff). Два независимых + потока: `events()` (те же `CommonEvent`) и `connectionStatus()` + (`Connecting`/`Connected`/`Disconnected`/`Failed`) — статус НЕ мешается + с основным потоком событий. См. ниже. `Agent` — `AutoCloseable`; `agent.close()` закрывает HttpClient. Не нужно вручную создавать `HttpClient` и накатывать на него JSON/Bearer-плагины. @@ -36,6 +41,30 @@ dependencies { } ``` +## Что клиент хранит локально (persistence) + +Либа **не** имеет `SettingsRepository` / `Config` — это намеренно: UI-фреймворки хранят настройки по-разному (JSON-файл, Keychain, Android DataStore, NSUserDefaults, ...). Либа не навязывает формат, но клиент должен сериализовать у себя минимум: + +| Поле | Что это | Где взять | +|---|---|---| +| `clientId` (параметр `id` в `AgentikAgent`) | Идентичность клиента в логах сервера (X-Client-Id header). Не user-id в агенте, не device-id — это **произвольная строка клиента**, обычно `-`. Сервер использует для log multiplexing и не интерпретирует. | Генерируется один раз при первом запуске (`UUID.randomUUID().toString()`) и сохраняется. Никогда не меняется. | +| `baseUrl` | URL сервера (`http://host:8080/agentik`). Должен включать path-prefix фасада, не только хост. | Из настроек пользователя / дефолт | +| `token` | Bearer-токен. `null` = анонимный доступ (если сервер разрешает). | Из настроек пользователя / secure-storage | + +Опционально (для UX): `engineFactory` — обычно compile-time выбор по платформе (`CIO` JVM/Native, `OkHttp` JVM, `Darwin` iOS/macOS). + +Минимальный JSON для UI, который хранит в файле: + +```json +{ + "clientId": "my-android-app-550e8400-e29b-41d4-a716-446655440000", + "baseUrl": "https://agent.example.com/agentik", + "token": "s3cret" +} +``` + +⚠️ `clientId` **генерируется один раз** при установке и больше не меняется — иначе сломается log multiplexing на сервере. + ## Быстрый старт: свой клиент за 5 минут Один self-contained пример: создаём агента, открываем диалог, @@ -322,7 +351,49 @@ UI-обновление списка — отдельная задача, реш ``` Покрывают: JSON-парсинг `Event`-ов, SSE-стрим, recovery после разрыва, -401/404. +401/404, reconnect-cycle `ReconnectingOutbox` (4 кейса: успех / обрыв + +reconnect / exhausted attempts → Failed / close → cancel). + +## Auto-reconnect для живого outbox + +Базовый `OutboxStore.events(after)` — cold SSE-стрим, при обрыве (мобильная +сеть, рестарт сервера) клиент сам должен реконнектиться с `after = lastEventDate`. +Это повторяется в каждом клиенте. `ReconnectingOutbox` берёт это на себя: + +```kotlin +val recon = ReconnectingOutbox( + outbox = agent.outbox, // или HttpEventStore + scope = myScreenScope, + policy = BackoffPolicy.Default, // 1s → 2s → ... → 30s, ±20% jitter +) + +scope.launch { recon.events(Instant.DISTANT_PAST).collect { handle(it) } } +scope.launch { + recon.connectionStatus().collect { status -> + when (status) { + is Connecting -> ui.showBanner("connecting...") + is Connected -> ui.hideBanner() + is Disconnected -> ui.showBanner("reconnecting in ${status.willRetryIn}…") + is Failed -> ui.showError(status.cause) + } + } +} + +// На выходе (например, navigation back): +recon.close() // отменяет background-loop, потоки терминируются +``` + +Два потока **независимы** — `events()` содержит только `CommonEvent`, +`connectionStatus()` содержит только `ConnectionStatus`. Никакого +"мешающего" `Connecting`/`Disconnected` в потоке событий. + +Параметры backoff (см. `BackoffPolicy`): +- `initial` / `max` — границы задержки +- `multiplier` — множитель на каждом шаге +- `jitter` — рандом-разброс (по умолчанию 20%) +- `maxAttempts` — лимит попыток; после — `Failed` + закрытие потока + +Если нужен фиксированный delay для тестов — `BackoffPolicy.Fixed(10.milliseconds, attempts = 3)`. ## Известное ограничение diff --git a/client/src/commonMain/kotlin/pw/binom/agentik/client/ReconnectingOutbox.kt b/client/src/commonMain/kotlin/pw/binom/agentik/client/ReconnectingOutbox.kt new file mode 100644 index 0000000..bb88273 --- /dev/null +++ b/client/src/commonMain/kotlin/pw/binom/agentik/client/ReconnectingOutbox.kt @@ -0,0 +1,239 @@ +package pw.binom.agentik.client + +import kotlinx.coroutines.CancellationException +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.Job +import kotlinx.coroutines.channels.BufferOverflow +import kotlinx.coroutines.currentCoroutineContext +import kotlinx.coroutines.delay +import kotlinx.coroutines.flow.Flow +import kotlinx.coroutines.flow.MutableSharedFlow +import kotlinx.coroutines.isActive +import kotlinx.coroutines.launch +import pw.binom.agentik.outbox.CommonEvent +import pw.binom.agentik.outbox.OutboxStore +import kotlin.concurrent.atomics.AtomicBoolean +import kotlin.concurrent.atomics.ExperimentalAtomicApi +import kotlin.math.min +import kotlin.random.Random +import kotlin.time.Duration +import kotlin.time.Duration.Companion.seconds +import kotlin.time.Instant + +/** + * Состояние подключения к удалённому [OutboxStore]. Эмитится через + * [ReconnectingOutbox.connectionStatus] — отдельным потоком, **не** + * смешивается с [ReconnectingOutbox.events]. + * + * Типичный цикл: + * ``` + * Connecting(1) → Connected → ... → Disconnected(reason, retryIn) → + * Connecting(2) → Connected → ... + * ``` + * При полном исчерпании попыток ([BackoffPolicy.maxAttempts]) — + * финальный [Failed]. + */ +sealed interface ConnectionStatus { + + /** Начата попытка подключения (включая первую — `attempt == 1`). */ + data class Connecting(val attempt: Int) : ConnectionStatus + + /** Получен первый event с сервера после [Connecting] / [Disconnected]. */ + data class Connected(val since: Instant) : ConnectionStatus + + /** + * Стрим оборвался (network error, server close, таймаут). [reason] — + * причина, `null` если штатное завершение. [willRetryIn] — через сколько + * будет следующая попытка (`null` если [Failed]). + */ + data class Disconnected( + val reason: Throwable?, + val willRetryIn: Duration?, + ) : ConnectionStatus + + /** + * Все попытки исчерпаны ([BackoffPolicy.maxAttempts]). Поток [events] + * закрывается после этого. Создатель [ReconnectingOutbox] должен + * решить, что делать — показать ошибку пользователю, пересоздать + * outbox и т.п. + */ + data class Failed(val cause: Throwable) : ConnectionStatus +} + +/** + * Политика backoff для [ReconnectingOutbox]. Параметры: + * + * - [initial] — задержка перед первой retry-попыткой. + * - [max] — потолок задержки (после серии умножений). + * - [multiplier] — множитель на каждом шаге (например, `2.0` → 1s, 2s, 4s, 8s, ...). + * - [jitter] — доля случайного разброса `[0, jitter]` от текущей задержки + * (например, `0.2` = ±20%). Снижает thundering-herd при массовом reconnect. + * - [maxAttempts] — лимит попыток. `Int.MAX_VALUE` = бесконечно. + */ +data class BackoffPolicy( + val initial: Duration = 1.seconds, + val max: Duration = 30.seconds, + val multiplier: Double = 2.0, + val maxAttempts: Int = Int.MAX_VALUE, + val jitter: Double = 0.2, +) { + init { + require(initial > Duration.ZERO) { "initial must be positive" } + require(max >= initial) { "max must be >= initial" } + require(multiplier >= 1.0) { "multiplier must be >= 1.0" } + require(maxAttempts >= 1) { "maxAttempts must be >= 1" } + require(jitter in 0.0..1.0) { "jitter must be in [0, 1]" } + } + + companion object { + /** 1s → 2s → 4s → ... → 30s, jitter ±20%, бесконечные попытки. */ + val Default: BackoffPolicy = BackoffPolicy() + + /** Только для тестов: фиксированные задержки без разброса. */ + fun Fixed(delay: Duration, attempts: Int = 3): BackoffPolicy = + BackoffPolicy( + initial = delay, + max = delay, + multiplier = 1.0, + maxAttempts = attempts, + jitter = 0.0, + ) + } +} + +/** + * Обёртка над [OutboxStore] с автоматическим reconnect при обрыве стрима. + * + * **Два независимых потока**: + * - [events] — `Flow`, тот же контракт что [OutboxStore.events], + * но с автоматическим переподключением через [BackoffPolicy]. Cursor + * (`lastSeen`) сохраняется между попытками — клиент не теряет события. + * - [connectionStatus] — `Flow`, **параллельный** поток + * lifecycle подключения. Не смешивается с [events]. + * + * ``` + * val outbox = ReconnectingOutbox(httpEventStore, scope) + * + * scope.launch { + * outbox.events(after = Instant.DISTANT_PAST).collect { e -> handle(e) } + * } + * scope.launch { + * outbox.connectionStatus().collect { s -> ui.showStatus(s) } + * } + * + * // На выходе: + * outbox.close() // отменяет background-loop, эмитит Cancelled-как-Disconnected + * ``` + * + * Создатель передаёт свой [scope] — жизненный цикл reconnect-цикла + * привязан к нему. Закрытие scope (или явный [close]) отменяет + * background-loop. После [close] оба flow терминируются. + */ +class ReconnectingOutbox( + private val outbox: OutboxStore, + private val scope: CoroutineScope, + private val policy: BackoffPolicy = BackoffPolicy.Default, + private val random: Random = Random.Default, +) : AutoCloseable { + + private val _events = MutableSharedFlow( + replay = 0, + extraBufferCapacity = 64, + onBufferOverflow = BufferOverflow.DROP_OLDEST, + ) + private val _status = MutableSharedFlow( + replay = 0, + extraBufferCapacity = 64, + onBufferOverflow = BufferOverflow.DROP_OLDEST, + ) + + @OptIn(ExperimentalAtomicApi::class) + private val started = AtomicBoolean(false) + private var job: Job? = null + + @Volatile + private var lastSeen: Instant? = null + + /** + * Live-события из [outbox] с авто-reconnect. [after] — начальный курсор; + * учитывается только при первом вызове (любом из [events] / + * [connectionStatus]). После reconnect курсор берётся из `date` + * последнего виденного события. + * + * Коллекторы независимы — каждый получает свою копию потока (shared). + * Медленный коллектор может пропускать события при переполнении буфера + * (`DROP_OLDEST`). + */ + fun events(after: Instant? = null): Flow { + ensureStarted(after) + return _events + } + + /** + * Lifecycle подключения: [ConnectionStatus.Connecting] / + * [ConnectionStatus.Connected] / [ConnectionStatus.Disconnected] / + * [ConnectionStatus.Failed]. **Не смешивается** с [events] — это + * отдельный поток для UI-индикации статуса сети. + */ + fun connectionStatus(): Flow { + ensureStarted(null) + return _status + } + + @OptIn(ExperimentalAtomicApi::class) + private fun ensureStarted(initialCursor: Instant?) { + if (!started.compareAndSet(false, true)) return + lastSeen = initialCursor + job = scope.launch { runLoop() } + } + + private suspend fun runLoop() { + var attempt = 0 + var connected = false + while (currentCoroutineContext().isActive) { + attempt++ + _status.emit(ConnectionStatus.Connecting(attempt)) + val error: Throwable? = try { + outbox.events(after = lastSeen).collect { event -> + lastSeen = event.date + _events.emit(event) + if (!connected) { + connected = true + _status.emit(ConnectionStatus.Connected(lastSeen!!)) + } + } + null + } catch (t: CancellationException) { + throw t + } catch (t: Throwable) { + t + } + connected = false + if (attempt >= policy.maxAttempts) { + _status.emit( + ConnectionStatus.Failed(error ?: RuntimeException("outbox flow ended normally")) + ) + return + } + val backoff = computeBackoff(attempt) + _status.emit(ConnectionStatus.Disconnected(error, backoff)) + delay(backoff) + } + } + + private fun computeBackoff(attempt: Int): Duration { + // attempt 1 → initial, 2 → initial * m, 3 → initial * m^2, ... + val base = (policy.initial.inWholeMilliseconds.toDouble() * + Math.pow(policy.multiplier, (attempt - 1).toDouble())) + .toLong() + val capped = min(base, policy.max.inWholeMilliseconds) + val jitterMs = (capped * policy.jitter * random.nextDouble()).toLong() + val finalMs = (capped + jitterMs).coerceAtLeast(1L) + return Duration.parse("${finalMs}ms") + } + + override fun close() { + job?.cancel() + job = null + } +} diff --git a/client/src/commonTest/kotlin/pw/binom/agentik/client/ReconnectingOutboxTest.kt b/client/src/commonTest/kotlin/pw/binom/agentik/client/ReconnectingOutboxTest.kt new file mode 100644 index 0000000..b0ed8f9 --- /dev/null +++ b/client/src/commonTest/kotlin/pw/binom/agentik/client/ReconnectingOutboxTest.kt @@ -0,0 +1,211 @@ +package pw.binom.agentik.client + +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.ExperimentalCoroutinesApi +import kotlinx.coroutines.Job +import kotlinx.coroutines.channels.Channel +import kotlinx.coroutines.flow.Flow +import kotlinx.coroutines.flow.emptyFlow +import kotlinx.coroutines.flow.flow +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.OutboxStore +import pw.binom.agentik.outbox.Event +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertNotNull +import kotlin.test.assertTrue +import kotlin.time.Duration +import kotlin.time.Instant + +/** + * In-memory [OutboxStore] для unit-тестов [ReconnectingOutbox]. + * + * Управление: + * - [push] — кладёт [CommonEvent] в очередь, флоу доставит. + * - [throwAtNextEvent] — следующий «тик» `events(after)` бросит этот Throwable + * (симулирует network error / stream break). + * + * Сигнатура [events] идентична боевой — её можно подменить боевым + * `HttpEventStore`, контракт один и тот же. + */ +internal class FakeOutbox : OutboxStore { + private sealed interface Msg { + data class Ev(val event: CommonEvent) : Msg + data class Err(val throwable: Throwable) : Msg + } + + private val channel = Channel(Channel.UNLIMITED) + + override fun events(after: Instant?): Flow = flow { + for (msg in channel) { + when (msg) { + is Msg.Err -> throw msg.throwable + is Msg.Ev -> emit(msg.event) + } + } + } + + fun push(event: CommonEvent) { channel.trySend(Msg.Ev(event)) } + fun throwAtNextEvent(t: Throwable) { channel.trySend(Msg.Err(t)) } + + override fun agentEvents(after: Instant?): Flow = emptyFlow() + override fun conversationEvents( + after: Instant?, + conversationId: String?, + ): Flow = emptyFlow() + override suspend fun earliestEventDate(): Instant = Instant.DISTANT_PAST + override fun close() { channel.close() } +} + +private fun testEvent(dateMs: Long): CommonEvent = + CommonEvent.Conversation( + date = Instant.fromEpochMilliseconds(dateMs), + conversationId = "test", + event = Event.End(date = Instant.fromEpochMilliseconds(dateMs)), + ) + +@OptIn(ExperimentalCoroutinesApi::class) +class ReconnectingOutboxTest { + + @Test + fun `first event after connect emits Connecting then Connected`() = runConnectionTest( + attempts = 5, + ) { ctx -> + val fake = ctx.fake + val status = ctx.statusLog + val events = ctx.eventsLog + + fake.push(testEvent(1000)) + ctx.advanceAndDrain(50) + + assertEquals(1, status.count { it is ConnectionStatus.Connecting && it.attempt == 1 }) + assertEquals(1, status.count { it is ConnectionStatus.Connected }) + assertEquals(1, events.size) + assertEquals(Instant.fromEpochMilliseconds(1000), events[0].date) + } + + @Test + fun `disconnect mid-stream triggers retry with backoff and resumes from last seen`() = + runConnectionTest(attempts = 5) { ctx -> + val fake = ctx.fake + val status = ctx.statusLog + val events = ctx.eventsLog + + fake.push(testEvent(1000)) + ctx.advanceAndDrain(50) + assertEquals(1, events.size) + + // Имитируем обрыв стрима после первого события. + fake.throwAtNextEvent(RuntimeException("simulated network error")) + ctx.advanceAndDrain(50) + + // После Disconnected должен прийти Connecting(2), затем Connected, + // затем новые события без дубля предыдущего. + val disconnectedIndex = status.indexOfFirst { it is ConnectionStatus.Disconnected } + val connecting2Index = status.indexOfFirst { + it is ConnectionStatus.Connecting && it.attempt == 2 + } + assertTrue(disconnectedIndex >= 0, "no Disconnected emitted, got: $status") + assertTrue(connecting2Index > disconnectedIndex, + "expected Connecting(2) after Disconnected, got: $status") + + // Push a new event with later date — cursor preserves lastSeen. + fake.push(testEvent(2000)) + ctx.advanceAndDrain(50) + + assertEquals(2, events.size) + assertEquals(Instant.fromEpochMilliseconds(1000), events[0].date) + assertEquals(Instant.fromEpochMilliseconds(2000), events[1].date) + } + + @Test + fun `exhausted attempts emits Failed and closes flow`() = runConnectionTest( + attempts = 3, + ) { ctx -> + val fake = ctx.fake + + // Каждая попытка connect бросает — все 3 попытки fail. + for (i in 0 until 3) { + fake.throwAtNextEvent(RuntimeException("server is dead #${i + 1}")) + ctx.advanceAndDrain(50) + } + + val failed = ctx.statusLog.filterIsInstance().firstOrNull() + assertNotNull(failed) { "expected Failed status, got: ${ctx.statusLog}" } + assertTrue(failed.cause is RuntimeException) + assertEquals(0, ctx.eventsLog.size) + } + + @Test + fun `close cancels background loop`() = runConnectionTest( + attempts = 5, + ) { ctx -> + val fake = ctx.fake + fake.push(testEvent(1000)) + ctx.advanceAndDrain(50) + assertEquals(1, ctx.eventsLog.size) + + ctx.recon.close() + ctx.advanceAndDrain(100) + + // После close запуск новых эмиссий не должен происходить. + val beforePush = ctx.eventsLog.size + fake.push(testEvent(2000)) + ctx.advanceAndDrain(100) + assertEquals(beforePush, ctx.eventsLog.size) + } + + private data class TestCtx( + val fake: FakeOutbox, + val recon: ReconnectingOutbox, + val statusLog: MutableList, + val eventsLog: MutableList, + val jobs: List, + val scope: CoroutineScope, + val advanceAndDrain: (Long) -> Unit, + ) + + /** + * Запускает [ReconnectingOutbox] с policy из `attempts` попыток по 10ms, + * сабскрайбит на оба потока в собирающие лист, и возвращает [TestCtx] + * с управляемым `advanceAndDrain(ms)` — прокрутить виртуальное время. + */ + @OptIn(ExperimentalCoroutinesApi::class) + private fun runConnectionTest( + attempts: Int, + block: suspend (TestCtx) -> Unit, + ) = runTest { + val policy = BackoffPolicy.Fixed( + delay = Duration.parse("10ms"), + attempts = attempts, + ) + val fake = FakeOutbox() + val recon = ReconnectingOutbox( + outbox = fake, + scope = this, + policy = policy, + ) + val statusLog = mutableListOf() + val eventsLog = mutableListOf() + val jobs = listOf( + launch { recon.connectionStatus().collect { statusLog.add(it) } }, + launch { recon.events().collect { eventsLog.add(it) } }, + ) + val advanceAndDrain: (Long) -> Unit = { ms -> + if (ms > 0) advanceTimeBy(ms) + runCurrent() + } + try { + TestCtx(fake, recon, statusLog, eventsLog, jobs, this, advanceAndDrain).also { block(it) } + } finally { + recon.close() + jobs.forEach { it.cancel() } + fake.close() + } + } +}