feat(journal-inmemory): add :journal-inmemory module with InMemoryJournalStore
ci / JVM build + tests (push) Failing after 1m2s

- Introduced a new `:journal-inmemory` module implementing `:journal-api` with an in-memory backend.
- Added `InMemoryJournalStore` for concurrent append, list, and clear operations using `Mutex` and `MutableList`.
- Use cases include tests, dev mode, embedded scenarios, and client-side in-process caching.
- Integrated the module into the project setup and documented usage in `:client/README.md`.
- Added comprehensive
This commit is contained in:
2026-09-21 03:45:40 +03:00
parent 818f022ba3
commit fb963bfb6b
5 changed files with 468 additions and 64 deletions
+291 -64
View File
@@ -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 фабрики. - `AgentikAgent(id, baseUrl, httpClient)` — entry-point. Возвращает `Agent`
- `AgentClient` — список и lifecycle диалогов. (тот же интерфейс, что в `:proto`).
- `ConversationClient` — `send()`, `events()`, `interrupt()`, - `Agent`: `createConversation` / `getConversation` / `getConversations` /
`getMessages()`, `rename()`, `close()`. `deleteConversation` / `journal` / `outbox`.
- Внутренний парсер SSE → `Flow<Event>`. - `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).
Решает: пишем нативный Kotlin-клиент, без curl/JS/Python boilerplate, `HttpClient` создаётся снаружи (выбор движка — на тебе: CIO, OkHttp,
с теми же типами, что и сервер. Один и тот же клиент работает на Darwin). Конфигурация (JSON + Bearer-токен) — через `applyAgentikDefaults`.
JVM, iOS, macOS, Linux, Windows.
## Где используется ## Подключение
- `:agentik-cli` — REPL.
- `:agentik-cli` — JVM/native CLI-клиент поверх `:client`.
- Любой внешний KMP-проект, который хочет встроить агента в свой UI.
## Как подключить
```kotlin ```kotlin
// build.gradle.kts // build.gradle.kts
kotlin { kotlin {
sourceSets.commonMain.dependencies { sourceSets.commonMain.dependencies {
api("pw.binom.agentik:client:0.1.0") api("pw.binom.agentik:client:0.1.0")
} // Опционально — только если будешь использовать `InMemoryJournalStore`
} // как клиентский кэш. Свой `MutableJournalStore` — не нужен.
api("pw.binom.agentik:journal-inmemory: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
} }
} }
``` ```
## Версии ## Быстрый старт: свой клиент за 5 минут
`gradle/libs.versions.toml` → `[versions] agentik-client`. Один self-contained пример: создаём HTTP-клиент, открываем диалог,
отправляем сообщение, печатаем streaming-ответ.
Поддерживает все KMP-таргеты, что и `:proto`.
## Примеры API
```kotlin ```kotlin
// список диалогов import pw.binom.agentik.client.AgentikAgent
agent.getConversations().collect { println(it.id to it.title) } 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-подписка на события отдельного диалога fun main() = runBlocking {
val sub = conversation.events(after = Instant.parse("2026-09-01T00:00:00Z")).collect { } // 1. HTTP-клиент. Движок выбираешь сам (CIO/OkHttp/Darwin).
val http = HttpClient(CIO) { applyAgentikDefaults(token = "s3cret") }
// прерывание текущего хода // 2. Agent — обёртка над :server фасадом.
conversation.interrupt() val agent = AgentikAgent(
id = "my-client",
baseUrl = "http://localhost:8080/agentik",
httpClient = http,
)
// история // 3. Открыть диалог, отправить сообщение.
conversation.getMessages(offset = 0).collect { msg -> val conv = agent.createConversation(temp = false)
when (msg) { conv.send(listOf(Content.Text("Привет")))
is Message.UserMessage -> println("user: ${msg.content}")
is Message.AssistantMessage -> println("assistant: ${msg.content}") // 4. Собирать streaming-ответ.
else -> Unit 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<Event>` и рендеришь как хочешь.
### Добавить локальный кэш истории (ещё 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<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` как образец).
- **Авторизация** — `applyAgentikDefaults(token = "...")` для Bearer;
для OAuth/что-то ещё — конфигурируй `HttpClient` сам.
## Тесты ## Тесты
``` ```
./gradlew :client:jvmTest ./gradlew :client:jvmTest
``` ```
Покрывают: JSON-парсинг Event'ов, SSE-стрим, recovery после разрыва, Покрывают: JSON-парсинг `Event`-ов, SSE-стрим, recovery после разрыва,
401/404. 401/404.
## Чего здесь НЕТ
- Никакого LLM-кода. Это просто клиент.
- Никакого persistent state. История хранится у сервера, клиент её
запрашивает через `getMessages` или подписывается через `events`.
## Текущий статус
Используется продакшеном. Бэкендом служит `:server` поверх `:standalone`,
но клиент совместим с любым сервером, который держит wire-контракт
`:server`.
## Известное ограничение ## Известное ограничение
SSE event-stream в не-TTY ssh-сессии (без `-tt`) закрывается на SSE event-stream в не-TTY ssh-сессии (без `-tt`) закрывается на
default-таймауте Ktor. Используйте либо ssh -tt, либо нативный default-таймауте Ktor. Используйте либо `ssh -tt`, либо нативный
terminal (TTY). Это upstream-особенность Ktor SSE. terminal (TTY). Это upstream-особенность Ktor SSE.
+30
View File
@@ -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)
}
}
}
@@ -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<MessageRecord>`. Для
* 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<MessageRecord> = 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<MessageRecord> = 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'а — снаружи.
}
}
@@ -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())
}
}
+1
View File
@@ -90,6 +90,7 @@ include(":outbox-api")
// KMP in-memory реализация MutableEventStore. ConcurrentLinkedDeque + TTL/size // KMP in-memory реализация MutableEventStore. ConcurrentLinkedDeque + TTL/size
// eviction. Для тестов, dev-режима, embedded-сценариев (Android core). // eviction. Для тестов, dev-режима, embedded-сценариев (Android core).
include(":outbox-inmemory") include(":outbox-inmemory")
include(":journal-inmemory")
//include(":event-store-in-memory") //include(":event-store-in-memory")
//include(":working-memory-api") //include(":working-memory-api")
include(":storage-inmemory") include(":storage-inmemory")