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.
This commit is contained in:
+71
-3
@@ -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<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() {}
|
||||
}
|
||||
```
|
||||
|
||||
## Тесты
|
||||
|
||||
```
|
||||
|
||||
@@ -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?)
|
||||
|
||||
Reference in New Issue
Block a user