2 Commits
10 .. 11

Author SHA1 Message Date
subochev 639c7d1748 docs(client): add «Кэш списка бесед» section + ios targets to :journal-inmemory
ci / JVM build + tests (push) Successful in 5m49s
release / Publish KMP libraries → caffeine Nexus (release) Successful in 36s
- :journal-inmemory — added iosX64/iosArm64/iosSimulatorArm64 to the
  target set so :storage-inmemory (which now depends on it) can build
  for iOS. Pure `MutableMap`+`Mutex` impl, no I/O, fully portable.

- client/README.md — new «Кэш списка бесед» section:
  - shows `agent.conversationStore` as the read-only entry point;
  - demonstrates the «remote.listFlow → local.upsert + outbox.agentEvents
    → local apply» pattern (Created/Renamed/Touched/Deleted);
  - notes that the cache is built into `AgentikAgent` by default;
  - mentions `wrapWithLocalConversationCache` for custom stores (SQLite/JSON);
  - points to `agentikHttpClient(...).raw` as the escape-hatch for clients
    that want direct HTTP.

Also adds `HttpConversationStore.kt` to the index (the file existed on
disk but wasn't `git add`ed in the previous commit).
2026-09-22 05:25:39 +03:00
subochev 5f0e0da361 feat(client): client-side caching for conversation list via agent.outbox events
Adds `agent.conversationStore` (read-only view on `conversation` table) to
the :proto Agent interface, plus `agent.renameConversation(id, title?)`
command. Client-side cache in :client is built from a snapshot
(`remote.listFlow(0)` → `local.upsert(...)`) + live updates via
`outbox.agentEvents()` (Created/Deleted/Renamed/Touched).

Changes:
- :journal-api — split `ConversationStore` (read-only: get/list) and
  `MutableConversationStore` (CRUD: upsert/delete/rename/touch);
  `ConversationStore` gained `listFlow` (cold-flow paging via `list`).
- :outbox-api — `AgentEvent.Touched(date, id, updatedAt)` event so
  client cache stays fresh after `send()` (which bumps `updatedAt`).
- :proto.Agent — added `conversationStore: ConversationStore` property,
  added `renameConversation(id, title?): Instant?` command, removed
  `getConversations(offset, limit)` (now: `conversationStore.list(...)`).
- :server — `GET /conversations` now returns `List<ConversationRecord>`
  (lightweight metadata, no handle/image-support flags); `PATCH
  /conversations/{id}` uses `agent.renameConversation` and returns
  the updated `ConversationRecord`.
- :journal-inmemory — expanded targets to jvm+macos+linux+mingw (matches
  :client); moved `InMemoryMutableConversationStore` here from
  :storage-inmemory so :client can use it without pulling ios targets.
- :storage-inmemory — depends on :journal-inmemory.
- :storage-ksqlite — pre-staged rename `KsqliteConversationStore` →
  `KsqliteMutableConversationStore` to match the new interface split.
- :standalone — `ChatAgent` exposes `conversationStore` as a read-only
  view of its `mutableConversationStore`; emits `AgentEvent.Touched`
  after each `send()` (after `conversationStore.touch(id, ts)`).
- :client — new `HttpConversationStore` (read-only HTTP impl);
  `AgentikAgent` wraps the agent with `wrapWithLocalConversationCache`
  so the client sees an in-memory cache (snapshot + outbox events)
  instead of direct HTTP. Cache scope + HttpClient + background job
  all cancelled in `agent.close()`.
- :client/README — new «Кэш списка бесед» section with the
  `listFlow → upsert` / `agentEvents → apply` pattern and a note that
  `conversationStore` is read-only (writes only via Agent commands).

All 96 jvmTest tasks green.
2026-09-22 05:21:01 +03:00
30 changed files with 534 additions and 135 deletions
@@ -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<pw.binom.agentik.outbox.CommonEvent>()
override fun agentEvents(after: Instant?) = emptyFlow<pw.binom.agentik.outbox.CommonEvent.Agent>()
override fun conversationEvents(after: Instant?, conversationId: String?) = emptyFlow<pw.binom.agentik.outbox.CommonEvent.Conversation>()
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<Conversation> =
conversations.toList()
override suspend fun renameConversation(id: String, title: String?): Instant? = null
}
/**
+120 -3
View File
@@ -285,6 +285,55 @@ session.scope.launch {
`rec is MessageRecord.UserMessage` для реплик пользователя,
`rec is MessageRecord.ToolCall` для отрисовки tool-call баббла, и т.п.
## Кэш списка бесед
`agent.conversationStore` — read-only view поверх `conversation`-таблицы
на сервере (`ConversationRecord` = id / title / isTemporal / createdAt /
updatedAt, без `Conversation` handle и без флагов image-support).
**Сценарий клиента:** показать список диалогов («как в Telegram»), чтобы
при открытии UI уже знал названия, не дёргал сервер лишний раз, и
моментально реагировал на создание/удаление/переименование в другой
вкладке.
Подход — тот же **«remote → local snapshot + live-events»**:
```kotlin
import pw.binom.agentik.client.AgentikAgent
import pw.binom.agentik.journal.ConversationRecord
import pw.binom.agentik.outbox.AgentEvent
import io.ktor.client.engine.cio.CIO
// `AgentikAgent` сам оборачивает HTTP-store в локальный кэш:
// remote.listFlow → local.upsert (snapshot)
// outbox.agentEvents → local.upsert / delete (live)
val agent = AgentikAgent(
id = "agentik",
baseUrl = "http://localhost:8080/agentik",
engineFactory = CIO,
)
// Кэш уже наполняется в фоне, читать можно сразу:
val all = agent.conversationStore.list(0, Int.MAX_VALUE)
all.forEach { rec -> println("${rec.id} ${rec.title ?: "(no title)"} ${rec.updatedAt}") }
// И наблюдать live-изменения (Created/Deleted/Renamed/Touched)
agent.outbox.agentEvents(kotlin.time.Instant.DISTANT_PAST).collect { ev ->
when (ev) {
is AgentEvent.Created -> println("+ ${ev.conversationId}")
is AgentEvent.Renamed -> println("~ ${ev.id} → ${ev.title}")
is AgentEvent.Touched -> println("↻ ${ev.id} (${ev.updatedAt})")
is AgentEvent.Deleted -> println("- ${ev.id}")
}
}
```
Если ты **не хочешь** встроенный кэш (например, тебе нужен прямой HTTP
для бэкенда-сервиса) — `agentikHttpClient(...).raw` оставлен как
escape-hatch. Сам `InMemoryMutableConversationStore` тоже доступен —
подмени его на свою реализацию через `wrapWithLocalConversationCache`,
если нужен SQLite/JSON-store.
## Стриминг live-ответа
Для streaming-рендера текущего хода подписывайся на `events()` и
@@ -337,13 +386,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<ConversationRecord> { /* 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() {}
}
```
## Тесты
```
+1
View File
@@ -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)
@@ -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<Conversation> {
val snapshots = httpClient.get("$agentUrl/conversations") {
parameter("offset", offset)
parameter("limit", limit)
}.body<List<ConversationSnapshot>>()
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<pw.binom.agentik.journal.ConversationRecord>()
return rec.updatedAt
}
override fun close() {
@@ -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<ConversationRecord> = localStore.list(offset, limit)
override fun close() {} // owned by outer close
}
override fun close() {
cacheScope.cancel()
runBlocking { syncJob.join() }
delegate.close()
}
}
@@ -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?)
@@ -0,0 +1,58 @@
package pw.binom.agentik.client
import io.ktor.client.HttpClient
import io.ktor.client.call.body
import io.ktor.client.request.get
import io.ktor.client.request.parameter
import io.ktor.http.HttpStatusCode
import pw.binom.agentik.journal.ConversationRecord
import pw.binom.agentik.journal.ConversationStore
/**
* HTTP-реализация [ConversationStore] (read-only metadata view),
* ходящая в `:server`-фасад.
*
* **Endpoint**: `GET {baseUrl}/conversations?offset=&limit=` —
* возвращает `List<ConversationRecord>` (id, title, isTemporal, createdAt,
* updatedAt) БЕЗ handle'ов и image-support флагов (это лёгкая проекция
* для UI-списка; handle берётся через `agent.getConversation(id)`).
*
* **Read-only**: запись в `conversation` table — только через команды
* `agent.createConversation / deleteConversation / renameConversation`.
*
* Клиентский кэш строится композицией `HttpConversationStore` (snapshot)
* + `agent.outbox.agentEvents(after)` (live deltas: Created/Deleted/
* Renamed/Touched) — см. `client/README.md` секция
* «Кэш списка бесед».
*/
internal class HttpConversationStore(
private val httpClient: HttpClient,
private val baseUrl: String,
) : ConversationStore {
private val agentUrl: String = baseUrl.trimEnd('/')
override suspend fun get(id: String): ConversationRecord? {
val response = httpClient.get("$agentUrl/conversations/$id")
if (response.status == HttpStatusCode.NotFound) return null
check(response.status == HttpStatusCode.OK) {
"conversationStore.get($id): server returned ${response.status}"
}
return response.body<ConversationRecord>()
}
override suspend fun list(offset: Int, limit: Int): List<ConversationRecord> {
val response = httpClient.get("$agentUrl/conversations") {
parameter("offset", offset)
parameter("limit", limit)
}
check(response.status == HttpStatusCode.OK) {
"conversationStore.list: server returned ${response.status}"
}
return response.body<List<ConversationRecord>>()
}
override fun close() {
// HttpClient закрывает владелец (AgentClient / AgentikAgent).
}
}
+1 -1
View File
@@ -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)
@@ -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<ConversationRecord>
/** Переименовать диалог; `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<ConversationRecord> = 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
}
}
@@ -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)
}
+10 -6
View File
@@ -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,13 @@ kotlin {
jvmToolchain(21)
jvm()
macosX64()
macosArm64()
iosX64()
iosArm64()
iosSimulatorArm64()
linuxX64()
linuxArm64()
mingwX64()
sourceSets {
@@ -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<String, ConversationRecord> = mutableMapOf()
private val mutex = Mutex()
@@ -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
}
@@ -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
}
@@ -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<Conversation>
/** Все диалоги, начиная с [offset], как поток: подгружает по [PAGE_SIZE] за раз. */
fun getConversations(offset: Int = 0): Flow<Conversation> = 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 {
@@ -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<RequestRename>()
val newUpdatedAt = agent.renameConversation(id, req.title)
if (newUpdatedAt == null) {
call.respond(HttpStatusCode.NotFound)
return@patch
}
val req = call.receive<RequestRename>()
c.rename(req.title)
call.respond(c.snapshot())
// Возвращаем обновлённый record (лёгкая метадата, не snapshot handle'а).
val rec = agent.conversationStore.get(id)!!
call.respond(rec)
}
post("/conversations/{id}/messages") {
@@ -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<pw.binom.agentik.journal.ConversationRecord>()
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<Conversation> = emptyList()
override suspend fun renameConversation(id: String, title: String?): Instant? = null
override fun close() {}
}
@@ -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,
@@ -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<String, ChatConversation> = 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<ProtoConversation> =
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,
@@ -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
@@ -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,
@@ -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,
@@ -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,
@@ -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,
+1
View File
@@ -21,6 +21,7 @@ kotlin {
sourceSets {
commonMain.dependencies {
api(project(":journal-api"))
api(project(":journal-inmemory"))
api(project(":reflection-api"))
api(project(":context-api"))
}
@@ -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(),
@@ -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()
@@ -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),
@@ -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
@@ -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),