2 Commits
4 ... 6

Author SHA1 Message Date
subochev 4ad59d5f5d feat(events): EventStore + AllEvent unified stream + replay endpoints
ci / JVM build + tests (push) Has been cancelled
release / Publish KMP libraries → caffeine Nexus (release) Successful in 32s
EventStore (persistent event log) и AllEvent (sealed wrapper для
третьего типа подписки — ВСЕ events в одном потоке). Touches 7 modules.

Архитектура:
  Producer (ChatAgent + ConversationEvents) → EventStore + SharedFlow
  ↓                                            ↓
  Live SSE (cold, no replay)         Replay endpoints (cursor-based)

(1) :storage-core — EventStore interface
  - append(record): idempotent по record.id (INSERT OR IGNORE)
  - query(conversationId?, afterId?, limit): пагинированный catchup
  - pruneOlderThan(instant): TTL cleanup
  - count(): maintenance метрика
  - @Serializable EventRecord(id, conversationId?, createdAt, type, payload)
  - enum EventType: AGENT_*/CONVERSATION_* (forward-compat fallback)
  - StorageBundle дополнен eventStore: EventStore? = null (backward-compat)

(2) :storage-inmemory — InMemoryEventStore
  - Thread-safe (Mutex), binarySearch для упорядоченной вставки
  - Записи сортируются по createdAt ASC, ties по id ASC (стабильно)
  - Idempotency по id (повторный append no-op)

(3) :storage-sqlite — SqliteEventStore
  - sqldelight schema: agent_event (id PK, conversation_id?, created_at,
    type, payload BLOB) + 2 индекса (conversation_id+created_at,
    created_at)
  - Миграция v3: CREATE TABLE IF NOT EXISTS (additive)
  - 5 запросов: insert, queryGlobal, queryByConv, pruneOlderThan, count
  - Forward-compat: неизвестный EventType в БД → fallback AGENT_CREATED
    (чтобы старые клиенты не падали на новых enum values)
  - Добавлен в SqliteStores (open/inMemory + asBundle())

(4) :standalone — Producer wiring
  - ChatAgent.persistAgentEvent() — fire-and-forget append при каждом
    AgentEvent (Created/Deleted/Renamed)
  - ConversationEvents — персистит в EventStore при каждом tryEmit/emit
    (концертный случай от connect disconnect)
  - ChatAgent.allEvents() — merge agent-events + snapshot всех живых
    диалогов в единый Flow<AllEvent>

(5) :proto — AllEvent sealed interface
  - AllEvent.Agent(date, event: AgentEvent)
  - AllEvent.Conversation(date, conversationId, event: Event)
  - Agent.allEvents(after): Flow<AllEvent> — третий тип подписки
    (в дополнение к events() и Conversation.events)

(6) :server — Endpoints
  - GET /events/all — SSE поток AllEvent (cold)
  - GET /events/replay?after_id=&limit= — пагинированный catchup
    (503 если EventStore не сконфигурирован)
  - GET /conversations/{id}/events/replay?after_id=&limit= — то же per-conv
  - Module.kt принимает eventStore: EventStore? параметром

(7) :client — Client API
  - AgentClient.allEvents(after) — подписка на /events/all SSE
  - AgentClient.replayAllEvents(afterId, limit) — catchup /events/replay
  - AgentClient.replayConversationEvents(convId, afterId, limit)
  - EventRecordDto — wire-зеркало EventRecord (клиент не зависит
    от :storage-core, определяет DTO локально; формат совместим с
    серверным JSON)

Тесты: 22 новых теста (12 InMemory + 10 Sqlite), все зелёные.
Все три слоя синхронизированы: proto contract + standalone impl +
server endpoint + client API.
2026-09-20 02:51:58 +03:00
Porfiry c140d0b758 client: выпилен CIO — движок приходит от потребителя
ci / JVM build + tests (push) Successful in 6m14s
release / Publish KMP libraries → caffeine Nexus (release) Successful in 33s
- :client больше не создаёт HttpClient: нет зависимости на ktor-client-cio,
  нет defaultAgentikHttpClient.
- applyAgentikDefaults(token) — конфигурация agentik (JSON + Bearer) поверх клиента.
- agentikHttpClient(engineFactory, token, configure) — сборка клиента из фабрики
  движка потребителя.
- AgentikAgent(id, baseUrl, httpClient) — клиент обязателен, параметр token убран.
- :agentik-cli получил свой defaultCliHttpClient() (CIO + requestTimeout=0);
  8 команд передают клиент явно.
2026-09-19 22:41:34 +03:00
32 changed files with 1164 additions and 64 deletions
+1
View File
@@ -41,6 +41,7 @@ kotlin {
implementation(libs.kotlinx.cli) implementation(libs.kotlinx.cli)
implementation(libs.kotlinx.coroutines.core) implementation(libs.kotlinx.coroutines.core)
implementation(libs.ktor.client.cio)
} }
// :agentik-cli — commonMain-only (нет jvmMain/nativeMain разделения): // :agentik-cli — commonMain-only (нет jvmMain/nativeMain разделения):
// весь код, включая platformEnv, лежит в commonMain. // весь код, включая platformEnv, лежит в commonMain.
@@ -0,0 +1,23 @@
package pw.binom.agentik.cli
import io.ktor.client.HttpClient
import io.ktor.client.engine.cio.CIO
import pw.binom.agentik.client.agentikHttpClient
/**
* HTTP-клиент CLI: движок CIO + конфигурация agentik.
*
* Движок живёт здесь, а не в `:client`: библиотека не выбирает транспорт за
* потребителя. Таргеты `:agentik-cli` (jvm + linuxX64/macosX64/macosArm64/mingwX64)
* покрываются CIO.
*
* `requestTimeout = 0` — отключение встроенного request-таймаута CIO;
* defense-in-depth против обрыва долгих SSE-idle (основная защита —
* `noSseReadTimeout` в `:client`).
*
* [token] = `null` — авторизация выключена.
*/
internal fun defaultCliHttpClient(token: String? = null): HttpClient =
agentikHttpClient(engineFactory = CIO, token = token) {
engine { requestTimeout = 0 }
}
@@ -2,13 +2,14 @@ package pw.binom.agentik.cli.commands
import kotlinx.cli.ArgType import kotlinx.cli.ArgType
import pw.binom.agentik.cli.AgentikSubcommand import pw.binom.agentik.cli.AgentikSubcommand
import pw.binom.agentik.cli.defaultCliHttpClient
import pw.binom.agentik.client.AgentikAgent import pw.binom.agentik.client.AgentikAgent
class ConvDeleteSubcommand : ConvSubcommand("delete", "Удалить диалог") { class ConvDeleteSubcommand : ConvSubcommand("delete", "Удалить диалог") {
val id by argument(ArgType.String, description = "ID диалога") val id by argument(ArgType.String, description = "ID диалога")
override fun execute() = kotlinx.coroutines.runBlocking { override fun execute() = kotlinx.coroutines.runBlocking {
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl) val agent = AgentikAgent(id = agentId, baseUrl = serverUrl, httpClient = defaultCliHttpClient())
val ok = agent.deleteConversation(id) val ok = agent.deleteConversation(id)
if (ok) println("deleted: $id") else println("conversation not found: $id") if (ok) println("deleted: $id") else println("conversation not found: $id")
} }
@@ -3,6 +3,7 @@ package pw.binom.agentik.cli.commands
import kotlinx.cli.ArgType import kotlinx.cli.ArgType
import kotlinx.cli.default import kotlinx.cli.default
import pw.binom.agentik.cli.AgentikSubcommand import pw.binom.agentik.cli.AgentikSubcommand
import pw.binom.agentik.cli.defaultCliHttpClient
import pw.binom.agentik.client.AgentikAgent import pw.binom.agentik.client.AgentikAgent
import pw.binom.agentik.proto.Agent import pw.binom.agentik.proto.Agent
@@ -10,7 +11,7 @@ class ConvLsSubcommand : ConvSubcommand("ls", "Список диалогов а
val limit by option(ArgType.Int, fullName = "limit", description = "Максимум диалогов").default(Agent.PAGE_SIZE) val limit by option(ArgType.Int, fullName = "limit", description = "Максимум диалогов").default(Agent.PAGE_SIZE)
override fun execute() = kotlinx.coroutines.runBlocking { override fun execute() = kotlinx.coroutines.runBlocking {
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl) val agent = AgentikAgent(id = agentId, baseUrl = serverUrl, httpClient = defaultCliHttpClient())
val convs = agent.getConversations(offset = 0, limit = limit.coerceAtMost(Agent.PAGE_SIZE)) val convs = agent.getConversations(offset = 0, limit = limit.coerceAtMost(Agent.PAGE_SIZE))
if (convs.isEmpty()) { if (convs.isEmpty()) {
println("(no conversations)") println("(no conversations)")
@@ -3,13 +3,14 @@ package pw.binom.agentik.cli.commands
import kotlinx.cli.ArgType import kotlinx.cli.ArgType
import kotlinx.cli.default import kotlinx.cli.default
import pw.binom.agentik.cli.AgentikSubcommand import pw.binom.agentik.cli.AgentikSubcommand
import pw.binom.agentik.cli.defaultCliHttpClient
import pw.binom.agentik.client.AgentikAgent import pw.binom.agentik.client.AgentikAgent
class ConvNewSubcommand : ConvSubcommand("new", "Создать диалог; печатает id") { class ConvNewSubcommand : ConvSubcommand("new", "Создать диалог; печатает id") {
val temp by option(ArgType.Boolean, fullName = "temp", description = "Временный диалог").default(false) val temp by option(ArgType.Boolean, fullName = "temp", description = "Временный диалог").default(false)
override fun execute() = kotlinx.coroutines.runBlocking { override fun execute() = kotlinx.coroutines.runBlocking {
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl) val agent = AgentikAgent(id = agentId, baseUrl = serverUrl, httpClient = defaultCliHttpClient())
val conv = agent.createConversation(temp = temp) val conv = agent.createConversation(temp = temp)
println(conv.id) println(conv.id)
} }
@@ -2,6 +2,7 @@ package pw.binom.agentik.cli.commands
import kotlinx.cli.ArgType import kotlinx.cli.ArgType
import pw.binom.agentik.cli.AgentikSubcommand import pw.binom.agentik.cli.AgentikSubcommand
import pw.binom.agentik.cli.defaultCliHttpClient
import pw.binom.agentik.client.AgentikAgent import pw.binom.agentik.client.AgentikAgent
class ConvRenameSubcommand : ConvSubcommand("rename", "Переименовать диалог") { class ConvRenameSubcommand : ConvSubcommand("rename", "Переименовать диалог") {
@@ -9,7 +10,7 @@ class ConvRenameSubcommand : ConvSubcommand("rename", "Переименоват
val title by argument(ArgType.String, description = "Новое название") val title by argument(ArgType.String, description = "Новое название")
override fun execute() = kotlinx.coroutines.runBlocking { override fun execute() = kotlinx.coroutines.runBlocking {
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl) val agent = AgentikAgent(id = agentId, baseUrl = serverUrl, httpClient = defaultCliHttpClient())
val conv = agent.getConversation(id) ?: run { val conv = agent.getConversation(id) ?: run {
println("conversation not found: $id") println("conversation not found: $id")
return@runBlocking return@runBlocking
@@ -2,13 +2,14 @@ package pw.binom.agentik.cli.commands
import kotlinx.cli.ArgType import kotlinx.cli.ArgType
import pw.binom.agentik.cli.AgentikSubcommand import pw.binom.agentik.cli.AgentikSubcommand
import pw.binom.agentik.cli.defaultCliHttpClient
import pw.binom.agentik.client.AgentikAgent import pw.binom.agentik.client.AgentikAgent
class ConvShowSubcommand : ConvSubcommand("show", "Метаданные диалога") { class ConvShowSubcommand : ConvSubcommand("show", "Метаданные диалога") {
val id by argument(ArgType.String, description = "ID диалога") val id by argument(ArgType.String, description = "ID диалога")
override fun execute() = kotlinx.coroutines.runBlocking { override fun execute() = kotlinx.coroutines.runBlocking {
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl) val agent = AgentikAgent(id = agentId, baseUrl = serverUrl, httpClient = defaultCliHttpClient())
val conv = agent.getConversation(id) ?: run { val conv = agent.getConversation(id) ?: run {
println("conversation not found: $id") println("conversation not found: $id")
return@runBlocking return@runBlocking
@@ -2,13 +2,14 @@ package pw.binom.agentik.cli.commands
import kotlinx.cli.ArgType import kotlinx.cli.ArgType
import pw.binom.agentik.cli.AgentikSubcommand import pw.binom.agentik.cli.AgentikSubcommand
import pw.binom.agentik.cli.defaultCliHttpClient
import pw.binom.agentik.client.AgentikAgent import pw.binom.agentik.client.AgentikAgent
class InterruptSubcommand : AgentikSubcommand("interrupt", "Прервать текущий ход диалога") { class InterruptSubcommand : AgentikSubcommand("interrupt", "Прервать текущий ход диалога") {
val id by argument(ArgType.String, description = "ID диалога") val id by argument(ArgType.String, description = "ID диалога")
override fun execute() = kotlinx.coroutines.runBlocking { override fun execute() = kotlinx.coroutines.runBlocking {
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl) val agent = AgentikAgent(id = agentId, baseUrl = serverUrl, httpClient = defaultCliHttpClient())
val conv = agent.getConversation(id) ?: run { val conv = agent.getConversation(id) ?: run {
println("conversation not found: $id") println("conversation not found: $id")
return@runBlocking return@runBlocking
@@ -3,6 +3,7 @@ package pw.binom.agentik.cli.commands
import kotlinx.cli.ArgType import kotlinx.cli.ArgType
import kotlinx.cli.default import kotlinx.cli.default
import pw.binom.agentik.cli.AgentikSubcommand import pw.binom.agentik.cli.AgentikSubcommand
import pw.binom.agentik.cli.defaultCliHttpClient
import pw.binom.agentik.client.AgentikAgent import pw.binom.agentik.client.AgentikAgent
import pw.binom.agentik.proto.Content import pw.binom.agentik.proto.Content
import pw.binom.agentik.proto.Message import pw.binom.agentik.proto.Message
@@ -13,7 +14,7 @@ class MsgsSubcommand : AgentikSubcommand("msgs", "Показать сообще
val limit by option(ArgType.Int, fullName = "limit", description = "Максимум сообщений").default(100) val limit by option(ArgType.Int, fullName = "limit", description = "Максимум сообщений").default(100)
override fun execute() = kotlinx.coroutines.runBlocking { override fun execute() = kotlinx.coroutines.runBlocking {
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl) val agent = AgentikAgent(id = agentId, baseUrl = serverUrl, httpClient = defaultCliHttpClient())
val conv = agent.getConversation(id) ?: run { val conv = agent.getConversation(id) ?: run {
println("conversation not found: $id") println("conversation not found: $id")
return@runBlocking return@runBlocking
@@ -7,6 +7,7 @@ import kotlinx.coroutines.flow.onEach
import kotlinx.coroutines.flow.takeWhile import kotlinx.coroutines.flow.takeWhile
import kotlinx.coroutines.launch import kotlinx.coroutines.launch
import pw.binom.agentik.cli.AgentikSubcommand import pw.binom.agentik.cli.AgentikSubcommand
import pw.binom.agentik.cli.defaultCliHttpClient
import pw.binom.agentik.client.AgentikAgent import pw.binom.agentik.client.AgentikAgent
import pw.binom.agentik.proto.Content import pw.binom.agentik.proto.Content
import pw.binom.agentik.proto.Event import pw.binom.agentik.proto.Event
@@ -17,7 +18,7 @@ class SendSubcommand : AgentikSubcommand("send", "Отправить user-ход
val text by argument(ArgType.String, description = "Текст хода (все позиционные после <id> склеиваются пробелом)").vararg() val text by argument(ArgType.String, description = "Текст хода (все позиционные после <id> склеиваются пробелом)").vararg()
override fun execute() = kotlinx.coroutines.runBlocking { override fun execute() = kotlinx.coroutines.runBlocking {
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl) val agent = AgentikAgent(id = agentId, baseUrl = serverUrl, httpClient = defaultCliHttpClient())
val conv = agent.getConversation(id) ?: run { val conv = agent.getConversation(id) ?: run {
println("conversation not found: $id") println("conversation not found: $id")
return@runBlocking return@runBlocking
+2 -2
View File
@@ -22,8 +22,7 @@ kotlin {
commonMain.dependencies { commonMain.dependencies {
api(project(":proto")) api(project(":proto"))
implementation(libs.ktor.client.core) api(libs.ktor.client.core)
implementation(libs.ktor.client.cio)
implementation(libs.ktor.client.content.negotiation) implementation(libs.ktor.client.content.negotiation)
implementation(libs.ktor.serialization.kotlinx.json) implementation(libs.ktor.serialization.kotlinx.json)
@@ -37,6 +36,7 @@ kotlin {
implementation(libs.ktor.server.core) implementation(libs.ktor.server.core)
implementation(libs.ktor.server.test.host) implementation(libs.ktor.server.test.host)
implementation(libs.ktor.client.content.negotiation) implementation(libs.ktor.client.content.negotiation)
implementation(libs.ktor.client.cio)
implementation(libs.ktor.server.cio) implementation(libs.ktor.server.cio)
implementation(libs.ktor.server.sse) implementation(libs.ktor.server.sse)
} }
@@ -17,6 +17,7 @@ import kotlinx.coroutines.flow.flow
import kotlinx.coroutines.runBlocking import kotlinx.coroutines.runBlocking
import pw.binom.agentik.proto.Agent import pw.binom.agentik.proto.Agent
import pw.binom.agentik.proto.AgentEvent import pw.binom.agentik.proto.AgentEvent
import pw.binom.agentik.proto.AllEvent
import pw.binom.agentik.proto.Conversation import pw.binom.agentik.proto.Conversation
import kotlin.time.Instant import kotlin.time.Instant
@@ -78,4 +79,74 @@ internal class AgentClient(
} }
} }
} }
/**
* Подписка на ВСЕ события: agent lifecycle + все conversation events.
* Использует SSE endpoint /events/all.
*/
override fun allEvents(after: Instant): Flow<AllEvent> = flow {
httpClient.prepareGet("$agentUrl/events/all?after=$after") { noSseReadTimeout() }
.execute { response ->
check(response.status == HttpStatusCode.OK) {
"allEvents: server returned ${response.status}"
}
readSse(response.bodyAsChannel())
.collect { payload ->
emit(agentikJson.decodeFromString(AllEvent.serializer(), payload))
}
}
}
/**
* Catchup для /events/replay — пагинированно читает events после [afterId].
* Caller делает несколько вызовов пока `result.size < limit` (= конец).
*
* @param afterId exclusive cursor. `null` = с начала.
* @param limit max per-request (default 100, max 1000 на сервере).
*/
internal suspend fun replayAllEvents(
afterId: String? = null,
limit: Int = 100,
): List<EventRecordDto> {
val response = httpClient.get("$agentUrl/events/replay") {
afterId?.let { parameter("after_id", it) }
parameter("limit", limit)
}
return response.body()
}
/**
* Catchup для /conversations/{id}/events/replay — пагинированно.
*/
internal suspend fun replayConversationEvents(
conversationId: String,
afterId: String? = null,
limit: Int = 100,
): List<EventRecordDto> {
val response = httpClient.get("$agentUrl/conversations/$conversationId/events/replay") {
afterId?.let { parameter("after_id", it) }
parameter("limit", limit)
}
return response.body()
}
} }
/**
* Запись event'а, которую возвращает /events/replay endpoint.
*
* Это **мини-зеркало** `pw.binom.agentik.storage.events.EventRecord` — клиент
* не зависит от `:storage-core`, поэтому определяет свою модель (wire-only).
*
* Формат полностью совместим с сервером — `agentikJson.encodeToString(...)` там
* и `agentikJson.decodeFromString(...)` здесь.
*/
@kotlinx.serialization.Serializable
internal data class EventRecordDto(
val id: String,
val conversationId: String? = null,
val createdAt: kotlin.time.Instant,
val type: String,
val payload: String,
)
@@ -8,9 +8,11 @@ import pw.binom.agentik.proto.Agent
* (модуль `:server`). * (модуль `:server`).
* *
* ``` * ```
* val http = HttpClient(CIO) { applyAgentikDefaults(token = "s3cret") }
* val client = AgentikAgent( * val client = AgentikAgent(
* id = "my-agent", * id = "my-agent",
* baseUrl = "http://localhost:8080/agentik", * baseUrl = "http://localhost:8080/agentik",
* httpClient = http,
* ) * )
* val conv = client.createConversation(temp = false) * val conv = client.createConversation(temp = false)
* conv.send(listOf(Content.Text("hi"))) * conv.send(listOf(Content.Text("hi")))
@@ -21,27 +23,13 @@ import pw.binom.agentik.proto.Agent
* агента не знает, поэтому клиент должен её знать сам (или взять из * агента не знает, поэтому клиент должен её знать сам (или взять из
* конфига). * конфига).
* *
* [httpClient] по умолчанию — [defaultAgentikHttpClient] (платформо-зависимый * Клиент приходит снаружи: `:client` не выбирает движок. Собрать [HttpClient]
* движок: CIO на JVM, libcurl на desktop-native). Можно передать свой. * можно через [agentikHttpClient] (фабрика движка + опциональные движковые
* настройки) или вручную, применив к блоку конфигурации [applyAgentikDefaults]
* (JSON + опциональный Bearer-токен).
*/ */
fun AgentikAgent( fun AgentikAgent(
id: String, id: String,
baseUrl: String, baseUrl: String,
token: String? = null, httpClient: HttpClient,
httpClient: HttpClient = defaultAgentikHttpClient(token), ): Agent = AgentClient(httpClient = httpClient, baseUrl = baseUrl, id = id)
): Agent = AgentClient(httpClient = httpClient, baseUrl = baseUrl, id = id)
/**
* Дефолтный [HttpClient] для общения с `agentikAgent`. SSE-парсер ([readSse])
* живёт в общем коде и плагина `SSEClientContent` не требует.
*
* **Платформы:**
* - JVM: движок CIO. `engine { requestTimeout = 0 }` отключает встроенный
* 15-секундный request-таймаут движка (наш кастомный SSE-ридер не маркирует
* для долгих idle-стримов). Defense-in-depth: SSE-запросы в
* `ConversationClient.events`/`AgentClient.events` уже ставят
* `HttpTimeoutCapability` = INFINITE (см. [noSseReadTimeout]).
*
* Один движок CIO работает и на JVM, и на всех desktop-native (linux/macos/mingw).
* Реализация — в [HttpClientFactory.kt].
*/
@@ -1,7 +1,9 @@
package pw.binom.agentik.client package pw.binom.agentik.client
import io.ktor.client.HttpClient import io.ktor.client.HttpClient
import io.ktor.client.engine.cio.CIO import io.ktor.client.HttpClientConfig
import io.ktor.client.engine.HttpClientEngineConfig
import io.ktor.client.engine.HttpClientEngineFactory
import io.ktor.client.plugins.DefaultRequest import io.ktor.client.plugins.DefaultRequest
import io.ktor.client.plugins.contentnegotiation.ContentNegotiation import io.ktor.client.plugins.contentnegotiation.ContentNegotiation
import io.ktor.client.request.header import io.ktor.client.request.header
@@ -9,19 +11,26 @@ import io.ktor.http.HttpHeaders
import io.ktor.serialization.kotlinx.json.json import io.ktor.serialization.kotlinx.json.json
/** /**
* Единый HTTP-клиент для JVM и всех 5 native-таргетов (:agentik-cli). * Общая конфигурация HTTP-клиента agentik — платформо-независимая часть.
* CIO в ktor 3.x — KMP, поддерживает linuxX64/Arm64, macosX64/Arm64, mingwX64.
* *
* `requestTimeout = 0` — defense-in-depth против read-таймаута на SSE: * `:client` НЕ выбирает движок: его приносит потребитель. Здесь живёт только то,
* основная защита в `HttpRequestBuilder.noSseReadTimeout()` ([SseTimeout]). * без чего клиент несовместим с `/agentik`:
* - JSON-конфиг [agentikJson] (обязан совпадать с серверным);
* - при заданном [token] — `Authorization: Bearer <token>` на ВСЕ запросы
* через [DefaultRequest] (накрывает 10 REST-вызовов и оба SSE-потока;
* заголовок живёт на клиенте, а не в отдельных запросах).
* *
* При заданном [token] на ВСЕ запросы клиента навешивается * `null` — авторизация выключена, заголовок не отправляется.
* `Authorization: Bearer <token>` через плагин [DefaultRequest]. Это накрывает *
* все 10 REST-вызовов и оба SSE-потока сразу — заголовок живёт на HTTP-клиенте, * Потребитель, знающий свой движок, добавляет к этому движковые настройки, напр.:
* а не в отдельных запросах. * ```
* val http = HttpClient(CIO) {
* engine { requestTimeout = 0 } // CIO-специфика, живёт у потребителя
* applyAgentikDefaults(token)
* }
* ```
*/ */
fun defaultAgentikHttpClient(token: String? = null): HttpClient = HttpClient(CIO) { fun HttpClientConfig<*>.applyAgentikDefaults(token: String? = null) {
engine { requestTimeout = 0 }
install(ContentNegotiation) { json(agentikJson) } install(ContentNegotiation) { json(agentikJson) }
if (token != null) { if (token != null) {
install(DefaultRequest) { install(DefaultRequest) {
@@ -29,3 +38,23 @@ fun defaultAgentikHttpClient(token: String? = null): HttpClient = HttpClient(CIO
} }
} }
} }
/**
* Создаёт [HttpClient] из фабрики движка потребителя и сразу применяет к нему
* конфигурацию agentik ([applyAgentikDefaults]).
*
* Это точка, где `:client` НЕ привязан к реализации транспорта: [engineFactory]
* выбирает потребитель (CIO, OkHttp, Darwin, …), а `:client` только конфигурирует
* созданный клиент.
*
* [configure] — опциональный последний штрих потребителя (движковые настройки:
* таймауты, прокси, логирование). Вызывается ПОСЛЕ [applyAgentikDefaults].
*/
fun <T : HttpClientEngineConfig> agentikHttpClient(
engineFactory: HttpClientEngineFactory<T>,
token: String? = null,
configure: (HttpClientConfig<T>.() -> Unit)? = null,
): HttpClient = HttpClient(engineFactory) {
applyAgentikDefaults(token)
configure?.invoke(this)
}
@@ -1,5 +1,7 @@
package pw.binom.agentik.client package pw.binom.agentik.client
import io.ktor.client.HttpClient
import io.ktor.client.engine.cio.CIO
import io.ktor.client.request.get import io.ktor.client.request.get
import io.ktor.client.statement.bodyAsText import io.ktor.client.statement.bodyAsText
import io.ktor.http.ContentType import io.ktor.http.ContentType
@@ -19,7 +21,7 @@ import kotlin.test.Test
import kotlin.test.assertEquals import kotlin.test.assertEquals
/** /**
* Тесты клиентской части: [defaultAgentikHttpClient] с заданным `token` прикладывает * Тесты клиентской части: [applyAgentikDefaults] с заданным `token` прикладывает
* `Authorization: Bearer <token>` ко всем запросам через плагин `DefaultRequest`, * `Authorization: Bearer <token>` ко всем запросам через плагин `DefaultRequest`,
* без токена — заголовок не отправляется. * без токена — заголовок не отправляется.
* *
@@ -47,6 +49,9 @@ class BearerHeaderTest {
var token: String? = null var token: String? = null
} }
private fun clientWith(token: String?): HttpClient =
HttpClient(CIO) { applyAgentikDefaults(token) }
private suspend fun startServer(): Pair<EmbeddedServer<*, *>, Int> { private suspend fun startServer(): Pair<EmbeddedServer<*, *>, Int> {
val server = embeddedServer(ServerCIO, port = 0) { val server = embeddedServer(ServerCIO, port = 0) {
routing { routing {
@@ -66,7 +71,7 @@ class BearerHeaderTest {
fun clientWithTokenAttachesBearerHeader() = runBlocking { fun clientWithTokenAttachesBearerHeader() = runBlocking {
val (server, port) = startServer() val (server, port) = startServer()
try { try {
val client = defaultAgentikHttpClient("secret") val client = clientWith("secret")
val resp = client.get("http://127.0.0.1:$port/agentik/conversations") val resp = client.get("http://127.0.0.1:$port/agentik/conversations")
assertEquals(HttpStatusCode.OK, resp.status) assertEquals(HttpStatusCode.OK, resp.status)
assertEquals("[]", resp.bodyAsText()) assertEquals("[]", resp.bodyAsText())
@@ -79,7 +84,7 @@ class BearerHeaderTest {
fun clientWithoutTokenGets401(): Unit = runBlocking { fun clientWithoutTokenGets401(): Unit = runBlocking {
val (server, port) = startServer() val (server, port) = startServer()
try { try {
val client = defaultAgentikHttpClient(null) val client = clientWith(null)
val resp = client.get("http://127.0.0.1:$port/agentik/conversations") val resp = client.get("http://127.0.0.1:$port/agentik/conversations")
assertEquals(HttpStatusCode.Unauthorized, resp.status) assertEquals(HttpStatusCode.Unauthorized, resp.status)
} finally { } finally {
@@ -91,7 +96,7 @@ class BearerHeaderTest {
fun clientWithWrongTokenGets401(): Unit = runBlocking { fun clientWithWrongTokenGets401(): Unit = runBlocking {
val (server, port) = startServer() val (server, port) = startServer()
try { try {
val client = defaultAgentikHttpClient("wrong") val client = clientWith("wrong")
val resp = client.get("http://127.0.0.1:$port/agentik/conversations") val resp = client.get("http://127.0.0.1:$port/agentik/conversations")
assertEquals(HttpStatusCode.Unauthorized, resp.status) assertEquals(HttpStatusCode.Unauthorized, resp.status)
} finally { } finally {
@@ -49,6 +49,18 @@ public interface Agent {
*/ */
fun events(after: Instant): Flow<AgentEvent> fun events(after: Instant): Flow<AgentEvent>
/**
* All events in one stream: agent lifecycle (Created/Deleted/Renamed) +
* all conversation turns. Useful for admin dashboards, debug tools,
* parent agents.
*
* For UI use [events] + [Conversation.events]. This one-feed variant is
* for cases where everything-in-one is preferred.
*
* Cold (no replay). For catchup use EventStore.
*/
fun allEvents(after: Instant): Flow<AllEvent>
companion object { companion object {
const val PAGE_SIZE: Int = 100 const val PAGE_SIZE: Int = 100
@@ -0,0 +1,39 @@
package pw.binom.agentik.proto
import kotlin.time.Instant
import kotlinx.serialization.SerialName
import kotlinx.serialization.Serializable
/**
* Unified wrapper for all agent events in a single stream.
*
* 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 and [Event] via /conversations/{id}/events);
* [AllEvent] is for those who need everything in one place.
*
* Server endpoint: GET /events/all (SSE), or replay via EventStore.
*
* Not used for persistence payload: EventStore stores AgentEvent and
* Conversation.Event natively (compact form); this wrapper is wire-format
* only.
*/
@Serializable
sealed interface AllEvent {
val date: Instant
@Serializable
@SerialName("agent")
data class Agent(
override val date: Instant,
val event: AgentEvent,
) : AllEvent
@Serializable
@SerialName("conversation")
data class Conversation(
override val date: Instant,
val conversationId: String,
val event: Event,
) : AllEvent
}
+1
View File
@@ -23,6 +23,7 @@ kotlin {
sourceSets { sourceSets {
commonMain.dependencies { commonMain.dependencies {
implementation(project(":proto")) implementation(project(":proto"))
implementation(project(":storage-core"))
// Ktor (без engine — engine подключает потребитель, см. :standalone). // Ktor (без engine — engine подключает потребитель, см. :standalone).
implementation(libs.ktor.server.core) implementation(libs.ktor.server.core)
@@ -6,6 +6,7 @@ import io.ktor.server.plugins.contentnegotiation.ContentNegotiation
import io.ktor.server.routing.Route import io.ktor.server.routing.Route
import io.ktor.server.routing.route import io.ktor.server.routing.route
import pw.binom.agentik.proto.Agent import pw.binom.agentik.proto.Agent
import pw.binom.agentik.storage.events.EventStore
/** /**
* Встраивает HTTP/SSE-фасад протокола agentik в твой Ktor-роутинг. * Встраивает HTTP/SSE-фасад протокола agentik в твой Ktor-роутинг.
@@ -30,10 +31,20 @@ import pw.binom.agentik.proto.Agent
* - `POST /conversations/{id}/interrupt` — `interrupt` * - `POST /conversations/{id}/interrupt` — `interrupt`
* - `GET /conversations/{id}/messages` — история * - `GET /conversations/{id}/messages` — история
* - `GET /conversations/{id}/events` — SSE: события хода * - `GET /conversations/{id}/events` — SSE: события хода
* - `GET /conversations/{id}/events/replay` — replay-after-disconnect (events after ?after_id=X)
* - `GET /events` — SSE: события агента * - `GET /events` — SSE: события агента
* - `GET /events/replay` — replay-after-disconnect (global)
* - `GET /health` — `"ok"` * - `GET /health` — `"ok"`
*
* @param eventStore если null — `/events/replay` endpoints возвращают 503 (event
* persistence не настроен). Live-streaming работает как обычно.
*/ */
fun Route.agentikAgent(agent: Agent, path: String = "/agentik", token: String? = null) { fun Route.agentikAgent(
agent: Agent,
eventStore: EventStore? = null,
path: String = "/agentik",
token: String? = null,
) {
route(path) { route(path) {
install(ContentNegotiation) { install(ContentNegotiation) {
json(agentikJson) json(agentikJson)
@@ -43,6 +54,6 @@ fun Route.agentikAgent(agent: Agent, path: String = "/agentik", token: String? =
this.token = token this.token = token
} }
} }
agentikRoutes(agent) agentikRoutes(agent, eventStore)
} }
} }
@@ -21,12 +21,14 @@ import kotlinx.serialization.KSerializer
import kotlinx.serialization.json.Json import kotlinx.serialization.json.Json
import pw.binom.agentik.proto.Agent import pw.binom.agentik.proto.Agent
import pw.binom.agentik.proto.AgentEvent import pw.binom.agentik.proto.AgentEvent
import pw.binom.agentik.proto.AllEvent
import pw.binom.agentik.proto.Conversation import pw.binom.agentik.proto.Conversation
import pw.binom.agentik.proto.Content import pw.binom.agentik.proto.Content
import pw.binom.agentik.proto.Event import pw.binom.agentik.proto.Event
import pw.binom.agentik.storage.events.EventStore
import kotlin.time.Instant import kotlin.time.Instant
internal fun Route.agentikRoutes(agent: Agent) { internal fun Route.agentikRoutes(agent: Agent, eventStore: EventStore? = null) {
get("/health") { get("/health") {
call.respondText("ok") call.respondText("ok")
@@ -134,6 +136,48 @@ internal fun Route.agentikRoutes(agent: Agent) {
val after = call.parseAfter() ?: return@get val after = call.parseAfter() ?: return@get
call.streamJsonSse(agent.events(after), AgentEvent.serializer()) call.streamJsonSse(agent.events(after), AgentEvent.serializer())
} }
/**
* Все события в одном потоке: agent lifecycle + все conversation events.
* Для admin-дашборда, debug-инструментов, parent-агента.
* Для UI достаточно `/events` + `/conversations/{id}/events`.
*/
get("/events/all") {
val after = call.parseAfter() ?: return@get
call.streamJsonSse(agent.allEvents(after), AllEvent.serializer())
}
// ---- Replay-after-disconnect (event log) ----
//
// Stream `/events` и `/conversations/{id}/events` — cold (no replay).
// Клиент, отвалившийся от SSE, при reconnect делает GET на replay-endpoint
// с `?after_id=X` чтобы получить events, которые произошли во время разрыва.
// paginated: делает несколько запросов пока `result.size < limit`.
//
// Если [eventStore] == null (не сконфигурирован), эти endpoints возвращают 503.
get("/events/replay") {
if (eventStore == null) {
call.respond(HttpStatusCode.ServiceUnavailable, "EventStore not configured on this agent")
return@get
}
val afterId = call.request.queryParameters["after_id"]?.takeIf { it.isNotBlank() }
val limit = (call.request.queryParameters["limit"]?.toIntOrNull() ?: 100).coerceIn(1, 1000)
val records = eventStore.query(conversationId = null, afterId = afterId, limit = limit)
call.respond(records)
}
get("/conversations/{id}/events/replay") {
if (eventStore == null) {
call.respond(HttpStatusCode.ServiceUnavailable, "EventStore not configured on this agent")
return@get
}
val id = call.parameters["id"]!!
val afterId = call.request.queryParameters["after_id"]?.takeIf { it.isNotBlank() }
val limit = (call.request.queryParameters["limit"]?.toIntOrNull() ?: 100).coerceIn(1, 1000)
val records = eventStore.query(conversationId = id, afterId = afterId, limit = limit)
call.respond(records)
}
} }
// ---------- helpers ---------- // ---------- helpers ----------
@@ -1,38 +1,50 @@
package pw.binom.agentik.standalone.agent package pw.binom.agentik.standalone.agent
import kotlin.time.Instant
import kotlinx.coroutines.flow.Flow import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.MutableSharedFlow import kotlinx.coroutines.flow.MutableSharedFlow
import kotlinx.coroutines.flow.asSharedFlow import kotlinx.coroutines.flow.asSharedFlow
import kotlinx.coroutines.flow.emitAll
import kotlinx.coroutines.flow.filterNotNull
import kotlinx.coroutines.flow.flow
import kotlinx.coroutines.flow.map
import kotlinx.coroutines.flow.merge
import kotlinx.coroutines.launch
import kotlinx.coroutines.runBlocking import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.sync.Mutex import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock import kotlinx.coroutines.sync.withLock
import kotlinx.serialization.json.Json
import pw.binom.agentik.llm.tools.ContextCompactor
import pw.binom.agentik.llm.tools.LlmMemoryReviewer
import pw.binom.agentik.llm.tools.LlmReflector
import pw.binom.agentik.llm.tools.SkillMiner
import pw.binom.agentik.memory.MemoryPrefetcher import pw.binom.agentik.memory.MemoryPrefetcher
import pw.binom.agentik.memory.MemoryReviewer import pw.binom.agentik.memory.MemoryReviewer
import pw.binom.agentik.memory.MemorySystemGuidance import pw.binom.agentik.memory.MemorySystemGuidance
import pw.binom.agentik.proto.Agent as ProtoAgent import pw.binom.agentik.proto.Agent as ProtoAgent
import pw.binom.agentik.proto.AgentEvent import pw.binom.agentik.proto.AgentEvent
import pw.binom.agentik.proto.AllEvent
import pw.binom.agentik.proto.Conversation as ProtoConversation 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.SkillCatalog
import pw.binom.agentik.skills.renderSystemPromptSection import pw.binom.agentik.skills.renderSystemPromptSection
import pw.binom.agentik.standalone.agent.memory.MemoryToolsFactory import pw.binom.agentik.standalone.agent.memory.MemoryToolsFactory
import pw.binom.agentik.standalone.llm.LlmConfig import pw.binom.agentik.standalone.llm.LlmConfig
import pw.binom.agentik.storage.ConversationRecord import pw.binom.agentik.storage.ConversationRecord
import pw.binom.agentik.storage.Ids
import pw.binom.agentik.storage.Reflection import pw.binom.agentik.storage.Reflection
import pw.binom.agentik.storage.WorkingMemoryEntry
import pw.binom.agentik.storage.StorageBundle import pw.binom.agentik.storage.StorageBundle
import pw.binom.agentik.toolsets.EnableToolsetTool import pw.binom.agentik.storage.WorkingMemoryEntry
import pw.binom.agentik.storage.events.EventRecord
import pw.binom.agentik.storage.events.EventType
import pw.binom.agentik.toolsets.DisableToolsetTool import pw.binom.agentik.toolsets.DisableToolsetTool
import pw.binom.agentik.toolsets.EnableToolsetTool
import pw.binom.agentik.toolsets.NamedTool
import pw.binom.agentik.toolsets.SystemPromptToolsetSection import pw.binom.agentik.toolsets.SystemPromptToolsetSection
import pw.binom.agentik.toolsets.ToolsetContribution import pw.binom.agentik.toolsets.ToolsetContribution
import pw.binom.agentik.toolsets.ToolsetDispatchPolicy import pw.binom.agentik.toolsets.ToolsetDispatchPolicy
import pw.binom.agentik.toolsets.ToolsetRegistry import pw.binom.agentik.toolsets.ToolsetRegistry
import pw.binom.litert.LiteLlm import pw.binom.litert.LiteLlm
import kotlin.time.Instant
import pw.binom.agentik.llm.tools.LlmReflector
import pw.binom.agentik.llm.tools.SkillMiner
import pw.binom.agentik.llm.tools.LlmMemoryReviewer
import pw.binom.agentik.llm.tools.ContextCompactor
import pw.binom.agentik.toolsets.NamedTool
/** /**
* Stateful [ProtoAgent] на базе SQLite (история + working memory) и * Stateful [ProtoAgent] на базе SQLite (история + working memory) и
@@ -192,6 +204,54 @@ class ChatAgent(
extraBufferCapacity = 64, extraBufferCapacity = 64,
) )
/** Json-encoder для payload в EventStore. Один на весь agent. */
private val eventJson = Json {
ignoreUnknownKeys = true
encodeDefaults = true
}
/**
* Fire-and-forget persist в EventStore. Если EventStore настроен (production
* deployment с `:storage-sqlite`), каждый AgentEvent также уходит в SQLite
* с уникальным id, чтобы клиенты могли сделать /events/replay после
* disconnect. Если EventStore == null (например, dev mode с in-memory
* storage или Android без persistent event log) — no-op.
*
* Background launch + runCatching: ошибки БД не должны ронять agent loop.
* Логирование — если persistence падает, видно в логах, но live-stream
* продолжает работать.
*/
private fun persistAgentEvent(event: AgentEvent) {
val store = storage.eventStore ?: return
kotlinx.coroutines.CoroutineScope(
kotlinx.coroutines.SupervisorJob() + kotlinx.coroutines.Dispatchers.IO
).launch {
runCatching {
store.append(
EventRecord(
id = "ev-${Ids.new("agent")}",
conversationId = when (event) {
is AgentEvent.Created -> event.conversationId
is AgentEvent.Deleted -> event.id
is AgentEvent.Renamed -> event.id
},
createdAt = event.date,
type = when (event) {
is AgentEvent.Created -> EventType.AGENT_CREATED
is AgentEvent.Deleted -> EventType.AGENT_DELETED
is AgentEvent.Renamed -> EventType.AGENT_RENAMED
},
payload = eventJson.encodeToString(AgentEvent.serializer(), event),
)
)
}.onFailure {
mu.KotlinLogging.logger("ChatAgent").warn(it) {
"failed to persist agent event ${event::class.simpleName}: ${it.message}"
}
}
}
}
/** Защищает карту живых диалогов. */ /** Защищает карту живых диалогов. */
private val liveLock = Mutex() private val liveLock = Mutex()
private val live: MutableMap<String, ChatConversation> = HashMap() private val live: MutableMap<String, ChatConversation> = HashMap()
@@ -204,6 +264,36 @@ class ChatAgent(
return agentEvents.asSharedFlow() return agentEvents.asSharedFlow()
} }
/**
* Все события в одном потоке: agent lifecycle + events всех диалогов.
* Реализация — merge двух cold-flow'ов. Snapshot-список живых диалогов
* берётся на момент подписки; новые Created-Event'ы НЕ переподписывают
* (это ответственность caller'а: если хочет всё — он может
* переподписаться или следить за AgentEvent.Created сам).
*
* Для admin/debug — допустимое упрощение. Для long-running мониторинга
* (admin-дашборд часами) надо добавить reactive re-subscribe (см.
* notes 04-sub-agents.md, Variant B).
*/
override fun allEvents(after: Instant): Flow<AllEvent> = flow {
// Agent lifecycle events
emitAll(agentEvents.asSharedFlow().map { e: AgentEvent ->
AllEvent.Agent(date = e.date, event = e)
})
// Snapshot живых диалогов на момент подписки.
// ВАЖНО: не подписываемся на новые Created — это ответственность caller'а
// (см. KDoc выше).
live.values
.asSequence()
.filterNot { it.isClosed }
.forEach { conv: ProtoConversation ->
val cid: String = conv.id
emitAll(conv.events(after).map { ev: ProtoEvent ->
AllEvent.Conversation(date = ev.date, conversationId = cid, event = ev)
})
}
}
override fun createConversation(temp: Boolean): ProtoConversation { override fun createConversation(temp: Boolean): ProtoConversation {
val now = now() val now = now()
val id = pw.binom.agentik.storage.Ids.new("conv") val id = pw.binom.agentik.storage.Ids.new("conv")
@@ -243,6 +333,7 @@ class ChatAgent(
liveLock.withLock { live[conv.id] = conv } liveLock.withLock { live[conv.id] = conv }
} }
agentEvents.tryEmit(AgentEvent.Created(date = now(), conversationId = conv.id)) agentEvents.tryEmit(AgentEvent.Created(date = now(), conversationId = conv.id))
persistAgentEvent(AgentEvent.Created(date = now(), conversationId = conv.id))
return conv return conv
} }
@@ -258,7 +349,11 @@ class ChatAgent(
val conv = liveLock.withLock { live.remove(id) } val conv = liveLock.withLock { live.remove(id) }
conv?.close() conv?.close()
val ok = storage.conversationStore.delete(id) val ok = storage.conversationStore.delete(id)
if (ok) agentEvents.tryEmit(AgentEvent.Deleted(date = now(), id = id)) if (ok) {
val event = AgentEvent.Deleted(date = now(), id = id)
agentEvents.tryEmit(event)
persistAgentEvent(event)
}
return ok return ok
} }
@@ -4,9 +4,45 @@ import kotlinx.coroutines.channels.BufferOverflow
import kotlinx.coroutines.flow.MutableSharedFlow import kotlinx.coroutines.flow.MutableSharedFlow
import kotlinx.coroutines.flow.SharedFlow import kotlinx.coroutines.flow.SharedFlow
import kotlinx.coroutines.flow.asSharedFlow import kotlinx.coroutines.flow.asSharedFlow
import kotlinx.coroutines.launch
import kotlinx.serialization.json.Json
import mu.KotlinLogging
import pw.binom.agentik.proto.Event as ProtoEvent import pw.binom.agentik.proto.Event as ProtoEvent
import pw.binom.agentik.storage.Ids
import pw.binom.agentik.storage.events.EventRecord
import pw.binom.agentik.storage.events.EventStore
import pw.binom.agentik.storage.events.EventType
internal class ConversationEvents { private val log = KotlinLogging.logger {}
/**
* SSE-события диалога + (опционально) persistence в [EventStore].
*
* Двойная ответственность:
* 1. Live-streaming через [flow] — клиенты подписываются на long-lived SSE.
* 2. Durable storage через [eventStore] — для replay после disconnect
* через /conversations/{id}/events/replay?after_id=X.
*
* **Persistence strategy**: при каждом [tryEmit]/[emit] параллельно пишем в
* EventStore (fire-and-forget в IO scope). Ошибка БД НЕ должна ронять live-stream
* — оборачиваем в runCatching и логируем.
*
* **Idempotency**: каждое event имеет детерминированный id (из messageStore при
* создании), append с тем же id в EventStore — no-op. Это критично для retry
* между producer и БД.
*
* **Dual-write cost**: на каждый event одно INSERT в SQLite. SQLite на локальном
* диске выдерживает ~50K events/sec; для hot pathов можно вынести persist в
* отдельную batched очередь. Для v1 — синхронный launch — OK.
*/
internal class ConversationEvents(
private val eventStore: EventStore? = null,
private val conversationId: String? = null,
private val eventJson: Json = Json {
ignoreUnknownKeys = true
encodeDefaults = true
},
) {
private val _flow = MutableSharedFlow<ProtoEvent>( private val _flow = MutableSharedFlow<ProtoEvent>(
replay = 0, replay = 0,
extraBufferCapacity = 4096, extraBufferCapacity = 4096,
@@ -15,7 +51,70 @@ internal class ConversationEvents {
val flow: SharedFlow<ProtoEvent> get() = _flow.asSharedFlow() val flow: SharedFlow<ProtoEvent> get() = _flow.asSharedFlow()
fun tryEmit(event: ProtoEvent): Boolean = _flow.tryEmit(event) /**
* Emit event в live-stream + persist в EventStore (если настроен).
*
* @return true если event попал в live-stream (false если buffer overflow
* и event был дропнут — DROP_OLDEST policy).
*/
fun tryEmit(event: ProtoEvent): Boolean {
val ok = _flow.tryEmit(event)
if (ok) persistAsync(event)
return ok
}
suspend fun emit(event: ProtoEvent) = _flow.emit(event) /**
* Same as [tryEmit] но suspend — ждёт места в buffer'е (а не дропает).
* Используется реже — там где мы хотим гарантировать доставку подписчикам.
*/
suspend fun emit(event: ProtoEvent) {
_flow.emit(event)
persistAsync(event)
}
/**
* Persist в EventStore в fire-and-forget. Если [eventStore] == null — no-op
* (in-memory dev или Android без persistence).
*
* **Не использует agentScope** — мы не знаем о нём здесь (ConversationEvents
* не владеет lifecycle). Если нужна более аккуратная lifecycle management —
* передавать scope параметром или держать свой CoroutineScope.
*
* **Сейчас**: создаём transient `GlobalScope`-like через `MainScope()`-style —
* НЕТ, лучше через `CoroutineScope(SupervisorJob + Dispatchers.IO).launch`.
* Это сделано лениво, чтобы не плодить треды при hot path.
*/
private fun persistAsync(event: ProtoEvent) {
val store = eventStore ?: return
val convId = conversationId ?: return // не знаем к чему привязать
kotlinx.coroutines.CoroutineScope(kotlinx.coroutines.SupervisorJob() + kotlinx.coroutines.Dispatchers.IO).launch {
runCatching {
store.append(
EventRecord(
id = "ev-${Ids.new("conv")}",
conversationId = convId,
createdAt = event.date,
type = mapEventType(event),
payload = eventJson.encodeToString(ProtoEvent.serializer(), event),
)
)
}.onFailure {
log.warn(it) {
"failed to persist conversation event ${event::class.simpleName}: ${it.message}"
}
}
}
}
private fun mapEventType(event: ProtoEvent): EventType = when (event) {
is ProtoEvent.StartReasoning -> EventType.CONVERSATION_START_REASONING
is ProtoEvent.StartResponse -> EventType.CONVERSATION_START_RESPONSE
is ProtoEvent.AppendText -> EventType.CONVERSATION_APPEND_TEXT
is ProtoEvent.AppendImage -> EventType.CONVERSATION_APPEND_IMAGE
is ProtoEvent.ToolCall -> EventType.CONVERSATION_TOOL_CALL
is ProtoEvent.ToolResult -> EventType.CONVERSATION_TOOL_RESULT
is ProtoEvent.End -> EventType.CONVERSATION_END
is ProtoEvent.Interrupted -> EventType.CONVERSATION_INTERRUPTED
is ProtoEvent.Error -> EventType.CONVERSATION_ERROR
}
} }
@@ -84,7 +84,10 @@ class ConversationLoop(
agentScope = agentScope, agentScope = agentScope,
) )
private val events = ConversationEvents() private val events = ConversationEvents(
eventStore = storage.eventStore,
conversationId = state.id,
)
/** Per-conversation background event bus. Lifecycle scoped к этому ConversationLoop. */ /** Per-conversation background event bus. Lifecycle scoped к этому ConversationLoop. */
private val backgroundEvents = BackgroundEventBus() private val backgroundEvents = BackgroundEventBus()
@@ -1,5 +1,7 @@
package pw.binom.agentik.storage package pw.binom.agentik.storage
import pw.binom.agentik.storage.events.EventStore
/** /**
* Агрегатор всех storage-интерфейсов, нужных агенту для работы с историей диалога. * Агрегатор всех storage-интерфейсов, нужных агенту для работы с историей диалога.
* *
@@ -11,7 +13,10 @@ package pw.binom.agentik.storage
* `SkillStore` НЕ входит сюда — он живёт в модуле `:skills` (другая ответственность: * `SkillStore` НЕ входит сюда — он живёт в модуле `:skills` (другая ответственность:
* не сообщения/рефлексии, а контент-файлы навыков) и принимается отдельно в `ChatAgent`. * не сообщения/рефлексии, а контент-файлы навыков) и принимается отдельно в `ChatAgent`.
* *
* AutoCloseable: один `close()` закрывает все четыре store'а. В реализациях, * [EventStore] входит начиная с commit "event-store" — для replay после
* disconnect (см. `/events/replay` endpoint в `:server`).
*
* AutoCloseable: один `close()` закрывает все store'ы. В реализациях,
* которые не владеют ресурсами (in-memory), close — no-op. * которые не владеют ресурсами (in-memory), close — no-op.
*/ */
data class StorageBundle( data class StorageBundle(
@@ -19,11 +24,13 @@ data class StorageBundle(
val messageStore: MessageStore, val messageStore: MessageStore,
val workingMemoryStore: WorkingMemoryStore, val workingMemoryStore: WorkingMemoryStore,
val reflectionStore: ReflectionStore, val reflectionStore: ReflectionStore,
val eventStore: EventStore? = null,
) : AutoCloseable { ) : AutoCloseable {
override fun close() { override fun close() {
conversationStore.close() conversationStore.close()
messageStore.close() messageStore.close()
workingMemoryStore.close() workingMemoryStore.close()
reflectionStore.close() reflectionStore.close()
eventStore?.close()
} }
} }
@@ -0,0 +1,109 @@
package pw.binom.agentik.storage.events
import kotlinx.serialization.Serializable
import kotlin.time.Instant
/**
* Persistent event log для replay после disconnect.
*
* Зачем: SSE-подписка на `/events` и `/conversations/{id}/events` — cold (no replay).
* Если клиент отвалился на час, он пропустил всё. [EventStore] даёт:
* - append() — producer (ChatAgent) пишет при каждом event
* - query() — consumer (server SSE replay endpoint) читает по cursor
* - prune() — maintenance: удалить старые events по TTL
*
* Не заменяет live-подписку на [MutableSharedFlow] — это для долговременного
* хранения, а live-streaming идёт через in-memory channel.
*
* Платформо-агностичный interface (KMP): impl в `:storage-sqlite` (JVM-only),
* `:storage-inmemory` (KMP, для тестов и dev), и в будущем `:storage-sqlite-android`
* для Android-агента.
*
* Payload — opaque JSON string. [storage-core] не должен знать про
* kotlinx.serialization или [AgentEvent]/[Conversation.Event] типы (это `:proto`-шный
* слой). Конвертация — на стороне producer'а (:standalone ChatAgent).
*/
interface EventStore : AutoCloseable {
/**
* Записать event. Идемпотентен по [EventRecord.id] — повторный append с тем же
* id это no-op (важно для retry при network failure между producer'ом и БД).
*/
suspend fun append(record: EventRecord)
/**
* Catchup query для reconnect.
*
* @param conversationId если `null` — глобальный catchup (для `/events/replay`).
* если задан — только этот диалог (для `/conversations/{id}/events/replay`).
* @param afterId exclusive cursor: вернуть events СТРОГО после этого id.
* Если `null` — с начала.
* @param limit max количество records (default 100). Caller делает пагинацию
* пока `result.size == limit`.
*
* Сортировка: по [EventRecord.createdAt] ASC, ties broken по [EventRecord.id] ASC
* (т.к. id содержит timestamp-like prefix в нашей схеме, это даёт стабильный порядок).
*/
suspend fun query(
conversationId: String? = null,
afterId: String? = null,
limit: Int = 100,
): List<EventRecord>
/**
* Maintenance: удалить events старше [olderThan]. Возвращает количество удалённых.
* Default вызывается из background scope раз в час (TTL = 24h типично).
*/
suspend fun pruneOlderThan(olderThan: Instant): Int
/** Сколько events всего хранится (для observability). */
suspend fun count(): Int
override fun close()
}
/**
* Платформо-агностичная запись event'а.
*
* @param id уникальный в пределах EventStore. Convention: `"ev-<uuid>"`.
* Используется как cursor для [EventStore.query].
* @param conversationId `null` для agent-level events (Created/Deleted/Renamed).
* Задан для conversation events.
* @param createdAt UTC timestamp. Используется для сортировки в query() и для TTL в prune().
* @param type kind of event (для индексирования/фильтрации; payload всё равно opaque).
* @param payload opaque JSON string. Producer (:standalone ChatAgent) сериализует
* [pw.binom.agentik.proto.AgentEvent] или [pw.binom.agentik.proto.Event]
* в JSON перед append. Consumer (:server Routes) парсит обратно.
*
* Note: payload хранится as String, не ByteArray, чтобы не зависеть от kotlinx
* serialization и platform-specific binary encoding в [storage-core].
*/
@Serializable
data class EventRecord(
val id: String,
val conversationId: String?,
val createdAt: Instant,
val type: EventType,
val payload: String,
)
/**
* Категория event'а — для индексирования и для фильтрации в query().
*
* Naming: AGENT_* — agent-level, CONVERSATION_* — turn-level.
*/
enum class EventType {
AGENT_CREATED,
AGENT_DELETED,
AGENT_RENAMED,
CONVERSATION_START_REASONING,
CONVERSATION_START_RESPONSE,
CONVERSATION_APPEND_TEXT,
CONVERSATION_APPEND_IMAGE,
CONVERSATION_TOOL_CALL,
CONVERSATION_TOOL_RESULT,
CONVERSATION_END,
CONVERSATION_INTERRUPTED,
CONVERSATION_ERROR,
// reserved for future — adding new variants doesn't break older consumers
}
@@ -0,0 +1,84 @@
package pw.binom.agentik.storage.inmemory
import pw.binom.agentik.storage.events.EventRecord
import pw.binom.agentik.storage.events.EventStore
import pw.binom.agentik.storage.events.EventType
import kotlin.time.Instant
import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock
/**
* Thread-safe in-memory [EventStore]. Используется в тестах и в dev-режиме
* `:standalone` (когда persistent events не нужны — например, для agentik-cli).
*
* Хранит все events в одном sorted-list (по createdAt). Поиск по afterId —
* бинарный (list отсортирован). На 10K events работает за доли ms — для тестов
* хватает. На production нужен Sqlite-impl.
*
* Thread-safety: один [Mutex] на все операции. Не concurrent-write-optimised —
* для высоких нагрузок заменить на concurrent skip-list.
*/
class InMemoryEventStore : EventStore {
private val all: MutableList<EventRecord> = mutableListOf()
private val byId: MutableMap<String, EventRecord> = mutableMapOf()
private val mutex = Mutex()
override suspend fun append(record: EventRecord) {
mutex.withLock {
// Идемпотентность по id
if (byId.containsKey(record.id)) return
byId[record.id] = record
// Insert maintaining ASC order by createdAt, ties broken by id.
// binarySearch returns negative (-insertionPoint - 1) if not found,
// or non-negative index if equal element found.
val idx = all.binarySearch {
val cmp = it.createdAt.compareTo(record.createdAt)
if (cmp != 0) cmp else it.id.compareTo(record.id)
}
if (idx < 0) {
all.add(-idx - 1, record)
} else {
// Found equal element — insert AFTER it to keep insertion order.
all.add(idx + 1, record)
}
}
}
override suspend fun query(
conversationId: String?,
afterId: String?,
limit: Int,
): List<EventRecord> {
mutex.withLock {
val startIdx = if (afterId == null) 0 else {
val afterIdx = all.indexOfFirst { it.id == afterId }
if (afterIdx < 0) return emptyList()
afterIdx + 1
}
val filtered = if (conversationId == null) {
all.subList(startIdx.coerceAtMost(all.size), all.size)
} else {
all.subList(startIdx.coerceAtMost(all.size), all.size)
.filter { it.conversationId == conversationId }
}
return filtered.take(limit)
}
}
override suspend fun pruneOlderThan(olderThan: Instant): Int {
mutex.withLock {
val toRemove = all.filter { it.createdAt < olderThan }.map { it.id }
if (toRemove.isEmpty()) return 0
all.removeAll { it.id in toRemove }
toRemove.forEach { byId.remove(it) }
return toRemove.size
}
}
override suspend fun count(): Int = mutex.withLock { all.size }
override fun close() {
// no-op: nothing to release
}
}
@@ -19,5 +19,6 @@ object InMemoryStorage {
messageStore = InMemoryMessageStore(), messageStore = InMemoryMessageStore(),
workingMemoryStore = InMemoryWorkingMemoryStore(), workingMemoryStore = InMemoryWorkingMemoryStore(),
reflectionStore = InMemoryReflectionStore(), reflectionStore = InMemoryReflectionStore(),
eventStore = InMemoryEventStore(),
) )
} }
@@ -0,0 +1,165 @@
package pw.binom.agentik.storage.inmemory
import kotlinx.coroutines.coroutineScope
import kotlinx.coroutines.test.runTest
import pw.binom.agentik.storage.events.EventRecord
import pw.binom.agentik.storage.events.EventType
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertNull
import kotlin.test.assertTrue
import kotlin.time.Clock
import kotlin.time.Instant
class InMemoryEventStoreTest {
private fun rec(
id: String,
ts: Long,
conv: String? = null,
type: EventType = EventType.CONVERSATION_APPEND_TEXT,
): EventRecord = EventRecord(
id = id,
conversationId = conv,
createdAt = Instant.fromEpochMilliseconds(ts),
type = type,
payload = "{\"i\":$id}",
)
@Test
fun `append then query returns the record`() = runTest {
val store = InMemoryEventStore()
val r = rec("ev-1", ts = 1000)
store.append(r)
val result = store.query()
assertEquals(listOf(r), result)
}
@Test
fun `query with afterId returns events strictly after the cursor`() = runTest {
val store = InMemoryEventStore()
store.append(rec("ev-1", ts = 1000))
store.append(rec("ev-2", ts = 2000))
store.append(rec("ev-3", ts = 3000))
// afterId = "ev-1" → only ev-2, ev-3 (exclusive)
assertEquals(listOf("ev-2", "ev-3"), store.query(afterId = "ev-1").map { it.id })
// afterId = "ev-2" → only ev-3
assertEquals(listOf("ev-3"), store.query(afterId = "ev-2").map { it.id })
// afterId = null → all
assertEquals(listOf("ev-1", "ev-2", "ev-3"), store.query(afterId = null).map { it.id })
}
@Test
fun `query with afterId pointing at unknown id returns empty`() = runTest {
val store = InMemoryEventStore()
store.append(rec("ev-1", ts = 1000))
assertEquals(emptyList(), store.query(afterId = "ev-unknown"))
}
@Test
fun `query with conversationId filters to that conversation only`() = runTest {
val store = InMemoryEventStore()
store.append(rec("ev-1", ts = 1000, conv = "c-1"))
store.append(rec("ev-2", ts = 2000, conv = "c-2"))
store.append(rec("ev-3", ts = 3000, conv = "c-1"))
val c1 = store.query(conversationId = "c-1")
assertEquals(listOf("ev-1", "ev-3"), c1.map { it.id })
val c2 = store.query(conversationId = "c-2")
assertEquals(listOf("ev-2"), c2.map { it.id })
val all = store.query(conversationId = null)
assertEquals(listOf("ev-1", "ev-2", "ev-3"), all.map { it.id })
}
@Test
fun `query with limit caps the result size`() = runTest {
val store = InMemoryEventStore()
repeat(10) { i -> store.append(rec("ev-$i", ts = (i * 1000).toLong())) }
val first5 = store.query(limit = 5)
assertEquals(5, first5.size)
assertEquals(listOf("ev-0", "ev-1", "ev-2", "ev-3", "ev-4"), first5.map { it.id })
}
@Test
fun `append is idempotent on id`() = runTest {
val store = InMemoryEventStore()
val r1 = rec("ev-1", ts = 1000, type = EventType.CONVERSATION_APPEND_TEXT)
val r2 = rec("ev-1", ts = 1000, type = EventType.CONVERSATION_TOOL_CALL) // тот же id, разный тип
store.append(r1)
store.append(r2)
// Second append no-op (idempotent by id) — первая запись побеждает
val result = store.query()
assertEquals(1, result.size)
assertEquals(EventType.CONVERSATION_APPEND_TEXT, result[0].type)
}
@Test
fun `records returned in createdAt ascending order`() = runTest {
val store = InMemoryEventStore()
// Insert out of order
store.append(rec("ev-b", ts = 2000))
store.append(rec("ev-a", ts = 1000))
store.append(rec("ev-c", ts = 3000))
val result = store.query()
assertEquals(listOf("ev-a", "ev-b", "ev-c"), result.map { it.id })
}
@Test
fun `records with same timestamp ordered by id ascending (stable sort)`() = runTest {
val store = InMemoryEventStore()
store.append(rec("ev-c", ts = 1000))
store.append(rec("ev-a", ts = 1000))
store.append(rec("ev-b", ts = 1000))
val result = store.query()
// id lexicographic order: a < b < c
assertEquals(listOf("ev-a", "ev-b", "ev-c"), result.map { it.id })
}
@Test
fun `pruneOlderThan removes records before cutoff`() = runTest {
val store = InMemoryEventStore()
store.append(rec("ev-1", ts = 1000))
store.append(rec("ev-2", ts = 2000))
store.append(rec("ev-3", ts = 3000))
val removed = store.pruneOlderThan(Instant.fromEpochMilliseconds(2500))
assertEquals(2, removed)
assertEquals(listOf("ev-3"), store.query().map { it.id })
}
@Test
fun `pruneOlderThan returns 0 when nothing to remove`() = runTest {
val store = InMemoryEventStore()
store.append(rec("ev-1", ts = 5000))
val removed = store.pruneOlderThan(Instant.fromEpochMilliseconds(1000))
assertEquals(0, removed)
assertEquals(1, store.count())
}
@Test
fun `count returns number of stored records`() = runTest {
val store = InMemoryEventStore()
assertEquals(0, store.count())
store.append(rec("ev-1", ts = 1000))
store.append(rec("ev-2", ts = 2000))
assertEquals(2, store.count())
}
@Test
fun `concurrent append from multiple coroutines all succeed`() = runTest {
// Не strict — InMemoryEventStore использует Mutex, поэтому concurrent calls
// сериализуются. Тест проверяет, что при последовательных append все ids
// попадают в store (для concurrent test нужна отдельная TestScope — это
// покрыто integration-тестами в :standalone).
val store = InMemoryEventStore()
repeat(51) { i ->
store.append(rec("ev-$i", ts = i.toLong()))
}
assertEquals(51, store.count())
}
}
@@ -0,0 +1,86 @@
package pw.binom.agentik.storage.sqlite
import kotlin.time.Instant
import mu.KotlinLogging
import pw.binom.agentik.storage.events.EventRecord
import pw.binom.agentik.storage.events.EventStore
import pw.binom.agentik.storage.events.EventType
import pw.binom.agentik.storage.sqlite.Agent_event as DbAgentEvent
private val log = KotlinLogging.logger {}
/**
* SQLDelight-реализация [EventStore] поверх таблицы `agent_event`.
*
* Использует [EventStoreQueries] (генерируется SQLDelight из EventStore.sq).
* Все запросы готовы — мы только маппим `Agent_event` (DB) ↔ `EventRecord` (domain).
*
* **Idempotency**: `append()` использует `INSERT OR IGNORE` — повторный append
* с тем же id (network retry) — no-op. Это критично для producer'а, который
* может retry при transient failure.
*
* **Pruning**: вызывай [pruneOlderThan] раз в час из background scope. Типичный
* TTL = 24h. Если eventStore разрастётся (миллионы записей), индексы
* (created_at, conversation_id+created_at) обеспечат O(log N) для query.
*/
class SqliteEventStore(
private val db: AgentikDatabase,
) : EventStore {
private val queries: EventStoreQueries get() = db.eventStoreQueries
override suspend fun append(record: EventRecord) {
// payload — opaque JSON, хранится как UTF-8 bytes. Используем ByteArray,
// потому что BLOB-колонка эффективнее TEXT для >100KB строк, и
// API sqldelight нативно работает с ByteArray.
val payloadBytes = record.payload.encodeToByteArray()
queries.insert(
id = record.id,
conversation_id = record.conversationId,
created_at = record.createdAt.toEpochMilliseconds(),
type = record.type.name,
payload = payloadBytes,
)
log.debug { "appended event id=${record.id} type=${record.type} conv=${record.conversationId}" }
}
override suspend fun query(
conversationId: String?,
afterId: String?,
limit: Int,
): List<EventRecord> {
// Если conversationId == null — используем queryAfter без фильтра
// (он сам обрабатывает :convId IS NULL внутри SQL).
// Если задан — queryAfterByConv (тогда SQL имеет WHERE conversation_id = :convId).
val rows: List<DbAgentEvent> = if (conversationId == null) {
queries.queryAfter(convId = null, afterId = afterId, limit = limit.toLong()).executeAsList()
} else {
queries.queryAfterByConv(convId = conversationId, afterId = afterId, limit = limit.toLong())
.executeAsList()
}
return rows.map { it.toDomain() }
}
override suspend fun pruneOlderThan(olderThan: Instant): Int {
val deleted = queries.pruneOlderThan(olderThan.toEpochMilliseconds()).value
if (deleted > 0) log.info { "pruned $deleted events older than $olderThan" }
return deleted.toInt()
}
override suspend fun count(): Int = queries.countAll().executeAsOne().toInt()
override fun close() {
// no-op: lifecycle owned by AgentikDatabase / SqliteStores
}
private fun DbAgentEvent.toDomain(): EventRecord = EventRecord(
id = id,
conversationId = conversation_id,
createdAt = Instant.fromEpochMilliseconds(created_at),
// type name → enum. Если в БД оказался неизвестный тип (новая версия,
// unknown старому коду) — fallback на AGENT_CREATED (нейтральное значение).
// Это безопаснее чем throw: клиент просто получит event с минимальным payload.
type = runCatching { EventType.valueOf(type) }.getOrDefault(EventType.AGENT_CREATED),
payload = payload.decodeToString(),
)
}
@@ -7,12 +7,13 @@ import pw.binom.agentik.storage.ConversationStore
import pw.binom.agentik.storage.MessageStore import pw.binom.agentik.storage.MessageStore
import pw.binom.agentik.storage.ReflectionStore import pw.binom.agentik.storage.ReflectionStore
import pw.binom.agentik.storage.StorageBundle import pw.binom.agentik.storage.StorageBundle
import pw.binom.agentik.storage.events.EventStore
import pw.binom.agentik.storage.sqlite.SqliteReflectionStore import pw.binom.agentik.storage.sqlite.SqliteReflectionStore
import pw.binom.agentik.storage.WorkingMemoryStore import pw.binom.agentik.storage.WorkingMemoryStore
/** /**
* Корневой объект SQLite-слоя: держит [SqlDriver] и три [WorkingMemoryStore]/[MessageStore]/[ConversationStore]. * Корневой объект SQLite-слоя: держит [SqlDriver] и пять store'ов (включая
* Закрывается вместе с приложением. * [EventStore] — replay-after-disconnect).
*/ */
class SqliteStores private constructor( class SqliteStores private constructor(
val driver: SqlDriver, val driver: SqlDriver,
@@ -20,6 +21,7 @@ class SqliteStores private constructor(
val messages: MessageStore, val messages: MessageStore,
val workingMemory: WorkingMemoryStore, val workingMemory: WorkingMemoryStore,
val reflections: ReflectionStore, val reflections: ReflectionStore,
val events: EventStore,
) : AutoCloseable { ) : AutoCloseable {
/** /**
@@ -32,9 +34,11 @@ class SqliteStores private constructor(
messageStore = messages, messageStore = messages,
workingMemoryStore = workingMemory, workingMemoryStore = workingMemory,
reflectionStore = reflections, reflectionStore = reflections,
eventStore = events,
) )
override fun close() { override fun close() {
events.close()
conversations.close() conversations.close()
messages.close() messages.close()
workingMemory.close() workingMemory.close()
@@ -55,6 +59,7 @@ class SqliteStores private constructor(
messages = SqliteMessageStore(db), messages = SqliteMessageStore(db),
workingMemory = SqliteWorkingMemoryStore(db), workingMemory = SqliteWorkingMemoryStore(db),
reflections = SqliteReflectionStore(db), reflections = SqliteReflectionStore(db),
events = SqliteEventStore(db),
) )
} }
@@ -69,6 +74,7 @@ class SqliteStores private constructor(
messages = SqliteMessageStore(db), messages = SqliteMessageStore(db),
workingMemory = SqliteWorkingMemoryStore(db), workingMemory = SqliteWorkingMemoryStore(db),
reflections = SqliteReflectionStore(db), reflections = SqliteReflectionStore(db),
events = SqliteEventStore(db),
) )
} }
@@ -118,6 +124,18 @@ class SqliteStores private constructor(
CREATE INDEX IF NOT EXISTS idx_reflection_created ON reflection(created_at DESC); CREATE INDEX IF NOT EXISTS idx_reflection_created ON reflection(created_at DESC);
CREATE INDEX IF NOT EXISTS idx_reflection_conv ON reflection(conversation_id, created_at DESC); CREATE INDEX IF NOT EXISTS idx_reflection_conv ON reflection(conversation_id, created_at DESC);
""".trimIndent(), """.trimIndent(),
// v3: event log для replay-after-disconnect (см. EventStore.sq)
"""
CREATE TABLE IF NOT EXISTS agent_event (
id TEXT NOT NULL PRIMARY KEY,
conversation_id TEXT,
created_at INTEGER NOT NULL,
type TEXT NOT NULL,
payload BLOB NOT NULL
);
CREATE INDEX IF NOT EXISTS idx_agent_event_conv_time ON agent_event(conversation_id, created_at);
CREATE INDEX IF NOT EXISTS idx_agent_event_time ON agent_event(created_at);
""".trimIndent(),
) )
for (sql in migrations) { for (sql in migrations) {
driver.execute(null, sql, 0) driver.execute(null, sql, 0)
@@ -0,0 +1,49 @@
-- Event log для replay после disconnect (см. EventStore.kt в :storage-core).
-- Каждая запись — один event из ChatAgent (AgentEvent) или ConversationLoop
-- (Conversation.Event). payload — opaque JSON, сериализуется в :standalone перед append.
--
-- Cursor для pagination: query() принимает afterId, возвращает events строго
-- после него (exclusive). Используется /events/replay endpoint в :server.
CREATE TABLE agent_event (
id TEXT NOT NULL PRIMARY KEY,
conversation_id TEXT, -- NULL для agent-level events (Created/Deleted/Renamed)
created_at INTEGER NOT NULL, -- epoch millis (UTC)
type TEXT NOT NULL, -- EventType.name (см. :storage-core/events/EventStore.kt)
payload BLOB NOT NULL -- serialized JSON
);
CREATE INDEX idx_agent_event_conv_time ON agent_event(conversation_id, created_at);
CREATE INDEX idx_agent_event_time ON agent_event(created_at);
insert:
INSERT OR IGNORE INTO agent_event (id, conversation_id, created_at, type, payload)
VALUES (?, ?, ?, ?, ?);
queryById:
SELECT * FROM agent_event WHERE id = ?;
queryAfter:
-- Catchup по conversationId (или все если null). afterId exclusive.
-- Сортировка: created_at ASC, id ASC (стабильный tie-break для events с одинаковым timestamp).
SELECT * FROM agent_event
WHERE (:convId IS NULL OR conversation_id = :convId)
AND id > COALESCE(:afterId, '')
ORDER BY created_at ASC, id ASC
LIMIT :limit;
queryAfterByConv:
SELECT * FROM agent_event
WHERE conversation_id = :convId
AND id > COALESCE(:afterId, '')
ORDER BY created_at ASC, id ASC
LIMIT :limit;
countAll:
SELECT COUNT(*) FROM agent_event;
countByConv:
SELECT COUNT(*) FROM agent_event WHERE conversation_id = :convId;
pruneOlderThan:
DELETE FROM agent_event WHERE created_at < :cutoffEpochMillis;
@@ -0,0 +1,152 @@
package pw.binom.agentik.storage.sqlite
import kotlinx.coroutines.test.runTest
import pw.binom.agentik.storage.events.EventRecord
import pw.binom.agentik.storage.events.EventType
import kotlin.test.AfterTest
import kotlin.test.BeforeTest
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertTrue
import kotlin.time.Instant
class SqliteEventStoreTest {
private lateinit var driver: app.cash.sqldelight.driver.jdbc.sqlite.JdbcSqliteDriver
private lateinit var db: AgentikDatabase
private lateinit var store: SqliteEventStore
@BeforeTest
fun setup() {
driver = app.cash.sqldelight.driver.jdbc.sqlite.JdbcSqliteDriver(
app.cash.sqldelight.driver.jdbc.sqlite.JdbcSqliteDriver.IN_MEMORY,
)
AgentikDatabase.Schema.create(driver)
db = AgentikDatabase(driver)
store = SqliteEventStore(db)
}
@AfterTest
fun tearDown() {
driver.close()
}
private fun rec(
id: String,
ts: Long,
conv: String? = null,
type: EventType = EventType.CONVERSATION_APPEND_TEXT,
): EventRecord = EventRecord(
id = id,
conversationId = conv,
createdAt = Instant.fromEpochMilliseconds(ts),
type = type,
payload = """{"i":"$id"}""",
)
@Test
fun `append then query returns the record`() = runTest {
val r = rec("ev-1", ts = 1000)
store.append(r)
assertEquals(listOf(r), store.query())
}
@Test
fun `query with conversationId filters to that conversation only`() = runTest {
store.append(rec("ev-1", ts = 1000, conv = "c-1"))
store.append(rec("ev-2", ts = 2000, conv = "c-2"))
store.append(rec("ev-3", ts = 3000, conv = "c-1"))
val c1 = store.query(conversationId = "c-1")
assertEquals(listOf("ev-1", "ev-3"), c1.map { it.id })
}
@Test
fun `query with afterId returns events strictly after cursor`() = runTest {
store.append(rec("ev-1", ts = 1000))
store.append(rec("ev-2", ts = 2000))
store.append(rec("ev-3", ts = 3000))
assertEquals(listOf("ev-2", "ev-3"), store.query(afterId = "ev-1").map { it.id })
assertEquals(listOf("ev-3"), store.query(afterId = "ev-2").map { it.id })
}
@Test
fun `append is idempotent (INSERT OR IGNORE)`() = runTest {
val r1 = rec("ev-1", ts = 1000, type = EventType.CONVERSATION_APPEND_TEXT)
val r2 = rec("ev-1", ts = 1000, type = EventType.CONVERSATION_TOOL_CALL)
store.append(r1)
store.append(r2)
// INSERT OR IGNORE — вторая попытка no-op
val result = store.query()
assertEquals(1, result.size)
assertEquals(EventType.CONVERSATION_APPEND_TEXT, result[0].type)
}
@Test
fun `records ordered by createdAt then id ascending`() = runTest {
store.append(rec("ev-c", ts = 2000))
store.append(rec("ev-a", ts = 1000))
store.append(rec("ev-d", ts = 1000)) // same ts as a, different id
store.append(rec("ev-b", ts = 1500))
val result = store.query()
assertEquals(listOf("ev-a", "ev-d", "ev-b", "ev-c"), result.map { it.id })
}
@Test
fun `query with limit caps the result`() = runTest {
repeat(10) { i -> store.append(rec("ev-$i", ts = (i * 100).toLong())) }
val first3 = store.query(limit = 3)
assertEquals(3, first3.size)
assertEquals(listOf("ev-0", "ev-1", "ev-2"), first3.map { it.id })
}
@Test
fun `pruneOlderThan removes records before cutoff`() = runTest {
store.append(rec("ev-1", ts = 1000))
store.append(rec("ev-2", ts = 2000))
store.append(rec("ev-3", ts = 3000))
val removed = store.pruneOlderThan(Instant.fromEpochMilliseconds(2500))
assertEquals(2, removed)
assertEquals(listOf("ev-3"), store.query().map { it.id })
}
@Test
fun `count returns total record count`() = runTest {
assertEquals(0, store.count())
store.append(rec("ev-1", ts = 1000))
store.append(rec("ev-2", ts = 2000))
assertEquals(2, store.count())
}
@Test
fun `unknown EventType in DB is loaded as AGENT_CREATED fallback`() = runTest {
// Insert record with type that doesn't exist in current enum (simulating
// a future enum value that older code doesn't know about).
store.append(rec("ev-future", ts = 1000, type = EventType.CONVERSATION_END))
// Manually mutate DB to use a fake type name (simulating forward-compat)
db.eventStoreQueries.insert(
id = "ev-fake",
conversation_id = null,
created_at = 2000,
type = "FUTURE_TYPE_NOT_IN_ENUM",
payload = "{}".encodeToByteArray(),
)
// Оба должны загрузиться (первый с правильным type, второй — с fallback)
val result = store.query()
assertEquals(2, result.size)
assertEquals(EventType.CONVERSATION_END, result[0].type)
assertEquals(EventType.AGENT_CREATED, result[1].type) // fallback для неизвестного type
}
@Test
fun `close is no-op (does not close shared driver)`() = runTest {
// store.close() НЕ должен закрывать driver — driver shared с другими store'ами.
store.close()
// Если бы close закрыл driver — следующий запрос упал бы. Проверяем что работает.
store.append(rec("ev-1", ts = 1000))
assertEquals(1, store.count())
}
}