diff --git a/client/README.md b/client/README.md index b4be149..87f7cd1 100644 --- a/client/README.md +++ b/client/README.md @@ -1,101 +1,328 @@ -# `:client` — Ktor-клиент к `:server`/`:proto` (KMP, jvm + native) +# `:client` — Ktor-клиент к `:server` (KMP, jvm + native) -## Что это +Тонкий HTTP-клиент к `:server`-фасаду + локальные примитивы, чтобы +собирать свои клиенты (UI, CLI, parent-агенты, A2A-bridge) без бойлерплейта +про HTTP, JSON, SSE и lifecycle `Conversation`. -Ktor client (`io.ktor.client.HttpClient` + `ContentNegotiation(json) + -Sse`), превращающий HTTP/SSE-фасад `:server` в `Agent`/`Conversation` -интерфейсы `:proto`: +## Что есть -- `AgentikAgent(id, baseUrl)` — entry-point фабрики. -- `AgentClient` — список и lifecycle диалогов. -- `ConversationClient` — `send()`, `events()`, `interrupt()`, - `getMessages()`, `rename()`, `close()`. -- Внутренний парсер SSE → `Flow`. +- `AgentikAgent(id, baseUrl, httpClient)` — entry-point. Возвращает `Agent` + (тот же интерфейс, что в `:proto`). +- `Agent`: `createConversation` / `getConversation` / `getConversations` / + `deleteConversation` / `journal` / `outbox`. +- `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). -Решает: пишем нативный Kotlin-клиент, без curl/JS/Python boilerplate, -с теми же типами, что и сервер. Один и тот же клиент работает на -JVM, iOS, macOS, Linux, Windows. +`HttpClient` создаётся снаружи (выбор движка — на тебе: CIO, OkHttp, +Darwin). Конфигурация (JSON + Bearer-токен) — через `applyAgentikDefaults`. -## Где используется - -- `:agentik-cli` — REPL. -- `:agentik-cli` — JVM/native CLI-клиент поверх `:client`. -- Любой внешний KMP-проект, который хочет встроить агента в свой UI. - -## Как подключить +## Подключение ```kotlin // build.gradle.kts kotlin { sourceSets.commonMain.dependencies { api("pw.binom.agentik:client:0.1.0") - } -} - -// ваш код: -val agent = AgentikAgent(id = "agentik", baseUrl = "http://192.168.76.166:8080/agentik") -val conv = agent.createConversation(title = "test") -conv.send(listOf(Content.Text("hello"))).collect { event -> - when (event) { - is Event.AppendText -> print(event.body) - is Event.End -> println("\n--- end ---") - is Event.Error -> error("agent error: ${event.message}") - else -> Unit + // Опционально — только если будешь использовать `InMemoryJournalStore` + // как клиентский кэш. Свой `MutableJournalStore` — не нужен. + api("pw.binom.agentik:journal-inmemory:0.1.0") } } ``` -## Версии +## Быстрый старт: свой клиент за 5 минут -`gradle/libs.versions.toml` → `[versions] agentik-client`. - -Поддерживает все KMP-таргеты, что и `:proto`. - -## Примеры API +Один self-contained пример: создаём HTTP-клиент, открываем диалог, +отправляем сообщение, печатаем streaming-ответ. ```kotlin -// список диалогов -agent.getConversations().collect { println(it.id to it.title) } +import pw.binom.agentik.client.AgentikAgent +import pw.binom.agentik.client.applyAgentikDefaults +import pw.binom.agentik.proto.Content +import pw.binom.agentik.proto.Event +import io.ktor.client.HttpClient +import io.ktor.client.engine.cio.CIO +import kotlinx.coroutines.runBlocking +import kotlin.time.Clock -// live-подписка на события отдельного диалога -val sub = conversation.events(after = Instant.parse("2026-09-01T00:00:00Z")).collect { } +fun main() = runBlocking { + // 1. HTTP-клиент. Движок выбираешь сам (CIO/OkHttp/Darwin). + val http = HttpClient(CIO) { applyAgentikDefaults(token = "s3cret") } -// прерывание текущего хода -conversation.interrupt() + // 2. Agent — обёртка над :server фасадом. + val agent = AgentikAgent( + id = "my-client", + baseUrl = "http://localhost:8080/agentik", + httpClient = http, + ) -// история -conversation.getMessages(offset = 0).collect { msg -> - when (msg) { - is Message.UserMessage -> println("user: ${msg.content}") - is Message.AssistantMessage -> println("assistant: ${msg.content}") - else -> Unit + // 3. Открыть диалог, отправить сообщение. + val conv = agent.createConversation(temp = false) + conv.send(listOf(Content.Text("Привет"))) + + // 4. Собирать 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 + } + } + + // 5. Чистый shutdown. + conv.close() + http.close() +} +``` + +**Это весь клиент.** `:server` сам хранит историю, контекст, события. +Ты только получаешь типизированный `Flow` и рендеришь как хочешь. + +### Добавить локальный кэш истории (ещё 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` — `applyAgentikDefaults` регистрирует + `agentikJson` и `InstantSerializer`. +- SSE-парсер — `readSse()` внутри `:client`. +- Cursor-менеджмент для `listFlow` — дефолтная имплементация в + `JournalStore.listFlow` сама пагинирует. +- Lifecycle подписок на `events()` — `Conversation.close()` отменяет SSE-job. +- Bearer-токен в каждом запросе — `applyAgentikDefaults(token = ...)` инжектит + один раз на весь `HttpClient`. + +### Что нужно написать самому + +- UI-рендеринг `Event`'ов — это твоё (Compose/HTML/CLI). +- Диалог с пользователем — ввод текста, отображение кнопок и т.п. +- Persist кэша между запусками (если нужно) — замени `InMemoryJournalStore` + на свой `MutableJournalStore` (см. `:journal-ksqlite` как пример). + + +## Базовый пример: send + collect events + +```kotlin +import pw.binom.agentik.client.AgentikAgent +import pw.binom.agentik.client.applyAgentikDefaults +import pw.binom.agentik.proto.Content +import pw.binom.agentik.proto.Event +import io.ktor.client.HttpClient +import io.ktor.client.engine.cio.CIO + +val http = HttpClient(CIO) { applyAgentikDefaults(token = "s3cret") } +val agent = AgentikAgent( + id = "agentik", + baseUrl = "http://localhost:8080/agentik", + httpClient = http, +) + +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` как образец). +- **Авторизация** — `applyAgentikDefaults(token = "...")` для Bearer; + для OAuth/что-то ещё — конфигурируй `HttpClient` сам. + ## Тесты ``` ./gradlew :client:jvmTest ``` -Покрывают: JSON-парсинг Event'ов, SSE-стрим, recovery после разрыва, +Покрывают: JSON-парсинг `Event`-ов, SSE-стрим, recovery после разрыва, 401/404. -## Чего здесь НЕТ - -- Никакого LLM-кода. Это просто клиент. -- Никакого persistent state. История хранится у сервера, клиент её - запрашивает через `getMessages` или подписывается через `events`. - -## Текущий статус - -Используется продакшеном. Бэкендом служит `:server` поверх `:standalone`, -но клиент совместим с любым сервером, который держит wire-контракт -`:server`. - ## Известное ограничение SSE event-stream в не-TTY ssh-сессии (без `-tt`) закрывается на -default-таймауте Ktor. Используйте либо ssh -tt, либо нативный +default-таймауте Ktor. Используйте либо `ssh -tt`, либо нативный terminal (TTY). Это upstream-особенность Ktor SSE. diff --git a/journal-inmemory/build.gradle.kts b/journal-inmemory/build.gradle.kts new file mode 100644 index 0000000..2d12009 --- /dev/null +++ b/journal-inmemory/build.gradle.kts @@ -0,0 +1,30 @@ +plugins { + alias(libs.plugins.kotlin.multiplatform) +} + +// KMP-реализация [MutableJournalStore] на `MutableList` + `Mutex` — для +// тестов, dev-режима, embedded-сценариев (Android core, CLI, in-process кэш +// в клиенте) и как образец для своей реализации. +// +// `list` фильтрует по `conversationId`+`createdAt>after` и сортирует +// по `createdAt ASC`. Paging — поверх отфильтрованного списка. +// +// Зависимости: только `:journal-api`. Никакого I/O — pure in-memory. + +kotlin { + jvmToolchain(21) + + jvm() + linuxX64() + mingwX64() + + sourceSets { + commonMain.dependencies { + api(project(":journal-api")) + } + commonTest.dependencies { + implementation(kotlin("test")) + implementation(libs.kotlinx.coroutines.test) + } + } +} diff --git a/journal-inmemory/src/commonMain/kotlin/pw/binom/agentik/journal/inmemory/InMemoryJournalStore.kt b/journal-inmemory/src/commonMain/kotlin/pw/binom/agentik/journal/inmemory/InMemoryJournalStore.kt new file mode 100644 index 0000000..2a6768f --- /dev/null +++ b/journal-inmemory/src/commonMain/kotlin/pw/binom/agentik/journal/inmemory/InMemoryJournalStore.kt @@ -0,0 +1,68 @@ +package pw.binom.agentik.journal.inmemory + +import kotlinx.coroutines.sync.Mutex +import kotlinx.coroutines.sync.withLock +import pw.binom.agentik.journal.JournalStore +import pw.binom.agentik.journal.MessageRecord +import pw.binom.agentik.journal.MutableJournalStore +import kotlin.time.Instant + +/** + * Простая in-memory [MutableJournalStore] для тестов, dev-режима и + * клиентских in-process кэшей. + * + * **Thread-safety**: `Mutex` поверх `MutableList`. Для + * embedded/CLI сценариев достаточно; для hot-path на сервере используйте + * [pw.binom.agentik.journal.ksqlite.KsqliteJournalStore]. + * + * **Контракт `list`**: возвращает подмножество с + * `conversationId == conversationId && createdAt > after`, отсортированное + * по `createdAt ASC`. `offset/limit` — paging поверх отфильтрованного списка. + * + * **Очистка**: [clear] сбрасывает кэш (например, когда диалог удалён + * на сервере). [close] — no-op. + * + * Типичный кэш-паттерн в клиенте: + * ``` + * val local = InMemoryJournalStore() + * val remote = HttpJournalStore(httpClient, baseUrl) + * // backfill + кэширование: + * remote.listFlow(convId, Instant.DISTANT_PAST).collect { local.append(it) } + * // после этого `local.list(convId, after, offset, limit)` отдаёт из кэша. + * ``` + */ +class InMemoryJournalStore : MutableJournalStore { + + private val mutex = Mutex() + private val records: MutableList = mutableListOf() + + override suspend fun append(record: MessageRecord): Unit = mutex.withLock { + records.add(record) + } + + override suspend fun list( + conversationId: String, + after: Instant, + offset: Int, + limit: Int, + ): List = mutex.withLock { + records.asSequence() + .filter { it.conversationId == conversationId && it.createdAt > after } + .sortedBy { it.createdAt } + .drop(offset) + .take(limit) + .toList() + } + + /** Сбросить кэш (например, когда диалог удалён). */ + suspend fun clear(): Unit = mutex.withLock { + records.clear() + } + + /** Сколько записей сейчас в кэше. Для тестов/диагностики. */ + suspend fun size(): Int = mutex.withLock { records.size } + + override fun close() { + // no-op: lifecycle HttpClient'а — снаружи. + } +} diff --git a/journal-inmemory/src/commonTest/kotlin/pw/binom/agentik/journal/inmemory/InMemoryJournalStoreTest.kt b/journal-inmemory/src/commonTest/kotlin/pw/binom/agentik/journal/inmemory/InMemoryJournalStoreTest.kt new file mode 100644 index 0000000..456c54e --- /dev/null +++ b/journal-inmemory/src/commonTest/kotlin/pw/binom/agentik/journal/inmemory/InMemoryJournalStoreTest.kt @@ -0,0 +1,78 @@ +package pw.binom.agentik.journal.inmemory + +import kotlinx.coroutines.test.runTest +import pw.binom.agentik.journal.MessageRecord +import kotlin.time.Duration.Companion.seconds +import kotlin.time.Instant +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertTrue + +class InMemoryJournalStoreTest { + + private fun userMsg(id: String, convId: String, text: String, at: Instant) = + MessageRecord.UserMessage( + id = id, + conversationId = convId, + content = listOf(pw.binom.agentik.journal.Content.Text(text)), + createdAt = at, + ) + + @Test + fun `append then list returns records sorted by createdAt ASC`() = runTest { + val store = InMemoryJournalStore() + val t0 = Instant.parse("2026-09-21T10:00:00Z") + store.append(userMsg("m1", "c1", "first", t0)) + store.append(userMsg("m2", "c1", "second", t0 + 1.seconds)) + store.append(userMsg("m3", "c1", "third", t0 + 2.seconds)) + + val all = store.list("c1", Instant.DISTANT_PAST, 0, 100) + assertEquals(3, all.size) + assertEquals(listOf("m1", "m2", "m3"), all.map { it.id }) + } + + @Test + fun `list filters by conversationId`() = runTest { + val store = InMemoryJournalStore() + val t0 = Instant.parse("2026-09-21T10:00:00Z") + store.append(userMsg("m1", "c1", "a", t0)) + store.append(userMsg("m2", "c2", "b", t0 + 1.seconds)) + store.append(userMsg("m3", "c1", "c", t0 + 2.seconds)) + + assertEquals(2, store.list("c1", Instant.DISTANT_PAST, 0, 100).size) + assertEquals(1, store.list("c2", Instant.DISTANT_PAST, 0, 100).size) + } + + @Test + fun `list filters by after cursor`() = runTest { + val store = InMemoryJournalStore() + val t0 = Instant.parse("2026-09-21T10:00:00Z") + store.append(userMsg("m1", "c1", "a", t0)) + store.append(userMsg("m2", "c1", "b", t0 + 10.seconds)) + store.append(userMsg("m3", "c1", "c", t0 + 20.seconds)) + + val afterT0 = store.list("c1", t0, 0, 100) + assertEquals(listOf("m2", "m3"), afterT0.map { it.id }) + } + + @Test + fun `list applies offset and limit`() = runTest { + val store = InMemoryJournalStore() + val t0 = Instant.parse("2026-09-21T10:00:00Z") + repeat(10) { i -> store.append(userMsg("m$i", "c1", "x", t0 + i.seconds)) } + + val page = store.list("c1", Instant.DISTANT_PAST, offset = 3, limit = 4) + assertEquals(listOf("m3", "m4", "m5", "m6"), page.map { it.id }) + } + + @Test + fun `clear empties the cache`() = runTest { + val store = InMemoryJournalStore() + val t0 = Instant.parse("2026-09-21T10:00:00Z") + store.append(userMsg("m1", "c1", "x", t0)) + assertEquals(1, store.size()) + store.clear() + assertEquals(0, store.size()) + assertTrue(store.list("c1", Instant.DISTANT_PAST, 0, 100).isEmpty()) + } +} diff --git a/settings.gradle.kts b/settings.gradle.kts index 4a42db8..9fd16c6 100644 --- a/settings.gradle.kts +++ b/settings.gradle.kts @@ -90,6 +90,7 @@ include(":outbox-api") // KMP in-memory реализация MutableEventStore. ConcurrentLinkedDeque + TTL/size // eviction. Для тестов, dev-режима, embedded-сценариев (Android core). include(":outbox-inmemory") +include(":journal-inmemory") //include(":event-store-in-memory") //include(":working-memory-api") include(":storage-inmemory")