Compare commits
3 Commits
c0a933d251
..
11
| Author | SHA1 | Date | |
|---|---|---|---|
| 639c7d1748 | |||
| 5f0e0da361 | |||
| acb4ee6186 |
@@ -3,6 +3,7 @@ package pw.binom.agentik.tui
|
|||||||
import kotlinx.coroutines.flow.Flow
|
import kotlinx.coroutines.flow.Flow
|
||||||
import kotlinx.coroutines.flow.MutableSharedFlow
|
import kotlinx.coroutines.flow.MutableSharedFlow
|
||||||
import kotlinx.coroutines.flow.emptyFlow
|
import kotlinx.coroutines.flow.emptyFlow
|
||||||
|
import pw.binom.agentik.journal.ConversationStore
|
||||||
import pw.binom.agentik.journal.JournalStore
|
import pw.binom.agentik.journal.JournalStore
|
||||||
import pw.binom.agentik.outbox.OutboxStore
|
import pw.binom.agentik.outbox.OutboxStore
|
||||||
import pw.binom.agentik.proto.Agent
|
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 journal: JournalStore = error("journal not used in TuiBackend tests")
|
||||||
override val outbox: OutboxStore = object : OutboxStore {
|
override val outbox: OutboxStore = object : OutboxStore {
|
||||||
override fun events(after: Instant?) = emptyFlow<pw.binom.agentik.outbox.CommonEvent>()
|
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 suspend fun earliestEventDate(): Instant = Instant.DISTANT_PAST
|
||||||
override fun close() {}
|
override fun close() {}
|
||||||
}
|
}
|
||||||
|
override val conversationStore: ConversationStore = error("conversationStore not used in TuiBackend tests")
|
||||||
|
|
||||||
override fun createConversation(temp: Boolean): Conversation {
|
override fun createConversation(temp: Boolean): Conversation {
|
||||||
createCount++
|
createCount++
|
||||||
@@ -49,8 +53,7 @@ internal class FakeAgent(
|
|||||||
override suspend fun deleteConversation(id: String): Boolean =
|
override suspend fun deleteConversation(id: String): Boolean =
|
||||||
conversations.removeAll { it.id == id }
|
conversations.removeAll { it.id == id }
|
||||||
|
|
||||||
override suspend fun getConversations(offset: Int, limit: Int): List<Conversation> =
|
override suspend fun renameConversation(id: String, title: String?): Instant? = null
|
||||||
conversations.toList()
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
|||||||
+120
-3
@@ -285,6 +285,55 @@ session.scope.launch {
|
|||||||
`rec is MessageRecord.UserMessage` для реплик пользователя,
|
`rec is MessageRecord.UserMessage` для реплик пользователя,
|
||||||
`rec is MessageRecord.ToolCall` для отрисовки tool-call баббла, и т.п.
|
`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-ответа
|
## Стриминг live-ответа
|
||||||
|
|
||||||
Для streaming-рендера текущего хода подписывайся на `events()` и
|
Для streaming-рендера текущего хода подписывайся на `events()` и
|
||||||
@@ -337,13 +386,81 @@ UI-обновление списка — отдельная задача, реш
|
|||||||
|
|
||||||
- **UI-рендеринг** — это твоя зона (Compose/HTML/etc.), `:client` только
|
- **UI-рендеринг** — это твоя зона (Compose/HTML/etc.), `:client` только
|
||||||
отдаёт типы и потоки.
|
отдаёт типы и потоки.
|
||||||
- **Персистентность кэша** — `InMemoryJournalStore` хранит в RAM. Для
|
- **Персистентность кэша** — `InMemoryJournalStore` и
|
||||||
диска пиши свой `MutableJournalStore` (см. `KsqliteJournalStore` в
|
`InMemoryMutableConversationStore` хранят в RAM. Для диска пиши свой
|
||||||
`:journal-ksqlite` как образец).
|
`MutableJournalStore` / `MutableConversationStore` (см. `KsqliteJournalStore`
|
||||||
|
в `:journal-ksqlite` как образец).
|
||||||
- **Нестандартные движковые настройки** — для `requestTimeout`,
|
- **Нестандартные движковые настройки** — для `requestTimeout`,
|
||||||
прокси и т.п. используй `agentikHttpClient(engineFactory, token)`
|
прокси и т.п. используй `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(":proto"))
|
||||||
api(project(":outbox-api"))
|
api(project(":outbox-api"))
|
||||||
api(project(":journal-api"))
|
api(project(":journal-api"))
|
||||||
|
implementation(project(":journal-inmemory"))
|
||||||
|
|
||||||
api(libs.ktor.client.core)
|
api(libs.ktor.client.core)
|
||||||
implementation(libs.ktor.client.content.negotiation)
|
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.delete
|
||||||
import io.ktor.client.request.get
|
import io.ktor.client.request.get
|
||||||
import io.ktor.client.request.parameter
|
import io.ktor.client.request.parameter
|
||||||
|
import io.ktor.client.request.patch
|
||||||
import io.ktor.client.request.post
|
import io.ktor.client.request.post
|
||||||
import io.ktor.client.request.setBody
|
import io.ktor.client.request.setBody
|
||||||
import io.ktor.client.statement.HttpResponse
|
|
||||||
import io.ktor.http.ContentType
|
import io.ktor.http.ContentType
|
||||||
import io.ktor.http.HttpStatusCode
|
import io.ktor.http.HttpStatusCode
|
||||||
import io.ktor.http.contentType
|
import io.ktor.http.contentType
|
||||||
import kotlinx.coroutines.runBlocking
|
import kotlinx.coroutines.runBlocking
|
||||||
|
import pw.binom.agentik.journal.ConversationStore
|
||||||
import pw.binom.agentik.journal.JournalStore
|
import pw.binom.agentik.journal.JournalStore
|
||||||
import pw.binom.agentik.outbox.OutboxStore
|
import pw.binom.agentik.outbox.OutboxStore
|
||||||
import pw.binom.agentik.proto.Agent
|
import pw.binom.agentik.proto.Agent
|
||||||
import pw.binom.agentik.proto.Conversation
|
import pw.binom.agentik.proto.Conversation
|
||||||
|
import kotlin.time.Instant
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* HTTP-реализация [Agent]. Ходит в `:server`-фасад, см. `agentikAgent(...)`.
|
* HTTP-реализация [Agent]. Ходит в `:server`-фасад, см. `agentikAgent(...)`.
|
||||||
*
|
*
|
||||||
* HttpClient создаётся внутри из переданного engine и закрывается в [close].
|
* HttpClient создаётся внутри из переданного engine и закрывается в [close].
|
||||||
*
|
*
|
||||||
* **Storage handles** ([journal], [outbox]) — read-only views на серверные
|
* **Storage handles** ([journal], [outbox], [conversationStore]) — read-only
|
||||||
* хранилища.
|
* views на серверные хранилища. Запись — только через команды
|
||||||
|
* [createConversation] / [deleteConversation] / [renameConversation].
|
||||||
*/
|
*/
|
||||||
internal class AgentClient(
|
internal class AgentClient(
|
||||||
override val id: String,
|
override val id: String,
|
||||||
@@ -35,6 +38,7 @@ internal class AgentClient(
|
|||||||
|
|
||||||
override val outbox: OutboxStore = HttpEventStore(httpClient = httpClient, baseUrl = agentUrl)
|
override val outbox: OutboxStore = HttpEventStore(httpClient = httpClient, baseUrl = agentUrl)
|
||||||
override val journal: JournalStore = HttpJournalStore(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 =
|
override fun createConversation(temp: Boolean): Conversation =
|
||||||
runBlocking {
|
runBlocking {
|
||||||
@@ -53,16 +57,18 @@ internal class AgentClient(
|
|||||||
}
|
}
|
||||||
|
|
||||||
override suspend fun deleteConversation(id: String): Boolean {
|
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
|
return response.status == HttpStatusCode.NoContent
|
||||||
}
|
}
|
||||||
|
|
||||||
override suspend fun getConversations(offset: Int, limit: Int): List<Conversation> {
|
override suspend fun renameConversation(id: String, title: String?): Instant? {
|
||||||
val snapshots = httpClient.get("$agentUrl/conversations") {
|
val response = httpClient.patch("$agentUrl/conversations/$id") {
|
||||||
parameter("offset", offset)
|
contentType(ContentType.Application.Json)
|
||||||
parameter("limit", limit)
|
setBody(RequestRename(title))
|
||||||
}.body<List<ConversationSnapshot>>()
|
}
|
||||||
return snapshots.map { ConversationClient(httpClient, agentUrl, it) }
|
if (response.status == HttpStatusCode.NotFound) return null
|
||||||
|
val rec = response.body<pw.binom.agentik.journal.ConversationRecord>()
|
||||||
|
return rec.updatedAt
|
||||||
}
|
}
|
||||||
|
|
||||||
override fun close() {
|
override fun close() {
|
||||||
|
|||||||
@@ -1,7 +1,20 @@
|
|||||||
package pw.binom.agentik.client
|
package pw.binom.agentik.client
|
||||||
|
|
||||||
import io.ktor.client.engine.HttpClientEngineFactory
|
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 pw.binom.agentik.proto.Agent
|
||||||
|
import kotlin.time.Instant
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Создаёт [Agent], который ходит в HTTP-фасад `agentikAgent` (модуль `:server`).
|
* Создаёт [Agent], который ходит в HTTP-фасад `agentikAgent` (модуль `:server`).
|
||||||
@@ -23,7 +36,7 @@ import pw.binom.agentik.proto.Agent
|
|||||||
* agent.outbox.conversationEvents(Instant.DISTANT_PAST, conv.id)
|
* agent.outbox.conversationEvents(Instant.DISTANT_PAST, conv.id)
|
||||||
* .map { it.event }
|
* .map { it.event }
|
||||||
* .collect { ... }
|
* .collect { ... }
|
||||||
* agent.close() // закрывает HttpClient
|
* agent.close() // закрывает HttpClient + локальный кэш
|
||||||
* ```
|
* ```
|
||||||
*
|
*
|
||||||
* ## Что клиент должен хранить локально (persistence)
|
* ## Что клиент должен хранить локально (persistence)
|
||||||
@@ -58,17 +71,103 @@ import pw.binom.agentik.proto.Agent
|
|||||||
* `clientId` генерируется один раз при первой установке (`UUID.randomUUID().toString()`)
|
* `clientId` генерируется один раз при первой установке (`UUID.randomUUID().toString()`)
|
||||||
* и больше не меняется — иначе сломается log multiplexing на сервере.
|
* и больше не меняется — иначе сломается 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()` закрывает
|
* **Lifecycle**: [Agent] — `AutoCloseable`. `agent.close()` закрывает
|
||||||
* HttpClient (идемпотентно). После этого `createConversation` /
|
* HttpClient + локальный кэш + background-coroutine (идемпотентно).
|
||||||
* `getConversation` etc. не определены.
|
* После этого `createConversation` / `getConversation` etc. не определены.
|
||||||
*/
|
*/
|
||||||
fun AgentikAgent(
|
fun AgentikAgent(
|
||||||
id: String,
|
id: String,
|
||||||
baseUrl: String,
|
baseUrl: String,
|
||||||
engineFactory: HttpClientEngineFactory<*>,
|
engineFactory: HttpClientEngineFactory<*>,
|
||||||
token: String? = null,
|
token: String? = null,
|
||||||
): Agent = AgentClient(
|
): Agent {
|
||||||
id = id,
|
val httpClient = agentikHttpClient(engineFactory = engineFactory, token = token)
|
||||||
baseUrl = baseUrl,
|
val client = AgentClient(id = id, baseUrl = baseUrl, httpClient = httpClient)
|
||||||
httpClient = agentikHttpClient(engineFactory = engineFactory, token = token),
|
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)
|
internal data class RequestCreateConversation(val temp: Boolean)
|
||||||
|
|
||||||
@Serializable
|
@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).
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -13,8 +13,10 @@ import kotlinx.coroutines.launch
|
|||||||
import pw.binom.agentik.outbox.CommonEvent
|
import pw.binom.agentik.outbox.CommonEvent
|
||||||
import pw.binom.agentik.outbox.OutboxStore
|
import pw.binom.agentik.outbox.OutboxStore
|
||||||
import kotlin.concurrent.atomics.AtomicBoolean
|
import kotlin.concurrent.atomics.AtomicBoolean
|
||||||
|
import kotlin.concurrent.atomics.AtomicReference
|
||||||
import kotlin.concurrent.atomics.ExperimentalAtomicApi
|
import kotlin.concurrent.atomics.ExperimentalAtomicApi
|
||||||
import kotlin.math.min
|
import kotlin.math.min
|
||||||
|
import kotlin.math.pow
|
||||||
import kotlin.random.Random
|
import kotlin.random.Random
|
||||||
import kotlin.time.Duration
|
import kotlin.time.Duration
|
||||||
import kotlin.time.Duration.Companion.seconds
|
import kotlin.time.Duration.Companion.seconds
|
||||||
@@ -151,8 +153,8 @@ class ReconnectingOutbox(
|
|||||||
private val started = AtomicBoolean(false)
|
private val started = AtomicBoolean(false)
|
||||||
private var job: Job? = null
|
private var job: Job? = null
|
||||||
|
|
||||||
@Volatile
|
@OptIn(ExperimentalAtomicApi::class)
|
||||||
private var lastSeen: Instant? = null
|
private val lastSeen: AtomicReference<Instant?> = AtomicReference(null)
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Live-события из [outbox] с авто-reconnect. [after] — начальный курсор;
|
* Live-события из [outbox] с авто-reconnect. [after] — начальный курсор;
|
||||||
@@ -183,10 +185,11 @@ class ReconnectingOutbox(
|
|||||||
@OptIn(ExperimentalAtomicApi::class)
|
@OptIn(ExperimentalAtomicApi::class)
|
||||||
private fun ensureStarted(initialCursor: Instant?) {
|
private fun ensureStarted(initialCursor: Instant?) {
|
||||||
if (!started.compareAndSet(false, true)) return
|
if (!started.compareAndSet(false, true)) return
|
||||||
lastSeen = initialCursor
|
lastSeen.store(initialCursor)
|
||||||
job = scope.launch { runLoop() }
|
job = scope.launch { runLoop() }
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@OptIn(ExperimentalAtomicApi::class)
|
||||||
private suspend fun runLoop() {
|
private suspend fun runLoop() {
|
||||||
var attempt = 0
|
var attempt = 0
|
||||||
var connected = false
|
var connected = false
|
||||||
@@ -194,12 +197,12 @@ class ReconnectingOutbox(
|
|||||||
attempt++
|
attempt++
|
||||||
_status.emit(ConnectionStatus.Connecting(attempt))
|
_status.emit(ConnectionStatus.Connecting(attempt))
|
||||||
val error: Throwable? = try {
|
val error: Throwable? = try {
|
||||||
outbox.events(after = lastSeen).collect { event ->
|
outbox.events(after = lastSeen.load()).collect { event ->
|
||||||
lastSeen = event.date
|
lastSeen.store(event.date)
|
||||||
_events.emit(event)
|
_events.emit(event)
|
||||||
if (!connected) {
|
if (!connected) {
|
||||||
connected = true
|
connected = true
|
||||||
_status.emit(ConnectionStatus.Connected(lastSeen!!))
|
_status.emit(ConnectionStatus.Connected(event.date))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
null
|
null
|
||||||
@@ -224,7 +227,7 @@ class ReconnectingOutbox(
|
|||||||
private fun computeBackoff(attempt: Int): Duration {
|
private fun computeBackoff(attempt: Int): Duration {
|
||||||
// attempt 1 → initial, 2 → initial * m, 3 → initial * m^2, ...
|
// attempt 1 → initial, 2 → initial * m, 3 → initial * m^2, ...
|
||||||
val base = (policy.initial.inWholeMilliseconds.toDouble() *
|
val base = (policy.initial.inWholeMilliseconds.toDouble() *
|
||||||
Math.pow(policy.multiplier, (attempt - 1).toDouble()))
|
policy.multiplier.pow((attempt - 1).toDouble()))
|
||||||
.toLong()
|
.toLong()
|
||||||
val capped = min(base, policy.max.inWholeMilliseconds)
|
val capped = min(base, policy.max.inWholeMilliseconds)
|
||||||
val jitterMs = (capped * policy.jitter * random.nextDouble()).toLong()
|
val jitterMs = (capped * policy.jitter * random.nextDouble()).toLong()
|
||||||
|
|||||||
+1
-1
@@ -194,7 +194,7 @@ suspend fun compact(dropFromOrderIdx: Long, conversationId: String): Long
|
|||||||
|
|
||||||
`compact` — атомарный «выбросить всё от `dropFromOrderIdx` и дальше, вставить новую синтетическую запись на следующий `order_idx`». Для v1 — просто `DELETE` от индекса (суммаризация появится в v2 вместе с LLM-вызовом для генерации текста).
|
`compact` — атомарный «выбросить всё от `dropFromOrderIdx` и дальше, вставить новую синтетическую запись на следующий `order_idx`». Для v1 — просто `DELETE` от индекса (суммаризация появится в v2 вместе с LLM-вызовом для генерации текста).
|
||||||
|
|
||||||
### `ConversationStore`
|
### `MutableConversationStore`
|
||||||
|
|
||||||
```kotlin
|
```kotlin
|
||||||
suspend fun upsert(record: ConversationRecord)
|
suspend fun upsert(record: ConversationRecord)
|
||||||
|
|||||||
@@ -1,27 +1,45 @@
|
|||||||
package pw.binom.agentik.journal
|
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 {
|
interface ConversationStore : AutoCloseable {
|
||||||
|
|
||||||
/** Создать или обновить snapshot диалога. */
|
|
||||||
suspend fun upsert(record: ConversationRecord)
|
|
||||||
|
|
||||||
/** Диалог по id, или `null`. */
|
/** Диалог по id, или `null`. */
|
||||||
suspend fun get(id: String): ConversationRecord?
|
suspend fun get(id: String): ConversationRecord?
|
||||||
|
|
||||||
/** Удалить диалог (вместе с его сообщениями и working memory). */
|
|
||||||
suspend fun delete(id: String): Boolean
|
|
||||||
|
|
||||||
/** Список диалогов, отсортированный по `updatedAt` DESC. */
|
/** Список диалогов, отсортированный по `updatedAt` DESC. */
|
||||||
suspend fun list(offset: Int, limit: Int): List<ConversationRecord>
|
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` диалога (например, после отправки сообщения). */
|
companion object {
|
||||||
suspend fun touch(id: String, now: Instant)
|
const val PAGE_SIZE: Int = 100
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
+21
@@ -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)
|
||||||
|
}
|
||||||
@@ -2,12 +2,10 @@ plugins {
|
|||||||
alias(libs.plugins.kotlin.multiplatform)
|
alias(libs.plugins.kotlin.multiplatform)
|
||||||
}
|
}
|
||||||
|
|
||||||
// KMP-реализация [MutableJournalStore] на `MutableList` + `Mutex` — для
|
// KMP-реализация [MutableJournalStore] и [MutableConversationStore] на
|
||||||
// тестов, dev-режима, embedded-сценариев (Android core, CLI, in-process кэш
|
// `MutableList`/`MutableMap` + `Mutex` — для тестов, dev-режима,
|
||||||
// в клиенте) и как образец для своей реализации.
|
// embedded-сценариев (Android core, CLI, in-process кэш в клиенте) и как
|
||||||
//
|
// образец для своей реализации.
|
||||||
// `list` фильтрует по `conversationId`+`createdAt>after` и сортирует
|
|
||||||
// по `createdAt ASC`. Paging — поверх отфильтрованного списка.
|
|
||||||
//
|
//
|
||||||
// Зависимости: только `:journal-api`. Никакого I/O — pure in-memory.
|
// Зависимости: только `:journal-api`. Никакого I/O — pure in-memory.
|
||||||
|
|
||||||
@@ -15,7 +13,13 @@ kotlin {
|
|||||||
jvmToolchain(21)
|
jvmToolchain(21)
|
||||||
|
|
||||||
jvm()
|
jvm()
|
||||||
|
macosX64()
|
||||||
|
macosArm64()
|
||||||
|
iosX64()
|
||||||
|
iosArm64()
|
||||||
|
iosSimulatorArm64()
|
||||||
linuxX64()
|
linuxX64()
|
||||||
|
linuxArm64()
|
||||||
mingwX64()
|
mingwX64()
|
||||||
|
|
||||||
sourceSets {
|
sourceSets {
|
||||||
|
|||||||
+17
-6
@@ -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.Mutex
|
||||||
import kotlinx.coroutines.sync.withLock
|
import kotlinx.coroutines.sync.withLock
|
||||||
import kotlin.time.Clock
|
|
||||||
import pw.binom.agentik.journal.ConversationRecord
|
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
|
import kotlin.time.Instant
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Thread-safe Map-импл [ConversationStore].
|
* Thread-safe Map-импл [MutableConversationStore] для клиентских
|
||||||
|
* in-process кэшей (и тестов/dev-режима).
|
||||||
*
|
*
|
||||||
* Использует `Mutex` для атомарности read-modify-write операций
|
* Использует `Mutex` для атомарности read-modify-write операций
|
||||||
* (rename, touch) — иначе два параллельных `rename` могут потерять обновления
|
* (rename, touch) — иначе два параллельных `rename` могут потерять обновления
|
||||||
* (lost-update race), что в SQLite невозможно из-за driver-level locking.
|
* (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,
|
private val clock: Clock = Clock.System,
|
||||||
) : ConversationStore {
|
) : MutableConversationStore {
|
||||||
|
|
||||||
private val byId: MutableMap<String, ConversationRecord> = mutableMapOf()
|
private val byId: MutableMap<String, ConversationRecord> = mutableMapOf()
|
||||||
private val mutex = Mutex()
|
private val mutex = Mutex()
|
||||||
+11
-11
@@ -1,4 +1,4 @@
|
|||||||
package pw.binom.agentik.storage.inmemory
|
package pw.binom.agentik.journal.inmemory
|
||||||
|
|
||||||
import pw.binom.agentik.journal.ConversationRecord
|
import pw.binom.agentik.journal.ConversationRecord
|
||||||
import kotlin.test.Test
|
import kotlin.test.Test
|
||||||
@@ -9,11 +9,11 @@ import kotlin.test.assertTrue
|
|||||||
import kotlin.time.Instant
|
import kotlin.time.Instant
|
||||||
import kotlinx.coroutines.test.runTest
|
import kotlinx.coroutines.test.runTest
|
||||||
|
|
||||||
class InMemoryConversationStoreTest {
|
class InMemoryMutableConversationStoreTest {
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
fun `upsert and get roundtrip preserves all fields`() = runTest {
|
fun `upsert and get roundtrip preserves all fields`() = runTest {
|
||||||
val store = InMemoryConversationStore()
|
val store = InMemoryMutableConversationStore()
|
||||||
val rec = ConversationRecord(
|
val rec = ConversationRecord(
|
||||||
id = "c1",
|
id = "c1",
|
||||||
title = "test",
|
title = "test",
|
||||||
@@ -28,13 +28,13 @@ class InMemoryConversationStoreTest {
|
|||||||
|
|
||||||
@Test
|
@Test
|
||||||
fun `get returns null for missing id`() = runTest {
|
fun `get returns null for missing id`() = runTest {
|
||||||
val store = InMemoryConversationStore()
|
val store = InMemoryMutableConversationStore()
|
||||||
assertNull(store.get("nope"))
|
assertNull(store.get("nope"))
|
||||||
}
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
fun `delete removes the record and returns true`() = runTest {
|
fun `delete removes the record and returns true`() = runTest {
|
||||||
val store = InMemoryConversationStore()
|
val store = InMemoryMutableConversationStore()
|
||||||
store.upsert(
|
store.upsert(
|
||||||
ConversationRecord(
|
ConversationRecord(
|
||||||
"c1", null, false,
|
"c1", null, false,
|
||||||
@@ -50,7 +50,7 @@ class InMemoryConversationStoreTest {
|
|||||||
|
|
||||||
@Test
|
@Test
|
||||||
fun `list sorts by updatedAt DESC and respects offset+limit`() = runTest {
|
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")
|
val t0 = Instant.parse("2026-09-15T10:00:00Z")
|
||||||
store.upsert(ConversationRecord("c1", null, false, t0, t0))
|
store.upsert(ConversationRecord("c1", null, false, t0, t0))
|
||||||
store.upsert(ConversationRecord("c2", null, false, t0, t0.plus(kotlin.time.Duration.parse("PT60S"))))
|
store.upsert(ConversationRecord("c2", null, false, t0, t0.plus(kotlin.time.Duration.parse("PT60S"))))
|
||||||
@@ -69,7 +69,7 @@ class InMemoryConversationStoreTest {
|
|||||||
|
|
||||||
@Test
|
@Test
|
||||||
fun `rename updates title and updatedAt returns new updatedAt`() = runTest {
|
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")
|
val t0 = Instant.parse("2026-09-15T10:00:00Z")
|
||||||
store.upsert(ConversationRecord("c1", null, false, t0, t0))
|
store.upsert(ConversationRecord("c1", null, false, t0, t0))
|
||||||
|
|
||||||
@@ -84,7 +84,7 @@ class InMemoryConversationStoreTest {
|
|||||||
|
|
||||||
@Test
|
@Test
|
||||||
fun `rename with null title clears it`() = runTest {
|
fun `rename with null title clears it`() = runTest {
|
||||||
val store = InMemoryConversationStore()
|
val store = InMemoryMutableConversationStore()
|
||||||
val t0 = Instant.parse("2026-09-15T10:00:00Z")
|
val t0 = Instant.parse("2026-09-15T10:00:00Z")
|
||||||
store.upsert(ConversationRecord("c1", "old", false, t0, t0))
|
store.upsert(ConversationRecord("c1", "old", false, t0, t0))
|
||||||
store.rename("c1", null)
|
store.rename("c1", null)
|
||||||
@@ -93,13 +93,13 @@ class InMemoryConversationStoreTest {
|
|||||||
|
|
||||||
@Test
|
@Test
|
||||||
fun `rename returns null for missing conversation`() = runTest {
|
fun `rename returns null for missing conversation`() = runTest {
|
||||||
val store = InMemoryConversationStore()
|
val store = InMemoryMutableConversationStore()
|
||||||
assertNull(store.rename("nope", "x"))
|
assertNull(store.rename("nope", "x"))
|
||||||
}
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
fun `touch bumps updatedAt without changing other fields`() = runTest {
|
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 t0 = Instant.parse("2026-09-15T10:00:00Z")
|
||||||
val t1 = Instant.parse("2026-09-15T10:01:00Z")
|
val t1 = Instant.parse("2026-09-15T10:01:00Z")
|
||||||
store.upsert(ConversationRecord("c1", "title", false, t0, t0))
|
store.upsert(ConversationRecord("c1", "title", false, t0, t0))
|
||||||
@@ -112,7 +112,7 @@ class InMemoryConversationStoreTest {
|
|||||||
|
|
||||||
@Test
|
@Test
|
||||||
fun `close is idempotent and does nothing`() {
|
fun `close is idempotent and does nothing`() {
|
||||||
val store = InMemoryConversationStore()
|
val store = InMemoryMutableConversationStore()
|
||||||
store.close()
|
store.close()
|
||||||
store.close() // должно быть no-op
|
store.close() // должно быть no-op
|
||||||
}
|
}
|
||||||
@@ -45,4 +45,13 @@ sealed interface AgentEvent {
|
|||||||
@Serializable
|
@Serializable
|
||||||
@SerialName("renamed")
|
@SerialName("renamed")
|
||||||
data class Renamed(override val date: Instant, val id: String, val title: String?) : AgentEvent
|
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
|
package pw.binom.agentik.proto
|
||||||
|
|
||||||
import kotlinx.coroutines.flow.Flow
|
import pw.binom.agentik.journal.ConversationStore
|
||||||
import kotlinx.coroutines.flow.flow
|
|
||||||
import pw.binom.agentik.journal.JournalStore
|
import pw.binom.agentik.journal.JournalStore
|
||||||
import pw.binom.agentik.outbox.OutboxStore
|
import pw.binom.agentik.outbox.OutboxStore
|
||||||
import kotlin.time.Instant
|
import kotlin.time.Instant
|
||||||
@@ -13,11 +12,16 @@ import kotlin.time.Instant
|
|||||||
* [createConversation] возвращает [Conversation], который сам хранит историю
|
* [createConversation] возвращает [Conversation], который сам хранит историю
|
||||||
* и которому отправляют ходы через [Conversation.send].
|
* и которому отправляют ходы через [Conversation.send].
|
||||||
*
|
*
|
||||||
* **Хранилища вынесены в [Agent.journal] и [Agent.outbox]**: оба read-only.
|
* **Хранилища вынесены в [Agent.journal], [Agent.outbox] и
|
||||||
* События больше НЕ часть [Agent] (раньше были `events()`/`allEvents()`) —
|
* [Agent.conversationStore]**: все три read-only views. События живут
|
||||||
* они теперь живут в [outbox] как `OutboxStore.events(after)` /
|
* в [outbox] как `OutboxStore.events(after)` / `outbox.agentEvents(after)`.
|
||||||
* `outbox.agentEvents(after)`. Это даёт единый путь для всех read-операций
|
* Это даёт единый путь для всех read-операций по хранилищу и убирает
|
||||||
* по хранилищу и убирает дублирование между протоколом и хранилищем.
|
* дублирование между протоколом и хранилищем.
|
||||||
|
*
|
||||||
|
* **Команды** (create / delete / rename) живут прямо на [Agent]. Они
|
||||||
|
* шлются клиентом и выполняются сервером — клиент **не** пишет в стор
|
||||||
|
* напрямую. Клиентский кэш [conversationStore] обновляется через
|
||||||
|
* `outbox.agentEvents()` (Created / Deleted / Renamed / Touched).
|
||||||
*/
|
*/
|
||||||
interface Agent : AutoCloseable {
|
interface Agent : AutoCloseable {
|
||||||
|
|
||||||
@@ -59,6 +63,22 @@ interface Agent : AutoCloseable {
|
|||||||
*/
|
*/
|
||||||
val outbox: OutboxStore
|
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-диалог с агентом. */
|
/** Создаёт новый stateful-диалог с агентом. */
|
||||||
fun createConversation(temp: Boolean): Conversation
|
fun createConversation(temp: Boolean): Conversation
|
||||||
|
|
||||||
@@ -68,19 +88,11 @@ interface Agent : AutoCloseable {
|
|||||||
/** Удаляет диалог. Возвращает `true`, если диалог существовал и удалён. */
|
/** Удаляет диалог. Возвращает `true`, если диалог существовал и удалён. */
|
||||||
suspend fun deleteConversation(id: String): Boolean
|
suspend fun deleteConversation(id: String): Boolean
|
||||||
|
|
||||||
/** Страница диалогов: не более [limit] штук, начиная с [offset]-го. */
|
/**
|
||||||
suspend fun getConversations(offset: Int, limit: Int): List<Conversation>
|
* Переименовывает диалог; `null` для сброса заголовка. Возвращает новый
|
||||||
|
* `updatedAt` или `null`, если диалог не найден.
|
||||||
/** Все диалоги, начиная с [offset], как поток: подгружает по [PAGE_SIZE] за раз. */
|
*/
|
||||||
fun getConversations(offset: Int = 0): Flow<Conversation> = flow {
|
suspend fun renameConversation(id: String, title: String?): Instant?
|
||||||
var skip = offset
|
|
||||||
while (true) {
|
|
||||||
val page = getConversations(skip, PAGE_SIZE)
|
|
||||||
if (page.isEmpty()) break
|
|
||||||
page.forEach { emit(it) }
|
|
||||||
skip += page.size
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
companion object {
|
companion object {
|
||||||
|
|
||||||
|
|||||||
@@ -38,7 +38,7 @@ internal fun Route.agentikRoutes(agent: Agent) {
|
|||||||
get("/conversations") {
|
get("/conversations") {
|
||||||
val offset = call.request.queryParameters["offset"]?.toIntOrNull() ?: 0
|
val offset = call.request.queryParameters["offset"]?.toIntOrNull() ?: 0
|
||||||
val limit = call.request.queryParameters["limit"]?.toIntOrNull() ?: Agent.PAGE_SIZE
|
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") {
|
post("/conversations") {
|
||||||
@@ -63,14 +63,15 @@ internal fun Route.agentikRoutes(agent: Agent) {
|
|||||||
|
|
||||||
patch("/conversations/{id}") {
|
patch("/conversations/{id}") {
|
||||||
val id = call.parameters["id"]!!
|
val id = call.parameters["id"]!!
|
||||||
val c = agent.getConversation(id)
|
val req = call.receive<RequestRename>()
|
||||||
if (c == null) {
|
val newUpdatedAt = agent.renameConversation(id, req.title)
|
||||||
|
if (newUpdatedAt == null) {
|
||||||
call.respond(HttpStatusCode.NotFound)
|
call.respond(HttpStatusCode.NotFound)
|
||||||
return@patch
|
return@patch
|
||||||
}
|
}
|
||||||
val req = call.receive<RequestRename>()
|
// Возвращаем обновлённый record (лёгкая метадата, не snapshot handle'а).
|
||||||
c.rename(req.title)
|
val rec = agent.conversationStore.get(id)!!
|
||||||
call.respond(c.snapshot())
|
call.respond(rec)
|
||||||
}
|
}
|
||||||
|
|
||||||
post("/conversations/{id}/messages") {
|
post("/conversations/{id}/messages") {
|
||||||
|
|||||||
@@ -13,6 +13,7 @@ import io.ktor.server.engine.embeddedServer
|
|||||||
import io.ktor.server.routing.routing
|
import io.ktor.server.routing.routing
|
||||||
import kotlinx.coroutines.flow.emptyFlow
|
import kotlinx.coroutines.flow.emptyFlow
|
||||||
import kotlinx.coroutines.runBlocking
|
import kotlinx.coroutines.runBlocking
|
||||||
|
import pw.binom.agentik.journal.ConversationStore
|
||||||
import pw.binom.agentik.journal.JournalStore
|
import pw.binom.agentik.journal.JournalStore
|
||||||
import pw.binom.agentik.journal.MessageRecord
|
import pw.binom.agentik.journal.MessageRecord
|
||||||
import pw.binom.agentik.outbox.CommonEvent
|
import pw.binom.agentik.outbox.CommonEvent
|
||||||
@@ -53,10 +54,15 @@ class BearerTokenTest {
|
|||||||
override suspend fun earliestEventDate(): Instant = Instant.DISTANT_PAST
|
override suspend fun earliestEventDate(): Instant = Instant.DISTANT_PAST
|
||||||
override fun close() {}
|
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 fun createConversation(temp: Boolean): Conversation = TODO("not needed by tests")
|
||||||
override suspend fun getConversation(id: String): Conversation? = null
|
override suspend fun getConversation(id: String): Conversation? = null
|
||||||
override suspend fun deleteConversation(id: String): Boolean = false
|
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() {}
|
override fun close() {}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -363,7 +363,7 @@ private fun runServer() {
|
|||||||
|
|
||||||
val agent = ChatAgent(
|
val agent = ChatAgent(
|
||||||
id = "agentik",
|
id = "agentik",
|
||||||
conversationStore = sqliteStores.conversations,
|
mutableConversationStore = sqliteStores.conversations,
|
||||||
messageStore = sqliteStores.messages,
|
messageStore = sqliteStores.messages,
|
||||||
workingMemoryStore = sqliteStores.workingMemory,
|
workingMemoryStore = sqliteStores.workingMemory,
|
||||||
reflectionStore = sqliteStores.reflections,
|
reflectionStore = sqliteStores.reflections,
|
||||||
|
|||||||
@@ -1,11 +1,6 @@
|
|||||||
package pw.binom.agentik.standalone.agent
|
package pw.binom.agentik.standalone.agent
|
||||||
|
|
||||||
import kotlin.time.Instant
|
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.runBlocking
|
||||||
import kotlinx.coroutines.sync.Mutex
|
import kotlinx.coroutines.sync.Mutex
|
||||||
import kotlinx.coroutines.sync.withLock
|
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.proto.Agent as ProtoAgent
|
||||||
import pw.binom.agentik.outbox.AgentEvent
|
import pw.binom.agentik.outbox.AgentEvent
|
||||||
import pw.binom.agentik.outbox.CommonEvent
|
import pw.binom.agentik.outbox.CommonEvent
|
||||||
import pw.binom.agentik.outbox.Event as ProtoEvent
|
|
||||||
import pw.binom.agentik.outbox.MutableOutboxStore
|
import pw.binom.agentik.outbox.MutableOutboxStore
|
||||||
import pw.binom.agentik.journal.JournalStore
|
import pw.binom.agentik.journal.JournalStore
|
||||||
import pw.binom.agentik.outbox.OutboxStore
|
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.standalone.llm.LlmConfig
|
||||||
import pw.binom.agentik.journal.ConversationRecord
|
import pw.binom.agentik.journal.ConversationRecord
|
||||||
import pw.binom.agentik.journal.ConversationStore
|
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.Reflection
|
||||||
import pw.binom.agentik.reflection.ReflectionStore
|
import pw.binom.agentik.reflection.ReflectionStore
|
||||||
import pw.binom.agentik.journal.MutableJournalStore
|
import pw.binom.agentik.journal.MutableJournalStore
|
||||||
@@ -65,7 +59,7 @@ import pw.binom.litert.LiteLlm
|
|||||||
*/
|
*/
|
||||||
class ChatAgent(
|
class ChatAgent(
|
||||||
override val id: String,
|
override val id: String,
|
||||||
private val conversationStore: ConversationStore,
|
private val mutableConversationStore: MutableConversationStore,
|
||||||
private val messageStore: MutableJournalStore,
|
private val messageStore: MutableJournalStore,
|
||||||
private val workingMemoryStore: ContextStore,
|
private val workingMemoryStore: ContextStore,
|
||||||
private val reflectionStore: ReflectionStore,
|
private val reflectionStore: ReflectionStore,
|
||||||
@@ -239,6 +233,16 @@ class ChatAgent(
|
|||||||
override val outbox: OutboxStore
|
override val outbox: OutboxStore
|
||||||
get() = eventStore
|
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 liveLock = Mutex()
|
||||||
private val live: MutableMap<String, ChatConversation> = HashMap()
|
private val live: MutableMap<String, ChatConversation> = HashMap()
|
||||||
@@ -262,12 +266,12 @@ class ChatAgent(
|
|||||||
// и не переживают рестарт агента (см. Memory #3709).
|
// и не переживают рестарт агента (см. Memory #3709).
|
||||||
if (!temp) {
|
if (!temp) {
|
||||||
runBlocking {
|
runBlocking {
|
||||||
conversationStore.upsert(rec)
|
mutableConversationStore.upsert(rec)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
val conv = ChatConversation(
|
val conv = ChatConversation(
|
||||||
record = rec,
|
record = rec,
|
||||||
conversationStore = conversationStore,
|
conversationStore = mutableConversationStore,
|
||||||
messageStore = messageStore,
|
messageStore = messageStore,
|
||||||
workingMemoryStore = workingMemoryStore,
|
workingMemoryStore = workingMemoryStore,
|
||||||
reflectionStore = reflectionStore,
|
reflectionStore = reflectionStore,
|
||||||
@@ -302,7 +306,7 @@ class ChatAgent(
|
|||||||
|
|
||||||
override suspend fun getConversation(id: String): ProtoConversation? {
|
override suspend fun getConversation(id: String): ProtoConversation? {
|
||||||
liveLock.withLock { live[id] }?.let { if (!it.isClosed) return it }
|
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 {
|
return newConversation(rec).also {
|
||||||
liveLock.withLock { live[id] = it }
|
liveLock.withLock { live[id] = it }
|
||||||
}
|
}
|
||||||
@@ -311,7 +315,7 @@ class ChatAgent(
|
|||||||
override suspend fun deleteConversation(id: String): Boolean {
|
override suspend fun deleteConversation(id: String): Boolean {
|
||||||
val conv = liveLock.withLock { live.remove(id) }
|
val conv = liveLock.withLock { live.remove(id) }
|
||||||
conv?.close()
|
conv?.close()
|
||||||
val ok = conversationStore.delete(id)
|
val ok = mutableConversationStore.delete(id)
|
||||||
if (ok) {
|
if (ok) {
|
||||||
val event = AgentEvent.Deleted(date = now(), id = id)
|
val event = AgentEvent.Deleted(date = now(), id = id)
|
||||||
eventStore.append(CommonEvent.Agent(date = now(), event = event))
|
eventStore.append(CommonEvent.Agent(date = now(), event = event))
|
||||||
@@ -319,17 +323,16 @@ class ChatAgent(
|
|||||||
return ok
|
return ok
|
||||||
}
|
}
|
||||||
|
|
||||||
override suspend fun getConversations(offset: Int, limit: Int): List<ProtoConversation> =
|
override suspend fun renameConversation(id: String, title: String?): Instant? {
|
||||||
conversationStore.list(offset = offset, limit = limit).map { rec ->
|
val newUpdatedAt = mutableConversationStore.rename(id, title) ?: return null
|
||||||
liveLock.withLock { live[rec.id] }
|
val event = AgentEvent.Renamed(date = newUpdatedAt, id = id, title = title)
|
||||||
?: newConversation(rec).also {
|
eventStore.append(CommonEvent.Agent(date = newUpdatedAt, event = event))
|
||||||
liveLock.withLock { live[rec.id] = it }
|
return newUpdatedAt
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
private fun newConversation(rec: ConversationRecord): ChatConversation = ChatConversation(
|
private fun newConversation(rec: ConversationRecord): ChatConversation = ChatConversation(
|
||||||
record = rec,
|
record = rec,
|
||||||
conversationStore = conversationStore,
|
conversationStore = mutableConversationStore,
|
||||||
messageStore = messageStore,
|
messageStore = messageStore,
|
||||||
workingMemoryStore = workingMemoryStore,
|
workingMemoryStore = workingMemoryStore,
|
||||||
reflectionStore = reflectionStore,
|
reflectionStore = reflectionStore,
|
||||||
|
|||||||
+14
-3
@@ -7,7 +7,6 @@ import kotlinx.coroutines.Job
|
|||||||
import kotlinx.coroutines.SupervisorJob
|
import kotlinx.coroutines.SupervisorJob
|
||||||
import kotlinx.coroutines.cancel
|
import kotlinx.coroutines.cancel
|
||||||
import kotlinx.coroutines.cancelAndJoin
|
import kotlinx.coroutines.cancelAndJoin
|
||||||
import kotlinx.coroutines.flow.Flow
|
|
||||||
import kotlinx.coroutines.launch
|
import kotlinx.coroutines.launch
|
||||||
import kotlinx.coroutines.runBlocking
|
import kotlinx.coroutines.runBlocking
|
||||||
import kotlinx.coroutines.sync.Mutex
|
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.skills.SkillStore
|
||||||
import pw.binom.agentik.journal.Content
|
import pw.binom.agentik.journal.Content
|
||||||
import pw.binom.agentik.journal.ConversationRecord
|
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.MessageContext
|
||||||
import pw.binom.agentik.journal.MessageOrigin
|
import pw.binom.agentik.journal.MessageOrigin
|
||||||
import pw.binom.agentik.journal.MessageRecord
|
import pw.binom.agentik.journal.MessageRecord
|
||||||
@@ -54,7 +53,7 @@ import pw.binom.agentik.toolsets.NamedTool
|
|||||||
|
|
||||||
class ConversationLoop(
|
class ConversationLoop(
|
||||||
record: ConversationRecord,
|
record: ConversationRecord,
|
||||||
private val conversationStore: ConversationStore,
|
private val conversationStore: MutableConversationStore,
|
||||||
private val messageStore: MutableJournalStore,
|
private val messageStore: MutableJournalStore,
|
||||||
private val workingMemoryStore: ContextStore,
|
private val workingMemoryStore: ContextStore,
|
||||||
private val reflectionStore: ReflectionStore?,
|
private val reflectionStore: ReflectionStore?,
|
||||||
@@ -440,6 +439,18 @@ class ConversationLoop(
|
|||||||
|
|
||||||
state.record = state.record.copy(updatedAt = assistantAt)
|
state.record = state.record.copy(updatedAt = assistantAt)
|
||||||
conversationStore.touch(id, 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
|
// BackgroundScheduler is event-driven — подписан на BackgroundEventBus
|
||||||
// (compaction/lifecycle/tool-failure events). Никаких interval-based
|
// (compaction/lifecycle/tool-failure events). Никаких interval-based
|
||||||
|
|||||||
+24
-7
@@ -62,7 +62,7 @@ class ChatAgentTest {
|
|||||||
skills: SkillCatalog = SkillCatalog.EMPTY,
|
skills: SkillCatalog = SkillCatalog.EMPTY,
|
||||||
): ChatAgent = ChatAgent(
|
): ChatAgent = ChatAgent(
|
||||||
id = "agentik",
|
id = "agentik",
|
||||||
conversationStore = sqliteStores.conversations,
|
mutableConversationStore = sqliteStores.conversations,
|
||||||
messageStore = sqliteStores.messages,
|
messageStore = sqliteStores.messages,
|
||||||
workingMemoryStore = sqliteStores.workingMemory,
|
workingMemoryStore = sqliteStores.workingMemory,
|
||||||
reflectionStore = sqliteStores.reflections,
|
reflectionStore = sqliteStores.reflections,
|
||||||
@@ -146,15 +146,32 @@ class ChatAgentTest {
|
|||||||
}
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
fun `getConversations returns all stored persistent conversations`() = runTest {
|
fun `conversationStore list returns all stored persistent conversations`() = runTest {
|
||||||
val agent = newAgent()
|
val agent = newAgent()
|
||||||
agent.createConversation(temp = false)
|
agent.createConversation(temp = false)
|
||||||
agent.createConversation(temp = true)
|
agent.createConversation(temp = true)
|
||||||
val list = agent.getConversations(0, 10)
|
val list = agent.conversationStore.list(0, 10)
|
||||||
// temp-беседы не персистятся, в списке только persistent
|
// temp-беседы не персистятся, в списке только persistent
|
||||||
assertEquals(1, list.size)
|
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
|
@Test
|
||||||
fun `deleteConversation removes conversation and data`() = runTest {
|
fun `deleteConversation removes conversation and data`() = runTest {
|
||||||
val agent = newAgent()
|
val agent = newAgent()
|
||||||
@@ -458,7 +475,7 @@ class ChatAgentTest {
|
|||||||
sqliteStores = KsqliteStores.open(dbPath)
|
sqliteStores = KsqliteStores.open(dbPath)
|
||||||
val agent1 = ChatAgent(
|
val agent1 = ChatAgent(
|
||||||
id = "agentik",
|
id = "agentik",
|
||||||
conversationStore = sqliteStores.conversations,
|
mutableConversationStore = sqliteStores.conversations,
|
||||||
messageStore = sqliteStores.messages,
|
messageStore = sqliteStores.messages,
|
||||||
workingMemoryStore = sqliteStores.workingMemory,
|
workingMemoryStore = sqliteStores.workingMemory,
|
||||||
reflectionStore = sqliteStores.reflections,
|
reflectionStore = sqliteStores.reflections,
|
||||||
@@ -478,7 +495,7 @@ class ChatAgentTest {
|
|||||||
sqliteStores = KsqliteStores.open(dbPath)
|
sqliteStores = KsqliteStores.open(dbPath)
|
||||||
val agent2 = ChatAgent(
|
val agent2 = ChatAgent(
|
||||||
id = "agentik",
|
id = "agentik",
|
||||||
conversationStore = sqliteStores.conversations,
|
mutableConversationStore = sqliteStores.conversations,
|
||||||
messageStore = sqliteStores.messages,
|
messageStore = sqliteStores.messages,
|
||||||
workingMemoryStore = sqliteStores.workingMemory,
|
workingMemoryStore = sqliteStores.workingMemory,
|
||||||
reflectionStore = sqliteStores.reflections,
|
reflectionStore = sqliteStores.reflections,
|
||||||
@@ -500,7 +517,7 @@ class ChatAgentTest {
|
|||||||
sqliteStores = KsqliteStores.open(dbPath)
|
sqliteStores = KsqliteStores.open(dbPath)
|
||||||
val agent1 = ChatAgent(
|
val agent1 = ChatAgent(
|
||||||
id = "agentik",
|
id = "agentik",
|
||||||
conversationStore = sqliteStores.conversations,
|
mutableConversationStore = sqliteStores.conversations,
|
||||||
messageStore = sqliteStores.messages,
|
messageStore = sqliteStores.messages,
|
||||||
workingMemoryStore = sqliteStores.workingMemory,
|
workingMemoryStore = sqliteStores.workingMemory,
|
||||||
reflectionStore = sqliteStores.reflections,
|
reflectionStore = sqliteStores.reflections,
|
||||||
@@ -518,7 +535,7 @@ class ChatAgentTest {
|
|||||||
sqliteStores = KsqliteStores.open(dbPath)
|
sqliteStores = KsqliteStores.open(dbPath)
|
||||||
val agent2 = ChatAgent(
|
val agent2 = ChatAgent(
|
||||||
id = "agentik",
|
id = "agentik",
|
||||||
conversationStore = sqliteStores.conversations,
|
mutableConversationStore = sqliteStores.conversations,
|
||||||
messageStore = sqliteStores.messages,
|
messageStore = sqliteStores.messages,
|
||||||
workingMemoryStore = sqliteStores.workingMemory,
|
workingMemoryStore = sqliteStores.workingMemory,
|
||||||
reflectionStore = sqliteStores.reflections,
|
reflectionStore = sqliteStores.reflections,
|
||||||
|
|||||||
+1
-1
@@ -43,7 +43,7 @@ class ChatAgentToolsetsTest {
|
|||||||
val sqliteStores = KsqliteStores.inMemory("toolsets-${kotlin.random.Random.nextLong()}")
|
val sqliteStores = KsqliteStores.inMemory("toolsets-${kotlin.random.Random.nextLong()}")
|
||||||
val agent = ChatAgent(
|
val agent = ChatAgent(
|
||||||
id = "test-agent",
|
id = "test-agent",
|
||||||
conversationStore = sqliteStores.conversations,
|
mutableConversationStore = sqliteStores.conversations,
|
||||||
messageStore = sqliteStores.messages,
|
messageStore = sqliteStores.messages,
|
||||||
workingMemoryStore = sqliteStores.workingMemory,
|
workingMemoryStore = sqliteStores.workingMemory,
|
||||||
reflectionStore = sqliteStores.reflections,
|
reflectionStore = sqliteStores.reflections,
|
||||||
|
|||||||
+1
-1
@@ -63,7 +63,7 @@ class CompactionTest {
|
|||||||
val reviewer = if (memoryStore != null) KeywordMdReviewer() else null
|
val reviewer = if (memoryStore != null) KeywordMdReviewer() else null
|
||||||
return ChatAgent(
|
return ChatAgent(
|
||||||
id = "test",
|
id = "test",
|
||||||
conversationStore = sqliteStores.conversations,
|
mutableConversationStore = sqliteStores.conversations,
|
||||||
messageStore = sqliteStores.messages,
|
messageStore = sqliteStores.messages,
|
||||||
workingMemoryStore = sqliteStores.workingMemory,
|
workingMemoryStore = sqliteStores.workingMemory,
|
||||||
reflectionStore = sqliteStores.reflections,
|
reflectionStore = sqliteStores.reflections,
|
||||||
|
|||||||
+3
-3
@@ -76,7 +76,7 @@ class MemoryWiringTest {
|
|||||||
contextCompactor: ContextCompactor = EchoCompactor,
|
contextCompactor: ContextCompactor = EchoCompactor,
|
||||||
): ChatAgent = ChatAgent(
|
): ChatAgent = ChatAgent(
|
||||||
id = "agentik",
|
id = "agentik",
|
||||||
conversationStore = sqliteStores.conversations,
|
mutableConversationStore = sqliteStores.conversations,
|
||||||
messageStore = sqliteStores.messages,
|
messageStore = sqliteStores.messages,
|
||||||
workingMemoryStore = sqliteStores.workingMemory,
|
workingMemoryStore = sqliteStores.workingMemory,
|
||||||
reflectionStore = sqliteStores.reflections,
|
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 soulBody = "I am a helpful test persona. I always answer in one short line."
|
||||||
val agent = ChatAgent(
|
val agent = ChatAgent(
|
||||||
id = "agentik",
|
id = "agentik",
|
||||||
conversationStore = sqliteStores.conversations,
|
mutableConversationStore = sqliteStores.conversations,
|
||||||
messageStore = sqliteStores.messages,
|
messageStore = sqliteStores.messages,
|
||||||
workingMemoryStore = sqliteStores.workingMemory,
|
workingMemoryStore = sqliteStores.workingMemory,
|
||||||
reflectionStore = sqliteStores.reflections,
|
reflectionStore = sqliteStores.reflections,
|
||||||
@@ -143,7 +143,7 @@ class MemoryWiringTest {
|
|||||||
fun `soul body not added when null`() = runBlocking {
|
fun `soul body not added when null`() = runBlocking {
|
||||||
val agent = ChatAgent(
|
val agent = ChatAgent(
|
||||||
id = "agentik",
|
id = "agentik",
|
||||||
conversationStore = sqliteStores.conversations,
|
mutableConversationStore = sqliteStores.conversations,
|
||||||
messageStore = sqliteStores.messages,
|
messageStore = sqliteStores.messages,
|
||||||
workingMemoryStore = sqliteStores.workingMemory,
|
workingMemoryStore = sqliteStores.workingMemory,
|
||||||
reflectionStore = sqliteStores.reflections,
|
reflectionStore = sqliteStores.reflections,
|
||||||
|
|||||||
@@ -21,6 +21,7 @@ kotlin {
|
|||||||
sourceSets {
|
sourceSets {
|
||||||
commonMain.dependencies {
|
commonMain.dependencies {
|
||||||
api(project(":journal-api"))
|
api(project(":journal-api"))
|
||||||
|
api(project(":journal-inmemory"))
|
||||||
api(project(":reflection-api"))
|
api(project(":reflection-api"))
|
||||||
api(project(":context-api"))
|
api(project(":context-api"))
|
||||||
}
|
}
|
||||||
|
|||||||
+4
-3
@@ -2,7 +2,8 @@ package pw.binom.agentik.storage.inmemory
|
|||||||
|
|
||||||
import pw.binom.agentik.context.ContextStore
|
import pw.binom.agentik.context.ContextStore
|
||||||
import pw.binom.agentik.journal.MutableJournalStore
|
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 pw.binom.agentik.reflection.ReflectionStore
|
||||||
import kotlin.time.Clock
|
import kotlin.time.Clock
|
||||||
|
|
||||||
@@ -21,14 +22,14 @@ import kotlin.time.Clock
|
|||||||
*/
|
*/
|
||||||
object InMemoryStorage {
|
object InMemoryStorage {
|
||||||
data class Bundle(
|
data class Bundle(
|
||||||
val conversationStore: ConversationStore,
|
val conversationStore: MutableConversationStore,
|
||||||
val messageStore: MutableJournalStore,
|
val messageStore: MutableJournalStore,
|
||||||
val workingMemoryStore: ContextStore,
|
val workingMemoryStore: ContextStore,
|
||||||
val reflectionStore: ReflectionStore,
|
val reflectionStore: ReflectionStore,
|
||||||
)
|
)
|
||||||
|
|
||||||
fun create(clock: Clock = Clock.System): Bundle = Bundle(
|
fun create(clock: Clock = Clock.System): Bundle = Bundle(
|
||||||
conversationStore = InMemoryConversationStore(clock),
|
conversationStore = InMemoryMutableConversationStore(clock),
|
||||||
messageStore = InMemoryMessageStore(),
|
messageStore = InMemoryMessageStore(),
|
||||||
workingMemoryStore = InMemoryWorkingMemoryStore(),
|
workingMemoryStore = InMemoryWorkingMemoryStore(),
|
||||||
reflectionStore = InMemoryReflectionStore(),
|
reflectionStore = InMemoryReflectionStore(),
|
||||||
|
|||||||
+4
-4
@@ -3,7 +3,7 @@ package pw.binom.agentik.storage.ksqlite
|
|||||||
import kotlin.time.Clock
|
import kotlin.time.Clock
|
||||||
import kotlin.time.Instant
|
import kotlin.time.Instant
|
||||||
import pw.binom.agentik.journal.ConversationRecord
|
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 pw.binom.db.ksqlite.SQLiteConnection
|
||||||
import kotlinx.coroutines.Dispatchers
|
import kotlinx.coroutines.Dispatchers
|
||||||
import kotlinx.coroutines.sync.Mutex
|
import kotlinx.coroutines.sync.Mutex
|
||||||
@@ -11,15 +11,15 @@ import kotlinx.coroutines.sync.withLock
|
|||||||
import kotlinx.coroutines.withContext
|
import kotlinx.coroutines.withContext
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* ksqlite-реализация [ConversationStore]. Схема таблицы `conversation` живёт
|
* ksqlite-реализация [MutableConversationStore]. Схема таблицы `conversation` живёт
|
||||||
* в [Schema] (миграция через PRAGMA user_version) — этот класс только
|
* в [Schema] (миграция через PRAGMA user_version) — этот класс только
|
||||||
* готовит и выполняет SQL, ссылаясь на `Schema.COL_*` / `Schema.TABLE_*`.
|
* готовит и выполняет SQL, ссылаясь на `Schema.COL_*` / `Schema.TABLE_*`.
|
||||||
*/
|
*/
|
||||||
class KsqliteConversationStore(
|
class KsqliteMutableConversationStore(
|
||||||
private val connection: SQLiteConnection,
|
private val connection: SQLiteConnection,
|
||||||
private val messageStore: KsqliteMessageStore? = null,
|
private val messageStore: KsqliteMessageStore? = null,
|
||||||
private val workingMemoryStore: KsqliteWorkingMemoryStore? = null,
|
private val workingMemoryStore: KsqliteWorkingMemoryStore? = null,
|
||||||
) : ConversationStore {
|
) : MutableConversationStore {
|
||||||
|
|
||||||
private val mutex = Mutex()
|
private val mutex = Mutex()
|
||||||
|
|
||||||
+3
-3
@@ -2,7 +2,7 @@ package pw.binom.agentik.storage.ksqlite
|
|||||||
|
|
||||||
import pw.binom.agentik.context.ContextStore
|
import pw.binom.agentik.context.ContextStore
|
||||||
import pw.binom.agentik.journal.MutableJournalStore
|
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.agentik.reflection.ReflectionStore
|
||||||
import pw.binom.db.ksqlite.SQLiteConnection
|
import pw.binom.db.ksqlite.SQLiteConnection
|
||||||
|
|
||||||
@@ -17,7 +17,7 @@ import pw.binom.db.ksqlite.SQLiteConnection
|
|||||||
*/
|
*/
|
||||||
class KsqliteStores internal constructor(
|
class KsqliteStores internal constructor(
|
||||||
val connection: SQLiteConnection,
|
val connection: SQLiteConnection,
|
||||||
val conversations: ConversationStore,
|
val conversations: MutableConversationStore,
|
||||||
val messages: MutableJournalStore,
|
val messages: MutableJournalStore,
|
||||||
val workingMemory: ContextStore,
|
val workingMemory: ContextStore,
|
||||||
val reflections: ReflectionStore,
|
val reflections: ReflectionStore,
|
||||||
@@ -51,7 +51,7 @@ class KsqliteStores internal constructor(
|
|||||||
val working = KsqliteWorkingMemoryStore(conn)
|
val working = KsqliteWorkingMemoryStore(conn)
|
||||||
return KsqliteStores(
|
return KsqliteStores(
|
||||||
connection = conn,
|
connection = conn,
|
||||||
conversations = KsqliteConversationStore(conn, messages, working),
|
conversations = KsqliteMutableConversationStore(conn, messages, working),
|
||||||
messages = messages,
|
messages = messages,
|
||||||
workingMemory = working,
|
workingMemory = working,
|
||||||
reflections = KsqliteReflectionStore(conn),
|
reflections = KsqliteReflectionStore(conn),
|
||||||
|
|||||||
+1
-1
@@ -12,7 +12,7 @@ import kotlin.test.assertTrue
|
|||||||
import kotlin.time.Duration
|
import kotlin.time.Duration
|
||||||
import kotlin.time.Instant
|
import kotlin.time.Instant
|
||||||
|
|
||||||
class KsqliteConversationStoreTest {
|
class KsqliteMutableConversationStoreTest {
|
||||||
private lateinit var stores: KsqliteStores
|
private lateinit var stores: KsqliteStores
|
||||||
|
|
||||||
@BeforeTest
|
@BeforeTest
|
||||||
+1
-1
@@ -79,7 +79,7 @@ class SchemaMigrationTest {
|
|||||||
// constructor internal, тест в том же модуле и может его звать.
|
// constructor internal, тест в том же модуле и может его звать.
|
||||||
val stores = KsqliteStores(
|
val stores = KsqliteStores(
|
||||||
connection = conn,
|
connection = conn,
|
||||||
conversations = KsqliteConversationStore(conn),
|
conversations = KsqliteMutableConversationStore(conn),
|
||||||
messages = KsqliteMessageStore(conn),
|
messages = KsqliteMessageStore(conn),
|
||||||
workingMemory = KsqliteWorkingMemoryStore(conn),
|
workingMemory = KsqliteWorkingMemoryStore(conn),
|
||||||
reflections = KsqliteReflectionStore(conn),
|
reflections = KsqliteReflectionStore(conn),
|
||||||
|
|||||||
Reference in New Issue
Block a user