# `: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`) / `getMessages(after, offset, limit)` / `rename` / `interrupt` / `close`. - `HttpJournalStore` — `list(convId, after, offset, limit)` → `List` со всеми типами записей (User/Assistant/ToolCall/ToolResult/Error + tokens). - `HttpEventStore` — `events` / `agentEvents` / `conversationEvents` (SSE). - `ReconnectingOutbox(outbox, scope, policy)` — обёртка над `OutboxStore` с авто-reconnect при обрыве стрима (exponential backoff). Два независимых потока: `events()` (те же `CommonEvent`) и `connectionStatus()` (`Connecting`/`Connected`/`Disconnected`/`Failed`) — статус НЕ мешается с основным потоком событий. См. ниже. `Agent` — `AutoCloseable`; `agent.close()` закрывает HttpClient. Не нужно вручную создавать `HttpClient` и накатывать на него JSON/Bearer-плагины. ## Подключение ```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 — это **произвольная строка клиента**, обычно `-`. Сервер использует для 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` и рендеришь как хочешь. `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() 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`. ## Где `: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.