3 Commits

Author SHA1 Message Date
subochev c0a933d251 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).
2026-09-22 02:55:19 +03:00
subochev 29851c047a docs(client): document persistence contract for AgentikAgent consumers
Android-client review item 8: AgentikAgent(id, baseUrl, engineFactory,
token) constructor was well-documented per parameter but lacked guidance
on what client must persist locally. AgentSettingsRepository rejected —
UI frameworks persist settings differently (JSON file, Keychain, Android
DataStore, NSUserDefaults), lib doesn't impose format.

KDoc on AgentikAgent now contains table of {clientId, baseUrl, token}
with where each comes from and the critical constraint that clientId
must be generated once on first install (UUID.randomUUID().toString())
and never changed — otherwise log multiplexing on the server breaks.
client/README.md 'persistence' section mirrors this for offline reading
with a minimal JSON example.

KDoc-only change. No code, no API surface.
2026-09-22 02:55:08 +03:00
subochev c9995b263e refactor(protocol): add toolName to ToolResult, remove proto typealiases, rename id→toolCallId, drop Conversation.events()
Three protocol-level changes from Android-client review (items 1-3, 5-6):

1) toolName denormalization in ToolResult (3 layers):
   - :outbox-api/Event.ToolResult: +toolName: String? = null
   - :journal-api/MessageRecord.ToolResult: +toolName: String? = null
   - :proto/Message.ToolResult: +toolName: String? = null
   - :storage-ksqlite, :journal-ksqlite ResultPayload codec: +toolName
   - :standalone/ToolDispatcher, ConversationLoop: thread toolName = call.name
   Nullable + default = backward-compat for already-persisted histories
   and existing clients.

2) Drop proto/Event.kt, AgentEvent.kt, CommonEvent.kt typealiases.
   is proto.Event.End failed with 'Unresolved reference End' (alias
   loses nested-class access). Use pw.binom.agentik.outbox.{Event,
   AgentEvent, CommonEvent} directly everywhere — :proto already has
   api(:outbox-api), the package is visible to consumers, no shim
   needed. 21 files rewired, 3 files deleted.

3) Rename Event.ToolResult.id → toolCallId (option B per user).
   In :outbox-api Event.ToolResult.id == Event.ToolCall.id (one value,
   one name); the persistent journal keeps MessageRecord.ToolResult.id
   as its own PK + toolCallId as FK to the call — different semantics,
   left untouched. Fixed ToolDispatcher bug: emitted id = resultId
   while KDoc claimed id == ToolCall.id; now emits toolCallId = callId.

4) Remove Conversation.events() from :proto; OutboxStore is sole event source.
   Conversation is a pure per-conversation abstraction (send/getMessages/
   rename/close). Live events only via agent.outbox.conversationEvents/
   agentEvents/events. HTTP route /conversations/{id}/events stays for
   wire-compat but routes through outbox internally (map { it.event }).

jvmTest green (95 tasks).
2026-09-22 02:55:02 +03:00
28 changed files with 685 additions and 122 deletions
@@ -27,9 +27,10 @@ class SendSubcommand : AgentikSubcommand("send", "Отправить user-ход
// Подписываемся на поток событий ДО send: события, отправленные // Подписываемся на поток событий ДО send: события, отправленные
// до подписки, не реплеятся (shared-flow без replay). // до подписки, не реплеятся (shared-flow без replay).
val eventsJob = launch { val eventsJob = launch {
conv.events(Instant.DISTANT_PAST) agent.outbox.conversationEvents(Instant.DISTANT_PAST, conv.id)
// onEach печатает и терминальный event, takeWhile лишь // onEach печатает и терминальный event, takeWhile лишь
// завершает сбор после него. // завершает сбор после него.
.map { it.event }
.onEach { ev -> emit(ev) } .onEach { ev -> emit(ev) }
.takeWhile { ev -> !isTerminal(ev) } .takeWhile { ev -> !isTerminal(ev) }
.collect { } .collect { }
@@ -53,7 +54,7 @@ class SendSubcommand : AgentikSubcommand("send", "Отправить user-ход
is Event.AppendText -> println("event AppendText ${escape(ev.body)}") is Event.AppendText -> println("event AppendText ${escape(ev.body)}")
is Event.AppendImage -> println("event AppendImage <${ev.body.size}B ${ev.mime}>") is Event.AppendImage -> println("event AppendImage <${ev.body.size}B ${ev.mime}>")
is Event.ToolCall -> println("event ToolCall ${ev.id} ${ev.toolName} ${escape(ev.toolArgs)}") is Event.ToolCall -> println("event ToolCall ${ev.id} ${ev.toolName} ${escape(ev.toolArgs)}")
is Event.ToolResult -> println("event ToolResult ${ev.id} ${escape(ev.result ?: "")}") is Event.ToolResult -> println("event ToolResult ${ev.toolCallId} ${escape(ev.result ?: "")}")
is Event.End -> println("event End") is Event.End -> println("event End")
is Event.Interrupted -> println("event Interrupted") is Event.Interrupted -> println("event Interrupted")
is Event.Error -> println("event Error ${ev.code ?: ""} ${escape(ev.message)}") is Event.Error -> println("event Error ${ev.code ?: ""} ${escape(ev.message)}")
@@ -7,7 +7,7 @@ import kotlinx.coroutines.launch
import pw.binom.agentik.proto.Agent import pw.binom.agentik.proto.Agent
import pw.binom.agentik.proto.Content import pw.binom.agentik.proto.Content
import pw.binom.agentik.proto.Conversation import pw.binom.agentik.proto.Conversation
import pw.binom.agentik.proto.Event import pw.binom.agentik.outbox.Event
import kotlin.coroutines.CoroutineContext import kotlin.coroutines.CoroutineContext
import kotlin.time.Instant import kotlin.time.Instant
@@ -92,12 +92,12 @@ internal class TuiBackend(
} }
/** /**
* Подписывается на [Conversation.events] и перенаправляет их в [state]. * Подписывается на `outbox.conversationEvents(after, conv.id)` и перенаправляет их в [state].
*/ */
private fun subscribeEvents(conv: Conversation, from: Instant) { private fun subscribeEvents(conv: Conversation, from: Instant) {
eventsJob?.cancel() eventsJob?.cancel()
eventsJob = scope.launch { eventsJob = scope.launch {
conv.events(from).collect { ev -> dispatch(ev) } agent.outbox.conversationEvents(from, conv.id).collect { ce -> dispatch(ce.event) }
} }
} }
@@ -8,7 +8,7 @@ import pw.binom.agentik.outbox.OutboxStore
import pw.binom.agentik.proto.Agent import pw.binom.agentik.proto.Agent
import pw.binom.agentik.proto.Content import pw.binom.agentik.proto.Content
import pw.binom.agentik.proto.Conversation import pw.binom.agentik.proto.Conversation
import pw.binom.agentik.proto.Event import pw.binom.agentik.outbox.Event
import pw.binom.agentik.proto.Message import pw.binom.agentik.proto.Message
import pw.binom.agentik.proto.MessageContext import pw.binom.agentik.proto.MessageContext
import kotlin.time.Instant import kotlin.time.Instant
@@ -31,7 +31,7 @@ internal class FakeAgent(
// emptyFlow, journal — error-on-access (никто не должен его трогать). // emptyFlow, journal — error-on-access (никто не должен его трогать).
override val journal: JournalStore = error("journal not used in TuiBackend tests") override val journal: JournalStore = error("journal not used in TuiBackend tests")
override val outbox: OutboxStore = object : OutboxStore { override val outbox: OutboxStore = object : OutboxStore {
override fun events(after: Instant?) = emptyFlow<pw.binom.agentik.proto.CommonEvent>() override fun events(after: Instant?) = emptyFlow<pw.binom.agentik.outbox.CommonEvent>()
override suspend fun earliestEventDate(): Instant = Instant.DISTANT_PAST override suspend fun earliestEventDate(): Instant = Instant.DISTANT_PAST
override fun close() {} override fun close() {}
} }
@@ -4,7 +4,7 @@ import kotlinx.coroutines.ExperimentalCoroutinesApi
import kotlinx.coroutines.test.runCurrent import kotlinx.coroutines.test.runCurrent
import kotlinx.coroutines.test.runTest import kotlinx.coroutines.test.runTest
import pw.binom.agentik.proto.Content import pw.binom.agentik.proto.Content
import pw.binom.agentik.proto.Event import pw.binom.agentik.outbox.Event
import kotlin.test.Test import kotlin.test.Test
import kotlin.test.assertEquals import kotlin.test.assertEquals
import kotlin.test.assertFalse import kotlin.test.assertFalse
@@ -170,7 +170,7 @@ class TuiBackendTest {
runCurrent() runCurrent()
val now = kotlin.time.Clock.System.now() val now = kotlin.time.Clock.System.now()
conv.emit(Event.ToolCall(date = now, id = "1", title = null, toolName = "echo", toolArgs = """{"x":1}""")) conv.emit(Event.ToolCall(date = now, id = "1", title = null, toolName = "echo", toolArgs = """{"x":1}"""))
conv.emit(Event.ToolResult(date = now, id = "1", result = "ok")) conv.emit(Event.ToolResult(date = now, toolCallId = "1", result = "ok"))
runCurrent() runCurrent()
val toolMsgs = state.messages.value.filterIsInstance<TuiMessage.ToolCall>() val toolMsgs = state.messages.value.filterIsInstance<TuiMessage.ToolCall>()
+72 -1
View File
@@ -16,6 +16,11 @@
- `HttpJournalStore` — `list(convId, after, offset, limit)` → `List<MessageRecord>` - `HttpJournalStore` — `list(convId, after, offset, limit)` → `List<MessageRecord>`
со всеми типами записей (User/Assistant/ToolCall/ToolResult/Error + tokens). со всеми типами записей (User/Assistant/ToolCall/ToolResult/Error + tokens).
- `HttpEventStore` — `events` / `agentEvents` / `conversationEvents` (SSE). - `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. Не нужно `Agent` — `AutoCloseable`; `agent.close()` закрывает HttpClient. Не нужно
вручную создавать `HttpClient` и накатывать на него JSON/Bearer-плагины. вручную создавать `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 минут ## Быстрый старт: свой клиент за 5 минут
Один self-contained пример: создаём агента, открываем диалог, Один self-contained пример: создаём агента, открываем диалог,
@@ -322,7 +351,49 @@ UI-обновление списка — отдельная задача, реш
``` ```
Покрывают: JSON-парсинг `Event`-ов, SSE-стрим, recovery после разрыва, Покрывают: 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)`.
## Известное ограничение ## Известное ограничение
@@ -20,12 +20,43 @@ import pw.binom.agentik.proto.Agent
* ) * )
* val conv = agent.createConversation(temp = false) * val conv = agent.createConversation(temp = false)
* conv.send(listOf(Content.Text("hi"))) * conv.send(listOf(Content.Text("hi")))
* conv.events(Instant.DISTANT_PAST).collect { ... } * agent.outbox.conversationEvents(Instant.DISTANT_PAST, conv.id)
* .map { it.event }
* .collect { ... }
* agent.close() // закрывает HttpClient * agent.close() // закрывает HttpClient
* ``` * ```
* *
* [id] пробрасывается в `Agent.id` — сервер про идентичность агента не знает, * ## Что клиент должен хранить локально (persistence)
* поэтому клиент должен её знать сам (или взять из конфига). *
* Либа **не** имеет `SettingsRepository` / `Config` — это намеренно:
* UI-фреймворки хранят настройки по-разному (JSON-файл, Keychain,
* `SharedPreferences`, Android DataStore, NSUserDefaults, ...). Либа
* не навязывает формат, но вот минимальный набор, который клиент должен
* сериализовать у себя, чтобы пережить перезапуск:
*
* | Поле | Что это | Где взять |
* |---|---|---|
* | `id` | Идентичность клиента в логах сервера (X-Client-Id header). Не user-id в агенте, не device-id, а произвольная строка клиента — обычно `<app-name>-<installation-uuid>`. Сервер использует для log multiplexing и не интерпретирует. | Генерируется клиентом при первом запуске, сохраняется локально |
* | `baseUrl` | URL сервера (`http://host:8080/agentik`). Должен включать path-prefix фасада, не только хост. | Из настроек пользователя / дефолт |
* | `token` | Bearer-токен. `null` = анонимный доступ (если сервер разрешает). | Из настроек пользователя / secure-storage |
*
* Опционально (для UX):
* | Поле | Зачем |
* |---|---|
* | `engineFactory` | Зависит от платформы (`CIO` JVM/Native, `OkHttp` JVM, `Darwin` iOS/macOS). Выбор — обычно compile-time. |
*
* Пример минимального persistence-файла (для UI, который хранит JSON):
*
* ```json
* {
* "clientId": "my-android-app-550e8400-e29b-41d4-a716-446655440000",
* "baseUrl": "https://agent.example.com/agentik",
* "token": "s3cret"
* }
* ```
*
* `clientId` генерируется один раз при первой установке (`UUID.randomUUID().toString()`)
* и больше не меняется — иначе сломается log multiplexing на сервере.
* *
* **Lifecycle**: [Agent] — `AutoCloseable`. `agent.close()` закрывает * **Lifecycle**: [Agent] — `AutoCloseable`. `agent.close()` закрывает
* HttpClient (идемпотентно). После этого `createConversation` / * HttpClient (идемпотентно). После этого `createConversation` /
@@ -6,18 +6,12 @@ import io.ktor.client.request.get
import io.ktor.client.request.parameter import io.ktor.client.request.parameter
import io.ktor.client.request.patch import io.ktor.client.request.patch
import io.ktor.client.request.post import io.ktor.client.request.post
import io.ktor.client.request.prepareGet
import io.ktor.client.request.setBody import io.ktor.client.request.setBody
import io.ktor.client.statement.bodyAsChannel
import io.ktor.http.ContentType import io.ktor.http.ContentType
import io.ktor.http.HttpStatusCode
import io.ktor.http.contentType import io.ktor.http.contentType
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.flow
import kotlinx.serialization.Serializable import kotlinx.serialization.Serializable
import pw.binom.agentik.proto.Content import pw.binom.agentik.proto.Content
import pw.binom.agentik.proto.Conversation import pw.binom.agentik.proto.Conversation
import pw.binom.agentik.proto.Event
import pw.binom.agentik.proto.Message import pw.binom.agentik.proto.Message
import pw.binom.agentik.proto.MessageContext import pw.binom.agentik.proto.MessageContext
import kotlin.time.Instant import kotlin.time.Instant
@@ -72,22 +66,6 @@ internal class ConversationClient(
httpClient.post("$convUrl/interrupt") httpClient.post("$convUrl/interrupt")
} }
override fun events(after: Instant): Flow<Event> = flow {
// prepareGet + execute (а не get) обязателен: `get` дожидается полного
// тела ответа, а SSE-поток не заканчивается никогда — вызов висел бы
// вечно. `execute` отдаёт HttpResponse со стриминговым bodyAsChannel.
httpClient.prepareGet("$convUrl/events?after=$after") { noSseReadTimeout() }
.execute { response ->
check(response.status == HttpStatusCode.OK) {
"events: server returned ${response.status}"
}
readSse(response.bodyAsChannel())
.collect { payload ->
emit(agentikJson.decodeFromString(Event.serializer(), payload))
}
}
}
override suspend fun getMessages(after: Instant, offset: Int, limit: Int): List<Message> = override suspend fun getMessages(after: Instant, offset: Int, limit: Int): List<Message> =
httpClient.get("$convUrl/messages") { httpClient.get("$convUrl/messages") {
parameter("after", after.toString()) parameter("after", after.toString())
@@ -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()
}
}
}
@@ -55,6 +55,14 @@ sealed interface MessageRecord {
override val id: String, override val id: String,
override val conversationId: String, override val conversationId: String,
val toolCallId: String, val toolCallId: String,
/**
* Имя тула, денормализованное из соответствующего `MessageRecord.ToolCall.toolName`.
* Денормализация экономна (одна строка в SQLite) и снимает с UI
* необходимость сопоставления `toolCallId → toolName`. `null` —
* безопасный backfill для записей до миграции или для сиротливых
* результатов без предшествующего `ToolCall`.
*/
val toolName: String? = null,
val result: String?, val result: String?,
override val createdAt: Instant, override val createdAt: Instant,
) : MessageRecord ) : MessageRecord
@@ -30,7 +30,7 @@ internal fun encodeRecord(record: MessageRecord): Pair<String, String> = when (r
) )
is MessageRecord.ToolResult -> "tool_result" to Json.encodeToString( is MessageRecord.ToolResult -> "tool_result" to Json.encodeToString(
ResultPayload.serializer(), ResultPayload.serializer(),
ResultPayload(toolCallId = record.toolCallId, result = record.result), ResultPayload(toolCallId = record.toolCallId, toolName = record.toolName, result = record.result),
) )
is MessageRecord.Error -> "error" to Json.encodeToString( is MessageRecord.Error -> "error" to Json.encodeToString(
ErrorPayload.serializer(), ErrorPayload.serializer(),
@@ -59,7 +59,7 @@ internal fun SQLiteResultSet.toMessageRecord(json: Json): MessageRecord {
} }
"tool_result" -> { "tool_result" -> {
val p = Json.decodeFromString(ResultPayload.serializer(), payload) val p = Json.decodeFromString(ResultPayload.serializer(), payload)
MessageRecord.ToolResult(id = id, conversationId = convId, toolCallId = p.toolCallId, result = p.result, createdAt = createdAt) MessageRecord.ToolResult(id = id, conversationId = convId, toolCallId = p.toolCallId, toolName = p.toolName, result = p.result, createdAt = createdAt)
} }
"error" -> { "error" -> {
val p = Json.decodeFromString(ErrorPayload.serializer(), payload) val p = Json.decodeFromString(ErrorPayload.serializer(), payload)
@@ -72,8 +72,18 @@ internal fun SQLiteResultSet.toMessageRecord(json: Json): MessageRecord {
@kotlinx.serialization.Serializable @kotlinx.serialization.Serializable
internal data class CallPayload(val name: String, val title: String?, val argsJson: String) internal data class CallPayload(val name: String, val title: String?, val argsJson: String)
/**
* Тулрезалт-сериализация для SQLite. [toolName] денормализован из
* соответствующего `ToolCall.name` для упрощения UI (нет нужды в
* локальной `Map<id, name>`). Nullable с дефолтом — старые записи
* без поля десериализуются как `null`.
*/
@kotlinx.serialization.Serializable @kotlinx.serialization.Serializable
internal data class ResultPayload(val toolCallId: String, val result: String?) internal data class ResultPayload(
val toolCallId: String,
val toolName: String? = null,
val result: String?,
)
@kotlinx.serialization.Serializable @kotlinx.serialization.Serializable
internal data class ErrorPayload(val message: String, val code: String?) internal data class ErrorPayload(val message: String, val code: String?)
@@ -14,9 +14,9 @@ import kotlin.time.Instant
* идентична `OutboxStore.events`: поток **не реплеит** прошлое, для бэкфилла * идентична `OutboxStore.events`: поток **не реплеит** прошлое, для бэкфилла
* используются `Agent.getConversations` / `getConversation`. * используются `Agent.getConversations` / `getConversation`.
* *
* **История**: раньше жил в `:proto` (как `pw.binom.agentik.proto.AgentEvent`). * **История**: до 2026-09-21 жил в `:proto` как `pw.binom.agentik.proto.AgentEvent`;
* После миграции в `:outbox-api` — `:proto.AgentEvent` стал typealias'ом, * typealias удалён 2026-09-21 (стирал nested-типы в `is`/`when`) — потребители
* backward-compat для существующих импортов сохранён. * импортируют напрямую из `pw.binom.agentik.outbox.AgentEvent`.
*/ */
@Serializable @Serializable
sealed interface AgentEvent { sealed interface AgentEvent {
@@ -9,17 +9,17 @@ import kotlinx.serialization.Serializable
* *
* Useful for admin dashboards, debug tools, parent agents: one subscription * Useful for admin dashboards, debug tools, parent agents: one subscription
* instead of N+1. For regular UI use two separate SSE feeds * instead of N+1. For regular UI use two separate SSE feeds
* ([AgentEvent] via `/events` и `Event` via `/conversations/{id}/events`); * ([AgentEvent] via `/events` и [Event] via `/conversations/{id}/events`);
* [CommonEvent] — for those who need everything in one place. * [CommonEvent] — for those who need everything in one place.
* *
* Server endpoint: `GET /events/all` (SSE), or replay via `OutboxStore.events(after)`. * Server endpoint: `GET /events/all` (SSE), or replay via `OutboxStore.events(after)`.
* *
* **История**: раньше жил в `:proto` (как `pw.binom.agentik.proto.CommonEvent`). * **История**: до 2026-09-21 жил в `:proto` как `pw.binom.agentik.proto.CommonEvent`;
* После миграции в `:outbox-api` — `:proto.CommonEvent` стал typealias'ом, * при миграции в `:outbox-api` был оставлен typealias в `:proto` для backward-compat,
* backward-compat для существующих импортов сохранён. `CommonEvent.Conversation` * но он стирал nested-типы (`CommonEvent.Agent`, `CommonEvent.Conversation`),
* ссылается на [Event] (тоже в `:outbox-api` теперь) — раньше был * что ломало `is CommonEvent.Agent` на стороне клиента. Typealias'ы
* `pw.binom.agentik.proto.Event`, теперь это `pw.binom.agentik.outbox.Event` * `Event`/`AgentEvent`/`CommonEvent` из `:proto` удалены — потребители
* (он тоже typealias-нут в `:proto.Event`). * импортируют напрямую из `pw.binom.agentik.outbox.*`.
*/ */
@Serializable @Serializable
sealed interface CommonEvent { sealed interface CommonEvent {
@@ -15,9 +15,12 @@ import kotlin.time.Instant
* `StartReasoning?` → `StartResponse(TEXT|IMAGE)` → ...контент... → `End` | `Interrupted` | `Error`. * `StartReasoning?` → `StartResponse(TEXT|IMAGE)` → ...контент... → `End` | `Interrupted` | `Error`.
* `StartReasoning` может отсутствовать, если агент не показывал рассуждения. * `StartReasoning` может отсутствовать, если агент не показывал рассуждения.
* *
* **История**: раньше жил в `:proto` (как `pw.binom.agentik.proto.Event`). * **История**: до 2026-09-21 жил в `:proto` как `pw.binom.agentik.proto.Event`;
* После миграции в `:outbox-api` — `:proto.Event` стал typealias'ом, * при миграции в `:outbox-api` был оставлен typealias в `:proto` для
* backward-compat для существующих импортов сохранён. * backward-compat, но он стирал nested-типы (`Event.End`, `Event.ToolCall`,
* `Event.ToolResult`), что ломало `is Event.End` на стороне клиента.
* Typealias удалён 2026-09-21 — потребители импортируют напрямую из
* `pw.binom.agentik.outbox.Event`.
*/ */
@Serializable @Serializable
sealed interface Event { sealed interface Event {
@@ -75,12 +78,26 @@ sealed interface Event {
/** /**
* Результат вызова тула. Приходит целиком после завершения исполнения. * Результат вызова тула. Приходит целиком после завершения исполнения.
* [id] совпадает с [ToolCall.id], к которому относится результат, и *
* с id `Message.ToolResult` в истории. * [toolCallId] = id [ToolCall], к которому относится результат, и
* `MessageRecord.ToolResult.toolCallId` в истории. Один Call → один Result,
* пара `(date, toolCallId)` уникальна — отдельный `id` в live-событии
* не нужен (PK живёт в персистентном журнале).
*
* [toolName] денормализован из соответствующего [ToolCall.toolName] —
* UI рендерит имя тула без локальной `Map<toolCallId, name>` и без риска
* «Result пришёл до Call». `null` допустим для backfill'а старых
* записей, у которых поле отсутствует, или теоретического случая
* Result без предшествующего Call (orphan).
*/ */
@Serializable @Serializable
@SerialName("tool_result") @SerialName("tool_result")
data class ToolResult(override val date: Instant, val id: String, val result: String?) : Event data class ToolResult(
override val date: Instant,
val toolCallId: String,
val toolName: String? = null,
val result: String?,
) : Event
/** /**
* Ошибка хода. После неё поток завершается; дальнейшие события могут * Ошибка хода. После неё поток завершается; дальнейшие события могут
@@ -1,10 +0,0 @@
package pw.binom.agentik.proto
/**
* Backward-compat typealias: `AgentEvent` теперь живёт в `:outbox-api`
* (логически принадлежит сущности outbox, не wire-протоколу `:proto`).
*
* Существующие импорты `pw.binom.agentik.proto.AgentEvent` продолжают
* работать транспарентно. Использовать typealias в новом коде.
*/
typealias AgentEvent = pw.binom.agentik.outbox.AgentEvent
@@ -1,9 +0,0 @@
package pw.binom.agentik.proto
/**
* Backward-compat typealias: `CommonEvent` теперь живёт в `:outbox-api`.
*
* Существующие импорты `pw.binom.agentik.proto.CommonEvent` продолжают
* работать транспарентно.
*/
typealias CommonEvent = pw.binom.agentik.outbox.CommonEvent
@@ -8,6 +8,19 @@ import kotlin.time.Instant
* Stateful-диалог клиента и [Agent]. Хранит собственную историю: на каждый * Stateful-диалог клиента и [Agent]. Хранит собственную историю: на каждый
* [send] агенту не нужно пересылать транскрипт — он уже живёт внутри * [send] агенту не нужно пересылать транскрипт — он уже живёт внутри
* [Conversation]. * [Conversation].
*
* **Live-события** диалога (turn stream: StartReasoning / AppendText / End /
* ToolCall / ToolResult / ...) НЕ часть этого интерфейса — единственный
* источник live-событий это [pw.binom.agentik.outbox.OutboxStore].
* Подписаться на события конкретного диалога:
* ```
* agent.outbox.conversationEvents(after = lastSeen, conversationId = id)
* .map { it.event }
* .collect { e -> ... }
* ```
* Для cross-conversation view (admin / parent-agent / debug):
* `agent.outbox.events(after)`. Для lifecycle агента (created/deleted/renamed):
* `agent.outbox.agentEvents(after)`.
*/ */
interface Conversation : AutoCloseable { interface Conversation : AutoCloseable {
val id: String val id: String
@@ -33,7 +46,7 @@ interface Conversation : AutoCloseable {
/** /**
* Ставит новый user-ход в очередь. Возвращает управление сразу — поток * Ставит новый user-ход в очередь. Возвращает управление сразу — поток
* событий ответа приходит через [events]. * событий ответа приходит через `agent.outbox.conversationEvents(...)`.
* *
* Если в момент вызова выполняется другой ход, новый встаёт в очередь * Если в момент вызова выполняется другой ход, новый встаёт в очередь
* за ним. Чтобы отменить текущий — вызови [interrupt] перед [send]. * за ним. Чтобы отменить текущий — вызови [interrupt] перед [send].
@@ -48,25 +61,13 @@ interface Conversation : AutoCloseable {
/** /**
* Прерывает текущий исполняемый ход (best-effort: LLM-stream прибивается, * Прерывает текущий исполняемый ход (best-effort: LLM-stream прибивается,
* in-flight tool может доехать или отвалиться). В [events] эмитится * in-flight tool может доехать или отвалиться). В `agent.outbox.conversationEvents`
* [Event.Interrupted], затем может начаться следующий ход из очереди. * эмитится `Event.Interrupted`, затем может начаться следующий ход из очереди.
* *
* Если хода нет — no-op. * Если хода нет — no-op.
*/ */
suspend fun interrupt() suspend fun interrupt()
/**
* Live-подписка на всё, что происходит в диалоге, начиная с [after].
*
* **Не реплеит** события, произошедшие до [after] — для бэкфилла
* используй [getMessages]. Если [after] — момент последнего виденного
* клиентом события, поток продолжается «с того места».
*
* Подписки независимы: каждый вызов возвращает свой [Flow], отмена одного
* не влияет на других подписчиков и на сам диалог.
*/
fun events(after: Instant): Flow<Event>
/** Страница истории: не более [limit] сообщений после [after], начиная с [offset]-го. */ /** Страница истории: не более [limit] сообщений после [after], начиная с [offset]-го. */
suspend fun getMessages(after: Instant, offset: Int, limit: Int): List<Message> suspend fun getMessages(after: Instant, offset: Int, limit: Int): List<Message>
@@ -83,7 +84,7 @@ interface Conversation : AutoCloseable {
/** /**
* Освобождает ресурсы диалога (подписки, сетевые хэндлы). Идемпотентно. * Освобождает ресурсы диалога (подписки, сетевые хэндлы). Идемпотентно.
* После [close] дальнейшие вызовы [send]/[interrupt]/[events]/[getMessages]/[rename] не определены. * После [close] дальнейшие вызовы [send]/[interrupt]/[getMessages]/[rename] не определены.
*/ */
override fun close() override fun close()
@@ -1,11 +0,0 @@
package pw.binom.agentik.proto
/**
* Backward-compat typealias: `Event` теперь живёт в `:outbox-api`.
*
* Существующие импорты `pw.binom.agentik.proto.Event` продолжают
* работать транспарентно. `Conversation.events(after): Flow<Event>` в
* `:proto.Conversation` теперь фактически возвращает
* `pw.binom.agentik.outbox.Event` — тот же тип, другое имя.
*/
typealias Event = pw.binom.agentik.outbox.Event
@@ -45,9 +45,22 @@ sealed interface Message {
override val date: Instant override val date: Instant
) : Message ) : Message
/**
* Результат вызова тула. Приходит в историю `getMessages` после завершения хода.
* [id] совпадает с [ToolCall.id], к которому относится результат, и
* с id соответствующего `MessageRecord.ToolResult` в journal.
*
* [toolName] денормализован из [ToolCall.toolName] — UI рендерит
* имя тула в строке результата без отдельной `Map<id, name>`.
*/
@Serializable @Serializable
@SerialName("tool_result") @SerialName("tool_result")
class ToolResult(override val id: String, val result: String?, override val date: Instant) : Message class ToolResult(
override val id: String,
val toolName: String? = null,
val result: String?,
override val date: Instant,
) : Message
/** /**
* Ход завершился ошибкой (LLM, инициализация движка или иная отказоустойчивая * Ход завершился ошибкой (LLM, инициализация движка или иная отказоустойчивая
@@ -3,8 +3,8 @@ package pw.binom.agentik.server
import io.ktor.server.routing.Route import io.ktor.server.routing.Route
import io.ktor.server.routing.get import io.ktor.server.routing.get
import io.ktor.server.routing.route import io.ktor.server.routing.route
import pw.binom.agentik.outbox.CommonEvent
import pw.binom.agentik.outbox.OutboxStore import pw.binom.agentik.outbox.OutboxStore
import pw.binom.agentik.proto.CommonEvent
/** /**
* HTTP-фасад для [OutboxStore] (bounded-tail live event stream агента). * HTTP-фасад для [OutboxStore] (bounded-tail live event stream агента).
@@ -20,11 +20,11 @@ import kotlinx.coroutines.flow.map
import kotlinx.serialization.KSerializer import kotlinx.serialization.KSerializer
import kotlinx.serialization.json.Json import kotlinx.serialization.json.Json
import pw.binom.agentik.proto.Agent import pw.binom.agentik.proto.Agent
import pw.binom.agentik.proto.AgentEvent
import pw.binom.agentik.proto.CommonEvent
import pw.binom.agentik.proto.Conversation import pw.binom.agentik.proto.Conversation
import pw.binom.agentik.proto.Content import pw.binom.agentik.proto.Content
import pw.binom.agentik.proto.Event import pw.binom.agentik.outbox.AgentEvent
import pw.binom.agentik.outbox.CommonEvent
import pw.binom.agentik.outbox.Event
import kotlin.time.Instant import kotlin.time.Instant
internal fun Route.agentikRoutes(agent: Agent) { internal fun Route.agentikRoutes(agent: Agent) {
@@ -128,7 +128,10 @@ internal fun Route.agentikRoutes(agent: Agent) {
return@get return@get
} }
val after = call.parseAfter() ?: return@get val after = call.parseAfter() ?: return@get
call.streamJsonSse(c.events(after), Event.serializer()) // Live-источник событий — `OutboxStore` (единая точка истины);
// разворачиваем `CommonEvent.Conversation` → `Event` для совместимости
// wire-формата (клиент десериализует как `Event`, не как `CommonEvent.Conversation`).
call.streamJsonSse(agent.outbox.conversationEvents(after, id).map { it.event }, Event.serializer())
} }
get("/events") { get("/events") {
@@ -49,7 +49,8 @@ class A2aBridge(private val agent: Agent) : AgentHandler {
val reply = StringBuilder() val reply = StringBuilder()
val turnDone = CompletableDeferred<Unit>() val turnDone = CompletableDeferred<Unit>()
val subscription = async { val subscription = async {
conv.events(since).collect { e -> agent.outbox.conversationEvents(since, conv.id).collect { ce ->
val e = ce.event
when (e) { when (e) {
is Event.AppendText -> reply.append(e.body) is Event.AppendText -> reply.append(e.body)
is Event.End, is Event.Interrupted -> turnDone.complete(Unit) is Event.End, is Event.Interrupted -> turnDone.complete(Unit)
@@ -18,11 +18,11 @@ import pw.binom.agentik.memory.MemorySystemGuidance
import pw.binom.agentik.proto.Agent as ProtoAgent import pw.binom.agentik.proto.Agent as ProtoAgent
import pw.binom.agentik.outbox.AgentEvent import pw.binom.agentik.outbox.AgentEvent
import pw.binom.agentik.outbox.CommonEvent import pw.binom.agentik.outbox.CommonEvent
import pw.binom.agentik.outbox.Event as ProtoEvent
import pw.binom.agentik.outbox.MutableOutboxStore import pw.binom.agentik.outbox.MutableOutboxStore
import pw.binom.agentik.journal.JournalStore import pw.binom.agentik.journal.JournalStore
import pw.binom.agentik.outbox.OutboxStore import pw.binom.agentik.outbox.OutboxStore
import pw.binom.agentik.proto.Conversation as ProtoConversation import pw.binom.agentik.proto.Conversation as ProtoConversation
import pw.binom.agentik.proto.Event as ProtoEvent
import pw.binom.agentik.skills.SkillCatalog import pw.binom.agentik.skills.SkillCatalog
import pw.binom.agentik.skills.renderSystemPromptSection import pw.binom.agentik.skills.renderSystemPromptSection
import pw.binom.agentik.standalone.agent.memory.MemoryToolsFactory import pw.binom.agentik.standalone.agent.memory.MemoryToolsFactory
@@ -3,15 +3,15 @@ package pw.binom.agentik.standalone.agent
import kotlinx.coroutines.runBlocking import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.flow.Flow import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.map import kotlinx.coroutines.flow.map
import pw.binom.agentik.outbox.MutableOutboxStore
import pw.binom.agentik.outbox.CommonEvent import pw.binom.agentik.outbox.CommonEvent
import pw.binom.agentik.proto.Event as ProtoEvent import pw.binom.agentik.outbox.Event
import pw.binom.agentik.outbox.MutableOutboxStore
internal class ConversationEvents( internal class ConversationEvents(
private val globalEventStore: MutableOutboxStore, private val globalEventStore: MutableOutboxStore,
private val conversationId: String, private val conversationId: String,
) { ) {
fun tryEmit(event: ProtoEvent): Boolean { fun tryEmit(event: Event): Boolean {
runBlocking { runBlocking {
globalEventStore.append( globalEventStore.append(
CommonEvent.Conversation( CommonEvent.Conversation(
@@ -24,7 +24,7 @@ internal class ConversationEvents(
return true return true
} }
fun events(after: kotlin.time.Instant?): Flow<ProtoEvent> = fun events(after: kotlin.time.Instant?): Flow<Event> =
globalEventStore.conversationEvents(after = after, conversationId = conversationId) globalEventStore.conversationEvents(after = after, conversationId = conversationId)
.map { it.event } .map { it.event }
} }
@@ -220,9 +220,6 @@ class ConversationLoop(
toolDispatcher.currentToolJob?.cancel() toolDispatcher.currentToolJob?.cancel()
} }
override fun events(after: Instant): Flow<ProtoEvent> =
events.events(after)
override suspend fun getMessages(after: Instant, offset: Int, limit: Int): List<ProtoMessage> = override suspend fun getMessages(after: Instant, offset: Int, limit: Int): List<ProtoMessage> =
messageStore.list(conversationId = id, after = after, offset = offset, limit = limit) messageStore.list(conversationId = id, after = after, offset = offset, limit = limit)
.map { it.toProto() } .map { it.toProto() }
@@ -568,6 +565,7 @@ internal fun MessageRecord.toProto(): ProtoMessage = when (this) {
is MessageRecord.ToolResult -> ProtoMessage.ToolResult( is MessageRecord.ToolResult -> ProtoMessage.ToolResult(
id = id, id = id,
date = createdAt, date = createdAt,
toolName = toolName,
result = result, result = result,
) )
is MessageRecord.Error -> ProtoMessage.Error( is MessageRecord.Error -> ProtoMessage.Error(
@@ -93,7 +93,7 @@ internal class ToolDispatcher(
} }
val resultAt = now() val resultAt = now()
events.tryEmit(ProtoEvent.ToolResult(date = resultAt, id = resultId, result = resultText)) events.tryEmit(ProtoEvent.ToolResult(date = resultAt, toolCallId = callId, toolName = call.name, result = resultText))
// Эмитим background event — другие компоненты (BackgroundScheduler) // Эмитим background event — другие компоненты (BackgroundScheduler)
// решают, делать ли что-то. Cancellation = not a failure (не эмитим Failed). // решают, делать ли что-то. Cancellation = not a failure (не эмитим Failed).
@@ -110,6 +110,7 @@ internal class ToolDispatcher(
id = resultId, id = resultId,
conversationId = state.id, conversationId = state.id,
toolCallId = callId, toolCallId = callId,
toolName = call.name,
result = resultText, result = resultText,
createdAt = resultAt, createdAt = resultAt,
), ),
@@ -319,7 +319,7 @@ class ChatAgentTest {
val events = mutableListOf<ProtoEvent>() val events = mutableListOf<ProtoEvent>()
val job = launch(start = kotlinx.coroutines.CoroutineStart.UNDISPATCHED) { val job = launch(start = kotlinx.coroutines.CoroutineStart.UNDISPATCHED) {
conv.events(Instant.DISTANT_PAST).collect { events.add(it) } agent.outbox.conversationEvents(Instant.DISTANT_PAST, conv.id).collect { events.add(it.event) }
} }
conv.send(listOf(Content.Text("hi"))) conv.send(listOf(Content.Text("hi")))
delay(50) delay(50)
@@ -356,7 +356,7 @@ class ChatAgentTest {
// отправки событий подписка ничего не увидит. // отправки событий подписка ничего не увидит.
val events = mutableListOf<ProtoEvent>() val events = mutableListOf<ProtoEvent>()
val eventsJob = launch(start = kotlinx.coroutines.CoroutineStart.UNDISPATCHED) { val eventsJob = launch(start = kotlinx.coroutines.CoroutineStart.UNDISPATCHED) {
conv.events(Instant.DISTANT_PAST).collect { events.add(it) } agent.outbox.conversationEvents(Instant.DISTANT_PAST, conv.id).collect { events.add(it.event) }
} }
val sendJob = launch { val sendJob = launch {
@@ -414,7 +414,7 @@ class ChatAgentTest {
// Подписываемся ДО send — SharedFlow без replay // Подписываемся ДО send — SharedFlow без replay
val events = mutableListOf<ProtoEvent>() val events = mutableListOf<ProtoEvent>()
val eventsJob = launch(start = kotlinx.coroutines.CoroutineStart.UNDISPATCHED) { val eventsJob = launch(start = kotlinx.coroutines.CoroutineStart.UNDISPATCHED) {
conv.events(Instant.DISTANT_PAST).collect { events.add(it) } agent.outbox.conversationEvents(Instant.DISTANT_PAST, conv.id).collect { events.add(it.event) }
} }
val sendJob = launch { val sendJob = launch {
@@ -30,7 +30,7 @@ internal fun encodeRecord(record: MessageRecord): Pair<String, String> = when (r
) )
is MessageRecord.ToolResult -> "tool_result" to Json.encodeToString( is MessageRecord.ToolResult -> "tool_result" to Json.encodeToString(
ResultPayload.serializer(), ResultPayload.serializer(),
ResultPayload(toolCallId = record.toolCallId, result = record.result), ResultPayload(toolCallId = record.toolCallId, toolName = record.toolName, result = record.result),
) )
is MessageRecord.Error -> "error" to Json.encodeToString( is MessageRecord.Error -> "error" to Json.encodeToString(
ErrorPayload.serializer(), ErrorPayload.serializer(),
@@ -59,7 +59,7 @@ internal fun SQLiteResultSet.toMessageRecord(json: Json): MessageRecord {
} }
"tool_result" -> { "tool_result" -> {
val p = Json.decodeFromString(ResultPayload.serializer(), payload) val p = Json.decodeFromString(ResultPayload.serializer(), payload)
MessageRecord.ToolResult(id = id, conversationId = convId, toolCallId = p.toolCallId, result = p.result, createdAt = createdAt) MessageRecord.ToolResult(id = id, conversationId = convId, toolCallId = p.toolCallId, toolName = p.toolName, result = p.result, createdAt = createdAt)
} }
"error" -> { "error" -> {
val p = Json.decodeFromString(ErrorPayload.serializer(), payload) val p = Json.decodeFromString(ErrorPayload.serializer(), payload)
@@ -72,8 +72,18 @@ internal fun SQLiteResultSet.toMessageRecord(json: Json): MessageRecord {
@kotlinx.serialization.Serializable @kotlinx.serialization.Serializable
internal data class CallPayload(val name: String, val title: String?, val argsJson: String) internal data class CallPayload(val name: String, val title: String?, val argsJson: String)
/**
* Тулрезалт-сериализация для SQLite. [toolName] денормализован из
* соответствующего `ToolCall.name` для упрощения UI (нет нужды в
* локальной `Map<id, name>`). Nullable с дефолтом — старые записи
* без поля десериализуются как `null`.
*/
@kotlinx.serialization.Serializable @kotlinx.serialization.Serializable
internal data class ResultPayload(val toolCallId: String, val result: String?) internal data class ResultPayload(
val toolCallId: String,
val toolName: String? = null,
val result: String?,
)
@kotlinx.serialization.Serializable @kotlinx.serialization.Serializable
internal data class ErrorPayload(val message: String, val code: String?) internal data class ErrorPayload(val message: String, val code: String?)