Files
agentik/client/README.md
T
subochev c0a933d251
ci / JVM build + tests (push) Successful in 6m3s
release / Publish KMP libraries → caffeine Nexus (release) Failing after 23s
feat(client): add ReconnectingOutbox with parallel connectionStatus flow
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

403 lines
18 KiB
Markdown
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# `: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.proto.Content
import pw.binom.agentik.proto.Event
import io.ktor.client.engine.cio.CIO
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, если не нужен
)
// 2. Открыть диалог, отправить сообщение.
val conv = agent.createConversation(temp = false)
conv.send(listOf(Content.Text("Привет")))
// 3. Собирать streaming-ответ.
conv.events(after = Clock.System.now()).collect { ev ->
when (ev) {
is Event.StartResponse -> println("[start]")
is Event.AppendText -> print(ev.body)
is Event.End -> println("[end]")
is Event.Error -> println("[error: ${ev.message}]")
else -> Unit
}
}
// 4. Чистый shutdown.
conv.close()
agent.close()
}
```
**Это весь клиент.** `:server` сам хранит историю, контекст, события.
Ты только получаешь типизированный `Flow<Event>` и рендеришь как хочешь.
`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.proto.Content
import pw.binom.agentik.proto.Event
import io.ktor.client.engine.cio.CIO
val agent = AgentikAgent(
id = "agentik",
baseUrl = "http://localhost:8080/agentik",
engineFactory = CIO,
)
val conv = agent.createConversation(temp = false)
conv.send(listOf(Content.Text("Привет, расскажи про себя")))
conv.events(after = kotlin.time.Clock.System.now()).collect { ev ->
when (ev) {
is Event.AppendText -> print(ev.body) // streaming чанки
is Event.End -> println("\n--- end ---")
is Event.Error -> error("agent error: ${ev.message}")
else -> Unit
}
}
```
## История с локальным кэшем
Главный паттерн: **клиент держит свой `MutableJournalStore` и периодически
(или разово) синхронизирует с удалённым через `listFlow`**. Дальше всё
чтение истории — из локального кэша.
`InMemoryJournalStore` — это `MutableJournalStore`, ты можешь реализовать
свой (например с персистентностью в SQLite/JSON/whatever) — главное чтобы
реализовывал интерфейс.
```kotlin
import pw.binom.agentik.journal.inmemory.InMemoryJournalStore
import pw.binom.agentik.proto.Content
import pw.binom.agentik.proto.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: на каждом `End` хода просим у сервера новые записи.
scope.launch {
agent.getConversation(conversationId)!!.events(Instant.DISTANT_PAST).collect { ev ->
if (ev is Event.End) {
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 баббла, и т.п.
## Стриминг live-ответа
Для streaming-рендера текущего хода подписывайся на `events()` и
собирай `Event.AppendText`-чанки в свой буфер. Это **не идёт в кэш** —
только для UI-feedback во время хода. После `End` хода запись уже
появится в кэше через refresh-блок выше.
```kotlin
import pw.binom.agentik.proto.Event
agent.getConversation(convId)!!.events(Instant.DISTANT_PAST).collect { ev ->
when (ev) {
is Event.StartResponse -> println("[start]")
is Event.AppendText -> print(ev.body)
is Event.AppendImage -> showImage(ev.body)
is Event.ToolCall -> println("[tool: ${ev.toolName}]")
is Event.ToolResult -> println("[result]")
is Event.End -> println("[end]")
is Event.Error -> println("[error: ${ev.message}]")
else -> Unit
}
}
```
## Прерывание хода
```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` хранит в RAM. Для
диска пиши свой `MutableJournalStore` (см. `KsqliteJournalStore` в
`:journal-ksqlite` как образец).
- **Нестандартные движковые настройки** — для `requestTimeout`,
прокси и т.п. используй `agentikHttpClient(engineFactory, token)`
напрямую.
## Тесты
```
./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.