feat(client): add ReconnectingOutbox with parallel connectionStatus flow
ci / JVM build + tests (push) Successful in 6m3s
release / Publish KMP libraries → caffeine Nexus (release) Failing after 23s

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<CommonEvent> with auto-reconnect, cursor
    (lastSeen) preserved across retries, so client never loses events
  - connectionStatus(): Flow<ConnectionStatus> = 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).
This commit is contained in:
2026-09-22 02:55:19 +03:00
parent 29851c047a
commit c0a933d251
3 changed files with 522 additions and 1 deletions
+72 -1
View File
@@ -16,6 +16,11 @@
- `HttpJournalStore` — `list(convId, after, offset, limit)` → `List<MessageRecord>`
со всеми типами записей (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 — это **произвольная строка клиента**, обычно `<app-name>-<installation-uuid>`. Сервер использует для 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)`.
## Известное ограничение
@@ -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<CommonEvent>`, тот же контракт что [OutboxStore.events],
* но с автоматическим переподключением через [BackoffPolicy]. Cursor
* (`lastSeen`) сохраняется между попытками — клиент не теряет события.
* - [connectionStatus] — `Flow<ConnectionStatus>`, **параллельный** поток
* 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<CommonEvent>(
replay = 0,
extraBufferCapacity = 64,
onBufferOverflow = BufferOverflow.DROP_OLDEST,
)
private val _status = MutableSharedFlow<ConnectionStatus>(
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<CommonEvent> {
ensureStarted(after)
return _events
}
/**
* Lifecycle подключения: [ConnectionStatus.Connecting] /
* [ConnectionStatus.Connected] / [ConnectionStatus.Disconnected] /
* [ConnectionStatus.Failed]. **Не смешивается** с [events] — это
* отдельный поток для UI-индикации статуса сети.
*/
fun connectionStatus(): Flow<ConnectionStatus> {
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
}
}
@@ -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<Msg>(Channel.UNLIMITED)
override fun events(after: Instant?): Flow<CommonEvent> = 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<CommonEvent.Agent> = emptyFlow()
override fun conversationEvents(
after: Instant?,
conversationId: String?,
): Flow<CommonEvent.Conversation> = 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<ConnectionStatus.Failed>().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<ConnectionStatus>,
val eventsLog: MutableList<CommonEvent>,
val jobs: List<Job>,
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<ConnectionStatus>()
val eventsLog = mutableListOf<CommonEvent>()
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()
}
}
}