Files
agentik/client/README.md
T
subochev f946186ef5
ci / JVM build + tests (push) Failing after 54s
release / Publish KMP libraries → caffeine Nexus (release) Failing after 5s
feat(client): refactor AgentikAgent to manage its own HttpClient
- `AgentikAgent` now accepts `engineFactory` and an optional `token` to create an internal `HttpClient`, handling all configuration (JSON, Bearer).
- Removed `applyAgentikDefaults` and replaced it with `agentikHttpClient` for `HttpClient` creation with consistent settings.
- Updated `Agent` to implement `AutoCloseable`, ensuring proper resource closure with `agent.close()`.
- Adjusted tests, docs, and examples to align with the new `AgentikAgent` API.
2026-09-21 03:53:35 +03:00

14 KiB
Raw Permalink Blame History

: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).

Agent — AutoCloseable; agent.close() закрывает HttpClient. Не нужно вручную создавать HttpClient и накатывать на него JSON/Bearer-плагины.

Подключение

// 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")
}

Быстрый старт: свой клиент за 5 минут

Один self-contained пример: создаём агента, открываем диалог, отправляем сообщение, печатаем streaming-ответ.

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 строки)

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

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) — главное чтобы реализовывал интерфейс.

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-блок выше.

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
    }
}

Прерывание хода

agent.getConversation(convId)!!.interrupt()

Multi-conversation

Один Agent, много ChatSession:

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.

Известное ограничение

SSE event-stream в не-TTY ssh-сессии (без -tt) закрывается на default-таймауте Ktor. Используйте либо ssh -tt, либо нативный terminal (TTY). Это upstream-особенность Ktor SSE.