diff --git a/agentik-tui/src/commonTest/kotlin/pw/binom/agentik/tui/FakeAgent.kt b/agentik-tui/src/commonTest/kotlin/pw/binom/agentik/tui/FakeAgent.kt index a555af7..0f626f3 100644 --- a/agentik-tui/src/commonTest/kotlin/pw/binom/agentik/tui/FakeAgent.kt +++ b/agentik-tui/src/commonTest/kotlin/pw/binom/agentik/tui/FakeAgent.kt @@ -3,6 +3,7 @@ package pw.binom.agentik.tui import kotlinx.coroutines.flow.Flow import kotlinx.coroutines.flow.MutableSharedFlow import kotlinx.coroutines.flow.emptyFlow +import pw.binom.agentik.journal.ConversationStore import pw.binom.agentik.journal.JournalStore import pw.binom.agentik.outbox.OutboxStore import pw.binom.agentik.proto.Agent @@ -32,9 +33,12 @@ internal class FakeAgent( override val journal: JournalStore = error("journal not used in TuiBackend tests") override val outbox: OutboxStore = object : OutboxStore { override fun events(after: Instant?) = emptyFlow() + override fun agentEvents(after: Instant?) = emptyFlow() + override fun conversationEvents(after: Instant?, conversationId: String?) = emptyFlow() override suspend fun earliestEventDate(): Instant = Instant.DISTANT_PAST override fun close() {} } + override val conversationStore: ConversationStore = error("conversationStore not used in TuiBackend tests") override fun createConversation(temp: Boolean): Conversation { createCount++ @@ -49,8 +53,7 @@ internal class FakeAgent( override suspend fun deleteConversation(id: String): Boolean = conversations.removeAll { it.id == id } - override suspend fun getConversations(offset: Int, limit: Int): List = - conversations.toList() + override suspend fun renameConversation(id: String, title: String?): Instant? = null } /** diff --git a/client/README.md b/client/README.md index 4ffc871..ce49c54 100644 --- a/client/README.md +++ b/client/README.md @@ -337,13 +337,81 @@ UI-обновление списка — отдельная задача, реш - **UI-рендеринг** — это твоя зона (Compose/HTML/etc.), `:client` только отдаёт типы и потоки. -- **Персистентность кэша** — `InMemoryJournalStore` хранит в RAM. Для - диска пиши свой `MutableJournalStore` (см. `KsqliteJournalStore` в - `:journal-ksqlite` как образец). +- **Персистентность кэша** — `InMemoryJournalStore` и + `InMemoryMutableConversationStore` хранят в RAM. Для диска пиши свой + `MutableJournalStore` / `MutableConversationStore` (см. `KsqliteJournalStore` + в `:journal-ksqlite` как образец). - **Нестандартные движковые настройки** — для `requestTimeout`, прокси и т.п. используй `agentikHttpClient(engineFactory, token)` напрямую. +## Кэш списка бесед + +`agent.conversationStore`, который видит клиент — это **локальный кэш**, +а не прямой HTTP. Внутри `AgentikAgent` (в `wrapWithLocalConversationCache`) +лежит `InMemoryMutableConversationStore`, синхронизированный с сервером: + +1. **Seed при старте**: один snapshot через `remote.listFlow(0)` → заливаем + в `localStore.upsert(...)`. +2. **Live-обновления**: подписка на `outbox.agentEvents(after)`: + - `Created(id)` → `remote.get(id)` → `local.upsert(record)` + - `Deleted(id)` → `local.delete(id)` + - `Renamed(id, title)` → `local.rename(id, title)` + - `Touched(id, updatedAt)` → `local.touch(id, updatedAt)` + +UI читает `agent.conversationStore.list(0, PAGE_SIZE)` — мгновенно, без +HTTP, в т.ч. оффлайн. Список бесед всегда свежий: сервер эмитит +`AgentEvent.Created` / `Deleted` / `Renamed` / `Touched` в свой outbox, +клиент видит их через SSE и применяет к локальной копии. + +**Команды** (создать / переименовать / удалить) идут через `agent`: + +```kotlin +// Создать новую беседу: +val conv = agent.createConversation(temp = false) // → POST /conversations + // → server эмитит Created + // → client cache получает Created + // → UI увидит её в списке +// Переименовать: +agent.renameConversation(conv.id, "Новый заголовок") // → PATCH /conversations/{id} + // → server эмитит Renamed + // → client cache обновляет title +// Удалить: +agent.deleteConversation(conv.id) // → DELETE /conversations/{id} + // → server эмитит Deleted + // → client cache удаляет запись +``` + +`conversationStore` доступен **только для чтения**. Это read-only projection +на серверную таблицу `conversation` (id + title + timestamps). Для активной +работы (send / interrupt) получай handle через `agent.getConversation(id)`. + +**Никогда не пиши в `conversationStore` напрямую.** Все модификации — +командами `agent.createConversation / deleteConversation / renameConversation`. + +### Если хочется своего cache-импла + +`InMemoryMutableConversationStore` подходит для 99% случаев — Map + +Mutex, KMP, тесты зелёные. Если нужен диск (cold-start восстановление +после перезапуска) — реализуй свой `MutableConversationStore` поверх +SQLite/Room/Core Data, см. `KsqliteMutableConversationStore` в +`:storage-ksqlite` как образец. + +```kotlin +import pw.binom.agentik.journal.MutableConversationStore +import pw.binom.agentik.journal.ConversationRecord + +class MySqliteConversationStore(db: MyDb) : MutableConversationStore { + override suspend fun upsert(record: ConversationRecord) { /* INSERT OR REPLACE */ } + override suspend fun get(id: String): ConversationRecord? { /* SELECT */ } + override suspend fun list(offset: Int, limit: Int): List { /* SELECT ORDER BY updatedAt DESC */ } + override suspend fun delete(id: String): Boolean { /* DELETE */ } + override suspend fun rename(id: String, title: String?): Instant? { /* UPDATE + bump updatedAt */ } + override suspend fun touch(id: String, now: Instant) { /* UPDATE updatedAt */ } + override fun close() {} +} +``` + ## Тесты ``` diff --git a/client/build.gradle.kts b/client/build.gradle.kts index 36c5389..06df030 100644 --- a/client/build.gradle.kts +++ b/client/build.gradle.kts @@ -23,6 +23,7 @@ kotlin { api(project(":proto")) api(project(":outbox-api")) api(project(":journal-api")) + implementation(project(":journal-inmemory")) api(libs.ktor.client.core) implementation(libs.ktor.client.content.negotiation) diff --git a/client/src/commonMain/kotlin/pw/binom/agentik/client/AgentClient.kt b/client/src/commonMain/kotlin/pw/binom/agentik/client/AgentClient.kt index cea1f79..ddaf3e9 100644 --- a/client/src/commonMain/kotlin/pw/binom/agentik/client/AgentClient.kt +++ b/client/src/commonMain/kotlin/pw/binom/agentik/client/AgentClient.kt @@ -5,25 +5,28 @@ import io.ktor.client.call.body import io.ktor.client.request.delete import io.ktor.client.request.get import io.ktor.client.request.parameter +import io.ktor.client.request.patch import io.ktor.client.request.post import io.ktor.client.request.setBody -import io.ktor.client.statement.HttpResponse import io.ktor.http.ContentType import io.ktor.http.HttpStatusCode import io.ktor.http.contentType import kotlinx.coroutines.runBlocking +import pw.binom.agentik.journal.ConversationStore import pw.binom.agentik.journal.JournalStore import pw.binom.agentik.outbox.OutboxStore import pw.binom.agentik.proto.Agent import pw.binom.agentik.proto.Conversation +import kotlin.time.Instant /** * HTTP-реализация [Agent]. Ходит в `:server`-фасад, см. `agentikAgent(...)`. * * HttpClient создаётся внутри из переданного engine и закрывается в [close]. * - * **Storage handles** ([journal], [outbox]) — read-only views на серверные - * хранилища. + * **Storage handles** ([journal], [outbox], [conversationStore]) — read-only + * views на серверные хранилища. Запись — только через команды + * [createConversation] / [deleteConversation] / [renameConversation]. */ internal class AgentClient( override val id: String, @@ -35,6 +38,7 @@ internal class AgentClient( override val outbox: OutboxStore = HttpEventStore(httpClient = httpClient, baseUrl = agentUrl) override val journal: JournalStore = HttpJournalStore(httpClient = httpClient, baseUrl = agentUrl) + override val conversationStore: ConversationStore = HttpConversationStore(httpClient = httpClient, baseUrl = agentUrl) override fun createConversation(temp: Boolean): Conversation = runBlocking { @@ -53,16 +57,18 @@ internal class AgentClient( } override suspend fun deleteConversation(id: String): Boolean { - val response: HttpResponse = httpClient.delete("$agentUrl/conversations/$id") + val response = httpClient.delete("$agentUrl/conversations/$id") return response.status == HttpStatusCode.NoContent } - override suspend fun getConversations(offset: Int, limit: Int): List { - val snapshots = httpClient.get("$agentUrl/conversations") { - parameter("offset", offset) - parameter("limit", limit) - }.body>() - return snapshots.map { ConversationClient(httpClient, agentUrl, it) } + override suspend fun renameConversation(id: String, title: String?): Instant? { + val response = httpClient.patch("$agentUrl/conversations/$id") { + contentType(ContentType.Application.Json) + setBody(RequestRename(title)) + } + if (response.status == HttpStatusCode.NotFound) return null + val rec = response.body() + return rec.updatedAt } override fun close() { diff --git a/client/src/commonMain/kotlin/pw/binom/agentik/client/AgentikAgent.kt b/client/src/commonMain/kotlin/pw/binom/agentik/client/AgentikAgent.kt index e2c6052..82d4c79 100644 --- a/client/src/commonMain/kotlin/pw/binom/agentik/client/AgentikAgent.kt +++ b/client/src/commonMain/kotlin/pw/binom/agentik/client/AgentikAgent.kt @@ -1,7 +1,20 @@ package pw.binom.agentik.client import io.ktor.client.engine.HttpClientEngineFactory +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.Job +import kotlinx.coroutines.SupervisorJob +import kotlinx.coroutines.cancel +import kotlinx.coroutines.launch +import kotlinx.coroutines.runBlocking +import pw.binom.agentik.journal.ConversationRecord +import pw.binom.agentik.journal.ConversationStore +import pw.binom.agentik.journal.MutableConversationStore +import pw.binom.agentik.journal.inmemory.InMemoryMutableConversationStore +import pw.binom.agentik.outbox.AgentEvent import pw.binom.agentik.proto.Agent +import kotlin.time.Instant /** * Создаёт [Agent], который ходит в HTTP-фасад `agentikAgent` (модуль `:server`). @@ -23,7 +36,7 @@ import pw.binom.agentik.proto.Agent * agent.outbox.conversationEvents(Instant.DISTANT_PAST, conv.id) * .map { it.event } * .collect { ... } - * agent.close() // закрывает HttpClient + * agent.close() // закрывает HttpClient + локальный кэш * ``` * * ## Что клиент должен хранить локально (persistence) @@ -58,17 +71,103 @@ import pw.binom.agentik.proto.Agent * `clientId` генерируется один раз при первой установке (`UUID.randomUUID().toString()`) * и больше не меняется — иначе сломается log multiplexing на сервере. * + * ## Локальный кэш списка бесед + * + * [conversationStore], который видит клиент — это **кэш**, не прямой HTTP. + * Внутри лежит [InMemoryMutableConversationStore], который: + * 1. На старте делает snapshot через `remote.listFlow(0)` → `local.upsert(...)`. + * 2. Подписывается на `outbox.agentEvents(after)` → для каждого + * [AgentEvent.Created] / `Deleted` / `Renamed` / `Touched` применяет + * соответствующий `upsert/delete/rename/touch` к локальной копии. + * + * UI читает `agent.conversationStore.list(0, PAGE_SIZE)` — мгновенно, + * без HTTP, в т.ч. оффлайн. Команды (create/delete/rename) идут + * через [Agent] и **не** через `conversationStore` (он read-only). + * * **Lifecycle**: [Agent] — `AutoCloseable`. `agent.close()` закрывает - * HttpClient (идемпотентно). После этого `createConversation` / - * `getConversation` etc. не определены. + * HttpClient + локальный кэш + background-coroutine (идемпотентно). + * После этого `createConversation` / `getConversation` etc. не определены. */ fun AgentikAgent( id: String, baseUrl: String, engineFactory: HttpClientEngineFactory<*>, token: String? = null, -): Agent = AgentClient( - id = id, - baseUrl = baseUrl, - httpClient = agentikHttpClient(engineFactory = engineFactory, token = token), -) +): Agent { + val httpClient = agentikHttpClient(engineFactory = engineFactory, token = token) + val client = AgentClient(id = id, baseUrl = baseUrl, httpClient = httpClient) + return wrapWithLocalConversationCache(client, scopeClient = client) +} + +/** + * Оборачивает [Agent] так, что [Agent.conversationStore] становится + * локальным in-memory кэшем, синхронизированным с удалённым стором + * через outbox-события. + * + * - **Seed**: при создании делает один snapshot через + * `remote.listFlow(0)` и заливает в [InMemoryMutableConversationStore]. + * - **Live**: подписка на `agent.outbox.agentEvents(after)` применяет + * `Created` / `Deleted` / `Renamed` / `Touched` к локальному кэшу. + * + * Возвращает обёртку, у которой переопределён только [Agent.conversationStore] + * (на read-only projection локального [InMemoryMutableConversationStore]). + * Остальные методы [Agent] — delegated в [delegate]. + */ +private fun wrapWithLocalConversationCache( + delegate: Agent, + scopeClient: Agent, +): Agent = object : Agent by delegate { + + private val localStore: MutableConversationStore = InMemoryMutableConversationStore() + private val cacheScope: CoroutineScope = CoroutineScope(SupervisorJob() + Dispatchers.Default) + private val syncJob: Job + + init { + // Делаем cacheStore read-only view на localStore. + // (Через вложенный класс — см. ниже.) + // Запускаем seed + live-refresh параллельно. + syncJob = cacheScope.launch { + // 1. seed — snapshot всех текущих бесед с сервера + try { + delegate.conversationStore.listFlow(offset = 0, pageSize = ConversationStore.PAGE_SIZE) + .collect { rec -> localStore.upsert(rec) } + } catch (_: Throwable) { + // seed может упасть (offline / 5xx) — не критично, + // live-источник всё равно догонит при первом событии. + } + + // 2. live — применяем outbox-события. + // Используем `first()` для knownId после Created — потом отписываемся, + // потому что Created нужно вытянуть полный record через `remote.get(id)`. + // Renamed/Touched меняют локальную копию без round-trip. + delegate.outbox.agentEvents(after = Instant.DISTANT_PAST).collect { ce -> + when (val ev = ce.event) { + is AgentEvent.Created -> { + // Created не несёт title/timestamps — нужно сходить в remote. + val rec = delegate.conversationStore.get(ev.conversationId) + if (rec != null) localStore.upsert(rec) + } + is AgentEvent.Deleted -> localStore.delete(ev.id) + is AgentEvent.Renamed -> localStore.rename(ev.id, ev.title) + is AgentEvent.Touched -> localStore.touch(ev.id, ev.updatedAt) + } + } + } + } + + /** + * Read-only projection локального кэша — клиент через него только + * читает (`get` / `list` / `listFlow`). + */ + override val conversationStore: ConversationStore = object : ConversationStore { + override suspend fun get(id: String): ConversationRecord? = localStore.get(id) + override suspend fun list(offset: Int, limit: Int): List = localStore.list(offset, limit) + override fun close() {} // owned by outer close + } + + override fun close() { + cacheScope.cancel() + runBlocking { syncJob.join() } + delegate.close() + } +} diff --git a/client/src/commonMain/kotlin/pw/binom/agentik/client/Dto.kt b/client/src/commonMain/kotlin/pw/binom/agentik/client/Dto.kt index db9eabf..2bc615a 100644 --- a/client/src/commonMain/kotlin/pw/binom/agentik/client/Dto.kt +++ b/client/src/commonMain/kotlin/pw/binom/agentik/client/Dto.kt @@ -23,4 +23,4 @@ data class ConversationSnapshot( internal data class RequestCreateConversation(val temp: Boolean) @Serializable -internal data class RequestRename(val title: String) +internal data class RequestRename(val title: String?) diff --git a/docs/STANDALONE.md b/docs/STANDALONE.md index 9aacd80..87bcdc3 100644 --- a/docs/STANDALONE.md +++ b/docs/STANDALONE.md @@ -194,7 +194,7 @@ suspend fun compact(dropFromOrderIdx: Long, conversationId: String): Long `compact` — атомарный «выбросить всё от `dropFromOrderIdx` и дальше, вставить новую синтетическую запись на следующий `order_idx`». Для v1 — просто `DELETE` от индекса (суммаризация появится в v2 вместе с LLM-вызовом для генерации текста). -### `ConversationStore` +### `MutableConversationStore` ```kotlin suspend fun upsert(record: ConversationRecord) diff --git a/journal-api/src/commonMain/kotlin/pw/binom/agentik/journal/ConversationStore.kt b/journal-api/src/commonMain/kotlin/pw/binom/agentik/journal/ConversationStore.kt index f16ec26..22e31f9 100644 --- a/journal-api/src/commonMain/kotlin/pw/binom/agentik/journal/ConversationStore.kt +++ b/journal-api/src/commonMain/kotlin/pw/binom/agentik/journal/ConversationStore.kt @@ -1,27 +1,45 @@ package pw.binom.agentik.journal -import kotlin.time.Instant +import kotlinx.coroutines.flow.Flow +import kotlinx.coroutines.flow.flow /** - * CRUD по таблице `conversation`. + * Read-only view of the `conversation` table (CRUD-операции находятся + * в [MutableConversationStore] и используются только внутри ChatAgent). + * + * Клиенты видят [ConversationStore] через [pw.binom.agentik.proto.Agent.conversationStore] + * (по аналогии с `journal` / `outbox`) и строят свой локальный кэш: + * - **seed** через [list] (snapshot страницы) или [listFlow] (cold-flow paging); + * - **live-refresh** через `outbox.agentEvents()` — Created / Deleted / + * Renamed / Touched. + * + * Запись в хранилище **не** делается клиентом — только команды + * `agent.createConversation / deleteConversation / renameConversation`. */ interface ConversationStore : AutoCloseable { - /** Создать или обновить snapshot диалога. */ - suspend fun upsert(record: ConversationRecord) - /** Диалог по id, или `null`. */ suspend fun get(id: String): ConversationRecord? - /** Удалить диалог (вместе с его сообщениями и working memory). */ - suspend fun delete(id: String): Boolean - /** Список диалогов, отсортированный по `updatedAt` DESC. */ suspend fun list(offset: Int, limit: Int): List - /** Переименовать диалог; `null` для сброса заголовка. Возвращает новый `updatedAt` или `null`, если не найден. */ - suspend fun rename(id: String, title: String?): Instant? + /** + * Cold-flow paging через [list]. Default-реализация делает N+1 round-trip + * (по странице через `list()` пока не получит короткую страницу). Для + * HTTP-импл — это лишние round-trip'ы; реализация может переопределить. + */ + fun listFlow(offset: Int = 0, pageSize: Int = PAGE_SIZE): Flow = flow { + var skip = offset + while (true) { + val page = list(skip, pageSize) + if (page.isEmpty()) break + page.forEach { emit(it) } + skip += page.size + } + } - /** Обновить `updatedAt` диалога (например, после отправки сообщения). */ - suspend fun touch(id: String, now: Instant) + companion object { + const val PAGE_SIZE: Int = 100 + } } diff --git a/journal-api/src/commonMain/kotlin/pw/binom/agentik/journal/MutableConversationStore.kt b/journal-api/src/commonMain/kotlin/pw/binom/agentik/journal/MutableConversationStore.kt new file mode 100644 index 0000000..a4110f0 --- /dev/null +++ b/journal-api/src/commonMain/kotlin/pw/binom/agentik/journal/MutableConversationStore.kt @@ -0,0 +1,21 @@ +package pw.binom.agentik.journal + +import kotlin.time.Instant + +/** + * CRUD по таблице `conversation`. + */ +interface MutableConversationStore : ConversationStore { + + /** Создать или обновить snapshot диалога. */ + suspend fun upsert(record: ConversationRecord) + + /** Удалить диалог (вместе с его сообщениями и working memory). */ + suspend fun delete(id: String): Boolean + + /** Переименовать диалог; `null` для сброса заголовка. Возвращает новый `updatedAt` или `null`, если не найден. */ + suspend fun rename(id: String, title: String?): Instant? + + /** Обновить `updatedAt` диалога (например, после отправки сообщения). */ + suspend fun touch(id: String, now: Instant) +} diff --git a/journal-inmemory/build.gradle.kts b/journal-inmemory/build.gradle.kts index 2d12009..41e9f7a 100644 --- a/journal-inmemory/build.gradle.kts +++ b/journal-inmemory/build.gradle.kts @@ -2,12 +2,10 @@ plugins { alias(libs.plugins.kotlin.multiplatform) } -// KMP-реализация [MutableJournalStore] на `MutableList` + `Mutex` — для -// тестов, dev-режима, embedded-сценариев (Android core, CLI, in-process кэш -// в клиенте) и как образец для своей реализации. -// -// `list` фильтрует по `conversationId`+`createdAt>after` и сортирует -// по `createdAt ASC`. Paging — поверх отфильтрованного списка. +// KMP-реализация [MutableJournalStore] и [MutableConversationStore] на +// `MutableList`/`MutableMap` + `Mutex` — для тестов, dev-режима, +// embedded-сценариев (Android core, CLI, in-process кэш в клиенте) и как +// образец для своей реализации. // // Зависимости: только `:journal-api`. Никакого I/O — pure in-memory. @@ -15,7 +13,10 @@ kotlin { jvmToolchain(21) jvm() + macosX64() + macosArm64() linuxX64() + linuxArm64() mingwX64() sourceSets { diff --git a/storage-inmemory/src/commonMain/kotlin/pw/binom/agentik/storage/inmemory/InMemoryConversationStore.kt b/journal-inmemory/src/commonMain/kotlin/pw/binom/agentik/journal/inmemory/InMemoryMutableConversationStore.kt similarity index 72% rename from storage-inmemory/src/commonMain/kotlin/pw/binom/agentik/storage/inmemory/InMemoryConversationStore.kt rename to journal-inmemory/src/commonMain/kotlin/pw/binom/agentik/journal/inmemory/InMemoryMutableConversationStore.kt index 02d4062..13e65d0 100644 --- a/storage-inmemory/src/commonMain/kotlin/pw/binom/agentik/storage/inmemory/InMemoryConversationStore.kt +++ b/journal-inmemory/src/commonMain/kotlin/pw/binom/agentik/journal/inmemory/InMemoryMutableConversationStore.kt @@ -1,22 +1,33 @@ -package pw.binom.agentik.storage.inmemory +package pw.binom.agentik.journal.inmemory import kotlinx.coroutines.sync.Mutex import kotlinx.coroutines.sync.withLock -import kotlin.time.Clock import pw.binom.agentik.journal.ConversationRecord -import pw.binom.agentik.journal.ConversationStore +import pw.binom.agentik.journal.MutableConversationStore +import kotlin.time.Clock import kotlin.time.Instant /** - * Thread-safe Map-импл [ConversationStore]. + * Thread-safe Map-импл [MutableConversationStore] для клиентских + * in-process кэшей (и тестов/dev-режима). * * Использует `Mutex` для атомарности read-modify-write операций * (rename, touch) — иначе два параллельных `rename` могут потерять обновления * (lost-update race), что в SQLite невозможно из-за driver-level locking. + * + * **Сортировка**: `list()` сортирует по `updatedAt DESC`. + * + * **Типичный кэш-паттерн в клиенте** (см. `client/README.md`): + * ``` + * val local = InMemoryMutableConversationStore() + * // seed: remote.listFlow → local.upsert + * // live-refresh: outbox.agentEvents → local.upsert/delete/rename/touch + * // UI: local.list(0, PAGE_SIZE) + * ``` */ -class InMemoryConversationStore( +class InMemoryMutableConversationStore( private val clock: Clock = Clock.System, -) : ConversationStore { +) : MutableConversationStore { private val byId: MutableMap = mutableMapOf() private val mutex = Mutex() diff --git a/storage-inmemory/src/commonTest/kotlin/pw/binom/agentik/storage/inmemory/InMemoryConversationStoreTest.kt b/journal-inmemory/src/commonTest/kotlin/pw/binom/agentik/journal/inmemory/InMemoryMutableConversationStoreTest.kt similarity index 86% rename from storage-inmemory/src/commonTest/kotlin/pw/binom/agentik/storage/inmemory/InMemoryConversationStoreTest.kt rename to journal-inmemory/src/commonTest/kotlin/pw/binom/agentik/journal/inmemory/InMemoryMutableConversationStoreTest.kt index c7a2113..be24b8b 100644 --- a/storage-inmemory/src/commonTest/kotlin/pw/binom/agentik/storage/inmemory/InMemoryConversationStoreTest.kt +++ b/journal-inmemory/src/commonTest/kotlin/pw/binom/agentik/journal/inmemory/InMemoryMutableConversationStoreTest.kt @@ -1,4 +1,4 @@ -package pw.binom.agentik.storage.inmemory +package pw.binom.agentik.journal.inmemory import pw.binom.agentik.journal.ConversationRecord import kotlin.test.Test @@ -9,11 +9,11 @@ import kotlin.test.assertTrue import kotlin.time.Instant import kotlinx.coroutines.test.runTest -class InMemoryConversationStoreTest { +class InMemoryMutableConversationStoreTest { @Test fun `upsert and get roundtrip preserves all fields`() = runTest { - val store = InMemoryConversationStore() + val store = InMemoryMutableConversationStore() val rec = ConversationRecord( id = "c1", title = "test", @@ -28,13 +28,13 @@ class InMemoryConversationStoreTest { @Test fun `get returns null for missing id`() = runTest { - val store = InMemoryConversationStore() + val store = InMemoryMutableConversationStore() assertNull(store.get("nope")) } @Test fun `delete removes the record and returns true`() = runTest { - val store = InMemoryConversationStore() + val store = InMemoryMutableConversationStore() store.upsert( ConversationRecord( "c1", null, false, @@ -50,7 +50,7 @@ class InMemoryConversationStoreTest { @Test fun `list sorts by updatedAt DESC and respects offset+limit`() = runTest { - val store = InMemoryConversationStore() + val store = InMemoryMutableConversationStore() val t0 = Instant.parse("2026-09-15T10:00:00Z") store.upsert(ConversationRecord("c1", null, false, t0, t0)) store.upsert(ConversationRecord("c2", null, false, t0, t0.plus(kotlin.time.Duration.parse("PT60S")))) @@ -69,7 +69,7 @@ class InMemoryConversationStoreTest { @Test fun `rename updates title and updatedAt returns new updatedAt`() = runTest { - val store = InMemoryConversationStore() + val store = InMemoryMutableConversationStore() val t0 = Instant.parse("2026-09-15T10:00:00Z") store.upsert(ConversationRecord("c1", null, false, t0, t0)) @@ -84,7 +84,7 @@ class InMemoryConversationStoreTest { @Test fun `rename with null title clears it`() = runTest { - val store = InMemoryConversationStore() + val store = InMemoryMutableConversationStore() val t0 = Instant.parse("2026-09-15T10:00:00Z") store.upsert(ConversationRecord("c1", "old", false, t0, t0)) store.rename("c1", null) @@ -93,13 +93,13 @@ class InMemoryConversationStoreTest { @Test fun `rename returns null for missing conversation`() = runTest { - val store = InMemoryConversationStore() + val store = InMemoryMutableConversationStore() assertNull(store.rename("nope", "x")) } @Test fun `touch bumps updatedAt without changing other fields`() = runTest { - val store = InMemoryConversationStore() + val store = InMemoryMutableConversationStore() val t0 = Instant.parse("2026-09-15T10:00:00Z") val t1 = Instant.parse("2026-09-15T10:01:00Z") store.upsert(ConversationRecord("c1", "title", false, t0, t0)) @@ -112,7 +112,7 @@ class InMemoryConversationStoreTest { @Test fun `close is idempotent and does nothing`() { - val store = InMemoryConversationStore() + val store = InMemoryMutableConversationStore() store.close() store.close() // должно быть no-op } diff --git a/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/AgentEvent.kt b/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/AgentEvent.kt index f0475ac..d878b98 100644 --- a/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/AgentEvent.kt +++ b/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/AgentEvent.kt @@ -45,4 +45,13 @@ sealed interface AgentEvent { @Serializable @SerialName("renamed") data class Renamed(override val date: Instant, val id: String, val title: String?) : AgentEvent + + /** + * Обновлён `updatedAt` диалога (после `send()` или другого события, + * бампнувшего активность). Клиентский кэш [ConversationStore] может + * применить этот event для пересортировки списка. + */ + @Serializable + @SerialName("touched") + data class Touched(override val date: Instant, val id: String, val updatedAt: Instant) : AgentEvent } diff --git a/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Agent.kt b/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Agent.kt index b7501e4..2b6d29e 100644 --- a/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Agent.kt +++ b/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Agent.kt @@ -1,7 +1,6 @@ package pw.binom.agentik.proto -import kotlinx.coroutines.flow.Flow -import kotlinx.coroutines.flow.flow +import pw.binom.agentik.journal.ConversationStore import pw.binom.agentik.journal.JournalStore import pw.binom.agentik.outbox.OutboxStore import kotlin.time.Instant @@ -13,11 +12,16 @@ import kotlin.time.Instant * [createConversation] возвращает [Conversation], который сам хранит историю * и которому отправляют ходы через [Conversation.send]. * - * **Хранилища вынесены в [Agent.journal] и [Agent.outbox]**: оба read-only. - * События больше НЕ часть [Agent] (раньше были `events()`/`allEvents()`) — - * они теперь живут в [outbox] как `OutboxStore.events(after)` / - * `outbox.agentEvents(after)`. Это даёт единый путь для всех read-операций - * по хранилищу и убирает дублирование между протоколом и хранилищем. + * **Хранилища вынесены в [Agent.journal], [Agent.outbox] и + * [Agent.conversationStore]**: все три read-only views. События живут + * в [outbox] как `OutboxStore.events(after)` / `outbox.agentEvents(after)`. + * Это даёт единый путь для всех read-операций по хранилищу и убирает + * дублирование между протоколом и хранилищем. + * + * **Команды** (create / delete / rename) живут прямо на [Agent]. Они + * шлются клиентом и выполняются сервером — клиент **не** пишет в стор + * напрямую. Клиентский кэш [conversationStore] обновляется через + * `outbox.agentEvents()` (Created / Deleted / Renamed / Touched). */ interface Agent : AutoCloseable { @@ -59,6 +63,22 @@ interface Agent : AutoCloseable { */ val outbox: OutboxStore + /** + * Read-only view на `conversation` table (id + title + timestamps). + * + * Используется HTTP-фасадом `:server` для endpoint'а + * `GET /{path}/conversations?offset=&limit=` — внешние клиенты + * получают лёгкую метадату (без handle'ов и image-support флагов) + * для рендера списка диалогов. Для активной работы (send / interrupt) + * клиент отдельно получает handle через [getConversation]. + * + * **Read-only**: write-доступ только через `MutableConversationStore` + * внутри ChatAgent, не через [Agent] interface. Клиент модифицирует + * диалоги командами: [createConversation] / [deleteConversation] / + * [renameConversation]. + */ + val conversationStore: ConversationStore + /** Создаёт новый stateful-диалог с агентом. */ fun createConversation(temp: Boolean): Conversation @@ -68,19 +88,11 @@ interface Agent : AutoCloseable { /** Удаляет диалог. Возвращает `true`, если диалог существовал и удалён. */ suspend fun deleteConversation(id: String): Boolean - /** Страница диалогов: не более [limit] штук, начиная с [offset]-го. */ - suspend fun getConversations(offset: Int, limit: Int): List - - /** Все диалоги, начиная с [offset], как поток: подгружает по [PAGE_SIZE] за раз. */ - fun getConversations(offset: Int = 0): Flow = flow { - var skip = offset - while (true) { - val page = getConversations(skip, PAGE_SIZE) - if (page.isEmpty()) break - page.forEach { emit(it) } - skip += page.size - } - } + /** + * Переименовывает диалог; `null` для сброса заголовка. Возвращает новый + * `updatedAt` или `null`, если диалог не найден. + */ + suspend fun renameConversation(id: String, title: String?): Instant? companion object { diff --git a/server/src/commonMain/kotlin/pw/binom/agentik/server/Routes.kt b/server/src/commonMain/kotlin/pw/binom/agentik/server/Routes.kt index 9ca5692..d679680 100644 --- a/server/src/commonMain/kotlin/pw/binom/agentik/server/Routes.kt +++ b/server/src/commonMain/kotlin/pw/binom/agentik/server/Routes.kt @@ -38,7 +38,7 @@ internal fun Route.agentikRoutes(agent: Agent) { get("/conversations") { val offset = call.request.queryParameters["offset"]?.toIntOrNull() ?: 0 val limit = call.request.queryParameters["limit"]?.toIntOrNull() ?: Agent.PAGE_SIZE - call.respond(agent.getConversations(offset, limit).map { it.snapshot() }) + call.respond(agent.conversationStore.list(offset, limit)) } post("/conversations") { @@ -63,14 +63,15 @@ internal fun Route.agentikRoutes(agent: Agent) { patch("/conversations/{id}") { val id = call.parameters["id"]!! - val c = agent.getConversation(id) - if (c == null) { + val req = call.receive() + val newUpdatedAt = agent.renameConversation(id, req.title) + if (newUpdatedAt == null) { call.respond(HttpStatusCode.NotFound) return@patch } - val req = call.receive() - c.rename(req.title) - call.respond(c.snapshot()) + // Возвращаем обновлённый record (лёгкая метадата, не snapshot handle'а). + val rec = agent.conversationStore.get(id)!! + call.respond(rec) } post("/conversations/{id}/messages") { diff --git a/server/src/commonTest/kotlin/pw/binom/agentik/server/BearerTokenTest.kt b/server/src/commonTest/kotlin/pw/binom/agentik/server/BearerTokenTest.kt index 944e483..b058252 100644 --- a/server/src/commonTest/kotlin/pw/binom/agentik/server/BearerTokenTest.kt +++ b/server/src/commonTest/kotlin/pw/binom/agentik/server/BearerTokenTest.kt @@ -13,6 +13,7 @@ import io.ktor.server.engine.embeddedServer import io.ktor.server.routing.routing import kotlinx.coroutines.flow.emptyFlow import kotlinx.coroutines.runBlocking +import pw.binom.agentik.journal.ConversationStore import pw.binom.agentik.journal.JournalStore import pw.binom.agentik.journal.MessageRecord import pw.binom.agentik.outbox.CommonEvent @@ -53,10 +54,15 @@ class BearerTokenTest { override suspend fun earliestEventDate(): Instant = Instant.DISTANT_PAST override fun close() {} } + override val conversationStore: ConversationStore = object : ConversationStore { + override suspend fun get(id: String) = null + override suspend fun list(offset: Int, limit: Int) = emptyList() + override fun close() {} + } override fun createConversation(temp: Boolean): Conversation = TODO("not needed by tests") override suspend fun getConversation(id: String): Conversation? = null override suspend fun deleteConversation(id: String): Boolean = false - override suspend fun getConversations(offset: Int, limit: Int): List = emptyList() + override suspend fun renameConversation(id: String, title: String?): Instant? = null override fun close() {} } diff --git a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/Main.kt b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/Main.kt index 63b67b6..2852c1e 100644 --- a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/Main.kt +++ b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/Main.kt @@ -363,7 +363,7 @@ private fun runServer() { val agent = ChatAgent( id = "agentik", - conversationStore = sqliteStores.conversations, + mutableConversationStore = sqliteStores.conversations, messageStore = sqliteStores.messages, workingMemoryStore = sqliteStores.workingMemory, reflectionStore = sqliteStores.reflections, diff --git a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ChatAgent.kt b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ChatAgent.kt index ff4a486..28e4810 100644 --- a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ChatAgent.kt +++ b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ChatAgent.kt @@ -1,11 +1,6 @@ package pw.binom.agentik.standalone.agent import kotlin.time.Instant -import kotlinx.coroutines.flow.Flow -import kotlinx.coroutines.flow.emitAll -import kotlinx.coroutines.flow.flow -import kotlinx.coroutines.flow.map -import kotlinx.coroutines.launch import kotlinx.coroutines.runBlocking import kotlinx.coroutines.sync.Mutex import kotlinx.coroutines.sync.withLock @@ -18,7 +13,6 @@ import pw.binom.agentik.memory.MemorySystemGuidance import pw.binom.agentik.proto.Agent as ProtoAgent import pw.binom.agentik.outbox.AgentEvent import pw.binom.agentik.outbox.CommonEvent -import pw.binom.agentik.outbox.Event as ProtoEvent import pw.binom.agentik.outbox.MutableOutboxStore import pw.binom.agentik.journal.JournalStore import pw.binom.agentik.outbox.OutboxStore @@ -29,7 +23,7 @@ import pw.binom.agentik.standalone.agent.memory.MemoryToolsFactory import pw.binom.agentik.standalone.llm.LlmConfig import pw.binom.agentik.journal.ConversationRecord import pw.binom.agentik.journal.ConversationStore -import pw.binom.agentik.journal.Ids +import pw.binom.agentik.journal.MutableConversationStore import pw.binom.agentik.reflection.Reflection import pw.binom.agentik.reflection.ReflectionStore import pw.binom.agentik.journal.MutableJournalStore @@ -65,7 +59,7 @@ import pw.binom.litert.LiteLlm */ class ChatAgent( override val id: String, - private val conversationStore: ConversationStore, + private val mutableConversationStore: MutableConversationStore, private val messageStore: MutableJournalStore, private val workingMemoryStore: ContextStore, private val reflectionStore: ReflectionStore, @@ -239,6 +233,16 @@ class ChatAgent( override val outbox: OutboxStore get() = eventStore + /** + * Read-only view of [mutableConversationStore] для HTTP-фасада в `:server` + * (`GET /conversations` → список [ConversationRecord] для UI). + * + * Сужение с [MutableConversationStore] на [ConversationStore] тривиальна + * через interface-наследование (read-only projection). + */ + override val conversationStore: ConversationStore + get() = mutableConversationStore + /** Защищает карту живых диалогов. */ private val liveLock = Mutex() private val live: MutableMap = HashMap() @@ -262,12 +266,12 @@ class ChatAgent( // и не переживают рестарт агента (см. Memory #3709). if (!temp) { runBlocking { - conversationStore.upsert(rec) + mutableConversationStore.upsert(rec) } } val conv = ChatConversation( record = rec, - conversationStore = conversationStore, + conversationStore = mutableConversationStore, messageStore = messageStore, workingMemoryStore = workingMemoryStore, reflectionStore = reflectionStore, @@ -302,7 +306,7 @@ class ChatAgent( override suspend fun getConversation(id: String): ProtoConversation? { liveLock.withLock { live[id] }?.let { if (!it.isClosed) return it } - val rec = conversationStore.get(id) ?: return null + val rec = mutableConversationStore.get(id) ?: return null return newConversation(rec).also { liveLock.withLock { live[id] = it } } @@ -311,7 +315,7 @@ class ChatAgent( override suspend fun deleteConversation(id: String): Boolean { val conv = liveLock.withLock { live.remove(id) } conv?.close() - val ok = conversationStore.delete(id) + val ok = mutableConversationStore.delete(id) if (ok) { val event = AgentEvent.Deleted(date = now(), id = id) eventStore.append(CommonEvent.Agent(date = now(), event = event)) @@ -319,17 +323,16 @@ class ChatAgent( return ok } - override suspend fun getConversations(offset: Int, limit: Int): List = - conversationStore.list(offset = offset, limit = limit).map { rec -> - liveLock.withLock { live[rec.id] } - ?: newConversation(rec).also { - liveLock.withLock { live[rec.id] = it } - } - } + override suspend fun renameConversation(id: String, title: String?): Instant? { + val newUpdatedAt = mutableConversationStore.rename(id, title) ?: return null + val event = AgentEvent.Renamed(date = newUpdatedAt, id = id, title = title) + eventStore.append(CommonEvent.Agent(date = newUpdatedAt, event = event)) + return newUpdatedAt + } private fun newConversation(rec: ConversationRecord): ChatConversation = ChatConversation( record = rec, - conversationStore = conversationStore, + conversationStore = mutableConversationStore, messageStore = messageStore, workingMemoryStore = workingMemoryStore, reflectionStore = reflectionStore, diff --git a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ConversationLoop.kt b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ConversationLoop.kt index 621b583..8af2106 100644 --- a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ConversationLoop.kt +++ b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ConversationLoop.kt @@ -7,7 +7,6 @@ import kotlinx.coroutines.Job import kotlinx.coroutines.SupervisorJob import kotlinx.coroutines.cancel import kotlinx.coroutines.cancelAndJoin -import kotlinx.coroutines.flow.Flow import kotlinx.coroutines.launch import kotlinx.coroutines.runBlocking import kotlinx.coroutines.sync.Mutex @@ -31,7 +30,7 @@ import pw.binom.agentik.proto.MessageContext as ProtoMessageContext import pw.binom.agentik.skills.SkillStore import pw.binom.agentik.journal.Content import pw.binom.agentik.journal.ConversationRecord -import pw.binom.agentik.journal.ConversationStore +import pw.binom.agentik.journal.MutableConversationStore import pw.binom.agentik.journal.MessageContext import pw.binom.agentik.journal.MessageOrigin import pw.binom.agentik.journal.MessageRecord @@ -54,7 +53,7 @@ import pw.binom.agentik.toolsets.NamedTool class ConversationLoop( record: ConversationRecord, - private val conversationStore: ConversationStore, + private val conversationStore: MutableConversationStore, private val messageStore: MutableJournalStore, private val workingMemoryStore: ContextStore, private val reflectionStore: ReflectionStore?, @@ -440,6 +439,18 @@ class ConversationLoop( state.record = state.record.copy(updatedAt = assistantAt) conversationStore.touch(id, assistantAt) + if (!state.isTemporal) { + eventStore.append( + pw.binom.agentik.outbox.CommonEvent.Agent( + date = assistantAt, + event = pw.binom.agentik.outbox.AgentEvent.Touched( + date = assistantAt, + id = id, + updatedAt = assistantAt, + ), + ) + ) + } // BackgroundScheduler is event-driven — подписан на BackgroundEventBus // (compaction/lifecycle/tool-failure events). Никаких interval-based diff --git a/standalone/src/commonTest/kotlin/pw/binom/agentik/standalone/agent/ChatAgentTest.kt b/standalone/src/commonTest/kotlin/pw/binom/agentik/standalone/agent/ChatAgentTest.kt index 61a5537..9cf0747 100644 --- a/standalone/src/commonTest/kotlin/pw/binom/agentik/standalone/agent/ChatAgentTest.kt +++ b/standalone/src/commonTest/kotlin/pw/binom/agentik/standalone/agent/ChatAgentTest.kt @@ -62,7 +62,7 @@ class ChatAgentTest { skills: SkillCatalog = SkillCatalog.EMPTY, ): ChatAgent = ChatAgent( id = "agentik", - conversationStore = sqliteStores.conversations, + mutableConversationStore = sqliteStores.conversations, messageStore = sqliteStores.messages, workingMemoryStore = sqliteStores.workingMemory, reflectionStore = sqliteStores.reflections, @@ -146,15 +146,32 @@ class ChatAgentTest { } @Test - fun `getConversations returns all stored persistent conversations`() = runTest { + fun `conversationStore list returns all stored persistent conversations`() = runTest { val agent = newAgent() agent.createConversation(temp = false) agent.createConversation(temp = true) - val list = agent.getConversations(0, 10) + val list = agent.conversationStore.list(0, 10) // temp-беседы не персистятся, в списке только persistent assertEquals(1, list.size) } + @Test + fun `renameConversation updates title and emits event`() = runTest { + val agent = newAgent() + val conv = agent.createConversation(temp = false) + val newTitle = "Renamed!" + val updatedAt = agent.renameConversation(conv.id, newTitle) + assertNotNull(updatedAt) + val rec = agent.conversationStore.get(conv.id) + assertEquals(newTitle, rec?.title) + } + + @Test + fun `renameConversation returns null for unknown id`() = runTest { + val agent = newAgent() + assertNull(agent.renameConversation("nope", "x")) + } + @Test fun `deleteConversation removes conversation and data`() = runTest { val agent = newAgent() @@ -458,7 +475,7 @@ class ChatAgentTest { sqliteStores = KsqliteStores.open(dbPath) val agent1 = ChatAgent( id = "agentik", - conversationStore = sqliteStores.conversations, + mutableConversationStore = sqliteStores.conversations, messageStore = sqliteStores.messages, workingMemoryStore = sqliteStores.workingMemory, reflectionStore = sqliteStores.reflections, @@ -478,7 +495,7 @@ class ChatAgentTest { sqliteStores = KsqliteStores.open(dbPath) val agent2 = ChatAgent( id = "agentik", - conversationStore = sqliteStores.conversations, + mutableConversationStore = sqliteStores.conversations, messageStore = sqliteStores.messages, workingMemoryStore = sqliteStores.workingMemory, reflectionStore = sqliteStores.reflections, @@ -500,7 +517,7 @@ class ChatAgentTest { sqliteStores = KsqliteStores.open(dbPath) val agent1 = ChatAgent( id = "agentik", - conversationStore = sqliteStores.conversations, + mutableConversationStore = sqliteStores.conversations, messageStore = sqliteStores.messages, workingMemoryStore = sqliteStores.workingMemory, reflectionStore = sqliteStores.reflections, @@ -518,7 +535,7 @@ class ChatAgentTest { sqliteStores = KsqliteStores.open(dbPath) val agent2 = ChatAgent( id = "agentik", - conversationStore = sqliteStores.conversations, + mutableConversationStore = sqliteStores.conversations, messageStore = sqliteStores.messages, workingMemoryStore = sqliteStores.workingMemory, reflectionStore = sqliteStores.reflections, diff --git a/standalone/src/commonTest/kotlin/pw/binom/agentik/standalone/agent/ChatAgentToolsetsTest.kt b/standalone/src/commonTest/kotlin/pw/binom/agentik/standalone/agent/ChatAgentToolsetsTest.kt index 6121b6c..83a7a33 100644 --- a/standalone/src/commonTest/kotlin/pw/binom/agentik/standalone/agent/ChatAgentToolsetsTest.kt +++ b/standalone/src/commonTest/kotlin/pw/binom/agentik/standalone/agent/ChatAgentToolsetsTest.kt @@ -43,7 +43,7 @@ class ChatAgentToolsetsTest { val sqliteStores = KsqliteStores.inMemory("toolsets-${kotlin.random.Random.nextLong()}") val agent = ChatAgent( id = "test-agent", - conversationStore = sqliteStores.conversations, + mutableConversationStore = sqliteStores.conversations, messageStore = sqliteStores.messages, workingMemoryStore = sqliteStores.workingMemory, reflectionStore = sqliteStores.reflections, diff --git a/standalone/src/commonTest/kotlin/pw/binom/agentik/standalone/agent/CompactionTest.kt b/standalone/src/commonTest/kotlin/pw/binom/agentik/standalone/agent/CompactionTest.kt index 70b63a3..011b89e 100644 --- a/standalone/src/commonTest/kotlin/pw/binom/agentik/standalone/agent/CompactionTest.kt +++ b/standalone/src/commonTest/kotlin/pw/binom/agentik/standalone/agent/CompactionTest.kt @@ -63,7 +63,7 @@ class CompactionTest { val reviewer = if (memoryStore != null) KeywordMdReviewer() else null return ChatAgent( id = "test", - conversationStore = sqliteStores.conversations, + mutableConversationStore = sqliteStores.conversations, messageStore = sqliteStores.messages, workingMemoryStore = sqliteStores.workingMemory, reflectionStore = sqliteStores.reflections, diff --git a/standalone/src/commonTest/kotlin/pw/binom/agentik/standalone/agent/MemoryWiringTest.kt b/standalone/src/commonTest/kotlin/pw/binom/agentik/standalone/agent/MemoryWiringTest.kt index 680bdbb..021c883 100644 --- a/standalone/src/commonTest/kotlin/pw/binom/agentik/standalone/agent/MemoryWiringTest.kt +++ b/standalone/src/commonTest/kotlin/pw/binom/agentik/standalone/agent/MemoryWiringTest.kt @@ -76,7 +76,7 @@ class MemoryWiringTest { contextCompactor: ContextCompactor = EchoCompactor, ): ChatAgent = ChatAgent( id = "agentik", - conversationStore = sqliteStores.conversations, + mutableConversationStore = sqliteStores.conversations, messageStore = sqliteStores.messages, workingMemoryStore = sqliteStores.workingMemory, reflectionStore = sqliteStores.reflections, @@ -115,7 +115,7 @@ class MemoryWiringTest { val soulBody = "I am a helpful test persona. I always answer in one short line." val agent = ChatAgent( id = "agentik", - conversationStore = sqliteStores.conversations, + mutableConversationStore = sqliteStores.conversations, messageStore = sqliteStores.messages, workingMemoryStore = sqliteStores.workingMemory, reflectionStore = sqliteStores.reflections, @@ -143,7 +143,7 @@ class MemoryWiringTest { fun `soul body not added when null`() = runBlocking { val agent = ChatAgent( id = "agentik", - conversationStore = sqliteStores.conversations, + mutableConversationStore = sqliteStores.conversations, messageStore = sqliteStores.messages, workingMemoryStore = sqliteStores.workingMemory, reflectionStore = sqliteStores.reflections, diff --git a/storage-inmemory/build.gradle.kts b/storage-inmemory/build.gradle.kts index 6ebf3a0..9f0ff4c 100644 --- a/storage-inmemory/build.gradle.kts +++ b/storage-inmemory/build.gradle.kts @@ -21,6 +21,7 @@ kotlin { sourceSets { commonMain.dependencies { api(project(":journal-api")) + api(project(":journal-inmemory")) api(project(":reflection-api")) api(project(":context-api")) } diff --git a/storage-inmemory/src/commonMain/kotlin/pw/binom/agentik/storage/inmemory/InMemoryStorage.kt b/storage-inmemory/src/commonMain/kotlin/pw/binom/agentik/storage/inmemory/InMemoryStorage.kt index 3612eee..439c79d 100644 --- a/storage-inmemory/src/commonMain/kotlin/pw/binom/agentik/storage/inmemory/InMemoryStorage.kt +++ b/storage-inmemory/src/commonMain/kotlin/pw/binom/agentik/storage/inmemory/InMemoryStorage.kt @@ -2,7 +2,8 @@ package pw.binom.agentik.storage.inmemory import pw.binom.agentik.context.ContextStore import pw.binom.agentik.journal.MutableJournalStore -import pw.binom.agentik.journal.ConversationStore +import pw.binom.agentik.journal.MutableConversationStore +import pw.binom.agentik.journal.inmemory.InMemoryMutableConversationStore import pw.binom.agentik.reflection.ReflectionStore import kotlin.time.Clock @@ -21,14 +22,14 @@ import kotlin.time.Clock */ object InMemoryStorage { data class Bundle( - val conversationStore: ConversationStore, + val conversationStore: MutableConversationStore, val messageStore: MutableJournalStore, val workingMemoryStore: ContextStore, val reflectionStore: ReflectionStore, ) fun create(clock: Clock = Clock.System): Bundle = Bundle( - conversationStore = InMemoryConversationStore(clock), + conversationStore = InMemoryMutableConversationStore(clock), messageStore = InMemoryMessageStore(), workingMemoryStore = InMemoryWorkingMemoryStore(), reflectionStore = InMemoryReflectionStore(), diff --git a/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteConversationStore.kt b/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteMutableConversationStore.kt similarity index 97% rename from storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteConversationStore.kt rename to storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteMutableConversationStore.kt index 9c0f1dd..c898353 100644 --- a/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteConversationStore.kt +++ b/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteMutableConversationStore.kt @@ -3,7 +3,7 @@ package pw.binom.agentik.storage.ksqlite import kotlin.time.Clock import kotlin.time.Instant import pw.binom.agentik.journal.ConversationRecord -import pw.binom.agentik.journal.ConversationStore +import pw.binom.agentik.journal.MutableConversationStore import pw.binom.db.ksqlite.SQLiteConnection import kotlinx.coroutines.Dispatchers import kotlinx.coroutines.sync.Mutex @@ -11,15 +11,15 @@ import kotlinx.coroutines.sync.withLock import kotlinx.coroutines.withContext /** - * ksqlite-реализация [ConversationStore]. Схема таблицы `conversation` живёт + * ksqlite-реализация [MutableConversationStore]. Схема таблицы `conversation` живёт * в [Schema] (миграция через PRAGMA user_version) — этот класс только * готовит и выполняет SQL, ссылаясь на `Schema.COL_*` / `Schema.TABLE_*`. */ -class KsqliteConversationStore( +class KsqliteMutableConversationStore( private val connection: SQLiteConnection, private val messageStore: KsqliteMessageStore? = null, private val workingMemoryStore: KsqliteWorkingMemoryStore? = null, -) : ConversationStore { +) : MutableConversationStore { private val mutex = Mutex() diff --git a/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteStores.kt b/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteStores.kt index a98d0eb..78692d1 100644 --- a/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteStores.kt +++ b/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteStores.kt @@ -2,7 +2,7 @@ package pw.binom.agentik.storage.ksqlite import pw.binom.agentik.context.ContextStore import pw.binom.agentik.journal.MutableJournalStore -import pw.binom.agentik.journal.ConversationStore +import pw.binom.agentik.journal.MutableConversationStore import pw.binom.agentik.reflection.ReflectionStore import pw.binom.db.ksqlite.SQLiteConnection @@ -17,7 +17,7 @@ import pw.binom.db.ksqlite.SQLiteConnection */ class KsqliteStores internal constructor( val connection: SQLiteConnection, - val conversations: ConversationStore, + val conversations: MutableConversationStore, val messages: MutableJournalStore, val workingMemory: ContextStore, val reflections: ReflectionStore, @@ -51,7 +51,7 @@ class KsqliteStores internal constructor( val working = KsqliteWorkingMemoryStore(conn) return KsqliteStores( connection = conn, - conversations = KsqliteConversationStore(conn, messages, working), + conversations = KsqliteMutableConversationStore(conn, messages, working), messages = messages, workingMemory = working, reflections = KsqliteReflectionStore(conn), diff --git a/storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteConversationStoreTest.kt b/storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteMutableConversationStoreTest.kt similarity index 98% rename from storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteConversationStoreTest.kt rename to storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteMutableConversationStoreTest.kt index aa01424..ae1ba43 100644 --- a/storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteConversationStoreTest.kt +++ b/storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteMutableConversationStoreTest.kt @@ -12,7 +12,7 @@ import kotlin.test.assertTrue import kotlin.time.Duration import kotlin.time.Instant -class KsqliteConversationStoreTest { +class KsqliteMutableConversationStoreTest { private lateinit var stores: KsqliteStores @BeforeTest diff --git a/storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/SchemaMigrationTest.kt b/storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/SchemaMigrationTest.kt index e00d18a..90b0baf 100644 --- a/storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/SchemaMigrationTest.kt +++ b/storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/SchemaMigrationTest.kt @@ -79,7 +79,7 @@ class SchemaMigrationTest { // constructor internal, тест в том же модуле и может его звать. val stores = KsqliteStores( connection = conn, - conversations = KsqliteConversationStore(conn), + conversations = KsqliteMutableConversationStore(conn), messages = KsqliteMessageStore(conn), workingMemory = KsqliteWorkingMemoryStore(conn), reflections = KsqliteReflectionStore(conn),