564 lines
27 KiB
Markdown
564 lines
27 KiB
Markdown
# `:client` — Ktor-клиент к `:server` (KMP, jvm + native)
|
||
|
||
Тонкий HTTP-клиент к `:server`-фасаду + локальные примитивы, чтобы
|
||
собирать свои клиенты (UI, CLI, parent-агенты, A2A-bridge) без бойлерплейта
|
||
про HTTP, JSON, SSE и lifecycle `Conversation`.
|
||
|
||
## Что есть
|
||
|
||
- `AgentikAgent(id, baseUrl, engineFactory, token?)` — entry-point. Возвращает
|
||
`Agent` (тот же интерфейс, что в `:proto`). HttpClient создаётся внутри
|
||
из переданной `engineFactory` (`CIO`, `OkHttp`, `Darwin`).
|
||
- `Agent`: `createConversation` / `getConversation` / `getConversations` /
|
||
`deleteConversation` / `journal` / `outbox` / `close`.
|
||
- `Conversation`: `send(content, context?)` / `events(after)` (SSE `Flow<Event>`)
|
||
/ `getMessages(after, offset, limit)` / `rename` / `interrupt` / `close`.
|
||
- `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-плагины.
|
||
|
||
## Подключение
|
||
|
||
```kotlin
|
||
// build.gradle.kts
|
||
dependencies {
|
||
api("pw.binom.agentik:client:0.1.0")
|
||
// Движок — на твой выбор (один из):
|
||
implementation("io.ktor:ktor-client-cio:3.x") // JVM/Native
|
||
implementation("io.ktor:ktor-client-okhttp:3.x") // JVM
|
||
implementation("io.ktor:ktor-client-darwin:3.x") // iOS/macOS
|
||
// Опционально — только если будешь использовать `InMemoryJournalStore`
|
||
// как клиентский кэш. Свой `MutableJournalStore` — не нужен.
|
||
api("pw.binom.agentik:journal-inmemory:0.1.0")
|
||
}
|
||
```
|
||
|
||
## Что клиент хранит локально (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 пример: создаём агента, открываем диалог,
|
||
отправляем сообщение, печатаем streaming-ответ.
|
||
|
||
```kotlin
|
||
import pw.binom.agentik.client.AgentikAgent
|
||
import pw.binom.agentik.content.Content
|
||
import pw.binom.agentik.outbox.Event
|
||
import pw.binom.agentik.outbox.OnlineEvent
|
||
import io.ktor.client.engine.cio.CIO
|
||
import kotlinx.coroutines.launch
|
||
import kotlinx.coroutines.runBlocking
|
||
import kotlin.time.Clock
|
||
|
||
fun main() = runBlocking {
|
||
// 1. Agent — обёртка над :server фасадом. HttpClient создаётся внутри.
|
||
val agent = AgentikAgent(
|
||
id = "my-client",
|
||
baseUrl = "http://localhost:8080/agentik",
|
||
engineFactory = CIO,
|
||
token = "s3cret", // или null, если не нужен
|
||
)
|
||
|
||
val conv = agent.createConversation(temp = false)
|
||
|
||
// 2. Два независимых потока событий диалога:
|
||
// durable (outbox) — целые события, с курсором после переподключения;
|
||
// online (OnlineOutbox) — стриминг ответа, только live (без курсора).
|
||
launch {
|
||
agent.outbox.conversationEvents(after = Clock.System.now(), conversationId = conv.id)
|
||
.collect { ce ->
|
||
when (val ev = ce.event) {
|
||
is Event.AssistantMessage -> println("[answer ready: ${ev.content}]")
|
||
is Event.Interrupted -> println("[interrupted]")
|
||
is Event.Error -> println("[error: ${ev.message}]")
|
||
else -> Unit
|
||
}
|
||
}
|
||
}
|
||
launch {
|
||
agent.onlineOutbox.onlineEvents(conv.id).collect { ev ->
|
||
when (ev) {
|
||
is OnlineEvent.Working -> println("[working]")
|
||
is OnlineEvent.AppendText -> print(ev.body)
|
||
is OnlineEvent.StartResponse -> println("[start]")
|
||
is OnlineEvent.End -> println("\n[end]")
|
||
else -> Unit
|
||
}
|
||
}
|
||
}
|
||
|
||
// 3. Отправить ход (fire-and-forget — ответ придёт по подпискам выше).
|
||
conv.send(listOf(Content.Text("Привет")))
|
||
|
||
// 4. Чистый shutdown.
|
||
conv.close()
|
||
agent.close()
|
||
}
|
||
```
|
||
|
||
**Это весь клиент.** `:server` сам хранит историю, контекст, события.
|
||
Ты только получаешь два типизированных `Flow` и рендеришь как хочешь.
|
||
|
||
> **Durable vs online.** `Event` (в `agent.outbox`) — «целые» события, их
|
||
> можно перезапросить по курсору `after`. `OnlineEvent` (в
|
||
> `agent.onlineOutbox`) — поток стриминга (`Working`/`End`/`AppendText`/
|
||
> `AppendImage`), **никогда не сохраняется** и не реплеится: потерянный при
|
||
> обрыве фрагмент невосстановим, но целый ответ всегда придёт durable-
|
||
> `Event.AssistantMessage` и/или ляжет в journal.
|
||
|
||
`HttpClient`, `applyAgentikDefaults`, выбор engine'а — всё скрыто
|
||
внутри `AgentikAgent`. Один вызов — один готовый `Agent`.
|
||
|
||
### Добавить локальный кэш истории (ещё 4 строки)
|
||
|
||
```kotlin
|
||
import pw.binom.agentik.journal.inmemory.InMemoryJournalStore
|
||
import kotlin.time.Instant
|
||
|
||
// Свой кэш. Хочешь SQLite/JSON/etc. — реализуй MutableJournalStore сам.
|
||
val cache = InMemoryJournalStore()
|
||
|
||
// Backfill + live-refresh в одном фоне:
|
||
launch {
|
||
agent.journal.listFlow(conv.id, Instant.DISTANT_PAST)
|
||
.collect { cache.append(it) }
|
||
}
|
||
|
||
// История — теперь из кэша, без HTTP:
|
||
val history = cache.list(conv.id, Instant.DISTANT_PAST, 0, Int.MAX_VALUE)
|
||
history.forEach { rec ->
|
||
when (rec) {
|
||
is pw.binom.agentik.journal.MessageRecord.UserMessage -> print("user> ${rec.content.text()}")
|
||
is pw.binom.agentik.journal.MessageRecord.AssistantMessage -> print("agent> ${rec.content.text()}")
|
||
is pw.binom.agentik.journal.MessageRecord.ToolCall -> print("[tool: ${rec.toolName}]")
|
||
is pw.binom.agentik.journal.MessageRecord.ToolResult -> print("[result]")
|
||
is pw.binom.agentik.journal.MessageRecord.Error -> print("[error: ${rec.message}]")
|
||
}
|
||
}
|
||
```
|
||
|
||
Шаблон "remote.listFlow → local.append" работает с любым
|
||
`MutableJournalStore` (см. `:journal-api`). Это и есть кэширование
|
||
"без геморроя".
|
||
|
||
### Что вообще не нужно писать самому
|
||
|
||
- HTTP-сериализация `Event`/`Message` — `agentikHttpClient` регистрирует
|
||
`agentikJson` и `InstantSerializer`.
|
||
- SSE-парсер — `readSse()` внутри `:client`.
|
||
- Cursor-менеджмент для `listFlow` — дефолтная имплементация в
|
||
`JournalStore.listFlow` сама пагинирует.
|
||
- Lifecycle подписок на `events()` — `Conversation.close()` отменяет SSE-job.
|
||
- HTTP-клиент и Bearer — `AgentikAgent` создаёт `HttpClient(engineFactory)`
|
||
с Bearer'ом из `token=` под капотом; `agent.close()` его закрывает.
|
||
- Движковые настройки (requestTimeout и пр.) — `HttpClient(engineFactory) { ... }`
|
||
создаётся здесь; для нестандартных движковых настроек используй
|
||
`agentikHttpClient(engineFactory, token)` напрямую (он экспортирован).
|
||
|
||
### Что нужно написать самому
|
||
|
||
- UI-рендеринг `Event`'ов — это твоё (Compose/HTML/CLI).
|
||
- Диалог с пользователем — ввод текста, отображение кнопок и т.п.
|
||
- Persist кэша между запусками (если нужно) — замени `InMemoryJournalStore`
|
||
на свой `MutableJournalStore` (см. `:journal-ksqlite` как пример).
|
||
|
||
|
||
## Базовый пример: send + collect events
|
||
|
||
```kotlin
|
||
import pw.binom.agentik.client.AgentikAgent
|
||
import pw.binom.agentik.content.Content
|
||
import pw.binom.agentik.outbox.Event
|
||
import pw.binom.agentik.outbox.OnlineEvent
|
||
import io.ktor.client.engine.cio.CIO
|
||
import kotlinx.coroutines.launch
|
||
|
||
val agent = AgentikAgent(
|
||
id = "agentik",
|
||
baseUrl = "http://localhost:8080/agentik",
|
||
engineFactory = CIO,
|
||
)
|
||
|
||
val conv = agent.createConversation(temp = false)
|
||
|
||
// durable-поток (с курсором): terminal-события хода.
|
||
launch {
|
||
agent.outbox.conversationEvents(after = kotlin.time.Clock.System.now(), conversationId = conv.id)
|
||
.collect { ce ->
|
||
when (ce.event) {
|
||
is Event.AssistantMessage -> println("\n--- answer ready ---")
|
||
is Event.Error -> error("agent error: ${(ce.event as Event.Error).message}")
|
||
else -> Unit
|
||
}
|
||
}
|
||
}
|
||
// online-поток (live-only): стриминг ответа.
|
||
launch {
|
||
agent.onlineOutbox.onlineEvents(conv.id).collect { ev ->
|
||
when (ev) {
|
||
is OnlineEvent.AppendText -> print(ev.body) // streaming чанки
|
||
is OnlineEvent.End -> println("\n--- end ---")
|
||
else -> Unit
|
||
}
|
||
}
|
||
}
|
||
|
||
conv.send(listOf(Content.Text("Привет, расскажи про себя")))
|
||
```
|
||
|
||
## История с локальным кэшем
|
||
|
||
Главный паттерн: **клиент держит свой `MutableJournalStore` и периодически
|
||
(или разово) синхронизирует с удалённым через `listFlow`**. Дальше всё
|
||
чтение истории — из локального кэша.
|
||
|
||
`InMemoryJournalStore` — это `MutableJournalStore`, ты можешь реализовать
|
||
свой (например с персистентностью в SQLite/JSON/whatever) — главное чтобы
|
||
реализовывал интерфейс.
|
||
|
||
```kotlin
|
||
import pw.binom.agentik.journal.inmemory.InMemoryJournalStore
|
||
import pw.binom.agentik.content.Content
|
||
import pw.binom.agentik.outbox.Event
|
||
import kotlin.time.Instant
|
||
|
||
class ChatSession(
|
||
private val agent: pw.binom.agentik.proto.Agent,
|
||
val conversationId: String,
|
||
) : AutoCloseable {
|
||
|
||
// Локальный кэш. Замените InMemoryJournalStore на свой, если нужна
|
||
// персистентность (SQLite/JSON/etc.) — контракт `MutableJournalStore`
|
||
// (модуль `:journal-api`).
|
||
val cache = InMemoryJournalStore()
|
||
|
||
// Подписка на live-события этого диалога — будем обновлять кэш на `End`.
|
||
private val scope = kotlinx.coroutines.CoroutineScope(
|
||
kotlinx.coroutines.SupervisorJob() +
|
||
kotlinx.coroutines.Dispatchers.Default,
|
||
)
|
||
|
||
init {
|
||
// 1. Backfill: забираем всю историю разговора с сервера.
|
||
scope.launch {
|
||
agent.journal.listFlow(
|
||
conversationId = conversationId,
|
||
after = Instant.DISTANT_PAST,
|
||
).collect { cache.append(it) }
|
||
}
|
||
// 2. Live: на каждом завершённом ходе (durable AssistantMessage)
|
||
// просим у сервера новые записи.
|
||
scope.launch {
|
||
agent.outbox.conversationEvents(Instant.DISTANT_PAST, conversationId).collect { ce ->
|
||
if (ce.event is Event.AssistantMessage) {
|
||
val newest = cache.let {
|
||
// last-seen курсор — последний createdAt в кэше
|
||
it.list(conversationId, Instant.DISTANT_PAST, 0, 1).lastOrNull()?.createdAt
|
||
?: Instant.DISTANT_PAST
|
||
}
|
||
agent.journal.list(conversationId, newest, offset = 0, limit = 100)
|
||
.forEach { cache.append(it) }
|
||
}
|
||
}
|
||
}
|
||
}
|
||
|
||
fun history() = kotlinx.coroutines.runBlocking {
|
||
cache.list(conversationId, Instant.DISTANT_PAST, 0, Int.MAX_VALUE)
|
||
}
|
||
|
||
override fun close() {
|
||
scope.cancel()
|
||
}
|
||
}
|
||
|
||
// Использование:
|
||
val session = ChatSession(agent, conv.id)
|
||
|
||
// История — из кэша:
|
||
session.history().forEach { rec ->
|
||
when (rec) {
|
||
is MessageRecord.UserMessage -> println("user: ${rec.content.text()}")
|
||
is MessageRecord.AssistantMessage -> println("assistant: ${rec.content.text()}")
|
||
is MessageRecord.ToolCall -> println("tool-call: ${rec.toolName}")
|
||
is MessageRecord.ToolResult -> println("tool-result: ${rec.result}")
|
||
is MessageRecord.Error -> println("error: ${rec.message}")
|
||
}
|
||
}
|
||
|
||
// Отправить новое сообщение:
|
||
session.scope.launch {
|
||
agent.getConversation(conversationId)!!.send(listOf(Content.Text("Привет ещё раз")))
|
||
}
|
||
```
|
||
|
||
`InMemoryJournalStore` отдаёт `MessageRecord` со всем payload'ом
|
||
(текст + tool-call/tool-result + tokens). UI сам решает что показать —
|
||
`rec is MessageRecord.UserMessage` для реплик пользователя,
|
||
`rec is MessageRecord.ToolCall` для отрисовки tool-call баббла, и т.п.
|
||
|
||
## Кэш списка бесед
|
||
|
||
`agent.conversationStore` — read-only view поверх `conversation`-таблицы
|
||
на сервере (`ConversationRecord` = id / title / isTemporal / createdAt /
|
||
updatedAt, без `Conversation` handle и без флагов image-support).
|
||
|
||
**Сценарий клиента:** показать список диалогов («как в Telegram»), чтобы
|
||
при открытии UI уже знал названия, не дёргал сервер лишний раз, и
|
||
моментально реагировал на создание/удаление/переименование в другой
|
||
вкладке.
|
||
|
||
Подход — тот же **«remote → local snapshot + live-events»**:
|
||
|
||
```kotlin
|
||
import pw.binom.agentik.client.AgentikAgent
|
||
import pw.binom.agentik.journal.ConversationRecord
|
||
import pw.binom.agentik.outbox.AgentEvent
|
||
import io.ktor.client.engine.cio.CIO
|
||
|
||
// `AgentikAgent` сам оборачивает HTTP-store в локальный кэш:
|
||
// remote.listFlow → local.upsert (snapshot)
|
||
// outbox.agentEvents → local.upsert / delete (live)
|
||
val agent = AgentikAgent(
|
||
id = "agentik",
|
||
baseUrl = "http://localhost:8080/agentik",
|
||
engineFactory = CIO,
|
||
)
|
||
|
||
// Кэш уже наполняется в фоне, читать можно сразу:
|
||
val all = agent.conversationStore.list(0, Int.MAX_VALUE)
|
||
all.forEach { rec -> println("${rec.id} ${rec.title ?: "(no title)"} ${rec.updatedAt}") }
|
||
|
||
// И наблюдать live-изменения (Created/Deleted/Renamed/Touched)
|
||
agent.outbox.agentEvents(kotlin.time.Instant.DISTANT_PAST).collect { ev ->
|
||
when (ev) {
|
||
is AgentEvent.Created -> println("+ ${ev.conversationId}")
|
||
is AgentEvent.Renamed -> println("~ ${ev.id} → ${ev.title}")
|
||
is AgentEvent.Touched -> println("↻ ${ev.id} (${ev.updatedAt})")
|
||
is AgentEvent.Deleted -> println("- ${ev.id}")
|
||
}
|
||
}
|
||
```
|
||
|
||
Если ты **не хочешь** встроенный кэш (например, тебе нужен прямой HTTP
|
||
для бэкенда-сервиса) — `agentikHttpClient(...).raw` оставлен как
|
||
escape-hatch. Сам `InMemoryMutableConversationStore` тоже доступен —
|
||
подмени его на свою реализацию через `wrapWithLocalConversationCache`,
|
||
если нужен SQLite/JSON-store.
|
||
|
||
## Стриминг live-ответа
|
||
|
||
Для streaming-рендера текущего хода подписывайся на `events()` и
|
||
собирай `Event.AppendText`-чанки в свой буфер. Это **не идёт в кэш** —
|
||
только для UI-feedback во время хода. После `End` хода запись уже
|
||
появится в кэше через refresh-блок выше.
|
||
|
||
```kotlin
|
||
import pw.binom.agentik.outbox.OnlineEvent
|
||
|
||
agent.onlineOutbox.onlineEvents(convId).collect { ev ->
|
||
when (ev) {
|
||
is OnlineEvent.Working -> println("[working]")
|
||
is OnlineEvent.StartResponse -> println("[start]")
|
||
is OnlineEvent.AppendText -> print(ev.body)
|
||
is OnlineEvent.AppendImage -> showImage(ev.body)
|
||
is OnlineEvent.End -> println("[end]")
|
||
else -> Unit
|
||
}
|
||
}
|
||
```
|
||
|
||
Инструментальные вызовы и целый ответ — durable-поток
|
||
(`agent.outbox.conversationEvents`): `Event.ToolCall`/`Event.ToolResult` и
|
||
`Event.AssistantMessage`/`Event.Interrupted`/`Event.Error`.
|
||
|
||
## Прерывание хода
|
||
|
||
```kotlin
|
||
agent.getConversation(convId)!!.interrupt()
|
||
```
|
||
|
||
## Multi-conversation
|
||
|
||
Один `Agent`, много `ChatSession`:
|
||
|
||
```kotlin
|
||
val sessions = mutableMapOf<String, ChatSession>()
|
||
|
||
fun open(convId: String): ChatSession =
|
||
sessions.getOrPut(convId) { ChatSession(agent, convId) }
|
||
|
||
fun close(convId: String) {
|
||
sessions.remove(convId)?.close()
|
||
}
|
||
```
|
||
|
||
Подписка на lifecycle диалогов (`agent.outbox.agentEvents(...)`) +
|
||
UI-обновление списка — отдельная задача, решается `Flow<CommonEvent.Agent>`.
|
||
|
||
## Где `:client` НЕ помогает
|
||
|
||
- **UI-рендеринг** — это твоя зона (Compose/HTML/etc.), `:client` только
|
||
отдаёт типы и потоки.
|
||
- **Персистентность кэша** — `InMemoryJournalStore` и
|
||
`InMemoryMutableConversationStore` хранят в RAM. Для диска пиши свой
|
||
`MutableJournalStore` / `MutableConversationStore` (см. `KsqliteJournalStore`
|
||
в `:journal-ksqlite` как образец).
|
||
- **Нестандартные движковые настройки** — для `requestTimeout`,
|
||
прокси и т.п. используй `agentikHttpClient(engineFactory, token)`
|
||
напрямую.
|
||
|
||
## Кэш списка бесед
|
||
|
||
`agent.conversationStore`, который видит клиент — это **локальный кэш**,
|
||
а не прямой HTTP. Внутри `AgentikAgent` (в `wrapWithLocalConversationCache`)
|
||
лежит `InMemoryMutableConversationStore`, синхронизированный с сервером:
|
||
|
||
1. **Seed при старте**: один snapshot через `remote.listFlow(0)` → заливаем
|
||
в `localStore.upsert(...)`.
|
||
2. **Live-обновления**: подписка на `outbox.agentEvents(after)`:
|
||
- `Created(id)` → `remote.get(id)` → `local.upsert(record)`
|
||
- `Deleted(id)` → `local.delete(id)`
|
||
- `Renamed(id, title)` → `local.rename(id, title)`
|
||
- `Touched(id, updatedAt)` → `local.touch(id, updatedAt)`
|
||
|
||
UI читает `agent.conversationStore.list(0, PAGE_SIZE)` — мгновенно, без
|
||
HTTP, в т.ч. оффлайн. Список бесед всегда свежий: сервер эмитит
|
||
`AgentEvent.Created` / `Deleted` / `Renamed` / `Touched` в свой outbox,
|
||
клиент видит их через SSE и применяет к локальной копии.
|
||
|
||
**Команды** (создать / переименовать / удалить) идут через `agent`:
|
||
|
||
```kotlin
|
||
// Создать новую беседу:
|
||
val conv = agent.createConversation(temp = false) // → POST /conversations
|
||
// → server эмитит Created
|
||
// → client cache получает Created
|
||
// → UI увидит её в списке
|
||
// Переименовать:
|
||
agent.renameConversation(conv.id, "Новый заголовок") // → PATCH /conversations/{id}
|
||
// → server эмитит Renamed
|
||
// → client cache обновляет title
|
||
// Удалить:
|
||
agent.deleteConversation(conv.id) // → DELETE /conversations/{id}
|
||
// → server эмитит Deleted
|
||
// → client cache удаляет запись
|
||
```
|
||
|
||
`conversationStore` доступен **только для чтения**. Это read-only projection
|
||
на серверную таблицу `conversation` (id + title + timestamps). Для активной
|
||
работы (send / interrupt) получай handle через `agent.getConversation(id)`.
|
||
|
||
**Никогда не пиши в `conversationStore` напрямую.** Все модификации —
|
||
командами `agent.createConversation / deleteConversation / renameConversation`.
|
||
|
||
### Если хочется своего cache-импла
|
||
|
||
`InMemoryMutableConversationStore` подходит для 99% случаев — Map +
|
||
Mutex, KMP, тесты зелёные. Если нужен диск (cold-start восстановление
|
||
после перезапуска) — реализуй свой `MutableConversationStore` поверх
|
||
SQLite/Room/Core Data, см. `KsqliteMutableConversationStore` в
|
||
`:journal-ksqlite` как образец.
|
||
|
||
```kotlin
|
||
import pw.binom.agentik.journal.MutableConversationStore
|
||
import pw.binom.agentik.journal.ConversationRecord
|
||
|
||
class MySqliteConversationStore(db: MyDb) : MutableConversationStore {
|
||
override suspend fun upsert(record: ConversationRecord) { /* INSERT OR REPLACE */ }
|
||
override suspend fun get(id: String): ConversationRecord? { /* SELECT */ }
|
||
override suspend fun list(offset: Int, limit: Int): List<ConversationRecord> { /* SELECT ORDER BY updatedAt DESC */ }
|
||
override suspend fun delete(id: String): Boolean { /* DELETE */ }
|
||
override suspend fun rename(id: String, title: String?): Instant? { /* UPDATE + bump updatedAt */ }
|
||
override suspend fun touch(id: String, now: Instant) { /* UPDATE updatedAt */ }
|
||
override fun close() {}
|
||
}
|
||
```
|
||
|
||
## Тесты
|
||
|
||
```
|
||
./gradlew :client:jvmTest
|
||
```
|
||
|
||
Покрывают: JSON-парсинг `Event`-ов, SSE-стрим, recovery после разрыва,
|
||
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)`.
|
||
|
||
## Известное ограничение
|
||
|
||
SSE event-stream в не-TTY ssh-сессии (без `-tt`) закрывается на
|
||
default-таймауте Ktor. Используйте либо `ssh -tt`, либо нативный
|
||
terminal (TTY). Это upstream-особенность Ktor SSE.
|