From 4ad59d5f5d8b3074c16a8d4f30cc653ee8696782 Mon Sep 17 00:00:00 2001 From: subochev Date: Sun, 20 Sep 2026 02:51:58 +0300 Subject: [PATCH] feat(events): EventStore + AllEvent unified stream + replay endpoints MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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 (5) :proto — AllEvent sealed interface - AllEvent.Agent(date, event: AgentEvent) - AllEvent.Conversation(date, conversationId, event: Event) - Agent.allEvents(after): Flow — третий тип подписки (в дополнение к 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. --- .../pw/binom/agentik/client/AgentClient.kt | 71 ++++++++ .../kotlin/pw/binom/agentik/proto/Agent.kt | 12 ++ .../kotlin/pw/binom/agentik/proto/AllEvent.kt | 39 +++++ server/build.gradle.kts | 1 + .../kotlin/pw/binom/agentik/server/Module.kt | 15 +- .../kotlin/pw/binom/agentik/server/Routes.kt | 46 ++++- .../agentik/standalone/agent/ChatAgent.kt | 113 +++++++++++- .../standalone/agent/ConversationEvents.kt | 105 ++++++++++- .../standalone/agent/ConversationLoop.kt | 5 +- .../pw/binom/agentik/storage/StorageBundle.kt | 9 +- .../agentik/storage/events/EventStore.kt | 109 ++++++++++++ .../storage/inmemory/InMemoryEventStore.kt | 84 +++++++++ .../storage/inmemory/InMemoryStorage.kt | 1 + .../inmemory/InMemoryEventStoreTest.kt | 165 ++++++++++++++++++ .../storage/sqlite/SqliteEventStore.kt | 86 +++++++++ .../agentik/storage/sqlite/SqliteStores.kt | 22 ++- .../agentik/storage/sqlite/EventStore.sq | 49 ++++++ .../storage/sqlite/SqliteEventStoreTest.kt | 152 ++++++++++++++++ 18 files changed, 1065 insertions(+), 19 deletions(-) create mode 100644 proto/src/commonMain/kotlin/pw/binom/agentik/proto/AllEvent.kt create mode 100644 storage-core/src/commonMain/kotlin/pw/binom/agentik/storage/events/EventStore.kt create mode 100644 storage-inmemory/src/commonMain/kotlin/pw/binom/agentik/storage/inmemory/InMemoryEventStore.kt create mode 100644 storage-inmemory/src/commonTest/kotlin/pw/binom/agentik/storage/inmemory/InMemoryEventStoreTest.kt create mode 100644 storage-sqlite/src/jvmMain/kotlin/pw/binom/agentik/storage/sqlite/SqliteEventStore.kt create mode 100644 storage-sqlite/src/jvmMain/sqldelight/pw/binom/agentik/storage/sqlite/EventStore.sq create mode 100644 storage-sqlite/src/jvmTest/kotlin/pw/binom/agentik/storage/sqlite/SqliteEventStoreTest.kt diff --git a/client/src/commonMain/kotlin/pw/binom/agentik/client/AgentClient.kt b/client/src/commonMain/kotlin/pw/binom/agentik/client/AgentClient.kt index 839c754..be49b3f 100644 --- a/client/src/commonMain/kotlin/pw/binom/agentik/client/AgentClient.kt +++ b/client/src/commonMain/kotlin/pw/binom/agentik/client/AgentClient.kt @@ -17,6 +17,7 @@ import kotlinx.coroutines.flow.flow import kotlinx.coroutines.runBlocking import pw.binom.agentik.proto.Agent import pw.binom.agentik.proto.AgentEvent +import pw.binom.agentik.proto.AllEvent import pw.binom.agentik.proto.Conversation 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 = 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 { + 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 { + 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, +) diff --git a/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Agent.kt b/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Agent.kt index 2265d71..bcdf196 100644 --- a/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Agent.kt +++ b/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Agent.kt @@ -49,6 +49,18 @@ public interface Agent { */ fun events(after: Instant): Flow + /** + * 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 + companion object { const val PAGE_SIZE: Int = 100 diff --git a/proto/src/commonMain/kotlin/pw/binom/agentik/proto/AllEvent.kt b/proto/src/commonMain/kotlin/pw/binom/agentik/proto/AllEvent.kt new file mode 100644 index 0000000..4ce16a5 --- /dev/null +++ b/proto/src/commonMain/kotlin/pw/binom/agentik/proto/AllEvent.kt @@ -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 +} diff --git a/server/build.gradle.kts b/server/build.gradle.kts index ebc506e..7b78f36 100644 --- a/server/build.gradle.kts +++ b/server/build.gradle.kts @@ -23,6 +23,7 @@ kotlin { sourceSets { commonMain.dependencies { implementation(project(":proto")) + implementation(project(":storage-core")) // Ktor (без engine — engine подключает потребитель, см. :standalone). implementation(libs.ktor.server.core) diff --git a/server/src/commonMain/kotlin/pw/binom/agentik/server/Module.kt b/server/src/commonMain/kotlin/pw/binom/agentik/server/Module.kt index 798550f..4ab26ab 100644 --- a/server/src/commonMain/kotlin/pw/binom/agentik/server/Module.kt +++ b/server/src/commonMain/kotlin/pw/binom/agentik/server/Module.kt @@ -6,6 +6,7 @@ import io.ktor.server.plugins.contentnegotiation.ContentNegotiation import io.ktor.server.routing.Route import io.ktor.server.routing.route import pw.binom.agentik.proto.Agent +import pw.binom.agentik.storage.events.EventStore /** * Встраивает HTTP/SSE-фасад протокола agentik в твой Ktor-роутинг. @@ -30,10 +31,20 @@ import pw.binom.agentik.proto.Agent * - `POST /conversations/{id}/interrupt` — `interrupt` * - `GET /conversations/{id}/messages` — история * - `GET /conversations/{id}/events` — SSE: события хода + * - `GET /conversations/{id}/events/replay` — replay-after-disconnect (events after ?after_id=X) * - `GET /events` — SSE: события агента + * - `GET /events/replay` — replay-after-disconnect (global) * - `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) { install(ContentNegotiation) { json(agentikJson) @@ -43,6 +54,6 @@ fun Route.agentikAgent(agent: Agent, path: String = "/agentik", token: String? = this.token = token } } - agentikRoutes(agent) + agentikRoutes(agent, eventStore) } } diff --git a/server/src/commonMain/kotlin/pw/binom/agentik/server/Routes.kt b/server/src/commonMain/kotlin/pw/binom/agentik/server/Routes.kt index d8d2672..1be914d 100644 --- a/server/src/commonMain/kotlin/pw/binom/agentik/server/Routes.kt +++ b/server/src/commonMain/kotlin/pw/binom/agentik/server/Routes.kt @@ -21,12 +21,14 @@ 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.AllEvent import pw.binom.agentik.proto.Conversation import pw.binom.agentik.proto.Content import pw.binom.agentik.proto.Event +import pw.binom.agentik.storage.events.EventStore import kotlin.time.Instant -internal fun Route.agentikRoutes(agent: Agent) { +internal fun Route.agentikRoutes(agent: Agent, eventStore: EventStore? = null) { get("/health") { call.respondText("ok") @@ -134,6 +136,48 @@ internal fun Route.agentikRoutes(agent: Agent) { val after = call.parseAfter() ?: return@get 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 ---------- diff --git a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ChatAgent.kt b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ChatAgent.kt index e377d4f..7bcf2e7 100644 --- a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ChatAgent.kt +++ b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ChatAgent.kt @@ -1,38 +1,50 @@ package pw.binom.agentik.standalone.agent +import kotlin.time.Instant import kotlinx.coroutines.flow.Flow import kotlinx.coroutines.flow.MutableSharedFlow 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.sync.Mutex 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.MemoryReviewer import pw.binom.agentik.memory.MemorySystemGuidance import pw.binom.agentik.proto.Agent as ProtoAgent 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.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.storage.ConversationRecord +import pw.binom.agentik.storage.Ids import pw.binom.agentik.storage.Reflection -import pw.binom.agentik.storage.WorkingMemoryEntry 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.EnableToolsetTool +import pw.binom.agentik.toolsets.NamedTool import pw.binom.agentik.toolsets.SystemPromptToolsetSection import pw.binom.agentik.toolsets.ToolsetContribution import pw.binom.agentik.toolsets.ToolsetDispatchPolicy import pw.binom.agentik.toolsets.ToolsetRegistry 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) и @@ -192,6 +204,54 @@ class ChatAgent( 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 live: MutableMap = HashMap() @@ -204,6 +264,36 @@ class ChatAgent( 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 = 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 { val now = now() val id = pw.binom.agentik.storage.Ids.new("conv") @@ -243,6 +333,7 @@ class ChatAgent( liveLock.withLock { live[conv.id] = conv } } agentEvents.tryEmit(AgentEvent.Created(date = now(), conversationId = conv.id)) + persistAgentEvent(AgentEvent.Created(date = now(), conversationId = conv.id)) return conv } @@ -258,7 +349,11 @@ class ChatAgent( val conv = liveLock.withLock { live.remove(id) } conv?.close() 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 } diff --git a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ConversationEvents.kt b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ConversationEvents.kt index 4f7ac0b..4347d11 100644 --- a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ConversationEvents.kt +++ b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ConversationEvents.kt @@ -4,9 +4,45 @@ import kotlinx.coroutines.channels.BufferOverflow import kotlinx.coroutines.flow.MutableSharedFlow import kotlinx.coroutines.flow.SharedFlow 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.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( replay = 0, extraBufferCapacity = 4096, @@ -15,7 +51,70 @@ internal class ConversationEvents { val flow: SharedFlow 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 + } } diff --git a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ConversationLoop.kt b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ConversationLoop.kt index 8375ace..909820c 100644 --- a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ConversationLoop.kt +++ b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ConversationLoop.kt @@ -84,7 +84,10 @@ class ConversationLoop( agentScope = agentScope, ) - private val events = ConversationEvents() + private val events = ConversationEvents( + eventStore = storage.eventStore, + conversationId = state.id, + ) /** Per-conversation background event bus. Lifecycle scoped к этому ConversationLoop. */ private val backgroundEvents = BackgroundEventBus() diff --git a/storage-core/src/commonMain/kotlin/pw/binom/agentik/storage/StorageBundle.kt b/storage-core/src/commonMain/kotlin/pw/binom/agentik/storage/StorageBundle.kt index 5dd9352..d68eaba 100644 --- a/storage-core/src/commonMain/kotlin/pw/binom/agentik/storage/StorageBundle.kt +++ b/storage-core/src/commonMain/kotlin/pw/binom/agentik/storage/StorageBundle.kt @@ -1,5 +1,7 @@ package pw.binom.agentik.storage +import pw.binom.agentik.storage.events.EventStore + /** * Агрегатор всех storage-интерфейсов, нужных агенту для работы с историей диалога. * @@ -11,7 +13,10 @@ package pw.binom.agentik.storage * `SkillStore` НЕ входит сюда — он живёт в модуле `:skills` (другая ответственность: * не сообщения/рефлексии, а контент-файлы навыков) и принимается отдельно в `ChatAgent`. * - * AutoCloseable: один `close()` закрывает все четыре store'а. В реализациях, + * [EventStore] входит начиная с commit "event-store" — для replay после + * disconnect (см. `/events/replay` endpoint в `:server`). + * + * AutoCloseable: один `close()` закрывает все store'ы. В реализациях, * которые не владеют ресурсами (in-memory), close — no-op. */ data class StorageBundle( @@ -19,11 +24,13 @@ data class StorageBundle( val messageStore: MessageStore, val workingMemoryStore: WorkingMemoryStore, val reflectionStore: ReflectionStore, + val eventStore: EventStore? = null, ) : AutoCloseable { override fun close() { conversationStore.close() messageStore.close() workingMemoryStore.close() reflectionStore.close() + eventStore?.close() } } diff --git a/storage-core/src/commonMain/kotlin/pw/binom/agentik/storage/events/EventStore.kt b/storage-core/src/commonMain/kotlin/pw/binom/agentik/storage/events/EventStore.kt new file mode 100644 index 0000000..2465bf5 --- /dev/null +++ b/storage-core/src/commonMain/kotlin/pw/binom/agentik/storage/events/EventStore.kt @@ -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 + + /** + * 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-"`. + * Используется как 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 +} diff --git a/storage-inmemory/src/commonMain/kotlin/pw/binom/agentik/storage/inmemory/InMemoryEventStore.kt b/storage-inmemory/src/commonMain/kotlin/pw/binom/agentik/storage/inmemory/InMemoryEventStore.kt new file mode 100644 index 0000000..407c3ce --- /dev/null +++ b/storage-inmemory/src/commonMain/kotlin/pw/binom/agentik/storage/inmemory/InMemoryEventStore.kt @@ -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 = mutableListOf() + private val byId: MutableMap = 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 { + 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 + } +} diff --git a/storage-inmemory/src/commonMain/kotlin/pw/binom/agentik/storage/inmemory/InMemoryStorage.kt b/storage-inmemory/src/commonMain/kotlin/pw/binom/agentik/storage/inmemory/InMemoryStorage.kt index d1a7d49..acb37ec 100644 --- a/storage-inmemory/src/commonMain/kotlin/pw/binom/agentik/storage/inmemory/InMemoryStorage.kt +++ b/storage-inmemory/src/commonMain/kotlin/pw/binom/agentik/storage/inmemory/InMemoryStorage.kt @@ -19,5 +19,6 @@ object InMemoryStorage { messageStore = InMemoryMessageStore(), workingMemoryStore = InMemoryWorkingMemoryStore(), reflectionStore = InMemoryReflectionStore(), + eventStore = InMemoryEventStore(), ) } diff --git a/storage-inmemory/src/commonTest/kotlin/pw/binom/agentik/storage/inmemory/InMemoryEventStoreTest.kt b/storage-inmemory/src/commonTest/kotlin/pw/binom/agentik/storage/inmemory/InMemoryEventStoreTest.kt new file mode 100644 index 0000000..9033b53 --- /dev/null +++ b/storage-inmemory/src/commonTest/kotlin/pw/binom/agentik/storage/inmemory/InMemoryEventStoreTest.kt @@ -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()) + } +} diff --git a/storage-sqlite/src/jvmMain/kotlin/pw/binom/agentik/storage/sqlite/SqliteEventStore.kt b/storage-sqlite/src/jvmMain/kotlin/pw/binom/agentik/storage/sqlite/SqliteEventStore.kt new file mode 100644 index 0000000..9484baa --- /dev/null +++ b/storage-sqlite/src/jvmMain/kotlin/pw/binom/agentik/storage/sqlite/SqliteEventStore.kt @@ -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 { + // Если conversationId == null — используем queryAfter без фильтра + // (он сам обрабатывает :convId IS NULL внутри SQL). + // Если задан — queryAfterByConv (тогда SQL имеет WHERE conversation_id = :convId). + val rows: List = 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(), + ) +} diff --git a/storage-sqlite/src/jvmMain/kotlin/pw/binom/agentik/storage/sqlite/SqliteStores.kt b/storage-sqlite/src/jvmMain/kotlin/pw/binom/agentik/storage/sqlite/SqliteStores.kt index a25805f..9756309 100644 --- a/storage-sqlite/src/jvmMain/kotlin/pw/binom/agentik/storage/sqlite/SqliteStores.kt +++ b/storage-sqlite/src/jvmMain/kotlin/pw/binom/agentik/storage/sqlite/SqliteStores.kt @@ -7,12 +7,13 @@ import pw.binom.agentik.storage.ConversationStore import pw.binom.agentik.storage.MessageStore import pw.binom.agentik.storage.ReflectionStore 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.WorkingMemoryStore /** - * Корневой объект SQLite-слоя: держит [SqlDriver] и три [WorkingMemoryStore]/[MessageStore]/[ConversationStore]. - * Закрывается вместе с приложением. + * Корневой объект SQLite-слоя: держит [SqlDriver] и пять store'ов (включая + * [EventStore] — replay-after-disconnect). */ class SqliteStores private constructor( val driver: SqlDriver, @@ -20,6 +21,7 @@ class SqliteStores private constructor( val messages: MessageStore, val workingMemory: WorkingMemoryStore, val reflections: ReflectionStore, + val events: EventStore, ) : AutoCloseable { /** @@ -32,9 +34,11 @@ class SqliteStores private constructor( messageStore = messages, workingMemoryStore = workingMemory, reflectionStore = reflections, + eventStore = events, ) override fun close() { + events.close() conversations.close() messages.close() workingMemory.close() @@ -55,6 +59,7 @@ class SqliteStores private constructor( messages = SqliteMessageStore(db), workingMemory = SqliteWorkingMemoryStore(db), reflections = SqliteReflectionStore(db), + events = SqliteEventStore(db), ) } @@ -69,6 +74,7 @@ class SqliteStores private constructor( messages = SqliteMessageStore(db), workingMemory = SqliteWorkingMemoryStore(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_conv ON reflection(conversation_id, created_at DESC); """.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) { driver.execute(null, sql, 0) diff --git a/storage-sqlite/src/jvmMain/sqldelight/pw/binom/agentik/storage/sqlite/EventStore.sq b/storage-sqlite/src/jvmMain/sqldelight/pw/binom/agentik/storage/sqlite/EventStore.sq new file mode 100644 index 0000000..6e0a867 --- /dev/null +++ b/storage-sqlite/src/jvmMain/sqldelight/pw/binom/agentik/storage/sqlite/EventStore.sq @@ -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; diff --git a/storage-sqlite/src/jvmTest/kotlin/pw/binom/agentik/storage/sqlite/SqliteEventStoreTest.kt b/storage-sqlite/src/jvmTest/kotlin/pw/binom/agentik/storage/sqlite/SqliteEventStoreTest.kt new file mode 100644 index 0000000..8a0a7d0 --- /dev/null +++ b/storage-sqlite/src/jvmTest/kotlin/pw/binom/agentik/storage/sqlite/SqliteEventStoreTest.kt @@ -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()) + } +}