Files
agentik/client
subochev acb4ee6186
ci / JVM build + tests (push) Successful in 5m48s
release / Publish KMP libraries → caffeine Nexus (release) Successful in 35s
fix(client): make ReconnectingOutbox KMP-native-safe (@Volatile→AtomicReference, Math.pow→kotlin.math.pow)
CI red on tag 10 release #1983: client:compileCommonMainKotlinMetadata
and :compileKotlinLinuxArm64 both failed with:
  e: ReconnectingOutbox.kt:154 Unresolved reference 'Volatile'
  e: ReconnectingOutbox.kt:227 Unresolved reference 'Math'

@Volatile is JVM-only annotation; java.lang.Math is JVM-only API. On
linuxArm64/macosArm64 they don't resolve.

Fix:
- @Volatile private var lastSeen: Instant? → AtomicReference<Instant?>
  (kotlin.concurrent.atomics, same module as the AtomicBoolean already
  used for ). .load() / .store() / @OptIn(ExperimentalAtomicApi::class).
- Math.pow(m, e) → m.pow(e) via kotlin.math.pow import.

commonMain stays KMP-clean; jvmTest green (95 tasks); linuxX64 / linuxArm64
/ mingwX64 / macosX64 / macosArm64 compile green.
2026-09-22 03:13:02 +03:00
..

: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-плагины.

Подключение

// 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, который хранит в файле:

{
  "clientId": "my-android-app-550e8400-e29b-41d4-a716-446655440000",
  "baseUrl": "https://agent.example.com/agentik",
  "token": "s3cret"
}

⚠️ clientId генерируется один раз при установке и больше не меняется — иначе сломается log multiplexing на сервере.

Быстрый старт: свой клиент за 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, reconnect-cycle ReconnectingOutbox (4 кейса: успех / обрыв + reconnect / exhausted attempts → Failed / close → cancel).

Auto-reconnect для живого outbox

Базовый OutboxStore.events(after) — cold SSE-стрим, при обрыве (мобильная сеть, рестарт сервера) клиент сам должен реконнектиться с after = lastEventDate. Это повторяется в каждом клиенте. ReconnectingOutbox берёт это на себя:

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.