Compare commits
6 Commits
7865eed836
...
11
| Author | SHA1 | Date | |
|---|---|---|---|
| 639c7d1748 | |||
| 5f0e0da361 | |||
| acb4ee6186 | |||
| c0a933d251 | |||
| 29851c047a | |||
| c9995b263e |
@@ -27,9 +27,10 @@ class SendSubcommand : AgentikSubcommand("send", "Отправить user-ход
|
||||
// Подписываемся на поток событий ДО send: события, отправленные
|
||||
// до подписки, не реплеятся (shared-flow без replay).
|
||||
val eventsJob = launch {
|
||||
conv.events(Instant.DISTANT_PAST)
|
||||
agent.outbox.conversationEvents(Instant.DISTANT_PAST, conv.id)
|
||||
// onEach печатает и терминальный event, takeWhile лишь
|
||||
// завершает сбор после него.
|
||||
.map { it.event }
|
||||
.onEach { ev -> emit(ev) }
|
||||
.takeWhile { ev -> !isTerminal(ev) }
|
||||
.collect { }
|
||||
@@ -53,7 +54,7 @@ class SendSubcommand : AgentikSubcommand("send", "Отправить user-ход
|
||||
is Event.AppendText -> println("event AppendText ${escape(ev.body)}")
|
||||
is Event.AppendImage -> println("event AppendImage <${ev.body.size}B ${ev.mime}>")
|
||||
is Event.ToolCall -> println("event ToolCall ${ev.id} ${ev.toolName} ${escape(ev.toolArgs)}")
|
||||
is Event.ToolResult -> println("event ToolResult ${ev.id} ${escape(ev.result ?: "")}")
|
||||
is Event.ToolResult -> println("event ToolResult ${ev.toolCallId} ${escape(ev.result ?: "")}")
|
||||
is Event.End -> println("event End")
|
||||
is Event.Interrupted -> println("event Interrupted")
|
||||
is Event.Error -> println("event Error ${ev.code ?: ""} ${escape(ev.message)}")
|
||||
|
||||
@@ -7,7 +7,7 @@ import kotlinx.coroutines.launch
|
||||
import pw.binom.agentik.proto.Agent
|
||||
import pw.binom.agentik.proto.Content
|
||||
import pw.binom.agentik.proto.Conversation
|
||||
import pw.binom.agentik.proto.Event
|
||||
import pw.binom.agentik.outbox.Event
|
||||
import kotlin.coroutines.CoroutineContext
|
||||
import kotlin.time.Instant
|
||||
|
||||
@@ -92,12 +92,12 @@ internal class TuiBackend(
|
||||
}
|
||||
|
||||
/**
|
||||
* Подписывается на [Conversation.events] и перенаправляет их в [state].
|
||||
* Подписывается на `outbox.conversationEvents(after, conv.id)` и перенаправляет их в [state].
|
||||
*/
|
||||
private fun subscribeEvents(conv: Conversation, from: Instant) {
|
||||
eventsJob?.cancel()
|
||||
eventsJob = scope.launch {
|
||||
conv.events(from).collect { ev -> dispatch(ev) }
|
||||
agent.outbox.conversationEvents(from, conv.id).collect { ce -> dispatch(ce.event) }
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -3,12 +3,13 @@ package pw.binom.agentik.tui
|
||||
import kotlinx.coroutines.flow.Flow
|
||||
import kotlinx.coroutines.flow.MutableSharedFlow
|
||||
import kotlinx.coroutines.flow.emptyFlow
|
||||
import pw.binom.agentik.journal.ConversationStore
|
||||
import pw.binom.agentik.journal.JournalStore
|
||||
import pw.binom.agentik.outbox.OutboxStore
|
||||
import pw.binom.agentik.proto.Agent
|
||||
import pw.binom.agentik.proto.Content
|
||||
import pw.binom.agentik.proto.Conversation
|
||||
import pw.binom.agentik.proto.Event
|
||||
import pw.binom.agentik.outbox.Event
|
||||
import pw.binom.agentik.proto.Message
|
||||
import pw.binom.agentik.proto.MessageContext
|
||||
import kotlin.time.Instant
|
||||
@@ -31,10 +32,13 @@ internal class FakeAgent(
|
||||
// emptyFlow, journal — error-on-access (никто не должен его трогать).
|
||||
override val journal: JournalStore = error("journal not used in TuiBackend tests")
|
||||
override val outbox: OutboxStore = object : OutboxStore {
|
||||
override fun events(after: Instant?) = emptyFlow<pw.binom.agentik.proto.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 fun close() {}
|
||||
}
|
||||
override val conversationStore: ConversationStore = error("conversationStore not used in TuiBackend tests")
|
||||
|
||||
override fun createConversation(temp: Boolean): Conversation {
|
||||
createCount++
|
||||
@@ -49,8 +53,7 @@ internal class FakeAgent(
|
||||
override suspend fun deleteConversation(id: String): Boolean =
|
||||
conversations.removeAll { it.id == id }
|
||||
|
||||
override suspend fun getConversations(offset: Int, limit: Int): List<Conversation> =
|
||||
conversations.toList()
|
||||
override suspend fun renameConversation(id: String, title: String?): Instant? = null
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -4,7 +4,7 @@ import kotlinx.coroutines.ExperimentalCoroutinesApi
|
||||
import kotlinx.coroutines.test.runCurrent
|
||||
import kotlinx.coroutines.test.runTest
|
||||
import pw.binom.agentik.proto.Content
|
||||
import pw.binom.agentik.proto.Event
|
||||
import pw.binom.agentik.outbox.Event
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertFalse
|
||||
@@ -170,7 +170,7 @@ class TuiBackendTest {
|
||||
runCurrent()
|
||||
val now = kotlin.time.Clock.System.now()
|
||||
conv.emit(Event.ToolCall(date = now, id = "1", title = null, toolName = "echo", toolArgs = """{"x":1}"""))
|
||||
conv.emit(Event.ToolResult(date = now, id = "1", result = "ok"))
|
||||
conv.emit(Event.ToolResult(date = now, toolCallId = "1", result = "ok"))
|
||||
runCurrent()
|
||||
|
||||
val toolMsgs = state.messages.value.filterIsInstance<TuiMessage.ToolCall>()
|
||||
|
||||
+192
-4
@@ -16,6 +16,11 @@
|
||||
- `HttpJournalStore` — `list(convId, after, offset, limit)` → `List<MessageRecord>`
|
||||
со всеми типами записей (User/Assistant/ToolCall/ToolResult/Error + tokens).
|
||||
- `HttpEventStore` — `events` / `agentEvents` / `conversationEvents` (SSE).
|
||||
- `ReconnectingOutbox(outbox, scope, policy)` — обёртка над `OutboxStore` с
|
||||
авто-reconnect при обрыве стрима (exponential backoff). Два независимых
|
||||
потока: `events()` (те же `CommonEvent`) и `connectionStatus()`
|
||||
(`Connecting`/`Connected`/`Disconnected`/`Failed`) — статус НЕ мешается
|
||||
с основным потоком событий. См. ниже.
|
||||
|
||||
`Agent` — `AutoCloseable`; `agent.close()` закрывает HttpClient. Не нужно
|
||||
вручную создавать `HttpClient` и накатывать на него JSON/Bearer-плагины.
|
||||
@@ -36,6 +41,30 @@ dependencies {
|
||||
}
|
||||
```
|
||||
|
||||
## Что клиент хранит локально (persistence)
|
||||
|
||||
Либа **не** имеет `SettingsRepository` / `Config` — это намеренно: UI-фреймворки хранят настройки по-разному (JSON-файл, Keychain, Android DataStore, NSUserDefaults, ...). Либа не навязывает формат, но клиент должен сериализовать у себя минимум:
|
||||
|
||||
| Поле | Что это | Где взять |
|
||||
|---|---|---|
|
||||
| `clientId` (параметр `id` в `AgentikAgent`) | Идентичность клиента в логах сервера (X-Client-Id header). Не user-id в агенте, не device-id — это **произвольная строка клиента**, обычно `<app-name>-<installation-uuid>`. Сервер использует для log multiplexing и не интерпретирует. | Генерируется один раз при первом запуске (`UUID.randomUUID().toString()`) и сохраняется. Никогда не меняется. |
|
||||
| `baseUrl` | URL сервера (`http://host:8080/agentik`). Должен включать path-prefix фасада, не только хост. | Из настроек пользователя / дефолт |
|
||||
| `token` | Bearer-токен. `null` = анонимный доступ (если сервер разрешает). | Из настроек пользователя / secure-storage |
|
||||
|
||||
Опционально (для UX): `engineFactory` — обычно compile-time выбор по платформе (`CIO` JVM/Native, `OkHttp` JVM, `Darwin` iOS/macOS).
|
||||
|
||||
Минимальный JSON для UI, который хранит в файле:
|
||||
|
||||
```json
|
||||
{
|
||||
"clientId": "my-android-app-550e8400-e29b-41d4-a716-446655440000",
|
||||
"baseUrl": "https://agent.example.com/agentik",
|
||||
"token": "s3cret"
|
||||
}
|
||||
```
|
||||
|
||||
⚠️ `clientId` **генерируется один раз** при установке и больше не меняется — иначе сломается log multiplexing на сервере.
|
||||
|
||||
## Быстрый старт: свой клиент за 5 минут
|
||||
|
||||
Один self-contained пример: создаём агента, открываем диалог,
|
||||
@@ -256,6 +285,55 @@ session.scope.launch {
|
||||
`rec is MessageRecord.UserMessage` для реплик пользователя,
|
||||
`rec is MessageRecord.ToolCall` для отрисовки tool-call баббла, и т.п.
|
||||
|
||||
## Кэш списка бесед
|
||||
|
||||
`agent.conversationStore` — read-only view поверх `conversation`-таблицы
|
||||
на сервере (`ConversationRecord` = id / title / isTemporal / createdAt /
|
||||
updatedAt, без `Conversation` handle и без флагов image-support).
|
||||
|
||||
**Сценарий клиента:** показать список диалогов («как в Telegram»), чтобы
|
||||
при открытии UI уже знал названия, не дёргал сервер лишний раз, и
|
||||
моментально реагировал на создание/удаление/переименование в другой
|
||||
вкладке.
|
||||
|
||||
Подход — тот же **«remote → local snapshot + live-events»**:
|
||||
|
||||
```kotlin
|
||||
import pw.binom.agentik.client.AgentikAgent
|
||||
import pw.binom.agentik.journal.ConversationRecord
|
||||
import pw.binom.agentik.outbox.AgentEvent
|
||||
import io.ktor.client.engine.cio.CIO
|
||||
|
||||
// `AgentikAgent` сам оборачивает HTTP-store в локальный кэш:
|
||||
// remote.listFlow → local.upsert (snapshot)
|
||||
// outbox.agentEvents → local.upsert / delete (live)
|
||||
val agent = AgentikAgent(
|
||||
id = "agentik",
|
||||
baseUrl = "http://localhost:8080/agentik",
|
||||
engineFactory = CIO,
|
||||
)
|
||||
|
||||
// Кэш уже наполняется в фоне, читать можно сразу:
|
||||
val all = agent.conversationStore.list(0, Int.MAX_VALUE)
|
||||
all.forEach { rec -> println("${rec.id} ${rec.title ?: "(no title)"} ${rec.updatedAt}") }
|
||||
|
||||
// И наблюдать live-изменения (Created/Deleted/Renamed/Touched)
|
||||
agent.outbox.agentEvents(kotlin.time.Instant.DISTANT_PAST).collect { ev ->
|
||||
when (ev) {
|
||||
is AgentEvent.Created -> println("+ ${ev.conversationId}")
|
||||
is AgentEvent.Renamed -> println("~ ${ev.id} → ${ev.title}")
|
||||
is AgentEvent.Touched -> println("↻ ${ev.id} (${ev.updatedAt})")
|
||||
is AgentEvent.Deleted -> println("- ${ev.id}")
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
Если ты **не хочешь** встроенный кэш (например, тебе нужен прямой HTTP
|
||||
для бэкенда-сервиса) — `agentikHttpClient(...).raw` оставлен как
|
||||
escape-hatch. Сам `InMemoryMutableConversationStore` тоже доступен —
|
||||
подмени его на свою реализацию через `wrapWithLocalConversationCache`,
|
||||
если нужен SQLite/JSON-store.
|
||||
|
||||
## Стриминг live-ответа
|
||||
|
||||
Для streaming-рендера текущего хода подписывайся на `events()` и
|
||||
@@ -308,13 +386,81 @@ UI-обновление списка — отдельная задача, реш
|
||||
|
||||
- **UI-рендеринг** — это твоя зона (Compose/HTML/etc.), `:client` только
|
||||
отдаёт типы и потоки.
|
||||
- **Персистентность кэша** — `InMemoryJournalStore` хранит в RAM. Для
|
||||
диска пиши свой `MutableJournalStore` (см. `KsqliteJournalStore` в
|
||||
`:journal-ksqlite` как образец).
|
||||
- **Персистентность кэша** — `InMemoryJournalStore` и
|
||||
`InMemoryMutableConversationStore` хранят в RAM. Для диска пиши свой
|
||||
`MutableJournalStore` / `MutableConversationStore` (см. `KsqliteJournalStore`
|
||||
в `:journal-ksqlite` как образец).
|
||||
- **Нестандартные движковые настройки** — для `requestTimeout`,
|
||||
прокси и т.п. используй `agentikHttpClient(engineFactory, token)`
|
||||
напрямую.
|
||||
|
||||
## Кэш списка бесед
|
||||
|
||||
`agent.conversationStore`, который видит клиент — это **локальный кэш**,
|
||||
а не прямой HTTP. Внутри `AgentikAgent` (в `wrapWithLocalConversationCache`)
|
||||
лежит `InMemoryMutableConversationStore`, синхронизированный с сервером:
|
||||
|
||||
1. **Seed при старте**: один snapshot через `remote.listFlow(0)` → заливаем
|
||||
в `localStore.upsert(...)`.
|
||||
2. **Live-обновления**: подписка на `outbox.agentEvents(after)`:
|
||||
- `Created(id)` → `remote.get(id)` → `local.upsert(record)`
|
||||
- `Deleted(id)` → `local.delete(id)`
|
||||
- `Renamed(id, title)` → `local.rename(id, title)`
|
||||
- `Touched(id, updatedAt)` → `local.touch(id, updatedAt)`
|
||||
|
||||
UI читает `agent.conversationStore.list(0, PAGE_SIZE)` — мгновенно, без
|
||||
HTTP, в т.ч. оффлайн. Список бесед всегда свежий: сервер эмитит
|
||||
`AgentEvent.Created` / `Deleted` / `Renamed` / `Touched` в свой outbox,
|
||||
клиент видит их через SSE и применяет к локальной копии.
|
||||
|
||||
**Команды** (создать / переименовать / удалить) идут через `agent`:
|
||||
|
||||
```kotlin
|
||||
// Создать новую беседу:
|
||||
val conv = agent.createConversation(temp = false) // → POST /conversations
|
||||
// → server эмитит Created
|
||||
// → client cache получает Created
|
||||
// → UI увидит её в списке
|
||||
// Переименовать:
|
||||
agent.renameConversation(conv.id, "Новый заголовок") // → PATCH /conversations/{id}
|
||||
// → server эмитит Renamed
|
||||
// → client cache обновляет title
|
||||
// Удалить:
|
||||
agent.deleteConversation(conv.id) // → DELETE /conversations/{id}
|
||||
// → server эмитит Deleted
|
||||
// → client cache удаляет запись
|
||||
```
|
||||
|
||||
`conversationStore` доступен **только для чтения**. Это read-only projection
|
||||
на серверную таблицу `conversation` (id + title + timestamps). Для активной
|
||||
работы (send / interrupt) получай handle через `agent.getConversation(id)`.
|
||||
|
||||
**Никогда не пиши в `conversationStore` напрямую.** Все модификации —
|
||||
командами `agent.createConversation / deleteConversation / renameConversation`.
|
||||
|
||||
### Если хочется своего cache-импла
|
||||
|
||||
`InMemoryMutableConversationStore` подходит для 99% случаев — Map +
|
||||
Mutex, KMP, тесты зелёные. Если нужен диск (cold-start восстановление
|
||||
после перезапуска) — реализуй свой `MutableConversationStore` поверх
|
||||
SQLite/Room/Core Data, см. `KsqliteMutableConversationStore` в
|
||||
`:storage-ksqlite` как образец.
|
||||
|
||||
```kotlin
|
||||
import pw.binom.agentik.journal.MutableConversationStore
|
||||
import pw.binom.agentik.journal.ConversationRecord
|
||||
|
||||
class MySqliteConversationStore(db: MyDb) : MutableConversationStore {
|
||||
override suspend fun upsert(record: ConversationRecord) { /* INSERT OR REPLACE */ }
|
||||
override suspend fun get(id: String): ConversationRecord? { /* SELECT */ }
|
||||
override suspend fun list(offset: Int, limit: Int): List<ConversationRecord> { /* SELECT ORDER BY updatedAt DESC */ }
|
||||
override suspend fun delete(id: String): Boolean { /* DELETE */ }
|
||||
override suspend fun rename(id: String, title: String?): Instant? { /* UPDATE + bump updatedAt */ }
|
||||
override suspend fun touch(id: String, now: Instant) { /* UPDATE updatedAt */ }
|
||||
override fun close() {}
|
||||
}
|
||||
```
|
||||
|
||||
## Тесты
|
||||
|
||||
```
|
||||
@@ -322,7 +468,49 @@ UI-обновление списка — отдельная задача, реш
|
||||
```
|
||||
|
||||
Покрывают: JSON-парсинг `Event`-ов, SSE-стрим, recovery после разрыва,
|
||||
401/404.
|
||||
401/404, reconnect-cycle `ReconnectingOutbox` (4 кейса: успех / обрыв +
|
||||
reconnect / exhausted attempts → Failed / close → cancel).
|
||||
|
||||
## Auto-reconnect для живого outbox
|
||||
|
||||
Базовый `OutboxStore.events(after)` — cold SSE-стрим, при обрыве (мобильная
|
||||
сеть, рестарт сервера) клиент сам должен реконнектиться с `after = lastEventDate`.
|
||||
Это повторяется в каждом клиенте. `ReconnectingOutbox` берёт это на себя:
|
||||
|
||||
```kotlin
|
||||
val recon = ReconnectingOutbox(
|
||||
outbox = agent.outbox, // или HttpEventStore
|
||||
scope = myScreenScope,
|
||||
policy = BackoffPolicy.Default, // 1s → 2s → ... → 30s, ±20% jitter
|
||||
)
|
||||
|
||||
scope.launch { recon.events(Instant.DISTANT_PAST).collect { handle(it) } }
|
||||
scope.launch {
|
||||
recon.connectionStatus().collect { status ->
|
||||
when (status) {
|
||||
is Connecting -> ui.showBanner("connecting...")
|
||||
is Connected -> ui.hideBanner()
|
||||
is Disconnected -> ui.showBanner("reconnecting in ${status.willRetryIn}…")
|
||||
is Failed -> ui.showError(status.cause)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// На выходе (например, navigation back):
|
||||
recon.close() // отменяет background-loop, потоки терминируются
|
||||
```
|
||||
|
||||
Два потока **независимы** — `events()` содержит только `CommonEvent`,
|
||||
`connectionStatus()` содержит только `ConnectionStatus`. Никакого
|
||||
"мешающего" `Connecting`/`Disconnected` в потоке событий.
|
||||
|
||||
Параметры backoff (см. `BackoffPolicy`):
|
||||
- `initial` / `max` — границы задержки
|
||||
- `multiplier` — множитель на каждом шаге
|
||||
- `jitter` — рандом-разброс (по умолчанию 20%)
|
||||
- `maxAttempts` — лимит попыток; после — `Failed` + закрытие потока
|
||||
|
||||
Если нужен фиксированный delay для тестов — `BackoffPolicy.Fixed(10.milliseconds, attempts = 3)`.
|
||||
|
||||
## Известное ограничение
|
||||
|
||||
|
||||
@@ -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`).
|
||||
@@ -20,24 +33,141 @@ import pw.binom.agentik.proto.Agent
|
||||
* )
|
||||
* val conv = agent.createConversation(temp = false)
|
||||
* conv.send(listOf(Content.Text("hi")))
|
||||
* conv.events(Instant.DISTANT_PAST).collect { ... }
|
||||
* agent.close() // закрывает HttpClient
|
||||
* agent.outbox.conversationEvents(Instant.DISTANT_PAST, conv.id)
|
||||
* .map { it.event }
|
||||
* .collect { ... }
|
||||
* agent.close() // закрывает HttpClient + локальный кэш
|
||||
* ```
|
||||
*
|
||||
* [id] пробрасывается в `Agent.id` — сервер про идентичность агента не знает,
|
||||
* поэтому клиент должен её знать сам (или взять из конфига).
|
||||
* ## Что клиент должен хранить локально (persistence)
|
||||
*
|
||||
* Либа **не** имеет `SettingsRepository` / `Config` — это намеренно:
|
||||
* UI-фреймворки хранят настройки по-разному (JSON-файл, Keychain,
|
||||
* `SharedPreferences`, Android DataStore, NSUserDefaults, ...). Либа
|
||||
* не навязывает формат, но вот минимальный набор, который клиент должен
|
||||
* сериализовать у себя, чтобы пережить перезапуск:
|
||||
*
|
||||
* | Поле | Что это | Где взять |
|
||||
* |---|---|---|
|
||||
* | `id` | Идентичность клиента в логах сервера (X-Client-Id header). Не user-id в агенте, не device-id, а произвольная строка клиента — обычно `<app-name>-<installation-uuid>`. Сервер использует для log multiplexing и не интерпретирует. | Генерируется клиентом при первом запуске, сохраняется локально |
|
||||
* | `baseUrl` | URL сервера (`http://host:8080/agentik`). Должен включать path-prefix фасада, не только хост. | Из настроек пользователя / дефолт |
|
||||
* | `token` | Bearer-токен. `null` = анонимный доступ (если сервер разрешает). | Из настроек пользователя / secure-storage |
|
||||
*
|
||||
* Опционально (для UX):
|
||||
* | Поле | Зачем |
|
||||
* |---|---|
|
||||
* | `engineFactory` | Зависит от платформы (`CIO` JVM/Native, `OkHttp` JVM, `Darwin` iOS/macOS). Выбор — обычно compile-time. |
|
||||
*
|
||||
* Пример минимального persistence-файла (для UI, который хранит JSON):
|
||||
*
|
||||
* ```json
|
||||
* {
|
||||
* "clientId": "my-android-app-550e8400-e29b-41d4-a716-446655440000",
|
||||
* "baseUrl": "https://agent.example.com/agentik",
|
||||
* "token": "s3cret"
|
||||
* }
|
||||
* ```
|
||||
*
|
||||
* `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()
|
||||
}
|
||||
}
|
||||
|
||||
@@ -6,18 +6,12 @@ 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.prepareGet
|
||||
import io.ktor.client.request.setBody
|
||||
import io.ktor.client.statement.bodyAsChannel
|
||||
import io.ktor.http.ContentType
|
||||
import io.ktor.http.HttpStatusCode
|
||||
import io.ktor.http.contentType
|
||||
import kotlinx.coroutines.flow.Flow
|
||||
import kotlinx.coroutines.flow.flow
|
||||
import kotlinx.serialization.Serializable
|
||||
import pw.binom.agentik.proto.Content
|
||||
import pw.binom.agentik.proto.Conversation
|
||||
import pw.binom.agentik.proto.Event
|
||||
import pw.binom.agentik.proto.Message
|
||||
import pw.binom.agentik.proto.MessageContext
|
||||
import kotlin.time.Instant
|
||||
@@ -72,22 +66,6 @@ internal class ConversationClient(
|
||||
httpClient.post("$convUrl/interrupt")
|
||||
}
|
||||
|
||||
override fun events(after: Instant): Flow<Event> = flow {
|
||||
// prepareGet + execute (а не get) обязателен: `get` дожидается полного
|
||||
// тела ответа, а SSE-поток не заканчивается никогда — вызов висел бы
|
||||
// вечно. `execute` отдаёт HttpResponse со стриминговым bodyAsChannel.
|
||||
httpClient.prepareGet("$convUrl/events?after=$after") { noSseReadTimeout() }
|
||||
.execute { response ->
|
||||
check(response.status == HttpStatusCode.OK) {
|
||||
"events: server returned ${response.status}"
|
||||
}
|
||||
readSse(response.bodyAsChannel())
|
||||
.collect { payload ->
|
||||
emit(agentikJson.decodeFromString(Event.serializer(), payload))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
override suspend fun getMessages(after: Instant, offset: Int, limit: Int): List<Message> =
|
||||
httpClient.get("$convUrl/messages") {
|
||||
parameter("after", after.toString())
|
||||
|
||||
@@ -23,4 +23,4 @@ data class ConversationSnapshot(
|
||||
internal data class RequestCreateConversation(val temp: Boolean)
|
||||
|
||||
@Serializable
|
||||
internal data class RequestRename(val title: String)
|
||||
internal data class RequestRename(val title: String?)
|
||||
|
||||
@@ -0,0 +1,58 @@
|
||||
package pw.binom.agentik.client
|
||||
|
||||
import io.ktor.client.HttpClient
|
||||
import io.ktor.client.call.body
|
||||
import io.ktor.client.request.get
|
||||
import io.ktor.client.request.parameter
|
||||
import io.ktor.http.HttpStatusCode
|
||||
import pw.binom.agentik.journal.ConversationRecord
|
||||
import pw.binom.agentik.journal.ConversationStore
|
||||
|
||||
/**
|
||||
* HTTP-реализация [ConversationStore] (read-only metadata view),
|
||||
* ходящая в `:server`-фасад.
|
||||
*
|
||||
* **Endpoint**: `GET {baseUrl}/conversations?offset=&limit=` —
|
||||
* возвращает `List<ConversationRecord>` (id, title, isTemporal, createdAt,
|
||||
* updatedAt) БЕЗ handle'ов и image-support флагов (это лёгкая проекция
|
||||
* для UI-списка; handle берётся через `agent.getConversation(id)`).
|
||||
*
|
||||
* **Read-only**: запись в `conversation` table — только через команды
|
||||
* `agent.createConversation / deleteConversation / renameConversation`.
|
||||
*
|
||||
* Клиентский кэш строится композицией `HttpConversationStore` (snapshot)
|
||||
* + `agent.outbox.agentEvents(after)` (live deltas: Created/Deleted/
|
||||
* Renamed/Touched) — см. `client/README.md` секция
|
||||
* «Кэш списка бесед».
|
||||
*/
|
||||
internal class HttpConversationStore(
|
||||
private val httpClient: HttpClient,
|
||||
private val baseUrl: String,
|
||||
) : ConversationStore {
|
||||
|
||||
private val agentUrl: String = baseUrl.trimEnd('/')
|
||||
|
||||
override suspend fun get(id: String): ConversationRecord? {
|
||||
val response = httpClient.get("$agentUrl/conversations/$id")
|
||||
if (response.status == HttpStatusCode.NotFound) return null
|
||||
check(response.status == HttpStatusCode.OK) {
|
||||
"conversationStore.get($id): server returned ${response.status}"
|
||||
}
|
||||
return response.body<ConversationRecord>()
|
||||
}
|
||||
|
||||
override suspend fun list(offset: Int, limit: Int): List<ConversationRecord> {
|
||||
val response = httpClient.get("$agentUrl/conversations") {
|
||||
parameter("offset", offset)
|
||||
parameter("limit", limit)
|
||||
}
|
||||
check(response.status == HttpStatusCode.OK) {
|
||||
"conversationStore.list: server returned ${response.status}"
|
||||
}
|
||||
return response.body<List<ConversationRecord>>()
|
||||
}
|
||||
|
||||
override fun close() {
|
||||
// HttpClient закрывает владелец (AgentClient / AgentikAgent).
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,242 @@
|
||||
package pw.binom.agentik.client
|
||||
|
||||
import kotlinx.coroutines.CancellationException
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.Job
|
||||
import kotlinx.coroutines.channels.BufferOverflow
|
||||
import kotlinx.coroutines.currentCoroutineContext
|
||||
import kotlinx.coroutines.delay
|
||||
import kotlinx.coroutines.flow.Flow
|
||||
import kotlinx.coroutines.flow.MutableSharedFlow
|
||||
import kotlinx.coroutines.isActive
|
||||
import kotlinx.coroutines.launch
|
||||
import pw.binom.agentik.outbox.CommonEvent
|
||||
import pw.binom.agentik.outbox.OutboxStore
|
||||
import kotlin.concurrent.atomics.AtomicBoolean
|
||||
import kotlin.concurrent.atomics.AtomicReference
|
||||
import kotlin.concurrent.atomics.ExperimentalAtomicApi
|
||||
import kotlin.math.min
|
||||
import kotlin.math.pow
|
||||
import kotlin.random.Random
|
||||
import kotlin.time.Duration
|
||||
import kotlin.time.Duration.Companion.seconds
|
||||
import kotlin.time.Instant
|
||||
|
||||
/**
|
||||
* Состояние подключения к удалённому [OutboxStore]. Эмитится через
|
||||
* [ReconnectingOutbox.connectionStatus] — отдельным потоком, **не**
|
||||
* смешивается с [ReconnectingOutbox.events].
|
||||
*
|
||||
* Типичный цикл:
|
||||
* ```
|
||||
* Connecting(1) → Connected → ... → Disconnected(reason, retryIn) →
|
||||
* Connecting(2) → Connected → ...
|
||||
* ```
|
||||
* При полном исчерпании попыток ([BackoffPolicy.maxAttempts]) —
|
||||
* финальный [Failed].
|
||||
*/
|
||||
sealed interface ConnectionStatus {
|
||||
|
||||
/** Начата попытка подключения (включая первую — `attempt == 1`). */
|
||||
data class Connecting(val attempt: Int) : ConnectionStatus
|
||||
|
||||
/** Получен первый event с сервера после [Connecting] / [Disconnected]. */
|
||||
data class Connected(val since: Instant) : ConnectionStatus
|
||||
|
||||
/**
|
||||
* Стрим оборвался (network error, server close, таймаут). [reason] —
|
||||
* причина, `null` если штатное завершение. [willRetryIn] — через сколько
|
||||
* будет следующая попытка (`null` если [Failed]).
|
||||
*/
|
||||
data class Disconnected(
|
||||
val reason: Throwable?,
|
||||
val willRetryIn: Duration?,
|
||||
) : ConnectionStatus
|
||||
|
||||
/**
|
||||
* Все попытки исчерпаны ([BackoffPolicy.maxAttempts]). Поток [events]
|
||||
* закрывается после этого. Создатель [ReconnectingOutbox] должен
|
||||
* решить, что делать — показать ошибку пользователю, пересоздать
|
||||
* outbox и т.п.
|
||||
*/
|
||||
data class Failed(val cause: Throwable) : ConnectionStatus
|
||||
}
|
||||
|
||||
/**
|
||||
* Политика backoff для [ReconnectingOutbox]. Параметры:
|
||||
*
|
||||
* - [initial] — задержка перед первой retry-попыткой.
|
||||
* - [max] — потолок задержки (после серии умножений).
|
||||
* - [multiplier] — множитель на каждом шаге (например, `2.0` → 1s, 2s, 4s, 8s, ...).
|
||||
* - [jitter] — доля случайного разброса `[0, jitter]` от текущей задержки
|
||||
* (например, `0.2` = ±20%). Снижает thundering-herd при массовом reconnect.
|
||||
* - [maxAttempts] — лимит попыток. `Int.MAX_VALUE` = бесконечно.
|
||||
*/
|
||||
data class BackoffPolicy(
|
||||
val initial: Duration = 1.seconds,
|
||||
val max: Duration = 30.seconds,
|
||||
val multiplier: Double = 2.0,
|
||||
val maxAttempts: Int = Int.MAX_VALUE,
|
||||
val jitter: Double = 0.2,
|
||||
) {
|
||||
init {
|
||||
require(initial > Duration.ZERO) { "initial must be positive" }
|
||||
require(max >= initial) { "max must be >= initial" }
|
||||
require(multiplier >= 1.0) { "multiplier must be >= 1.0" }
|
||||
require(maxAttempts >= 1) { "maxAttempts must be >= 1" }
|
||||
require(jitter in 0.0..1.0) { "jitter must be in [0, 1]" }
|
||||
}
|
||||
|
||||
companion object {
|
||||
/** 1s → 2s → 4s → ... → 30s, jitter ±20%, бесконечные попытки. */
|
||||
val Default: BackoffPolicy = BackoffPolicy()
|
||||
|
||||
/** Только для тестов: фиксированные задержки без разброса. */
|
||||
fun Fixed(delay: Duration, attempts: Int = 3): BackoffPolicy =
|
||||
BackoffPolicy(
|
||||
initial = delay,
|
||||
max = delay,
|
||||
multiplier = 1.0,
|
||||
maxAttempts = attempts,
|
||||
jitter = 0.0,
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Обёртка над [OutboxStore] с автоматическим reconnect при обрыве стрима.
|
||||
*
|
||||
* **Два независимых потока**:
|
||||
* - [events] — `Flow<CommonEvent>`, тот же контракт что [OutboxStore.events],
|
||||
* но с автоматическим переподключением через [BackoffPolicy]. Cursor
|
||||
* (`lastSeen`) сохраняется между попытками — клиент не теряет события.
|
||||
* - [connectionStatus] — `Flow<ConnectionStatus>`, **параллельный** поток
|
||||
* lifecycle подключения. Не смешивается с [events].
|
||||
*
|
||||
* ```
|
||||
* val outbox = ReconnectingOutbox(httpEventStore, scope)
|
||||
*
|
||||
* scope.launch {
|
||||
* outbox.events(after = Instant.DISTANT_PAST).collect { e -> handle(e) }
|
||||
* }
|
||||
* scope.launch {
|
||||
* outbox.connectionStatus().collect { s -> ui.showStatus(s) }
|
||||
* }
|
||||
*
|
||||
* // На выходе:
|
||||
* outbox.close() // отменяет background-loop, эмитит Cancelled-как-Disconnected
|
||||
* ```
|
||||
*
|
||||
* Создатель передаёт свой [scope] — жизненный цикл reconnect-цикла
|
||||
* привязан к нему. Закрытие scope (или явный [close]) отменяет
|
||||
* background-loop. После [close] оба flow терминируются.
|
||||
*/
|
||||
class ReconnectingOutbox(
|
||||
private val outbox: OutboxStore,
|
||||
private val scope: CoroutineScope,
|
||||
private val policy: BackoffPolicy = BackoffPolicy.Default,
|
||||
private val random: Random = Random.Default,
|
||||
) : AutoCloseable {
|
||||
|
||||
private val _events = MutableSharedFlow<CommonEvent>(
|
||||
replay = 0,
|
||||
extraBufferCapacity = 64,
|
||||
onBufferOverflow = BufferOverflow.DROP_OLDEST,
|
||||
)
|
||||
private val _status = MutableSharedFlow<ConnectionStatus>(
|
||||
replay = 0,
|
||||
extraBufferCapacity = 64,
|
||||
onBufferOverflow = BufferOverflow.DROP_OLDEST,
|
||||
)
|
||||
|
||||
@OptIn(ExperimentalAtomicApi::class)
|
||||
private val started = AtomicBoolean(false)
|
||||
private var job: Job? = null
|
||||
|
||||
@OptIn(ExperimentalAtomicApi::class)
|
||||
private val lastSeen: AtomicReference<Instant?> = AtomicReference(null)
|
||||
|
||||
/**
|
||||
* Live-события из [outbox] с авто-reconnect. [after] — начальный курсор;
|
||||
* учитывается только при первом вызове (любом из [events] /
|
||||
* [connectionStatus]). После reconnect курсор берётся из `date`
|
||||
* последнего виденного события.
|
||||
*
|
||||
* Коллекторы независимы — каждый получает свою копию потока (shared).
|
||||
* Медленный коллектор может пропускать события при переполнении буфера
|
||||
* (`DROP_OLDEST`).
|
||||
*/
|
||||
fun events(after: Instant? = null): Flow<CommonEvent> {
|
||||
ensureStarted(after)
|
||||
return _events
|
||||
}
|
||||
|
||||
/**
|
||||
* Lifecycle подключения: [ConnectionStatus.Connecting] /
|
||||
* [ConnectionStatus.Connected] / [ConnectionStatus.Disconnected] /
|
||||
* [ConnectionStatus.Failed]. **Не смешивается** с [events] — это
|
||||
* отдельный поток для UI-индикации статуса сети.
|
||||
*/
|
||||
fun connectionStatus(): Flow<ConnectionStatus> {
|
||||
ensureStarted(null)
|
||||
return _status
|
||||
}
|
||||
|
||||
@OptIn(ExperimentalAtomicApi::class)
|
||||
private fun ensureStarted(initialCursor: Instant?) {
|
||||
if (!started.compareAndSet(false, true)) return
|
||||
lastSeen.store(initialCursor)
|
||||
job = scope.launch { runLoop() }
|
||||
}
|
||||
|
||||
@OptIn(ExperimentalAtomicApi::class)
|
||||
private suspend fun runLoop() {
|
||||
var attempt = 0
|
||||
var connected = false
|
||||
while (currentCoroutineContext().isActive) {
|
||||
attempt++
|
||||
_status.emit(ConnectionStatus.Connecting(attempt))
|
||||
val error: Throwable? = try {
|
||||
outbox.events(after = lastSeen.load()).collect { event ->
|
||||
lastSeen.store(event.date)
|
||||
_events.emit(event)
|
||||
if (!connected) {
|
||||
connected = true
|
||||
_status.emit(ConnectionStatus.Connected(event.date))
|
||||
}
|
||||
}
|
||||
null
|
||||
} catch (t: CancellationException) {
|
||||
throw t
|
||||
} catch (t: Throwable) {
|
||||
t
|
||||
}
|
||||
connected = false
|
||||
if (attempt >= policy.maxAttempts) {
|
||||
_status.emit(
|
||||
ConnectionStatus.Failed(error ?: RuntimeException("outbox flow ended normally"))
|
||||
)
|
||||
return
|
||||
}
|
||||
val backoff = computeBackoff(attempt)
|
||||
_status.emit(ConnectionStatus.Disconnected(error, backoff))
|
||||
delay(backoff)
|
||||
}
|
||||
}
|
||||
|
||||
private fun computeBackoff(attempt: Int): Duration {
|
||||
// attempt 1 → initial, 2 → initial * m, 3 → initial * m^2, ...
|
||||
val base = (policy.initial.inWholeMilliseconds.toDouble() *
|
||||
policy.multiplier.pow((attempt - 1).toDouble()))
|
||||
.toLong()
|
||||
val capped = min(base, policy.max.inWholeMilliseconds)
|
||||
val jitterMs = (capped * policy.jitter * random.nextDouble()).toLong()
|
||||
val finalMs = (capped + jitterMs).coerceAtLeast(1L)
|
||||
return Duration.parse("${finalMs}ms")
|
||||
}
|
||||
|
||||
override fun close() {
|
||||
job?.cancel()
|
||||
job = null
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,211 @@
|
||||
package pw.binom.agentik.client
|
||||
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.ExperimentalCoroutinesApi
|
||||
import kotlinx.coroutines.Job
|
||||
import kotlinx.coroutines.channels.Channel
|
||||
import kotlinx.coroutines.flow.Flow
|
||||
import kotlinx.coroutines.flow.emptyFlow
|
||||
import kotlinx.coroutines.flow.flow
|
||||
import kotlinx.coroutines.launch
|
||||
import kotlinx.coroutines.test.advanceTimeBy
|
||||
import kotlinx.coroutines.test.runCurrent
|
||||
import kotlinx.coroutines.test.runTest
|
||||
import pw.binom.agentik.outbox.AgentEvent
|
||||
import pw.binom.agentik.outbox.CommonEvent
|
||||
import pw.binom.agentik.outbox.OutboxStore
|
||||
import pw.binom.agentik.outbox.Event
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertNotNull
|
||||
import kotlin.test.assertTrue
|
||||
import kotlin.time.Duration
|
||||
import kotlin.time.Instant
|
||||
|
||||
/**
|
||||
* In-memory [OutboxStore] для unit-тестов [ReconnectingOutbox].
|
||||
*
|
||||
* Управление:
|
||||
* - [push] — кладёт [CommonEvent] в очередь, флоу доставит.
|
||||
* - [throwAtNextEvent] — следующий «тик» `events(after)` бросит этот Throwable
|
||||
* (симулирует network error / stream break).
|
||||
*
|
||||
* Сигнатура [events] идентична боевой — её можно подменить боевым
|
||||
* `HttpEventStore`, контракт один и тот же.
|
||||
*/
|
||||
internal class FakeOutbox : OutboxStore {
|
||||
private sealed interface Msg {
|
||||
data class Ev(val event: CommonEvent) : Msg
|
||||
data class Err(val throwable: Throwable) : Msg
|
||||
}
|
||||
|
||||
private val channel = Channel<Msg>(Channel.UNLIMITED)
|
||||
|
||||
override fun events(after: Instant?): Flow<CommonEvent> = flow {
|
||||
for (msg in channel) {
|
||||
when (msg) {
|
||||
is Msg.Err -> throw msg.throwable
|
||||
is Msg.Ev -> emit(msg.event)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fun push(event: CommonEvent) { channel.trySend(Msg.Ev(event)) }
|
||||
fun throwAtNextEvent(t: Throwable) { channel.trySend(Msg.Err(t)) }
|
||||
|
||||
override fun agentEvents(after: Instant?): Flow<CommonEvent.Agent> = emptyFlow()
|
||||
override fun conversationEvents(
|
||||
after: Instant?,
|
||||
conversationId: String?,
|
||||
): Flow<CommonEvent.Conversation> = emptyFlow()
|
||||
override suspend fun earliestEventDate(): Instant = Instant.DISTANT_PAST
|
||||
override fun close() { channel.close() }
|
||||
}
|
||||
|
||||
private fun testEvent(dateMs: Long): CommonEvent =
|
||||
CommonEvent.Conversation(
|
||||
date = Instant.fromEpochMilliseconds(dateMs),
|
||||
conversationId = "test",
|
||||
event = Event.End(date = Instant.fromEpochMilliseconds(dateMs)),
|
||||
)
|
||||
|
||||
@OptIn(ExperimentalCoroutinesApi::class)
|
||||
class ReconnectingOutboxTest {
|
||||
|
||||
@Test
|
||||
fun `first event after connect emits Connecting then Connected`() = runConnectionTest(
|
||||
attempts = 5,
|
||||
) { ctx ->
|
||||
val fake = ctx.fake
|
||||
val status = ctx.statusLog
|
||||
val events = ctx.eventsLog
|
||||
|
||||
fake.push(testEvent(1000))
|
||||
ctx.advanceAndDrain(50)
|
||||
|
||||
assertEquals(1, status.count { it is ConnectionStatus.Connecting && it.attempt == 1 })
|
||||
assertEquals(1, status.count { it is ConnectionStatus.Connected })
|
||||
assertEquals(1, events.size)
|
||||
assertEquals(Instant.fromEpochMilliseconds(1000), events[0].date)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `disconnect mid-stream triggers retry with backoff and resumes from last seen`() =
|
||||
runConnectionTest(attempts = 5) { ctx ->
|
||||
val fake = ctx.fake
|
||||
val status = ctx.statusLog
|
||||
val events = ctx.eventsLog
|
||||
|
||||
fake.push(testEvent(1000))
|
||||
ctx.advanceAndDrain(50)
|
||||
assertEquals(1, events.size)
|
||||
|
||||
// Имитируем обрыв стрима после первого события.
|
||||
fake.throwAtNextEvent(RuntimeException("simulated network error"))
|
||||
ctx.advanceAndDrain(50)
|
||||
|
||||
// После Disconnected должен прийти Connecting(2), затем Connected,
|
||||
// затем новые события без дубля предыдущего.
|
||||
val disconnectedIndex = status.indexOfFirst { it is ConnectionStatus.Disconnected }
|
||||
val connecting2Index = status.indexOfFirst {
|
||||
it is ConnectionStatus.Connecting && it.attempt == 2
|
||||
}
|
||||
assertTrue(disconnectedIndex >= 0, "no Disconnected emitted, got: $status")
|
||||
assertTrue(connecting2Index > disconnectedIndex,
|
||||
"expected Connecting(2) after Disconnected, got: $status")
|
||||
|
||||
// Push a new event with later date — cursor preserves lastSeen.
|
||||
fake.push(testEvent(2000))
|
||||
ctx.advanceAndDrain(50)
|
||||
|
||||
assertEquals(2, events.size)
|
||||
assertEquals(Instant.fromEpochMilliseconds(1000), events[0].date)
|
||||
assertEquals(Instant.fromEpochMilliseconds(2000), events[1].date)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `exhausted attempts emits Failed and closes flow`() = runConnectionTest(
|
||||
attempts = 3,
|
||||
) { ctx ->
|
||||
val fake = ctx.fake
|
||||
|
||||
// Каждая попытка connect бросает — все 3 попытки fail.
|
||||
for (i in 0 until 3) {
|
||||
fake.throwAtNextEvent(RuntimeException("server is dead #${i + 1}"))
|
||||
ctx.advanceAndDrain(50)
|
||||
}
|
||||
|
||||
val failed = ctx.statusLog.filterIsInstance<ConnectionStatus.Failed>().firstOrNull()
|
||||
assertNotNull(failed) { "expected Failed status, got: ${ctx.statusLog}" }
|
||||
assertTrue(failed.cause is RuntimeException)
|
||||
assertEquals(0, ctx.eventsLog.size)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `close cancels background loop`() = runConnectionTest(
|
||||
attempts = 5,
|
||||
) { ctx ->
|
||||
val fake = ctx.fake
|
||||
fake.push(testEvent(1000))
|
||||
ctx.advanceAndDrain(50)
|
||||
assertEquals(1, ctx.eventsLog.size)
|
||||
|
||||
ctx.recon.close()
|
||||
ctx.advanceAndDrain(100)
|
||||
|
||||
// После close запуск новых эмиссий не должен происходить.
|
||||
val beforePush = ctx.eventsLog.size
|
||||
fake.push(testEvent(2000))
|
||||
ctx.advanceAndDrain(100)
|
||||
assertEquals(beforePush, ctx.eventsLog.size)
|
||||
}
|
||||
|
||||
private data class TestCtx(
|
||||
val fake: FakeOutbox,
|
||||
val recon: ReconnectingOutbox,
|
||||
val statusLog: MutableList<ConnectionStatus>,
|
||||
val eventsLog: MutableList<CommonEvent>,
|
||||
val jobs: List<Job>,
|
||||
val scope: CoroutineScope,
|
||||
val advanceAndDrain: (Long) -> Unit,
|
||||
)
|
||||
|
||||
/**
|
||||
* Запускает [ReconnectingOutbox] с policy из `attempts` попыток по 10ms,
|
||||
* сабскрайбит на оба потока в собирающие лист, и возвращает [TestCtx]
|
||||
* с управляемым `advanceAndDrain(ms)` — прокрутить виртуальное время.
|
||||
*/
|
||||
@OptIn(ExperimentalCoroutinesApi::class)
|
||||
private fun runConnectionTest(
|
||||
attempts: Int,
|
||||
block: suspend (TestCtx) -> Unit,
|
||||
) = runTest {
|
||||
val policy = BackoffPolicy.Fixed(
|
||||
delay = Duration.parse("10ms"),
|
||||
attempts = attempts,
|
||||
)
|
||||
val fake = FakeOutbox()
|
||||
val recon = ReconnectingOutbox(
|
||||
outbox = fake,
|
||||
scope = this,
|
||||
policy = policy,
|
||||
)
|
||||
val statusLog = mutableListOf<ConnectionStatus>()
|
||||
val eventsLog = mutableListOf<CommonEvent>()
|
||||
val jobs = listOf(
|
||||
launch { recon.connectionStatus().collect { statusLog.add(it) } },
|
||||
launch { recon.events().collect { eventsLog.add(it) } },
|
||||
)
|
||||
val advanceAndDrain: (Long) -> Unit = { ms ->
|
||||
if (ms > 0) advanceTimeBy(ms)
|
||||
runCurrent()
|
||||
}
|
||||
try {
|
||||
TestCtx(fake, recon, statusLog, eventsLog, jobs, this, advanceAndDrain).also { block(it) }
|
||||
} finally {
|
||||
recon.close()
|
||||
jobs.forEach { it.cancel() }
|
||||
fake.close()
|
||||
}
|
||||
}
|
||||
}
|
||||
+1
-1
@@ -194,7 +194,7 @@ suspend fun compact(dropFromOrderIdx: Long, conversationId: String): Long
|
||||
|
||||
`compact` — атомарный «выбросить всё от `dropFromOrderIdx` и дальше, вставить новую синтетическую запись на следующий `order_idx`». Для v1 — просто `DELETE` от индекса (суммаризация появится в v2 вместе с LLM-вызовом для генерации текста).
|
||||
|
||||
### `ConversationStore`
|
||||
### `MutableConversationStore`
|
||||
|
||||
```kotlin
|
||||
suspend fun upsert(record: ConversationRecord)
|
||||
|
||||
@@ -1,27 +1,45 @@
|
||||
package pw.binom.agentik.journal
|
||||
|
||||
import kotlin.time.Instant
|
||||
import kotlinx.coroutines.flow.Flow
|
||||
import kotlinx.coroutines.flow.flow
|
||||
|
||||
/**
|
||||
* CRUD по таблице `conversation`.
|
||||
* Read-only view of the `conversation` table (CRUD-операции находятся
|
||||
* в [MutableConversationStore] и используются только внутри ChatAgent).
|
||||
*
|
||||
* Клиенты видят [ConversationStore] через [pw.binom.agentik.proto.Agent.conversationStore]
|
||||
* (по аналогии с `journal` / `outbox`) и строят свой локальный кэш:
|
||||
* - **seed** через [list] (snapshot страницы) или [listFlow] (cold-flow paging);
|
||||
* - **live-refresh** через `outbox.agentEvents()` — Created / Deleted /
|
||||
* Renamed / Touched.
|
||||
*
|
||||
* Запись в хранилище **не** делается клиентом — только команды
|
||||
* `agent.createConversation / deleteConversation / renameConversation`.
|
||||
*/
|
||||
interface ConversationStore : AutoCloseable {
|
||||
|
||||
/** Создать или обновить snapshot диалога. */
|
||||
suspend fun upsert(record: ConversationRecord)
|
||||
|
||||
/** Диалог по id, или `null`. */
|
||||
suspend fun get(id: String): ConversationRecord?
|
||||
|
||||
/** Удалить диалог (вместе с его сообщениями и working memory). */
|
||||
suspend fun delete(id: String): Boolean
|
||||
|
||||
/** Список диалогов, отсортированный по `updatedAt` DESC. */
|
||||
suspend fun list(offset: Int, limit: Int): List<ConversationRecord>
|
||||
|
||||
/** Переименовать диалог; `null` для сброса заголовка. Возвращает новый `updatedAt` или `null`, если не найден. */
|
||||
suspend fun rename(id: String, title: String?): Instant?
|
||||
|
||||
/** Обновить `updatedAt` диалога (например, после отправки сообщения). */
|
||||
suspend fun touch(id: String, now: 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
|
||||
}
|
||||
}
|
||||
|
||||
companion object {
|
||||
const val PAGE_SIZE: Int = 100
|
||||
}
|
||||
}
|
||||
|
||||
@@ -55,6 +55,14 @@ sealed interface MessageRecord {
|
||||
override val id: String,
|
||||
override val conversationId: String,
|
||||
val toolCallId: String,
|
||||
/**
|
||||
* Имя тула, денормализованное из соответствующего `MessageRecord.ToolCall.toolName`.
|
||||
* Денормализация экономна (одна строка в SQLite) и снимает с UI
|
||||
* необходимость сопоставления `toolCallId → toolName`. `null` —
|
||||
* безопасный backfill для записей до миграции или для сиротливых
|
||||
* результатов без предшествующего `ToolCall`.
|
||||
*/
|
||||
val toolName: String? = null,
|
||||
val result: String?,
|
||||
override val createdAt: Instant,
|
||||
) : MessageRecord
|
||||
|
||||
+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)
|
||||
}
|
||||
|
||||
// KMP-реализация [MutableJournalStore] на `MutableList` + `Mutex` — для
|
||||
// тестов, dev-режима, embedded-сценариев (Android core, CLI, in-process кэш
|
||||
// в клиенте) и как образец для своей реализации.
|
||||
//
|
||||
// `list` фильтрует по `conversationId`+`createdAt>after` и сортирует
|
||||
// по `createdAt ASC`. Paging — поверх отфильтрованного списка.
|
||||
// KMP-реализация [MutableJournalStore] и [MutableConversationStore] на
|
||||
// `MutableList`/`MutableMap` + `Mutex` — для тестов, dev-режима,
|
||||
// embedded-сценариев (Android core, CLI, in-process кэш в клиенте) и как
|
||||
// образец для своей реализации.
|
||||
//
|
||||
// Зависимости: только `:journal-api`. Никакого I/O — pure in-memory.
|
||||
|
||||
@@ -15,7 +13,13 @@ kotlin {
|
||||
jvmToolchain(21)
|
||||
|
||||
jvm()
|
||||
macosX64()
|
||||
macosArm64()
|
||||
iosX64()
|
||||
iosArm64()
|
||||
iosSimulatorArm64()
|
||||
linuxX64()
|
||||
linuxArm64()
|
||||
mingwX64()
|
||||
|
||||
sourceSets {
|
||||
|
||||
+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.withLock
|
||||
import kotlin.time.Clock
|
||||
import pw.binom.agentik.journal.ConversationRecord
|
||||
import pw.binom.agentik.journal.ConversationStore
|
||||
import pw.binom.agentik.journal.MutableConversationStore
|
||||
import kotlin.time.Clock
|
||||
import kotlin.time.Instant
|
||||
|
||||
/**
|
||||
* Thread-safe Map-импл [ConversationStore].
|
||||
* Thread-safe Map-импл [MutableConversationStore] для клиентских
|
||||
* in-process кэшей (и тестов/dev-режима).
|
||||
*
|
||||
* Использует `Mutex` для атомарности read-modify-write операций
|
||||
* (rename, touch) — иначе два параллельных `rename` могут потерять обновления
|
||||
* (lost-update race), что в SQLite невозможно из-за driver-level locking.
|
||||
*
|
||||
* **Сортировка**: `list()` сортирует по `updatedAt DESC`.
|
||||
*
|
||||
* **Типичный кэш-паттерн в клиенте** (см. `client/README.md`):
|
||||
* ```
|
||||
* val local = InMemoryMutableConversationStore()
|
||||
* // seed: remote.listFlow → local.upsert
|
||||
* // live-refresh: outbox.agentEvents → local.upsert/delete/rename/touch
|
||||
* // UI: local.list(0, PAGE_SIZE)
|
||||
* ```
|
||||
*/
|
||||
class InMemoryConversationStore(
|
||||
class InMemoryMutableConversationStore(
|
||||
private val clock: Clock = Clock.System,
|
||||
) : ConversationStore {
|
||||
) : MutableConversationStore {
|
||||
|
||||
private val byId: MutableMap<String, ConversationRecord> = mutableMapOf()
|
||||
private val mutex = Mutex()
|
||||
+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 kotlin.test.Test
|
||||
@@ -9,11 +9,11 @@ import kotlin.test.assertTrue
|
||||
import kotlin.time.Instant
|
||||
import kotlinx.coroutines.test.runTest
|
||||
|
||||
class InMemoryConversationStoreTest {
|
||||
class InMemoryMutableConversationStoreTest {
|
||||
|
||||
@Test
|
||||
fun `upsert and get roundtrip preserves all fields`() = runTest {
|
||||
val store = InMemoryConversationStore()
|
||||
val store = InMemoryMutableConversationStore()
|
||||
val rec = ConversationRecord(
|
||||
id = "c1",
|
||||
title = "test",
|
||||
@@ -28,13 +28,13 @@ class InMemoryConversationStoreTest {
|
||||
|
||||
@Test
|
||||
fun `get returns null for missing id`() = runTest {
|
||||
val store = InMemoryConversationStore()
|
||||
val store = InMemoryMutableConversationStore()
|
||||
assertNull(store.get("nope"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `delete removes the record and returns true`() = runTest {
|
||||
val store = InMemoryConversationStore()
|
||||
val store = InMemoryMutableConversationStore()
|
||||
store.upsert(
|
||||
ConversationRecord(
|
||||
"c1", null, false,
|
||||
@@ -50,7 +50,7 @@ class InMemoryConversationStoreTest {
|
||||
|
||||
@Test
|
||||
fun `list sorts by updatedAt DESC and respects offset+limit`() = runTest {
|
||||
val store = InMemoryConversationStore()
|
||||
val store = InMemoryMutableConversationStore()
|
||||
val t0 = Instant.parse("2026-09-15T10:00:00Z")
|
||||
store.upsert(ConversationRecord("c1", null, false, t0, t0))
|
||||
store.upsert(ConversationRecord("c2", null, false, t0, t0.plus(kotlin.time.Duration.parse("PT60S"))))
|
||||
@@ -69,7 +69,7 @@ class InMemoryConversationStoreTest {
|
||||
|
||||
@Test
|
||||
fun `rename updates title and updatedAt returns new updatedAt`() = runTest {
|
||||
val store = InMemoryConversationStore()
|
||||
val store = InMemoryMutableConversationStore()
|
||||
val t0 = Instant.parse("2026-09-15T10:00:00Z")
|
||||
store.upsert(ConversationRecord("c1", null, false, t0, t0))
|
||||
|
||||
@@ -84,7 +84,7 @@ class InMemoryConversationStoreTest {
|
||||
|
||||
@Test
|
||||
fun `rename with null title clears it`() = runTest {
|
||||
val store = InMemoryConversationStore()
|
||||
val store = InMemoryMutableConversationStore()
|
||||
val t0 = Instant.parse("2026-09-15T10:00:00Z")
|
||||
store.upsert(ConversationRecord("c1", "old", false, t0, t0))
|
||||
store.rename("c1", null)
|
||||
@@ -93,13 +93,13 @@ class InMemoryConversationStoreTest {
|
||||
|
||||
@Test
|
||||
fun `rename returns null for missing conversation`() = runTest {
|
||||
val store = InMemoryConversationStore()
|
||||
val store = InMemoryMutableConversationStore()
|
||||
assertNull(store.rename("nope", "x"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `touch bumps updatedAt without changing other fields`() = runTest {
|
||||
val store = InMemoryConversationStore()
|
||||
val store = InMemoryMutableConversationStore()
|
||||
val t0 = Instant.parse("2026-09-15T10:00:00Z")
|
||||
val t1 = Instant.parse("2026-09-15T10:01:00Z")
|
||||
store.upsert(ConversationRecord("c1", "title", false, t0, t0))
|
||||
@@ -112,7 +112,7 @@ class InMemoryConversationStoreTest {
|
||||
|
||||
@Test
|
||||
fun `close is idempotent and does nothing`() {
|
||||
val store = InMemoryConversationStore()
|
||||
val store = InMemoryMutableConversationStore()
|
||||
store.close()
|
||||
store.close() // должно быть no-op
|
||||
}
|
||||
+13
-3
@@ -30,7 +30,7 @@ internal fun encodeRecord(record: MessageRecord): Pair<String, String> = when (r
|
||||
)
|
||||
is MessageRecord.ToolResult -> "tool_result" to Json.encodeToString(
|
||||
ResultPayload.serializer(),
|
||||
ResultPayload(toolCallId = record.toolCallId, result = record.result),
|
||||
ResultPayload(toolCallId = record.toolCallId, toolName = record.toolName, result = record.result),
|
||||
)
|
||||
is MessageRecord.Error -> "error" to Json.encodeToString(
|
||||
ErrorPayload.serializer(),
|
||||
@@ -59,7 +59,7 @@ internal fun SQLiteResultSet.toMessageRecord(json: Json): MessageRecord {
|
||||
}
|
||||
"tool_result" -> {
|
||||
val p = Json.decodeFromString(ResultPayload.serializer(), payload)
|
||||
MessageRecord.ToolResult(id = id, conversationId = convId, toolCallId = p.toolCallId, result = p.result, createdAt = createdAt)
|
||||
MessageRecord.ToolResult(id = id, conversationId = convId, toolCallId = p.toolCallId, toolName = p.toolName, result = p.result, createdAt = createdAt)
|
||||
}
|
||||
"error" -> {
|
||||
val p = Json.decodeFromString(ErrorPayload.serializer(), payload)
|
||||
@@ -72,8 +72,18 @@ internal fun SQLiteResultSet.toMessageRecord(json: Json): MessageRecord {
|
||||
@kotlinx.serialization.Serializable
|
||||
internal data class CallPayload(val name: String, val title: String?, val argsJson: String)
|
||||
|
||||
/**
|
||||
* Тулрезалт-сериализация для SQLite. [toolName] денормализован из
|
||||
* соответствующего `ToolCall.name` для упрощения UI (нет нужды в
|
||||
* локальной `Map<id, name>`). Nullable с дефолтом — старые записи
|
||||
* без поля десериализуются как `null`.
|
||||
*/
|
||||
@kotlinx.serialization.Serializable
|
||||
internal data class ResultPayload(val toolCallId: String, val result: String?)
|
||||
internal data class ResultPayload(
|
||||
val toolCallId: String,
|
||||
val toolName: String? = null,
|
||||
val result: String?,
|
||||
)
|
||||
|
||||
@kotlinx.serialization.Serializable
|
||||
internal data class ErrorPayload(val message: String, val code: String?)
|
||||
|
||||
@@ -14,9 +14,9 @@ import kotlin.time.Instant
|
||||
* идентична `OutboxStore.events`: поток **не реплеит** прошлое, для бэкфилла
|
||||
* используются `Agent.getConversations` / `getConversation`.
|
||||
*
|
||||
* **История**: раньше жил в `:proto` (как `pw.binom.agentik.proto.AgentEvent`).
|
||||
* После миграции в `:outbox-api` — `:proto.AgentEvent` стал typealias'ом,
|
||||
* backward-compat для существующих импортов сохранён.
|
||||
* **История**: до 2026-09-21 жил в `:proto` как `pw.binom.agentik.proto.AgentEvent`;
|
||||
* typealias удалён 2026-09-21 (стирал nested-типы в `is`/`when`) — потребители
|
||||
* импортируют напрямую из `pw.binom.agentik.outbox.AgentEvent`.
|
||||
*/
|
||||
@Serializable
|
||||
sealed interface AgentEvent {
|
||||
@@ -45,4 +45,13 @@ sealed interface AgentEvent {
|
||||
@Serializable
|
||||
@SerialName("renamed")
|
||||
data class Renamed(override val date: Instant, val id: String, val title: String?) : AgentEvent
|
||||
|
||||
/**
|
||||
* Обновлён `updatedAt` диалога (после `send()` или другого события,
|
||||
* бампнувшего активность). Клиентский кэш [ConversationStore] может
|
||||
* применить этот event для пересортировки списка.
|
||||
*/
|
||||
@Serializable
|
||||
@SerialName("touched")
|
||||
data class Touched(override val date: Instant, val id: String, val updatedAt: Instant) : AgentEvent
|
||||
}
|
||||
|
||||
@@ -9,17 +9,17 @@ import kotlinx.serialization.Serializable
|
||||
*
|
||||
* Useful for admin dashboards, debug tools, parent agents: one subscription
|
||||
* instead of N+1. For regular UI use two separate SSE feeds
|
||||
* ([AgentEvent] via `/events` и `Event` via `/conversations/{id}/events`);
|
||||
* ([AgentEvent] via `/events` и [Event] via `/conversations/{id}/events`);
|
||||
* [CommonEvent] — for those who need everything in one place.
|
||||
*
|
||||
* Server endpoint: `GET /events/all` (SSE), or replay via `OutboxStore.events(after)`.
|
||||
*
|
||||
* **История**: раньше жил в `:proto` (как `pw.binom.agentik.proto.CommonEvent`).
|
||||
* После миграции в `:outbox-api` — `:proto.CommonEvent` стал typealias'ом,
|
||||
* backward-compat для существующих импортов сохранён. `CommonEvent.Conversation`
|
||||
* ссылается на [Event] (тоже в `:outbox-api` теперь) — раньше был
|
||||
* `pw.binom.agentik.proto.Event`, теперь это `pw.binom.agentik.outbox.Event`
|
||||
* (он тоже typealias-нут в `:proto.Event`).
|
||||
* **История**: до 2026-09-21 жил в `:proto` как `pw.binom.agentik.proto.CommonEvent`;
|
||||
* при миграции в `:outbox-api` был оставлен typealias в `:proto` для backward-compat,
|
||||
* но он стирал nested-типы (`CommonEvent.Agent`, `CommonEvent.Conversation`),
|
||||
* что ломало `is CommonEvent.Agent` на стороне клиента. Typealias'ы
|
||||
* `Event`/`AgentEvent`/`CommonEvent` из `:proto` удалены — потребители
|
||||
* импортируют напрямую из `pw.binom.agentik.outbox.*`.
|
||||
*/
|
||||
@Serializable
|
||||
sealed interface CommonEvent {
|
||||
|
||||
@@ -15,9 +15,12 @@ import kotlin.time.Instant
|
||||
* `StartReasoning?` → `StartResponse(TEXT|IMAGE)` → ...контент... → `End` | `Interrupted` | `Error`.
|
||||
* `StartReasoning` может отсутствовать, если агент не показывал рассуждения.
|
||||
*
|
||||
* **История**: раньше жил в `:proto` (как `pw.binom.agentik.proto.Event`).
|
||||
* После миграции в `:outbox-api` — `:proto.Event` стал typealias'ом,
|
||||
* backward-compat для существующих импортов сохранён.
|
||||
* **История**: до 2026-09-21 жил в `:proto` как `pw.binom.agentik.proto.Event`;
|
||||
* при миграции в `:outbox-api` был оставлен typealias в `:proto` для
|
||||
* backward-compat, но он стирал nested-типы (`Event.End`, `Event.ToolCall`,
|
||||
* `Event.ToolResult`), что ломало `is Event.End` на стороне клиента.
|
||||
* Typealias удалён 2026-09-21 — потребители импортируют напрямую из
|
||||
* `pw.binom.agentik.outbox.Event`.
|
||||
*/
|
||||
@Serializable
|
||||
sealed interface Event {
|
||||
@@ -75,12 +78,26 @@ sealed interface Event {
|
||||
|
||||
/**
|
||||
* Результат вызова тула. Приходит целиком после завершения исполнения.
|
||||
* [id] совпадает с [ToolCall.id], к которому относится результат, и
|
||||
* с id `Message.ToolResult` в истории.
|
||||
*
|
||||
* [toolCallId] = id [ToolCall], к которому относится результат, и
|
||||
* `MessageRecord.ToolResult.toolCallId` в истории. Один Call → один Result,
|
||||
* пара `(date, toolCallId)` уникальна — отдельный `id` в live-событии
|
||||
* не нужен (PK живёт в персистентном журнале).
|
||||
*
|
||||
* [toolName] денормализован из соответствующего [ToolCall.toolName] —
|
||||
* UI рендерит имя тула без локальной `Map<toolCallId, name>` и без риска
|
||||
* «Result пришёл до Call». `null` допустим для backfill'а старых
|
||||
* записей, у которых поле отсутствует, или теоретического случая
|
||||
* Result без предшествующего Call (orphan).
|
||||
*/
|
||||
@Serializable
|
||||
@SerialName("tool_result")
|
||||
data class ToolResult(override val date: Instant, val id: String, val result: String?) : Event
|
||||
data class ToolResult(
|
||||
override val date: Instant,
|
||||
val toolCallId: String,
|
||||
val toolName: String? = null,
|
||||
val result: String?,
|
||||
) : Event
|
||||
|
||||
/**
|
||||
* Ошибка хода. После неё поток завершается; дальнейшие события могут
|
||||
|
||||
@@ -1,7 +1,6 @@
|
||||
package pw.binom.agentik.proto
|
||||
|
||||
import kotlinx.coroutines.flow.Flow
|
||||
import kotlinx.coroutines.flow.flow
|
||||
import pw.binom.agentik.journal.ConversationStore
|
||||
import pw.binom.agentik.journal.JournalStore
|
||||
import pw.binom.agentik.outbox.OutboxStore
|
||||
import kotlin.time.Instant
|
||||
@@ -13,11 +12,16 @@ import kotlin.time.Instant
|
||||
* [createConversation] возвращает [Conversation], который сам хранит историю
|
||||
* и которому отправляют ходы через [Conversation.send].
|
||||
*
|
||||
* **Хранилища вынесены в [Agent.journal] и [Agent.outbox]**: оба read-only.
|
||||
* События больше НЕ часть [Agent] (раньше были `events()`/`allEvents()`) —
|
||||
* они теперь живут в [outbox] как `OutboxStore.events(after)` /
|
||||
* `outbox.agentEvents(after)`. Это даёт единый путь для всех read-операций
|
||||
* по хранилищу и убирает дублирование между протоколом и хранилищем.
|
||||
* **Хранилища вынесены в [Agent.journal], [Agent.outbox] и
|
||||
* [Agent.conversationStore]**: все три read-only views. События живут
|
||||
* в [outbox] как `OutboxStore.events(after)` / `outbox.agentEvents(after)`.
|
||||
* Это даёт единый путь для всех read-операций по хранилищу и убирает
|
||||
* дублирование между протоколом и хранилищем.
|
||||
*
|
||||
* **Команды** (create / delete / rename) живут прямо на [Agent]. Они
|
||||
* шлются клиентом и выполняются сервером — клиент **не** пишет в стор
|
||||
* напрямую. Клиентский кэш [conversationStore] обновляется через
|
||||
* `outbox.agentEvents()` (Created / Deleted / Renamed / Touched).
|
||||
*/
|
||||
interface Agent : AutoCloseable {
|
||||
|
||||
@@ -59,6 +63,22 @@ interface Agent : AutoCloseable {
|
||||
*/
|
||||
val outbox: OutboxStore
|
||||
|
||||
/**
|
||||
* Read-only view на `conversation` table (id + title + timestamps).
|
||||
*
|
||||
* Используется HTTP-фасадом `:server` для endpoint'а
|
||||
* `GET /{path}/conversations?offset=&limit=` — внешние клиенты
|
||||
* получают лёгкую метадату (без handle'ов и image-support флагов)
|
||||
* для рендера списка диалогов. Для активной работы (send / interrupt)
|
||||
* клиент отдельно получает handle через [getConversation].
|
||||
*
|
||||
* **Read-only**: write-доступ только через `MutableConversationStore`
|
||||
* внутри ChatAgent, не через [Agent] interface. Клиент модифицирует
|
||||
* диалоги командами: [createConversation] / [deleteConversation] /
|
||||
* [renameConversation].
|
||||
*/
|
||||
val conversationStore: ConversationStore
|
||||
|
||||
/** Создаёт новый stateful-диалог с агентом. */
|
||||
fun createConversation(temp: Boolean): Conversation
|
||||
|
||||
@@ -68,19 +88,11 @@ interface Agent : AutoCloseable {
|
||||
/** Удаляет диалог. Возвращает `true`, если диалог существовал и удалён. */
|
||||
suspend fun deleteConversation(id: String): Boolean
|
||||
|
||||
/** Страница диалогов: не более [limit] штук, начиная с [offset]-го. */
|
||||
suspend fun getConversations(offset: Int, limit: Int): List<Conversation>
|
||||
|
||||
/** Все диалоги, начиная с [offset], как поток: подгружает по [PAGE_SIZE] за раз. */
|
||||
fun getConversations(offset: Int = 0): Flow<Conversation> = flow {
|
||||
var skip = offset
|
||||
while (true) {
|
||||
val page = getConversations(skip, PAGE_SIZE)
|
||||
if (page.isEmpty()) break
|
||||
page.forEach { emit(it) }
|
||||
skip += page.size
|
||||
}
|
||||
}
|
||||
/**
|
||||
* Переименовывает диалог; `null` для сброса заголовка. Возвращает новый
|
||||
* `updatedAt` или `null`, если диалог не найден.
|
||||
*/
|
||||
suspend fun renameConversation(id: String, title: String?): Instant?
|
||||
|
||||
companion object {
|
||||
|
||||
|
||||
@@ -1,10 +0,0 @@
|
||||
package pw.binom.agentik.proto
|
||||
|
||||
/**
|
||||
* Backward-compat typealias: `AgentEvent` теперь живёт в `:outbox-api`
|
||||
* (логически принадлежит сущности outbox, не wire-протоколу `:proto`).
|
||||
*
|
||||
* Существующие импорты `pw.binom.agentik.proto.AgentEvent` продолжают
|
||||
* работать транспарентно. Использовать typealias в новом коде.
|
||||
*/
|
||||
typealias AgentEvent = pw.binom.agentik.outbox.AgentEvent
|
||||
@@ -1,9 +0,0 @@
|
||||
package pw.binom.agentik.proto
|
||||
|
||||
/**
|
||||
* Backward-compat typealias: `CommonEvent` теперь живёт в `:outbox-api`.
|
||||
*
|
||||
* Существующие импорты `pw.binom.agentik.proto.CommonEvent` продолжают
|
||||
* работать транспарентно.
|
||||
*/
|
||||
typealias CommonEvent = pw.binom.agentik.outbox.CommonEvent
|
||||
@@ -8,6 +8,19 @@ import kotlin.time.Instant
|
||||
* Stateful-диалог клиента и [Agent]. Хранит собственную историю: на каждый
|
||||
* [send] агенту не нужно пересылать транскрипт — он уже живёт внутри
|
||||
* [Conversation].
|
||||
*
|
||||
* **Live-события** диалога (turn stream: StartReasoning / AppendText / End /
|
||||
* ToolCall / ToolResult / ...) НЕ часть этого интерфейса — единственный
|
||||
* источник live-событий это [pw.binom.agentik.outbox.OutboxStore].
|
||||
* Подписаться на события конкретного диалога:
|
||||
* ```
|
||||
* agent.outbox.conversationEvents(after = lastSeen, conversationId = id)
|
||||
* .map { it.event }
|
||||
* .collect { e -> ... }
|
||||
* ```
|
||||
* Для cross-conversation view (admin / parent-agent / debug):
|
||||
* `agent.outbox.events(after)`. Для lifecycle агента (created/deleted/renamed):
|
||||
* `agent.outbox.agentEvents(after)`.
|
||||
*/
|
||||
interface Conversation : AutoCloseable {
|
||||
val id: String
|
||||
@@ -33,7 +46,7 @@ interface Conversation : AutoCloseable {
|
||||
|
||||
/**
|
||||
* Ставит новый user-ход в очередь. Возвращает управление сразу — поток
|
||||
* событий ответа приходит через [events].
|
||||
* событий ответа приходит через `agent.outbox.conversationEvents(...)`.
|
||||
*
|
||||
* Если в момент вызова выполняется другой ход, новый встаёт в очередь
|
||||
* за ним. Чтобы отменить текущий — вызови [interrupt] перед [send].
|
||||
@@ -48,25 +61,13 @@ interface Conversation : AutoCloseable {
|
||||
|
||||
/**
|
||||
* Прерывает текущий исполняемый ход (best-effort: LLM-stream прибивается,
|
||||
* in-flight tool может доехать или отвалиться). В [events] эмитится
|
||||
* [Event.Interrupted], затем может начаться следующий ход из очереди.
|
||||
* in-flight tool может доехать или отвалиться). В `agent.outbox.conversationEvents`
|
||||
* эмитится `Event.Interrupted`, затем может начаться следующий ход из очереди.
|
||||
*
|
||||
* Если хода нет — no-op.
|
||||
*/
|
||||
suspend fun interrupt()
|
||||
|
||||
/**
|
||||
* Live-подписка на всё, что происходит в диалоге, начиная с [after].
|
||||
*
|
||||
* **Не реплеит** события, произошедшие до [after] — для бэкфилла
|
||||
* используй [getMessages]. Если [after] — момент последнего виденного
|
||||
* клиентом события, поток продолжается «с того места».
|
||||
*
|
||||
* Подписки независимы: каждый вызов возвращает свой [Flow], отмена одного
|
||||
* не влияет на других подписчиков и на сам диалог.
|
||||
*/
|
||||
fun events(after: Instant): Flow<Event>
|
||||
|
||||
/** Страница истории: не более [limit] сообщений после [after], начиная с [offset]-го. */
|
||||
suspend fun getMessages(after: Instant, offset: Int, limit: Int): List<Message>
|
||||
|
||||
@@ -83,7 +84,7 @@ interface Conversation : AutoCloseable {
|
||||
|
||||
/**
|
||||
* Освобождает ресурсы диалога (подписки, сетевые хэндлы). Идемпотентно.
|
||||
* После [close] дальнейшие вызовы [send]/[interrupt]/[events]/[getMessages]/[rename] не определены.
|
||||
* После [close] дальнейшие вызовы [send]/[interrupt]/[getMessages]/[rename] не определены.
|
||||
*/
|
||||
override fun close()
|
||||
|
||||
|
||||
@@ -1,11 +0,0 @@
|
||||
package pw.binom.agentik.proto
|
||||
|
||||
/**
|
||||
* Backward-compat typealias: `Event` теперь живёт в `:outbox-api`.
|
||||
*
|
||||
* Существующие импорты `pw.binom.agentik.proto.Event` продолжают
|
||||
* работать транспарентно. `Conversation.events(after): Flow<Event>` в
|
||||
* `:proto.Conversation` теперь фактически возвращает
|
||||
* `pw.binom.agentik.outbox.Event` — тот же тип, другое имя.
|
||||
*/
|
||||
typealias Event = pw.binom.agentik.outbox.Event
|
||||
@@ -45,9 +45,22 @@ sealed interface Message {
|
||||
override val date: Instant
|
||||
) : Message
|
||||
|
||||
/**
|
||||
* Результат вызова тула. Приходит в историю `getMessages` после завершения хода.
|
||||
* [id] совпадает с [ToolCall.id], к которому относится результат, и
|
||||
* с id соответствующего `MessageRecord.ToolResult` в journal.
|
||||
*
|
||||
* [toolName] денормализован из [ToolCall.toolName] — UI рендерит
|
||||
* имя тула в строке результата без отдельной `Map<id, name>`.
|
||||
*/
|
||||
@Serializable
|
||||
@SerialName("tool_result")
|
||||
class ToolResult(override val id: String, val result: String?, override val date: Instant) : Message
|
||||
class ToolResult(
|
||||
override val id: String,
|
||||
val toolName: String? = null,
|
||||
val result: String?,
|
||||
override val date: Instant,
|
||||
) : Message
|
||||
|
||||
/**
|
||||
* Ход завершился ошибкой (LLM, инициализация движка или иная отказоустойчивая
|
||||
|
||||
@@ -3,8 +3,8 @@ package pw.binom.agentik.server
|
||||
import io.ktor.server.routing.Route
|
||||
import io.ktor.server.routing.get
|
||||
import io.ktor.server.routing.route
|
||||
import pw.binom.agentik.outbox.CommonEvent
|
||||
import pw.binom.agentik.outbox.OutboxStore
|
||||
import pw.binom.agentik.proto.CommonEvent
|
||||
|
||||
/**
|
||||
* HTTP-фасад для [OutboxStore] (bounded-tail live event stream агента).
|
||||
|
||||
@@ -20,11 +20,11 @@ import kotlinx.coroutines.flow.map
|
||||
import kotlinx.serialization.KSerializer
|
||||
import kotlinx.serialization.json.Json
|
||||
import pw.binom.agentik.proto.Agent
|
||||
import pw.binom.agentik.proto.AgentEvent
|
||||
import pw.binom.agentik.proto.CommonEvent
|
||||
import pw.binom.agentik.proto.Conversation
|
||||
import pw.binom.agentik.proto.Content
|
||||
import pw.binom.agentik.proto.Event
|
||||
import pw.binom.agentik.outbox.AgentEvent
|
||||
import pw.binom.agentik.outbox.CommonEvent
|
||||
import pw.binom.agentik.outbox.Event
|
||||
import kotlin.time.Instant
|
||||
|
||||
internal fun Route.agentikRoutes(agent: Agent) {
|
||||
@@ -38,7 +38,7 @@ internal fun Route.agentikRoutes(agent: Agent) {
|
||||
get("/conversations") {
|
||||
val offset = call.request.queryParameters["offset"]?.toIntOrNull() ?: 0
|
||||
val limit = call.request.queryParameters["limit"]?.toIntOrNull() ?: Agent.PAGE_SIZE
|
||||
call.respond(agent.getConversations(offset, limit).map { it.snapshot() })
|
||||
call.respond(agent.conversationStore.list(offset, limit))
|
||||
}
|
||||
|
||||
post("/conversations") {
|
||||
@@ -63,14 +63,15 @@ internal fun Route.agentikRoutes(agent: Agent) {
|
||||
|
||||
patch("/conversations/{id}") {
|
||||
val id = call.parameters["id"]!!
|
||||
val c = agent.getConversation(id)
|
||||
if (c == null) {
|
||||
val req = call.receive<RequestRename>()
|
||||
val newUpdatedAt = agent.renameConversation(id, req.title)
|
||||
if (newUpdatedAt == null) {
|
||||
call.respond(HttpStatusCode.NotFound)
|
||||
return@patch
|
||||
}
|
||||
val req = call.receive<RequestRename>()
|
||||
c.rename(req.title)
|
||||
call.respond(c.snapshot())
|
||||
// Возвращаем обновлённый record (лёгкая метадата, не snapshot handle'а).
|
||||
val rec = agent.conversationStore.get(id)!!
|
||||
call.respond(rec)
|
||||
}
|
||||
|
||||
post("/conversations/{id}/messages") {
|
||||
@@ -128,7 +129,10 @@ internal fun Route.agentikRoutes(agent: Agent) {
|
||||
return@get
|
||||
}
|
||||
val after = call.parseAfter() ?: return@get
|
||||
call.streamJsonSse(c.events(after), Event.serializer())
|
||||
// Live-источник событий — `OutboxStore` (единая точка истины);
|
||||
// разворачиваем `CommonEvent.Conversation` → `Event` для совместимости
|
||||
// wire-формата (клиент десериализует как `Event`, не как `CommonEvent.Conversation`).
|
||||
call.streamJsonSse(agent.outbox.conversationEvents(after, id).map { it.event }, Event.serializer())
|
||||
}
|
||||
|
||||
get("/events") {
|
||||
|
||||
@@ -13,6 +13,7 @@ import io.ktor.server.engine.embeddedServer
|
||||
import io.ktor.server.routing.routing
|
||||
import kotlinx.coroutines.flow.emptyFlow
|
||||
import kotlinx.coroutines.runBlocking
|
||||
import pw.binom.agentik.journal.ConversationStore
|
||||
import pw.binom.agentik.journal.JournalStore
|
||||
import pw.binom.agentik.journal.MessageRecord
|
||||
import pw.binom.agentik.outbox.CommonEvent
|
||||
@@ -53,10 +54,15 @@ class BearerTokenTest {
|
||||
override suspend fun earliestEventDate(): Instant = Instant.DISTANT_PAST
|
||||
override fun close() {}
|
||||
}
|
||||
override val conversationStore: ConversationStore = object : ConversationStore {
|
||||
override suspend fun get(id: String) = null
|
||||
override suspend fun list(offset: Int, limit: Int) = emptyList<pw.binom.agentik.journal.ConversationRecord>()
|
||||
override fun close() {}
|
||||
}
|
||||
override fun createConversation(temp: Boolean): Conversation = TODO("not needed by tests")
|
||||
override suspend fun getConversation(id: String): Conversation? = null
|
||||
override suspend fun deleteConversation(id: String): Boolean = false
|
||||
override suspend fun getConversations(offset: Int, limit: Int): List<Conversation> = emptyList()
|
||||
override suspend fun renameConversation(id: String, title: String?): Instant? = null
|
||||
override fun close() {}
|
||||
}
|
||||
|
||||
|
||||
@@ -49,7 +49,8 @@ class A2aBridge(private val agent: Agent) : AgentHandler {
|
||||
val reply = StringBuilder()
|
||||
val turnDone = CompletableDeferred<Unit>()
|
||||
val subscription = async {
|
||||
conv.events(since).collect { e ->
|
||||
agent.outbox.conversationEvents(since, conv.id).collect { ce ->
|
||||
val e = ce.event
|
||||
when (e) {
|
||||
is Event.AppendText -> reply.append(e.body)
|
||||
is Event.End, is Event.Interrupted -> turnDone.complete(Unit)
|
||||
|
||||
@@ -363,7 +363,7 @@ private fun runServer() {
|
||||
|
||||
val agent = ChatAgent(
|
||||
id = "agentik",
|
||||
conversationStore = sqliteStores.conversations,
|
||||
mutableConversationStore = sqliteStores.conversations,
|
||||
messageStore = sqliteStores.messages,
|
||||
workingMemoryStore = sqliteStores.workingMemory,
|
||||
reflectionStore = sqliteStores.reflections,
|
||||
|
||||
@@ -1,11 +1,6 @@
|
||||
package pw.binom.agentik.standalone.agent
|
||||
|
||||
import kotlin.time.Instant
|
||||
import kotlinx.coroutines.flow.Flow
|
||||
import kotlinx.coroutines.flow.emitAll
|
||||
import kotlinx.coroutines.flow.flow
|
||||
import kotlinx.coroutines.flow.map
|
||||
import kotlinx.coroutines.launch
|
||||
import kotlinx.coroutines.runBlocking
|
||||
import kotlinx.coroutines.sync.Mutex
|
||||
import kotlinx.coroutines.sync.withLock
|
||||
@@ -22,14 +17,13 @@ import pw.binom.agentik.outbox.MutableOutboxStore
|
||||
import pw.binom.agentik.journal.JournalStore
|
||||
import pw.binom.agentik.outbox.OutboxStore
|
||||
import pw.binom.agentik.proto.Conversation as ProtoConversation
|
||||
import pw.binom.agentik.proto.Event as ProtoEvent
|
||||
import pw.binom.agentik.skills.SkillCatalog
|
||||
import pw.binom.agentik.skills.renderSystemPromptSection
|
||||
import pw.binom.agentik.standalone.agent.memory.MemoryToolsFactory
|
||||
import pw.binom.agentik.standalone.llm.LlmConfig
|
||||
import pw.binom.agentik.journal.ConversationRecord
|
||||
import pw.binom.agentik.journal.ConversationStore
|
||||
import pw.binom.agentik.journal.Ids
|
||||
import pw.binom.agentik.journal.MutableConversationStore
|
||||
import pw.binom.agentik.reflection.Reflection
|
||||
import pw.binom.agentik.reflection.ReflectionStore
|
||||
import pw.binom.agentik.journal.MutableJournalStore
|
||||
@@ -65,7 +59,7 @@ import pw.binom.litert.LiteLlm
|
||||
*/
|
||||
class ChatAgent(
|
||||
override val id: String,
|
||||
private val conversationStore: ConversationStore,
|
||||
private val mutableConversationStore: MutableConversationStore,
|
||||
private val messageStore: MutableJournalStore,
|
||||
private val workingMemoryStore: ContextStore,
|
||||
private val reflectionStore: ReflectionStore,
|
||||
@@ -239,6 +233,16 @@ class ChatAgent(
|
||||
override val outbox: OutboxStore
|
||||
get() = eventStore
|
||||
|
||||
/**
|
||||
* Read-only view of [mutableConversationStore] для HTTP-фасада в `:server`
|
||||
* (`GET /conversations` → список [ConversationRecord] для UI).
|
||||
*
|
||||
* Сужение с [MutableConversationStore] на [ConversationStore] тривиальна
|
||||
* через interface-наследование (read-only projection).
|
||||
*/
|
||||
override val conversationStore: ConversationStore
|
||||
get() = mutableConversationStore
|
||||
|
||||
/** Защищает карту живых диалогов. */
|
||||
private val liveLock = Mutex()
|
||||
private val live: MutableMap<String, ChatConversation> = HashMap()
|
||||
@@ -262,12 +266,12 @@ class ChatAgent(
|
||||
// и не переживают рестарт агента (см. Memory #3709).
|
||||
if (!temp) {
|
||||
runBlocking {
|
||||
conversationStore.upsert(rec)
|
||||
mutableConversationStore.upsert(rec)
|
||||
}
|
||||
}
|
||||
val conv = ChatConversation(
|
||||
record = rec,
|
||||
conversationStore = conversationStore,
|
||||
conversationStore = mutableConversationStore,
|
||||
messageStore = messageStore,
|
||||
workingMemoryStore = workingMemoryStore,
|
||||
reflectionStore = reflectionStore,
|
||||
@@ -302,7 +306,7 @@ class ChatAgent(
|
||||
|
||||
override suspend fun getConversation(id: String): ProtoConversation? {
|
||||
liveLock.withLock { live[id] }?.let { if (!it.isClosed) return it }
|
||||
val rec = conversationStore.get(id) ?: return null
|
||||
val rec = mutableConversationStore.get(id) ?: return null
|
||||
return newConversation(rec).also {
|
||||
liveLock.withLock { live[id] = it }
|
||||
}
|
||||
@@ -311,7 +315,7 @@ class ChatAgent(
|
||||
override suspend fun deleteConversation(id: String): Boolean {
|
||||
val conv = liveLock.withLock { live.remove(id) }
|
||||
conv?.close()
|
||||
val ok = conversationStore.delete(id)
|
||||
val ok = mutableConversationStore.delete(id)
|
||||
if (ok) {
|
||||
val event = AgentEvent.Deleted(date = now(), id = id)
|
||||
eventStore.append(CommonEvent.Agent(date = now(), event = event))
|
||||
@@ -319,17 +323,16 @@ class ChatAgent(
|
||||
return ok
|
||||
}
|
||||
|
||||
override suspend fun getConversations(offset: Int, limit: Int): List<ProtoConversation> =
|
||||
conversationStore.list(offset = offset, limit = limit).map { rec ->
|
||||
liveLock.withLock { live[rec.id] }
|
||||
?: newConversation(rec).also {
|
||||
liveLock.withLock { live[rec.id] = it }
|
||||
}
|
||||
override suspend fun renameConversation(id: String, title: String?): Instant? {
|
||||
val newUpdatedAt = mutableConversationStore.rename(id, title) ?: return null
|
||||
val event = AgentEvent.Renamed(date = newUpdatedAt, id = id, title = title)
|
||||
eventStore.append(CommonEvent.Agent(date = newUpdatedAt, event = event))
|
||||
return newUpdatedAt
|
||||
}
|
||||
|
||||
private fun newConversation(rec: ConversationRecord): ChatConversation = ChatConversation(
|
||||
record = rec,
|
||||
conversationStore = conversationStore,
|
||||
conversationStore = mutableConversationStore,
|
||||
messageStore = messageStore,
|
||||
workingMemoryStore = workingMemoryStore,
|
||||
reflectionStore = reflectionStore,
|
||||
|
||||
+4
-4
@@ -3,15 +3,15 @@ package pw.binom.agentik.standalone.agent
|
||||
import kotlinx.coroutines.runBlocking
|
||||
import kotlinx.coroutines.flow.Flow
|
||||
import kotlinx.coroutines.flow.map
|
||||
import pw.binom.agentik.outbox.MutableOutboxStore
|
||||
import pw.binom.agentik.outbox.CommonEvent
|
||||
import pw.binom.agentik.proto.Event as ProtoEvent
|
||||
import pw.binom.agentik.outbox.Event
|
||||
import pw.binom.agentik.outbox.MutableOutboxStore
|
||||
|
||||
internal class ConversationEvents(
|
||||
private val globalEventStore: MutableOutboxStore,
|
||||
private val conversationId: String,
|
||||
) {
|
||||
fun tryEmit(event: ProtoEvent): Boolean {
|
||||
fun tryEmit(event: Event): Boolean {
|
||||
runBlocking {
|
||||
globalEventStore.append(
|
||||
CommonEvent.Conversation(
|
||||
@@ -24,7 +24,7 @@ internal class ConversationEvents(
|
||||
return true
|
||||
}
|
||||
|
||||
fun events(after: kotlin.time.Instant?): Flow<ProtoEvent> =
|
||||
fun events(after: kotlin.time.Instant?): Flow<Event> =
|
||||
globalEventStore.conversationEvents(after = after, conversationId = conversationId)
|
||||
.map { it.event }
|
||||
}
|
||||
|
||||
+15
-6
@@ -7,7 +7,6 @@ import kotlinx.coroutines.Job
|
||||
import kotlinx.coroutines.SupervisorJob
|
||||
import kotlinx.coroutines.cancel
|
||||
import kotlinx.coroutines.cancelAndJoin
|
||||
import kotlinx.coroutines.flow.Flow
|
||||
import kotlinx.coroutines.launch
|
||||
import kotlinx.coroutines.runBlocking
|
||||
import kotlinx.coroutines.sync.Mutex
|
||||
@@ -31,7 +30,7 @@ import pw.binom.agentik.proto.MessageContext as ProtoMessageContext
|
||||
import pw.binom.agentik.skills.SkillStore
|
||||
import pw.binom.agentik.journal.Content
|
||||
import pw.binom.agentik.journal.ConversationRecord
|
||||
import pw.binom.agentik.journal.ConversationStore
|
||||
import pw.binom.agentik.journal.MutableConversationStore
|
||||
import pw.binom.agentik.journal.MessageContext
|
||||
import pw.binom.agentik.journal.MessageOrigin
|
||||
import pw.binom.agentik.journal.MessageRecord
|
||||
@@ -54,7 +53,7 @@ import pw.binom.agentik.toolsets.NamedTool
|
||||
|
||||
class ConversationLoop(
|
||||
record: ConversationRecord,
|
||||
private val conversationStore: ConversationStore,
|
||||
private val conversationStore: MutableConversationStore,
|
||||
private val messageStore: MutableJournalStore,
|
||||
private val workingMemoryStore: ContextStore,
|
||||
private val reflectionStore: ReflectionStore?,
|
||||
@@ -220,9 +219,6 @@ class ConversationLoop(
|
||||
toolDispatcher.currentToolJob?.cancel()
|
||||
}
|
||||
|
||||
override fun events(after: Instant): Flow<ProtoEvent> =
|
||||
events.events(after)
|
||||
|
||||
override suspend fun getMessages(after: Instant, offset: Int, limit: Int): List<ProtoMessage> =
|
||||
messageStore.list(conversationId = id, after = after, offset = offset, limit = limit)
|
||||
.map { it.toProto() }
|
||||
@@ -443,6 +439,18 @@ class ConversationLoop(
|
||||
|
||||
state.record = state.record.copy(updatedAt = assistantAt)
|
||||
conversationStore.touch(id, assistantAt)
|
||||
if (!state.isTemporal) {
|
||||
eventStore.append(
|
||||
pw.binom.agentik.outbox.CommonEvent.Agent(
|
||||
date = assistantAt,
|
||||
event = pw.binom.agentik.outbox.AgentEvent.Touched(
|
||||
date = assistantAt,
|
||||
id = id,
|
||||
updatedAt = assistantAt,
|
||||
),
|
||||
)
|
||||
)
|
||||
}
|
||||
|
||||
// BackgroundScheduler is event-driven — подписан на BackgroundEventBus
|
||||
// (compaction/lifecycle/tool-failure events). Никаких interval-based
|
||||
@@ -568,6 +576,7 @@ internal fun MessageRecord.toProto(): ProtoMessage = when (this) {
|
||||
is MessageRecord.ToolResult -> ProtoMessage.ToolResult(
|
||||
id = id,
|
||||
date = createdAt,
|
||||
toolName = toolName,
|
||||
result = result,
|
||||
)
|
||||
is MessageRecord.Error -> ProtoMessage.Error(
|
||||
|
||||
+2
-1
@@ -93,7 +93,7 @@ internal class ToolDispatcher(
|
||||
}
|
||||
|
||||
val resultAt = now()
|
||||
events.tryEmit(ProtoEvent.ToolResult(date = resultAt, id = resultId, result = resultText))
|
||||
events.tryEmit(ProtoEvent.ToolResult(date = resultAt, toolCallId = callId, toolName = call.name, result = resultText))
|
||||
|
||||
// Эмитим background event — другие компоненты (BackgroundScheduler)
|
||||
// решают, делать ли что-то. Cancellation = not a failure (не эмитим Failed).
|
||||
@@ -110,6 +110,7 @@ internal class ToolDispatcher(
|
||||
id = resultId,
|
||||
conversationId = state.id,
|
||||
toolCallId = callId,
|
||||
toolName = call.name,
|
||||
result = resultText,
|
||||
createdAt = resultAt,
|
||||
),
|
||||
|
||||
+27
-10
@@ -62,7 +62,7 @@ class ChatAgentTest {
|
||||
skills: SkillCatalog = SkillCatalog.EMPTY,
|
||||
): ChatAgent = ChatAgent(
|
||||
id = "agentik",
|
||||
conversationStore = sqliteStores.conversations,
|
||||
mutableConversationStore = sqliteStores.conversations,
|
||||
messageStore = sqliteStores.messages,
|
||||
workingMemoryStore = sqliteStores.workingMemory,
|
||||
reflectionStore = sqliteStores.reflections,
|
||||
@@ -146,15 +146,32 @@ class ChatAgentTest {
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `getConversations returns all stored persistent conversations`() = runTest {
|
||||
fun `conversationStore list returns all stored persistent conversations`() = runTest {
|
||||
val agent = newAgent()
|
||||
agent.createConversation(temp = false)
|
||||
agent.createConversation(temp = true)
|
||||
val list = agent.getConversations(0, 10)
|
||||
val list = agent.conversationStore.list(0, 10)
|
||||
// temp-беседы не персистятся, в списке только persistent
|
||||
assertEquals(1, list.size)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `renameConversation updates title and emits event`() = runTest {
|
||||
val agent = newAgent()
|
||||
val conv = agent.createConversation(temp = false)
|
||||
val newTitle = "Renamed!"
|
||||
val updatedAt = agent.renameConversation(conv.id, newTitle)
|
||||
assertNotNull(updatedAt)
|
||||
val rec = agent.conversationStore.get(conv.id)
|
||||
assertEquals(newTitle, rec?.title)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `renameConversation returns null for unknown id`() = runTest {
|
||||
val agent = newAgent()
|
||||
assertNull(agent.renameConversation("nope", "x"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `deleteConversation removes conversation and data`() = runTest {
|
||||
val agent = newAgent()
|
||||
@@ -319,7 +336,7 @@ class ChatAgentTest {
|
||||
|
||||
val events = mutableListOf<ProtoEvent>()
|
||||
val job = launch(start = kotlinx.coroutines.CoroutineStart.UNDISPATCHED) {
|
||||
conv.events(Instant.DISTANT_PAST).collect { events.add(it) }
|
||||
agent.outbox.conversationEvents(Instant.DISTANT_PAST, conv.id).collect { events.add(it.event) }
|
||||
}
|
||||
conv.send(listOf(Content.Text("hi")))
|
||||
delay(50)
|
||||
@@ -356,7 +373,7 @@ class ChatAgentTest {
|
||||
// отправки событий подписка ничего не увидит.
|
||||
val events = mutableListOf<ProtoEvent>()
|
||||
val eventsJob = launch(start = kotlinx.coroutines.CoroutineStart.UNDISPATCHED) {
|
||||
conv.events(Instant.DISTANT_PAST).collect { events.add(it) }
|
||||
agent.outbox.conversationEvents(Instant.DISTANT_PAST, conv.id).collect { events.add(it.event) }
|
||||
}
|
||||
|
||||
val sendJob = launch {
|
||||
@@ -414,7 +431,7 @@ class ChatAgentTest {
|
||||
// Подписываемся ДО send — SharedFlow без replay
|
||||
val events = mutableListOf<ProtoEvent>()
|
||||
val eventsJob = launch(start = kotlinx.coroutines.CoroutineStart.UNDISPATCHED) {
|
||||
conv.events(Instant.DISTANT_PAST).collect { events.add(it) }
|
||||
agent.outbox.conversationEvents(Instant.DISTANT_PAST, conv.id).collect { events.add(it.event) }
|
||||
}
|
||||
|
||||
val sendJob = launch {
|
||||
@@ -458,7 +475,7 @@ class ChatAgentTest {
|
||||
sqliteStores = KsqliteStores.open(dbPath)
|
||||
val agent1 = ChatAgent(
|
||||
id = "agentik",
|
||||
conversationStore = sqliteStores.conversations,
|
||||
mutableConversationStore = sqliteStores.conversations,
|
||||
messageStore = sqliteStores.messages,
|
||||
workingMemoryStore = sqliteStores.workingMemory,
|
||||
reflectionStore = sqliteStores.reflections,
|
||||
@@ -478,7 +495,7 @@ class ChatAgentTest {
|
||||
sqliteStores = KsqliteStores.open(dbPath)
|
||||
val agent2 = ChatAgent(
|
||||
id = "agentik",
|
||||
conversationStore = sqliteStores.conversations,
|
||||
mutableConversationStore = sqliteStores.conversations,
|
||||
messageStore = sqliteStores.messages,
|
||||
workingMemoryStore = sqliteStores.workingMemory,
|
||||
reflectionStore = sqliteStores.reflections,
|
||||
@@ -500,7 +517,7 @@ class ChatAgentTest {
|
||||
sqliteStores = KsqliteStores.open(dbPath)
|
||||
val agent1 = ChatAgent(
|
||||
id = "agentik",
|
||||
conversationStore = sqliteStores.conversations,
|
||||
mutableConversationStore = sqliteStores.conversations,
|
||||
messageStore = sqliteStores.messages,
|
||||
workingMemoryStore = sqliteStores.workingMemory,
|
||||
reflectionStore = sqliteStores.reflections,
|
||||
@@ -518,7 +535,7 @@ class ChatAgentTest {
|
||||
sqliteStores = KsqliteStores.open(dbPath)
|
||||
val agent2 = ChatAgent(
|
||||
id = "agentik",
|
||||
conversationStore = sqliteStores.conversations,
|
||||
mutableConversationStore = sqliteStores.conversations,
|
||||
messageStore = sqliteStores.messages,
|
||||
workingMemoryStore = sqliteStores.workingMemory,
|
||||
reflectionStore = sqliteStores.reflections,
|
||||
|
||||
+1
-1
@@ -43,7 +43,7 @@ class ChatAgentToolsetsTest {
|
||||
val sqliteStores = KsqliteStores.inMemory("toolsets-${kotlin.random.Random.nextLong()}")
|
||||
val agent = ChatAgent(
|
||||
id = "test-agent",
|
||||
conversationStore = sqliteStores.conversations,
|
||||
mutableConversationStore = sqliteStores.conversations,
|
||||
messageStore = sqliteStores.messages,
|
||||
workingMemoryStore = sqliteStores.workingMemory,
|
||||
reflectionStore = sqliteStores.reflections,
|
||||
|
||||
+1
-1
@@ -63,7 +63,7 @@ class CompactionTest {
|
||||
val reviewer = if (memoryStore != null) KeywordMdReviewer() else null
|
||||
return ChatAgent(
|
||||
id = "test",
|
||||
conversationStore = sqliteStores.conversations,
|
||||
mutableConversationStore = sqliteStores.conversations,
|
||||
messageStore = sqliteStores.messages,
|
||||
workingMemoryStore = sqliteStores.workingMemory,
|
||||
reflectionStore = sqliteStores.reflections,
|
||||
|
||||
+3
-3
@@ -76,7 +76,7 @@ class MemoryWiringTest {
|
||||
contextCompactor: ContextCompactor = EchoCompactor,
|
||||
): ChatAgent = ChatAgent(
|
||||
id = "agentik",
|
||||
conversationStore = sqliteStores.conversations,
|
||||
mutableConversationStore = sqliteStores.conversations,
|
||||
messageStore = sqliteStores.messages,
|
||||
workingMemoryStore = sqliteStores.workingMemory,
|
||||
reflectionStore = sqliteStores.reflections,
|
||||
@@ -115,7 +115,7 @@ class MemoryWiringTest {
|
||||
val soulBody = "I am a helpful test persona. I always answer in one short line."
|
||||
val agent = ChatAgent(
|
||||
id = "agentik",
|
||||
conversationStore = sqliteStores.conversations,
|
||||
mutableConversationStore = sqliteStores.conversations,
|
||||
messageStore = sqliteStores.messages,
|
||||
workingMemoryStore = sqliteStores.workingMemory,
|
||||
reflectionStore = sqliteStores.reflections,
|
||||
@@ -143,7 +143,7 @@ class MemoryWiringTest {
|
||||
fun `soul body not added when null`() = runBlocking {
|
||||
val agent = ChatAgent(
|
||||
id = "agentik",
|
||||
conversationStore = sqliteStores.conversations,
|
||||
mutableConversationStore = sqliteStores.conversations,
|
||||
messageStore = sqliteStores.messages,
|
||||
workingMemoryStore = sqliteStores.workingMemory,
|
||||
reflectionStore = sqliteStores.reflections,
|
||||
|
||||
@@ -21,6 +21,7 @@ kotlin {
|
||||
sourceSets {
|
||||
commonMain.dependencies {
|
||||
api(project(":journal-api"))
|
||||
api(project(":journal-inmemory"))
|
||||
api(project(":reflection-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.journal.MutableJournalStore
|
||||
import pw.binom.agentik.journal.ConversationStore
|
||||
import pw.binom.agentik.journal.MutableConversationStore
|
||||
import pw.binom.agentik.journal.inmemory.InMemoryMutableConversationStore
|
||||
import pw.binom.agentik.reflection.ReflectionStore
|
||||
import kotlin.time.Clock
|
||||
|
||||
@@ -21,14 +22,14 @@ import kotlin.time.Clock
|
||||
*/
|
||||
object InMemoryStorage {
|
||||
data class Bundle(
|
||||
val conversationStore: ConversationStore,
|
||||
val conversationStore: MutableConversationStore,
|
||||
val messageStore: MutableJournalStore,
|
||||
val workingMemoryStore: ContextStore,
|
||||
val reflectionStore: ReflectionStore,
|
||||
)
|
||||
|
||||
fun create(clock: Clock = Clock.System): Bundle = Bundle(
|
||||
conversationStore = InMemoryConversationStore(clock),
|
||||
conversationStore = InMemoryMutableConversationStore(clock),
|
||||
messageStore = InMemoryMessageStore(),
|
||||
workingMemoryStore = InMemoryWorkingMemoryStore(),
|
||||
reflectionStore = InMemoryReflectionStore(),
|
||||
|
||||
+4
-4
@@ -3,7 +3,7 @@ package pw.binom.agentik.storage.ksqlite
|
||||
import kotlin.time.Clock
|
||||
import kotlin.time.Instant
|
||||
import pw.binom.agentik.journal.ConversationRecord
|
||||
import pw.binom.agentik.journal.ConversationStore
|
||||
import pw.binom.agentik.journal.MutableConversationStore
|
||||
import pw.binom.db.ksqlite.SQLiteConnection
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.sync.Mutex
|
||||
@@ -11,15 +11,15 @@ import kotlinx.coroutines.sync.withLock
|
||||
import kotlinx.coroutines.withContext
|
||||
|
||||
/**
|
||||
* ksqlite-реализация [ConversationStore]. Схема таблицы `conversation` живёт
|
||||
* ksqlite-реализация [MutableConversationStore]. Схема таблицы `conversation` живёт
|
||||
* в [Schema] (миграция через PRAGMA user_version) — этот класс только
|
||||
* готовит и выполняет SQL, ссылаясь на `Schema.COL_*` / `Schema.TABLE_*`.
|
||||
*/
|
||||
class KsqliteConversationStore(
|
||||
class KsqliteMutableConversationStore(
|
||||
private val connection: SQLiteConnection,
|
||||
private val messageStore: KsqliteMessageStore? = null,
|
||||
private val workingMemoryStore: KsqliteWorkingMemoryStore? = null,
|
||||
) : ConversationStore {
|
||||
) : MutableConversationStore {
|
||||
|
||||
private val mutex = Mutex()
|
||||
|
||||
+3
-3
@@ -2,7 +2,7 @@ package pw.binom.agentik.storage.ksqlite
|
||||
|
||||
import pw.binom.agentik.context.ContextStore
|
||||
import pw.binom.agentik.journal.MutableJournalStore
|
||||
import pw.binom.agentik.journal.ConversationStore
|
||||
import pw.binom.agentik.journal.MutableConversationStore
|
||||
import pw.binom.agentik.reflection.ReflectionStore
|
||||
import pw.binom.db.ksqlite.SQLiteConnection
|
||||
|
||||
@@ -17,7 +17,7 @@ import pw.binom.db.ksqlite.SQLiteConnection
|
||||
*/
|
||||
class KsqliteStores internal constructor(
|
||||
val connection: SQLiteConnection,
|
||||
val conversations: ConversationStore,
|
||||
val conversations: MutableConversationStore,
|
||||
val messages: MutableJournalStore,
|
||||
val workingMemory: ContextStore,
|
||||
val reflections: ReflectionStore,
|
||||
@@ -51,7 +51,7 @@ class KsqliteStores internal constructor(
|
||||
val working = KsqliteWorkingMemoryStore(conn)
|
||||
return KsqliteStores(
|
||||
connection = conn,
|
||||
conversations = KsqliteConversationStore(conn, messages, working),
|
||||
conversations = KsqliteMutableConversationStore(conn, messages, working),
|
||||
messages = messages,
|
||||
workingMemory = working,
|
||||
reflections = KsqliteReflectionStore(conn),
|
||||
|
||||
+13
-3
@@ -30,7 +30,7 @@ internal fun encodeRecord(record: MessageRecord): Pair<String, String> = when (r
|
||||
)
|
||||
is MessageRecord.ToolResult -> "tool_result" to Json.encodeToString(
|
||||
ResultPayload.serializer(),
|
||||
ResultPayload(toolCallId = record.toolCallId, result = record.result),
|
||||
ResultPayload(toolCallId = record.toolCallId, toolName = record.toolName, result = record.result),
|
||||
)
|
||||
is MessageRecord.Error -> "error" to Json.encodeToString(
|
||||
ErrorPayload.serializer(),
|
||||
@@ -59,7 +59,7 @@ internal fun SQLiteResultSet.toMessageRecord(json: Json): MessageRecord {
|
||||
}
|
||||
"tool_result" -> {
|
||||
val p = Json.decodeFromString(ResultPayload.serializer(), payload)
|
||||
MessageRecord.ToolResult(id = id, conversationId = convId, toolCallId = p.toolCallId, result = p.result, createdAt = createdAt)
|
||||
MessageRecord.ToolResult(id = id, conversationId = convId, toolCallId = p.toolCallId, toolName = p.toolName, result = p.result, createdAt = createdAt)
|
||||
}
|
||||
"error" -> {
|
||||
val p = Json.decodeFromString(ErrorPayload.serializer(), payload)
|
||||
@@ -72,8 +72,18 @@ internal fun SQLiteResultSet.toMessageRecord(json: Json): MessageRecord {
|
||||
@kotlinx.serialization.Serializable
|
||||
internal data class CallPayload(val name: String, val title: String?, val argsJson: String)
|
||||
|
||||
/**
|
||||
* Тулрезалт-сериализация для SQLite. [toolName] денормализован из
|
||||
* соответствующего `ToolCall.name` для упрощения UI (нет нужды в
|
||||
* локальной `Map<id, name>`). Nullable с дефолтом — старые записи
|
||||
* без поля десериализуются как `null`.
|
||||
*/
|
||||
@kotlinx.serialization.Serializable
|
||||
internal data class ResultPayload(val toolCallId: String, val result: String?)
|
||||
internal data class ResultPayload(
|
||||
val toolCallId: String,
|
||||
val toolName: String? = null,
|
||||
val result: String?,
|
||||
)
|
||||
|
||||
@kotlinx.serialization.Serializable
|
||||
internal data class ErrorPayload(val message: String, val code: String?)
|
||||
|
||||
+1
-1
@@ -12,7 +12,7 @@ import kotlin.test.assertTrue
|
||||
import kotlin.time.Duration
|
||||
import kotlin.time.Instant
|
||||
|
||||
class KsqliteConversationStoreTest {
|
||||
class KsqliteMutableConversationStoreTest {
|
||||
private lateinit var stores: KsqliteStores
|
||||
|
||||
@BeforeTest
|
||||
+1
-1
@@ -79,7 +79,7 @@ class SchemaMigrationTest {
|
||||
// constructor internal, тест в том же модуле и может его звать.
|
||||
val stores = KsqliteStores(
|
||||
connection = conn,
|
||||
conversations = KsqliteConversationStore(conn),
|
||||
conversations = KsqliteMutableConversationStore(conn),
|
||||
messages = KsqliteMessageStore(conn),
|
||||
workingMemory = KsqliteWorkingMemoryStore(conn),
|
||||
reflections = KsqliteReflectionStore(conn),
|
||||
|
||||
Reference in New Issue
Block a user