diff --git a/agent-toolsets/build.gradle.kts b/agent-toolsets/build.gradle.kts index e110d07..3366eb7 100644 --- a/agent-toolsets/build.gradle.kts +++ b/agent-toolsets/build.gradle.kts @@ -21,8 +21,9 @@ kotlin { sourceSets { commonMain.dependencies { - api(project(":storage-bundle")) - // :storage-bundle transitively подтягивает :message-store-api и :working-memory-api. + api(project(":message-store-api")) + api(project(":message-log-api")) + api(project(":working-memory-api")) // litert-kmp: LiteTool интерфейс (sync describe/invoke) api(libs.litert.api) diff --git a/agent-toolsets/src/commonMain/kotlin/pw/binom/agentik/toolsets/ToolsetContext.kt b/agent-toolsets/src/commonMain/kotlin/pw/binom/agentik/toolsets/ToolsetContext.kt index af36424..ac2a60f 100644 --- a/agent-toolsets/src/commonMain/kotlin/pw/binom/agentik/toolsets/ToolsetContext.kt +++ b/agent-toolsets/src/commonMain/kotlin/pw/binom/agentik/toolsets/ToolsetContext.kt @@ -4,7 +4,7 @@ package pw.binom.agentik.toolsets * Контекст, который тулсеты получают при активации. * * В commit 4 — минимальный: логгер. Позже (commit 5+, если понадобится) сюда - * добавятся `StorageBundle`, `SkillStore` и пр., чтобы тулы внутри тулсета + * добавятся `SkillStore` и пр., чтобы тулы внутри тулсета * могли читать/писать сообщения и память. * * Если конкретному тулсету нужно больше, чем [Logger], он может объявить свой diff --git a/client/build.gradle.kts b/client/build.gradle.kts index 6729b3f..2a95945 100644 --- a/client/build.gradle.kts +++ b/client/build.gradle.kts @@ -21,6 +21,7 @@ kotlin { sourceSets { commonMain.dependencies { api(project(":proto")) + api(project(":event-store")) api(libs.ktor.client.core) implementation(libs.ktor.client.content.negotiation) 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 f62b4db..fcc9adf 100644 --- a/client/src/commonMain/kotlin/pw/binom/agentik/client/AgentClient.kt +++ b/client/src/commonMain/kotlin/pw/binom/agentik/client/AgentClient.kt @@ -15,6 +15,7 @@ import io.ktor.http.contentType import kotlinx.coroutines.flow.Flow import kotlinx.coroutines.flow.flow import kotlinx.coroutines.runBlocking +import pw.binom.agentik.eventStore.EventStore import pw.binom.agentik.proto.Agent import pw.binom.agentik.proto.AgentEvent import pw.binom.agentik.proto.CommonEvent @@ -38,6 +39,12 @@ internal class AgentClient( private val agentUrl: String = baseUrl.trimEnd('/') + /** + * Единый канал событий (lifecycle + per-conversation). Под капотом — + * [HttpEventStore]: каждый метод бьёт свой URL (см. KDoc). + */ + val eventStore: EventStore = HttpEventStore(httpClient = httpClient, baseUrl = agentUrl) + override fun createConversation(temp: Boolean): Conversation = runBlocking { val snapshot: ConversationSnapshot = httpClient.post("$agentUrl/conversations") { @@ -97,56 +104,4 @@ internal class AgentClient( } } } - - /** - * 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.messageStore.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/client/src/commonMain/kotlin/pw/binom/agentik/client/HttpEventStore.kt b/client/src/commonMain/kotlin/pw/binom/agentik/client/HttpEventStore.kt new file mode 100644 index 0000000..3f8359a --- /dev/null +++ b/client/src/commonMain/kotlin/pw/binom/agentik/client/HttpEventStore.kt @@ -0,0 +1,128 @@ +package pw.binom.agentik.client + +import io.ktor.client.HttpClient +import io.ktor.client.request.prepareGet +import io.ktor.client.statement.bodyAsChannel +import io.ktor.http.HttpStatusCode +import kotlinx.coroutines.flow.Flow +import kotlinx.coroutines.flow.flow +import pw.binom.agentik.eventStore.EventStore +import pw.binom.agentik.proto.AgentEvent +import pw.binom.agentik.proto.CommonEvent +import pw.binom.agentik.proto.Event +import kotlin.time.Clock +import kotlin.time.Instant + +/** + * HTTP-реализация [EventStore], ходящая в `:server`-фасад. + * + * **Хитрый план**: вместо того, чтобы все методы шли в один общий endpoint и + * фильтровали client-side ([EventStore.events]/[filterIsInstance]), эта + * реализация бьёт запросы по URL'ам в зависимости от того, какой класс + * событий нужен: + * - [events] → `GET /events/all` (полный поток CommonEvent) + * - [agentEvents] → `GET /events` (только lifecycle диалогов) + * - [conversationEvents] с `conversationId != null` → `GET /conversations/{id}/events` + * + * Так серверный фильтр (SQL `WHERE` или разные буферы) работает на своей стороне, + * а клиент получает ровно тот срез, который ему нужен, без лишнего трафика. + * + * Для [conversationEvents] с `conversationId == null` (события всех диалогов) + * fallback на default [EventStore.conversationEvents] — общий поток `/events/all` + * + фильтр client-side. Это редкий кейс (admin-дашборды), и оптимизировать его + * отдельно нерационально. + * + * [earliestEventDate] не имеет своего endpoint'а; возвращает `Clock.System.now()` + * (см. KDoc [EventStore.earliestEventDate] — для пустого буфера это и есть + * контрактное значение). Клиент, который полагался на gap detection через + * message store, продолжит работать — просто fallback никогда не сработает. + */ +internal class HttpEventStore( + private val httpClient: HttpClient, + private val baseUrl: String, +) : EventStore { + + private val agentUrl: String = baseUrl.trimEnd('/') + + override fun events(after: Instant?): Flow = flow { + val url = buildString { + append("$agentUrl/events/all") + if (after != null) append("?after=$after") + } + httpClient.prepareGet(url) { noSseReadTimeout() } + .execute { response -> + check(response.status == HttpStatusCode.OK) { + "events: server returned ${response.status}" + } + readSse(response.bodyAsChannel()) + .collect { payload -> + emit(agentikJson.decodeFromString(CommonEvent.serializer(), payload)) + } + } + } + + /** + * Override: идём в `/events` напрямую — сервер фильтрует только lifecycle-события. + * Default из [EventStore.agentEvents] читал бы `/events/all` + `filterIsInstance`. + */ + override fun agentEvents(after: Instant?): Flow = flow { + val url = buildString { + append("$agentUrl/events") + if (after != null) append("?after=$after") + } + httpClient.prepareGet(url) { noSseReadTimeout() } + .execute { response -> + check(response.status == HttpStatusCode.OK) { + "agentEvents: server returned ${response.status}" + } + readSse(response.bodyAsChannel()) + .collect { payload -> + val event = agentikJson.decodeFromString(AgentEvent.serializer(), payload) + emit(CommonEvent.Agent(date = event.date, event = event)) + } + } + } + + /** + * Override с `conversationId != null` — идём в `/conversations/{id}/events`. + * С `null` (события всех диалогов) — fallback на default impl из [EventStore]: + * общий `/events/all` + filter. + */ + override fun conversationEvents( + after: Instant?, + conversationId: String?, + ): Flow { + if (conversationId == null) { + return super.conversationEvents(after, null) + } + return flow { + val url = buildString { + append("$agentUrl/conversations/$conversationId/events") + if (after != null) append("?after=$after") + } + httpClient.prepareGet(url) { noSseReadTimeout() } + .execute { response -> + check(response.status == HttpStatusCode.OK) { + "conversationEvents: server returned ${response.status}" + } + readSse(response.bodyAsChannel()) + .collect { payload -> + val event = agentikJson.decodeFromString(Event.serializer(), payload) + emit(CommonEvent.Conversation(date = event.date, conversationId = conversationId, event = event)) + } + } + } + } + + /** + * У HTTP-варианта нет своего endpoint'а для earliest-event-date. + * Контракт [EventStore.earliestEventDate] для пустого буфера говорит + * "сейчас" — для HTTP-клиента буфер на нашей стороне всегда "пуст" + * (мы не держим своё состояние), поэтому возвращаем `Clock.System.now()`. + */ + override suspend fun earliestEventDate(): Instant = Clock.System.now() + + override fun close() { + // HttpClient закрывает владелец (AgentClient / AgentikAgent). + } +} diff --git a/event-store/build.gradle.kts b/event-store/build.gradle.kts index bce981d..342eadf 100644 --- a/event-store/build.gradle.kts +++ b/event-store/build.gradle.kts @@ -13,7 +13,10 @@ kotlin { jvmToolchain(21) jvm() + macosX64() + macosArm64() linuxX64() + linuxArm64() mingwX64() sourceSets { diff --git a/message-store-api/src/commonMain/kotlin/pw/binom/agentik/messageStore/events/EventStore.kt b/message-store-api/src/commonMain/kotlin/pw/binom/agentik/messageStore/events/EventStore.kt deleted file mode 100644 index 2caa695..0000000 --- a/message-store-api/src/commonMain/kotlin/pw/binom/agentik/messageStore/events/EventStore.kt +++ /dev/null @@ -1,109 +0,0 @@ -package pw.binom.agentik.messageStore.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/server/build.gradle.kts b/server/build.gradle.kts index e1bba06..ebc506e 100644 --- a/server/build.gradle.kts +++ b/server/build.gradle.kts @@ -23,8 +23,6 @@ kotlin { sourceSets { commonMain.dependencies { implementation(project(":proto")) - implementation(project(":message-store-api")) - implementation(project(":working-memory-api")) // 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 9f97060..29d4665 100644 --- a/server/src/commonMain/kotlin/pw/binom/agentik/server/Module.kt +++ b/server/src/commonMain/kotlin/pw/binom/agentik/server/Module.kt @@ -6,7 +6,6 @@ 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.messageStore.events.EventStore /** * Встраивает HTTP/SSE-фасад протокола agentik в твой Ktor-роутинг. @@ -30,18 +29,19 @@ import pw.binom.agentik.messageStore.events.EventStore * - `POST /conversations/{id}/messages` — `send` (202 Accepted) * - `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 /conversations/{id}/events` — SSE: события хода (catchup + live) + * - `GET /events` — SSE: события агента (catchup + live) + * - `GET /events/all` — SSE: всё в одном потоке * - `GET /health` — `"ok"` * - * @param eventStore если null — `/events/replay` endpoints возвращают 503 (event - * persistence не настроен). Live-streaming работает как обычно. + * Event-эндпоинты сами по себе дают catchup + live в одном Flow: клиент + * передаёт `?after=`, сервер сначала отдаёт буферизованные события + * с `date > after`, потом переключается на live. Отдельный replay-endpoint + * не нужен — для событий старше буфера клиент должен идти в + * `:message-log-api` (полный audit log). */ fun Route.agentikAgent( agent: Agent, - eventStore: EventStore? = null, path: String = "/agentik", token: String? = null, ) { @@ -54,6 +54,6 @@ fun Route.agentikAgent( this.token = token } } - agentikRoutes(agent, eventStore) + agentikRoutes(agent) } } 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 3da8183..1c2e2e8 100644 --- a/server/src/commonMain/kotlin/pw/binom/agentik/server/Routes.kt +++ b/server/src/commonMain/kotlin/pw/binom/agentik/server/Routes.kt @@ -24,10 +24,9 @@ import pw.binom.agentik.proto.CommonEvent import pw.binom.agentik.proto.Conversation import pw.binom.agentik.proto.Content import pw.binom.agentik.proto.Event -import pw.binom.agentik.messageStore.events.EventStore import kotlin.time.Instant -internal fun Route.agentikRoutes(agent: Agent, eventStore: EventStore? = null) { +internal fun Route.agentikRoutes(agent: Agent) { get("/health") { call.respondText("ok") @@ -145,38 +144,6 @@ internal fun Route.agentikRoutes(agent: Agent, eventStore: EventStore? = null) { val after = call.parseAfter() ?: return@get call.streamJsonSse(agent.allEvents(after), CommonEvent.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/settings.gradle.kts b/settings.gradle.kts index f8b639f..504ecc4 100644 --- a/settings.gradle.kts +++ b/settings.gradle.kts @@ -77,7 +77,6 @@ include(":event-store") // eviction. Для тестов, dev-режима, embedded-сценариев (Android core). include(":event-store-in-memory") include(":working-memory-api") -include(":storage-bundle") include(":storage-inmemory") // SQLDelight-реализация store'ов из :message-store-api и :working-memory-api. JVM-only // native драйверов для KMP вне JVM пока не публикует). Содержит 4 .sq-файла + diff --git a/standalone/build.gradle.kts b/standalone/build.gradle.kts index 8ab0fab..0f88189 100644 --- a/standalone/build.gradle.kts +++ b/standalone/build.gradle.kts @@ -32,9 +32,11 @@ kotlin { sourceSets { commonMain.dependencies { - api(project(":storage-bundle")) implementation(project(":proto")) implementation(project(":server")) + implementation(project(":message-store-api")) + implementation(project(":message-log-api")) + implementation(project(":working-memory-api")) // Commons implementation(libs.kotlinx.coroutines.core) diff --git a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/DebugRoutes.kt b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/DebugRoutes.kt index beaf990..6fc4af4 100644 --- a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/DebugRoutes.kt +++ b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/DebugRoutes.kt @@ -9,13 +9,15 @@ import io.ktor.server.routing.post import kotlinx.serialization.json.buildJsonObject import kotlinx.serialization.json.put import pw.binom.agentik.memory.ConversationTurn +import pw.binom.agentik.messageStore.ReflectionStore +import pw.binom.agentik.messageLog.MessageStore import pw.binom.agentik.proto.Agent import pw.binom.agentik.standalone.agent.ChatConversation import pw.binom.agentik.llm.tools.LlmReflector import pw.binom.agentik.llm.tools.SkillMiner import pw.binom.agentik.standalone.agent.memory.Curator -import pw.binom.agentik.storageBundle.StorageBundle import pw.binom.agentik.skills.SkillStore +import pw.binom.agentik.workingMemory.WorkingMemoryStore /** * Debug-эндпоинты для ручного триггерирования фоновых фич (без ожидания @@ -33,7 +35,9 @@ import pw.binom.agentik.skills.SkillStore */ internal fun Route.debugRoutes( agent: Agent, - storage: pw.binom.agentik.storageBundle.StorageBundle, + messageStore: MessageStore, + workingMemoryStore: WorkingMemoryStore, + reflectionStore: ReflectionStore, reflector: LlmReflector?, skillMiner: SkillMiner?, skillStore: SkillStore?, @@ -44,14 +48,14 @@ internal fun Route.debugRoutes( ?: return@post call.respondText("conversationId required", status = HttpStatusCode.BadRequest) val minerReflector = reflector if (minerReflector == null) return@post call.respondText("reflection disabled", status = HttpStatusCode.NotFound) - val turns = recentTurns(storage, convId, maxTurns = 6) + val turns = recentTurns(workingMemoryStore, convId, maxTurns = 6) if (turns.isEmpty()) return@post call.respondText("no turns in conversation", status = HttpStatusCode.NotFound) val reflection = minerReflector.reflect(turns) if (reflection == null) { call.respondText("""{"reflected":false,"reason":"unparseable LLM reply"}""", contentType = ContentType.Application.Json) } else { val stamped = reflection.copy(conversationId = convId) - storage.reflectionStore.insert(stamped) + reflectionStore.insert(stamped) call.respondText( buildJsonObject { put("reflected", true) @@ -71,7 +75,7 @@ internal fun Route.debugRoutes( val miner = skillMiner val store = skillStore if (miner == null || store == null) return@post call.respondText("skill mining disabled", status = HttpStatusCode.NotFound) - val turns = recentTurns(storage, convId, maxTurns = miner.maxTurns) + val turns = recentTurns(workingMemoryStore, convId, maxTurns = miner.maxTurns) if (turns.isEmpty()) return@post call.respondText("no turns in conversation", status = HttpStatusCode.NotFound) val mined = miner.mine(turns, store.catalog.skills) for (s in mined) store.upsert(s) @@ -110,7 +114,7 @@ internal fun Route.debugRoutes( get("/debug/tokens") { val convId = call.parameters["conversationId"] ?: return@get call.respondText("conversationId required", status = HttpStatusCode.BadRequest) - val stats = storage.messageStore.tokenStats(convId) + val stats = messageStore.tokenStats(convId) val json = buildJsonObject { put("conversationId", convId) put("turns", stats.turns.toString()) @@ -127,8 +131,8 @@ internal fun Route.debugRoutes( * (для debug-триггеров reflector/miner; та же логика, что у хуков * [ChatConversation]). */ -internal suspend fun recentTurns(storage: pw.binom.agentik.storageBundle.StorageBundle, conversationId: String, maxTurns: Int): List { - val rows = storage.workingMemoryStore.list(conversationId) +internal suspend fun recentTurns(workingMemoryStore: WorkingMemoryStore, conversationId: String, maxTurns: Int): List { + val rows = workingMemoryStore.list(conversationId) val pairs = mutableListOf() var pendingUser: String? = null for (row in rows) { diff --git a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/Main.kt b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/Main.kt index b9f4efd..8303efd 100644 --- a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/Main.kt +++ b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/Main.kt @@ -198,7 +198,7 @@ private fun runServer() { } val llm = config.llm.createLlm() - val storage = SqliteStores.open(dbPath = config.agent.dbPath).asBundle() + val sqliteStores = SqliteStores.open(dbPath = config.agent.dbPath) val mcpRegistry = McpRegistry.fromConfig(config.mcp) // Хранилище скилов: если skillsDir задан, читаем каталог + создаём @@ -319,12 +319,15 @@ private fun runServer() { if (config.reflection.interval > 0) LlmReflector(llm = llm) else null val recentReflections: List = if (config.reflection.topK > 0) kotlinx.coroutines.runBlocking { - storage.reflectionStore.listRecent(config.reflection.topK) + sqliteStores.reflections.listRecent(config.reflection.topK) } else emptyList() val agent = ChatAgent( id = "agentik", - storage = storage, + conversationStore = sqliteStores.conversations, + messageStore = sqliteStores.messages, + workingMemoryStore = sqliteStores.workingMemory, + reflectionStore = sqliteStores.reflections, llm = llm, llmConfig = config.llm, tools = mcpRegistry.namedTools, @@ -365,7 +368,9 @@ private fun runServer() { if (config.debug.endpoints) { debugRoutes( agent = agent, - storage = storage, + messageStore = sqliteStores.messages, + workingMemoryStore = sqliteStores.workingMemory, + reflectionStore = sqliteStores.reflections, reflector = reflector, skillMiner = skillMiner, skillStore = skillStore, @@ -403,14 +408,14 @@ private fun runServer() { // Token stats по существующим диалогам (агрегат на старте — каждая запись // парсится из payload_json, ну >100 turns и БД приличная — но в рамках // стартапа это терпимо). - val existingConvs = kotlinx.coroutines.runBlocking { storage.conversationStore.list(offset = 0, limit = 1000) } + val existingConvs = kotlinx.coroutines.runBlocking { sqliteStores.conversations.list(offset = 0, limit = 1000) } if (existingConvs.isNotEmpty()) { var totalTurns = 0 var totalIn = 0L var totalOut = 0L for (c in existingConvs) { if (c.isTemporal) continue - val s = kotlinx.coroutines.runBlocking { storage.messageStore.tokenStats(c.id) } + val s = kotlinx.coroutines.runBlocking { sqliteStores.messages.tokenStats(c.id) } totalTurns += s.turns totalIn += s.inputTokens totalOut += s.outputTokens @@ -422,7 +427,9 @@ private fun runServer() { Runtime.getRuntime().addShutdownHook(Thread { agent.close() mcpRegistry.close() - storage.close() + }) + Runtime.getRuntime().addShutdownHook(Thread { + sqliteStores.close() llm.close() memorySystem?.close() }) 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 5d87d4b..aeb0fbb 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 @@ -26,10 +26,12 @@ 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.messageStore.ConversationRecord +import pw.binom.agentik.messageStore.ConversationStore import pw.binom.agentik.messageStore.Ids import pw.binom.agentik.messageStore.Reflection -import pw.binom.agentik.messageStore.events.EventRecord -import pw.binom.agentik.messageStore.events.EventType +import pw.binom.agentik.messageStore.ReflectionStore +import pw.binom.agentik.messageLog.MessageStore +import pw.binom.agentik.workingMemory.WorkingMemoryStore import pw.binom.agentik.toolsets.DisableToolsetTool import pw.binom.agentik.toolsets.EnableToolsetTool import pw.binom.agentik.toolsets.NamedTool @@ -61,7 +63,10 @@ import pw.binom.litert.LiteLlm */ class ChatAgent( override val id: String, - private val storage: pw.binom.agentik.storageBundle.StorageBundle, + private val conversationStore: ConversationStore, + private val messageStore: MessageStore, + private val workingMemoryStore: WorkingMemoryStore, + private val reflectionStore: ReflectionStore, private val llm: LiteLlm, private val llmConfig: LlmConfig, private val tools: List = emptyList(), @@ -251,12 +256,15 @@ class ChatAgent( // и не переживают рестарт агента (см. Memory #3709). if (!temp) { runBlocking { - storage.conversationStore.upsert(rec) + conversationStore.upsert(rec) } } val conv = ChatConversation( record = rec, - storage = storage, + conversationStore = conversationStore, + messageStore = messageStore, + workingMemoryStore = workingMemoryStore, + reflectionStore = reflectionStore, eventStore = eventStore, llm = llm, systemPrompt = systemPrompt, @@ -268,7 +276,6 @@ class ChatAgent( contextWindow = contextWindow, compressionThreshold = compressionThreshold, contextCompactor = contextCompactor, - reflectionStore = storage.reflectionStore, reflector = reflector, skillMiner = skillMiner, skillMiningStore = skillStore, @@ -289,7 +296,7 @@ class ChatAgent( override suspend fun getConversation(id: String): ProtoConversation? { liveLock.withLock { live[id] }?.let { if (!it.isClosed) return it } - val rec = storage.conversationStore.get(id) ?: return null + val rec = conversationStore.get(id) ?: return null return newConversation(rec).also { liveLock.withLock { live[id] = it } } @@ -298,7 +305,7 @@ class ChatAgent( override suspend fun deleteConversation(id: String): Boolean { val conv = liveLock.withLock { live.remove(id) } conv?.close() - val ok = storage.conversationStore.delete(id) + val ok = conversationStore.delete(id) if (ok) { val event = AgentEvent.Deleted(date = now(), id = id) eventStore.append(CommonEvent.Agent(date = now(), event = event)) @@ -307,7 +314,7 @@ class ChatAgent( } override suspend fun getConversations(offset: Int, limit: Int): List = - storage.conversationStore.list(offset = offset, limit = limit).map { rec -> + conversationStore.list(offset = offset, limit = limit).map { rec -> liveLock.withLock { live[rec.id] } ?: newConversation(rec).also { liveLock.withLock { live[rec.id] = it } @@ -316,7 +323,10 @@ class ChatAgent( private fun newConversation(rec: ConversationRecord): ChatConversation = ChatConversation( record = rec, - storage = storage, + conversationStore = conversationStore, + messageStore = messageStore, + workingMemoryStore = workingMemoryStore, + reflectionStore = reflectionStore, eventStore = eventStore, llm = llm, systemPrompt = systemPrompt, @@ -328,11 +338,10 @@ class ChatAgent( contextWindow = contextWindow, compressionThreshold = compressionThreshold, contextCompactor = contextCompactor, - reflectionStore = storage.reflectionStore, - reflector = reflector, - skillMiner = skillMiner, - skillMiningStore = skillStore, - ) + reflector = reflector, + skillMiner = skillMiner, + skillMiningStore = skillStore, + ) override fun close() { runBlocking { 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 e202194..47d5eb5 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 @@ -36,7 +36,6 @@ import pw.binom.agentik.messageLog.MessageOrigin import pw.binom.agentik.messageLog.MessageRecord import pw.binom.agentik.messageLog.MessageStore import pw.binom.agentik.messageStore.ReflectionStore -import pw.binom.agentik.storageBundle.StorageBundle import pw.binom.agentik.messageLog.TurnTokens import pw.binom.agentik.workingMemory.WorkingMemoryEntry import pw.binom.agentik.workingMemory.WorkingMemoryStore @@ -55,7 +54,10 @@ import pw.binom.agentik.toolsets.NamedTool class ConversationLoop( record: ConversationRecord, - private val storage: StorageBundle, + private val conversationStore: ConversationStore, + private val messageStore: MessageStore, + private val workingMemoryStore: WorkingMemoryStore, + private val reflectionStore: ReflectionStore?, private val eventStore: pw.binom.agentik.eventStore.MutableEventStore, private val llm: LiteLlm, private val systemPrompt: String, @@ -67,7 +69,6 @@ class ConversationLoop( private val contextWindow: Int? = null, private val compressionThreshold: Double = 0.8, private val contextCompactor: ContextCompactor? = null, - private val reflectionStore: ReflectionStore? = null, private val reflector: LlmReflector? = null, private val skillMiner: SkillMiner? = null, private val skillMiningStore: SkillStore? = null, @@ -93,10 +94,6 @@ class ConversationLoop( /** Per-conversation background event bus. Lifecycle scoped к этому ConversationLoop. */ private val backgroundEvents = BackgroundEventBus() - private val conversationStore: ConversationStore get() = storage.conversationStore - private val messageStore: MessageStore get() = storage.messageStore - private val workingMemory: WorkingMemoryStore get() = storage.workingMemoryStore - private val toolsByName: MutableMap = tools.associateBy { it.name }.toMutableMap() private val contextBuilder = ContextBuilder(memoryPrefetcher = memoryPrefetcher) @@ -108,7 +105,7 @@ class ConversationLoop( contextCompactor = contextCompactor, memoryReviewer = memoryReviewer, memoryStoreForReview = memoryStoreForReview, - workingMemory = workingMemory, + workingMemory = workingMemoryStore, liteLlm = llm, systemPrompt = systemPrompt, backgroundEvents = backgroundEvents, @@ -128,7 +125,7 @@ class ConversationLoop( private val backgroundScheduler = BackgroundScheduler( state = state, - workingMemory = workingMemory, + workingMemory = workingMemoryStore, config = BackgroundConfig( memoryReviewer = memoryReviewer, memoryStore = memoryStoreForReview, @@ -179,7 +176,7 @@ class ConversationLoop( if (!state.isTemporal) { messageStore.append(userRecord) - workingMemory.append( + workingMemoryStore.append( conversationId = id, entry = WorkingMemoryEntry.User( sourceMessageId = userMessageId, @@ -412,7 +409,7 @@ class ConversationLoop( ) messageStore.append(assistantRecord) - workingMemory.append( + workingMemoryStore.append( conversationId = id, entry = WorkingMemoryEntry.Assistant( sourceMessageId = assistantId, @@ -422,7 +419,7 @@ class ConversationLoop( ) for (ex in toolExchanges) { - workingMemory.append( + workingMemoryStore.append( conversationId = id, entry = ex, now = assistantAt, diff --git a/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/ChatAgentTest.kt b/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/ChatAgentTest.kt index f63c9f6..d72961c 100644 --- a/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/ChatAgentTest.kt +++ b/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/ChatAgentTest.kt @@ -41,28 +41,31 @@ import pw.binom.agentik.toolsets.NamedTool class ChatAgentTest { - private lateinit var storage: pw.binom.agentik.storageBundle.StorageBundle + private lateinit var sqliteStores: pw.binom.agentik.storage.sqlite.SqliteStores private lateinit var fakeLlm: FakeLiteLlm @BeforeTest fun setup() { - storage = SqliteStores.inMemory().asBundle() + sqliteStores = SqliteStores.inMemory() fakeLlm = FakeLiteLlm() } @AfterTest fun tearDown() { - storage.close() + sqliteStores.close() } private fun newAgent( - storage: pw.binom.agentik.storageBundle.StorageBundle = this.storage, + sqliteStores: pw.binom.agentik.storage.sqlite.SqliteStores = this.sqliteStores, llm: LiteLlm = this.fakeLlm, tools: List = emptyList(), skills: SkillCatalog = SkillCatalog.EMPTY, ): ChatAgent = ChatAgent( id = "agentik", - storage = storage, + conversationStore = sqliteStores.conversations, + messageStore = sqliteStores.messages, + workingMemoryStore = sqliteStores.workingMemory, + reflectionStore = sqliteStores.reflections, llm = llm, llmConfig = LlmConfig( backend = LlmBackend.OPENAI, @@ -82,7 +85,7 @@ class ChatAgentTest { val agent = newAgent() val conv = agent.createConversation(temp = false) as ChatConversation - val wm = storage.workingMemoryStore.list(conv.id) + val wm = sqliteStores.workingMemory.list(conv.id) assertEquals(0, wm.size) // System prompt виден через LiteConversationConfig, который LLM получит // при первом send (см. `send passes system prompt and past history to LLM on first send`). @@ -159,7 +162,7 @@ class ChatAgentTest { val id = conv.id // добавим сообщение, чтобы потом убедиться, что каскад сработал - storage.messageStore.append( + sqliteStores.messages.append( pw.binom.agentik.messageLog.MessageRecord.UserMessage( id = "m1", conversationId = id, @@ -169,8 +172,8 @@ class ChatAgentTest { ) assertTrue(agent.deleteConversation(id)) assertNull(agent.getConversation(id)) - assertNull(storage.conversationStore.get(id)) - assertEquals(emptyList(), storage.messageStore.listAll(id)) + assertNull(sqliteStores.conversations.get(id)) + assertEquals(emptyList(), sqliteStores.messages.listAll(id)) } @Test @@ -188,7 +191,7 @@ class ChatAgentTest { conv.send(listOf(Content.Text("hi"))) // user message записан в audit + working memory - val msgs = storage.messageStore.listAll(conv.id) + val msgs = sqliteStores.messages.listAll(conv.id) assertEquals(2, msgs.size) assertEquals("hi", (msgs[0] as pw.binom.agentik.messageLog.MessageRecord.UserMessage).content.let { (it[0] as pw.binom.agentik.messageLog.Content.Text).body @@ -214,9 +217,9 @@ class ChatAgentTest { // В working_memory теперь НЕТ System-entries — только user + assistant. // Системный промт живёт в ChatConversation.systemPrompt и едет в LLM // через LiteConversationConfig.systemInstruction. - val wm1 = storage.workingMemoryStore.list(conv1.id) + val wm1 = sqliteStores.workingMemory.list(conv1.id) assertEquals(2, wm1.size) - val wm2 = storage.workingMemoryStore.list(conv2.id) + val wm2 = sqliteStores.workingMemory.list(conv2.id) assertEquals(2, wm2.size) } @@ -256,12 +259,12 @@ class ChatAgentTest { fakeLlm.reply = "first reply" conv.send(listOf(Content.Text("first user"))) // первый turn: WM = [user, assistant] - assertEquals(2, storage.workingMemoryStore.list(conv.id).size) + assertEquals(2, sqliteStores.workingMemory.list(conv.id).size) fakeLlm.reply = "second reply" conv.send(listOf(Content.Text("second user"))) // второй turn: WM должен вырасти до [user, assistant, user, assistant] - val wm = storage.workingMemoryStore.list(conv.id) + val wm = sqliteStores.workingMemory.list(conv.id) System.err.println("[TEST] wm.size=${wm.size}") wm.forEachIndexed { i, row -> System.err.println("[TEST] $i: ${row.entry::class.simpleName} id=${row.id}") } assertEquals(4, wm.size) @@ -325,7 +328,7 @@ class ChatAgentTest { assertTrue(events.any { it is ProtoEvent.Error && it.message == "boom from llm" }, "events=$events") assertTrue(events.any { it is ProtoEvent.End }, "events=$events") - val msgs = storage.messageStore.listAll(conv.id) + val msgs = sqliteStores.messages.listAll(conv.id) assertEquals(2, msgs.size) assertIs(msgs[0]) val err = assertIs(msgs[1]) @@ -370,12 +373,12 @@ class ChatAgentTest { eventsJob.cancel() // audit: только user (assistant не успел сгенериться) - val msgs = storage.messageStore.listAll(conv.id) + val msgs = sqliteStores.messages.listAll(conv.id) assertEquals(1, msgs.size) assertIs(msgs[0]) // working memory: только user (assistant skipped because пустой) - val wm = storage.workingMemoryStore.list(conv.id) + val wm = sqliteStores.workingMemory.list(conv.id) assertEquals(1, wm.size) assertTrue(wm[0].entry is WorkingMemoryEntry.User) @@ -427,14 +430,14 @@ class ChatAgentTest { eventsJob.cancel() // audit: user + toolcall + toolresult (tool выполнился), assistant может быть - val msgs = storage.messageStore.listAll(conv.id) + val msgs = sqliteStores.messages.listAll(conv.id) val toolResult = msgs.filterIsInstance().firstOrNull() assertNotNull(toolResult, "tool result должен быть в audit — tool выполнился нормально") val toolResultResult = toolResult!!.result!! assertTrue(toolResultResult.contains("echo"), "tool result содержит реальный ответ тулы: $toolResultResult") // working memory: user + tool_exchange - val wm = storage.workingMemoryStore.list(conv.id) + val wm = sqliteStores.workingMemory.list(conv.id) val exchanges = wm.mapNotNull { (it.entry as? WorkingMemoryEntry.ToolExchange) } assertEquals(1, exchanges.size) assertEquals("echo_tool", exchanges[0].toolName) @@ -450,12 +453,15 @@ class ChatAgentTest { @Test fun `temp conversation is not persisted across agent instances`() = runTest { // Поднимаем file-backed БД, создаём temp-беседу - storage.close() + sqliteStores.close() val dbPath = (System.getProperty("java.io.tmpdir") + "/agentik-test-${System.nanoTime()}.db") - storage = SqliteStores.open(dbPath).asBundle() + sqliteStores = SqliteStores.open(dbPath) val agent1 = ChatAgent( id = "agentik", - storage = storage, + conversationStore = sqliteStores.conversations, + messageStore = sqliteStores.messages, + workingMemoryStore = sqliteStores.workingMemory, + reflectionStore = sqliteStores.reflections, llm = FakeLiteLlm().also { fakeLlm = it }, llmConfig = LlmConfig( backend = pw.binom.agentik.standalone.llm.LlmBackend.OPENAI, @@ -468,11 +474,14 @@ class ChatAgentTest { assertNotNull(agent1.getConversation(tempId)) // Переоткрываем БД — temp-беседа не должна пережить рестарт - storage.close() - storage = SqliteStores.open(dbPath).asBundle() + sqliteStores.close() + sqliteStores = SqliteStores.open(dbPath) val agent2 = ChatAgent( id = "agentik", - storage = storage, + conversationStore = sqliteStores.conversations, + messageStore = sqliteStores.messages, + workingMemoryStore = sqliteStores.workingMemory, + reflectionStore = sqliteStores.reflections, llm = fakeLlm, llmConfig = LlmConfig( backend = pw.binom.agentik.standalone.llm.LlmBackend.OPENAI, @@ -486,12 +495,15 @@ class ChatAgentTest { @Test fun `non-temp conversation persists across agent instances`() = runTest { - storage.close() + sqliteStores.close() val dbPath = (System.getProperty("java.io.tmpdir") + "/agentik-test-${System.nanoTime()}.db") - storage = SqliteStores.open(dbPath).asBundle() + sqliteStores = SqliteStores.open(dbPath) val agent1 = ChatAgent( id = "agentik", - storage = storage, + conversationStore = sqliteStores.conversations, + messageStore = sqliteStores.messages, + workingMemoryStore = sqliteStores.workingMemory, + reflectionStore = sqliteStores.reflections, llm = FakeLiteLlm().also { fakeLlm = it }, llmConfig = LlmConfig( backend = pw.binom.agentik.standalone.llm.LlmBackend.OPENAI, @@ -502,11 +514,14 @@ class ChatAgentTest { val conv = agent1.createConversation(temp = false) val id = conv.id - storage.close() - storage = SqliteStores.open(dbPath).asBundle() + sqliteStores.close() + sqliteStores = SqliteStores.open(dbPath) val agent2 = ChatAgent( id = "agentik", - storage = storage, + conversationStore = sqliteStores.conversations, + messageStore = sqliteStores.messages, + workingMemoryStore = sqliteStores.workingMemory, + reflectionStore = sqliteStores.reflections, llm = fakeLlm, llmConfig = LlmConfig( backend = pw.binom.agentik.standalone.llm.LlmBackend.OPENAI, diff --git a/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/ChatAgentToolsetsTest.kt b/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/ChatAgentToolsetsTest.kt index 2714205..bf1efe2 100644 --- a/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/ChatAgentToolsetsTest.kt +++ b/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/ChatAgentToolsetsTest.kt @@ -39,11 +39,14 @@ class ChatAgentToolsetsTest { private fun newAgent( toolsets: List = emptyList(), - ): Pair { - val storage = SqliteStores.inMemory().asBundle() + ): Pair { + val sqliteStores = SqliteStores.inMemory() val agent = ChatAgent( id = "test-agent", - storage = storage, + conversationStore = sqliteStores.conversations, + messageStore = sqliteStores.messages, + workingMemoryStore = sqliteStores.workingMemory, + reflectionStore = sqliteStores.reflections, llm = stubLlm(), llmConfig = LlmConfig( backend = pw.binom.agentik.standalone.llm.LlmBackend.GOOGLE, @@ -52,7 +55,7 @@ class ChatAgentToolsetsTest { ), toolsets = toolsets, ) - return agent to storage + return agent to sqliteStores } @Test diff --git a/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/CompactionTest.kt b/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/CompactionTest.kt index fdadd3a..2ca6be3 100644 --- a/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/CompactionTest.kt +++ b/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/CompactionTest.kt @@ -40,18 +40,18 @@ import pw.binom.agentik.llm.tools.SummaryTurn */ class CompactionTest { - private lateinit var storage: pw.binom.agentik.storageBundle.StorageBundle + private lateinit var sqliteStores: pw.binom.agentik.storage.sqlite.SqliteStores private lateinit var fakeLlm: FakeLiteLlm @BeforeTest fun setup() { - storage = SqliteStores.inMemory().asBundle() + sqliteStores = SqliteStores.inMemory() fakeLlm = FakeLiteLlm() } @AfterTest fun tearDown() { - storage.close() + sqliteStores.close() } private fun newAgent( @@ -63,7 +63,10 @@ class CompactionTest { val reviewer = if (memoryStore != null) KeywordMdReviewer() else null return ChatAgent( id = "test", - storage = storage, + conversationStore = sqliteStores.conversations, + messageStore = sqliteStores.messages, + workingMemoryStore = sqliteStores.workingMemory, + reflectionStore = sqliteStores.reflections, llm = fakeLlm, llmConfig = LlmConfig( backend = LlmBackend.OPENAI, @@ -87,7 +90,7 @@ class CompactionTest { repeat(10) { conv.send(listOf(ProtoContent.Text("turn $it: ${"x".repeat(200)}"))) } - val wm = storage.workingMemoryStore.list(conv.id) + val wm = sqliteStores.workingMemory.list(conv.id) // Без compaction все ходы остаются в памяти (System + 10 user/assistant = 21 строк). val summaries = wm.filter { it.entry is pw.binom.agentik.workingMemory.WorkingMemoryEntry.Summary } assertEquals(0, summaries.size, "compaction must not run without contextWindow") @@ -99,7 +102,7 @@ class CompactionTest { val agent = newAgent(contextWindow = 10, compactor = null) val conv = agent.createConversation(temp = false) as ChatConversation conv.send(listOf(ProtoContent.Text("first"))) - val wm = storage.workingMemoryStore.list(conv.id) + val wm = sqliteStores.workingMemory.list(conv.id) // System + User + Assistant = 3. Без compactor — никаких Summary. val summaries = wm.filter { it.entry is pw.binom.agentik.workingMemory.WorkingMemoryEntry.Summary } assertEquals(0, summaries.size, "no compaction runs without compactor") @@ -118,7 +121,7 @@ class CompactionTest { // Compactor должен был быть вызван хотя бы раз. assertTrue(compactor.calls > 0, "compactor must be called at least once when above threshold") // В working memory должна появиться Summary. - val wm = storage.workingMemoryStore.list(conv.id) + val wm = sqliteStores.workingMemory.list(conv.id) val summaries = wm.filter { it.entry is pw.binom.agentik.workingMemory.WorkingMemoryEntry.Summary } assertTrue(summaries.isNotEmpty(), "at least one Summary entry should be present after compaction") // Summary-текст — то, что вернул наш compactor. @@ -158,7 +161,7 @@ class CompactionTest { conv.send(listOf(ProtoContent.Text("second turn"))) conv.send(listOf(ProtoContent.Text("third turn — long content ${"y".repeat(150)}"))) - val wm = storage.workingMemoryStore.list(conv.id) + val wm = sqliteStores.workingMemory.list(conv.id) // Должны быть: System + хотя бы один Summary + последние KEEP_RECENT_TURNS ходов. // KEEP_RECENT_TURNS = 4 → user/assistant последних двух ходов (third + second) могут быть не тронуты. val userAssistantCount = wm.count { diff --git a/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/MemoryWiringTest.kt b/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/MemoryWiringTest.kt index acdddf6..003268e 100644 --- a/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/MemoryWiringTest.kt +++ b/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/MemoryWiringTest.kt @@ -49,13 +49,13 @@ import kotlin.time.Instant */ class MemoryWiringTest { - private lateinit var storage: pw.binom.agentik.storageBundle.StorageBundle + private lateinit var sqliteStores: pw.binom.agentik.storage.sqlite.SqliteStores private lateinit var fakeLlm: FakeLiteLlm private lateinit var root: Path @BeforeTest fun setup() { - storage = SqliteStores.inMemory().asBundle() + sqliteStores = SqliteStores.inMemory() fakeLlm = FakeLiteLlm() root = Path(SystemTemporaryDirectory.toString(), "agentik-mem-${java.util.UUID.randomUUID()}") SystemFileSystem.createDirectories(root, mustCreate = true) @@ -63,7 +63,7 @@ class MemoryWiringTest { @AfterTest fun tearDown() { - storage.close() + sqliteStores.close() runCatching { SystemFileSystem.delete(root, mustExist = false) } } @@ -76,7 +76,10 @@ class MemoryWiringTest { contextCompactor: ContextCompactor = EchoCompactor, ): ChatAgent = ChatAgent( id = "agentik", - storage = storage, + conversationStore = sqliteStores.conversations, + messageStore = sqliteStores.messages, + workingMemoryStore = sqliteStores.workingMemory, + reflectionStore = sqliteStores.reflections, llm = fakeLlm, llmConfig = LlmConfig( backend = LlmBackend.OPENAI, @@ -112,7 +115,10 @@ class MemoryWiringTest { val soulBody = "I am a helpful test persona. I always answer in one short line." val agent = ChatAgent( id = "agentik", - storage = storage, + conversationStore = sqliteStores.conversations, + messageStore = sqliteStores.messages, + workingMemoryStore = sqliteStores.workingMemory, + reflectionStore = sqliteStores.reflections, llm = fakeLlm, llmConfig = LlmConfig( backend = LlmBackend.OPENAI, @@ -137,7 +143,10 @@ class MemoryWiringTest { fun `soul body not added when null`() = runBlocking { val agent = ChatAgent( id = "agentik", - storage = storage, + conversationStore = sqliteStores.conversations, + messageStore = sqliteStores.messages, + workingMemoryStore = sqliteStores.workingMemory, + reflectionStore = sqliteStores.reflections, llm = fakeLlm, llmConfig = LlmConfig( backend = LlmBackend.OPENAI, diff --git a/storage-bundle/build.gradle.kts b/storage-bundle/build.gradle.kts deleted file mode 100644 index a299fad..0000000 --- a/storage-bundle/build.gradle.kts +++ /dev/null @@ -1,24 +0,0 @@ -plugins { - alias(libs.plugins.kotlin.multiplatform) -} - -// Агрегатор всех 5 storage-интерфейсов (MessageStore + ReflectionStore + -// EventStore + ConversationStore + WorkingMemoryStore). Зависит от обоих API -// модулей. Нужен только server-side рантайму (`:standalone`, `:agentik-cli`, -// `:agent-toolsets`, будущий `:android-agent` core). Тонкие клиенты этот модуль -// НЕ подтягивают — они работают с одним только `:message-store-api`. - -kotlin { - jvmToolchain(21) - - jvm() - linuxX64() - mingwX64() - - sourceSets { - commonMain.dependencies { - api(project(":message-store-api")) - api(project(":working-memory-api")) - } - } -} diff --git a/storage-bundle/src/commonMain/kotlin/pw/binom/agentik/storageBundle/StorageBundle.kt b/storage-bundle/src/commonMain/kotlin/pw/binom/agentik/storageBundle/StorageBundle.kt deleted file mode 100644 index 570d7bd..0000000 --- a/storage-bundle/src/commonMain/kotlin/pw/binom/agentik/storageBundle/StorageBundle.kt +++ /dev/null @@ -1,42 +0,0 @@ -package pw.binom.agentik.storageBundle - -import pw.binom.agentik.messageStore.ConversationStore -import pw.binom.agentik.messageLog.MessageStore -import pw.binom.agentik.messageStore.ReflectionStore -import pw.binom.agentik.messageStore.events.EventStore -import pw.binom.agentik.workingMemory.WorkingMemoryStore - -/** - * Агрегатор всех storage-интерфейсов для полного агент-рантайма. - * - * Зависит от **обоих** API модулей — `:message-store-api` и `:working-memory-api`. - * Тонкие клиенты, которым нужен только audit log, могут не подтягивать этот модуль - * и работать с [MessageStore] / [ReflectionStore] / [EventStore] / [ConversationStore] - * напрямую. - * - * **Реализации агрегата**: - * - [SqliteStorageBundle] — JVM-only, SQLDelight + SQLite. В `:storage-sqlite`. - * - [InMemoryStorageBundle] — KMP, Map-based. В `:storage-inmemory`. - * - * `SkillStore` НЕ входит — он в `:skills` (контент-файлы навыков), принимается - * отдельно в `ChatAgent`. - * - * AutoCloseable: один `close()` закрывает все store'ы; в реализациях без - * ресурсов (in-memory) — no-op. - */ -data class StorageBundle( - val conversationStore: ConversationStore, - 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-inmemory/build.gradle.kts b/storage-inmemory/build.gradle.kts index e3605f0..f91ff85 100644 --- a/storage-inmemory/build.gradle.kts +++ b/storage-inmemory/build.gradle.kts @@ -22,7 +22,6 @@ kotlin { commonMain.dependencies { api(project(":message-store-api")) api(project(":message-log-api")) - api(project(":storage-bundle")) api(project(":working-memory-api")) } commonTest.dependencies { 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 deleted file mode 100644 index e114211..0000000 --- a/storage-inmemory/src/commonMain/kotlin/pw/binom/agentik/storage/inmemory/InMemoryEventStore.kt +++ /dev/null @@ -1,84 +0,0 @@ -package pw.binom.agentik.storage.inmemory - -import pw.binom.agentik.messageStore.events.EventRecord -import pw.binom.agentik.messageStore.events.EventStore -import pw.binom.agentik.messageStore.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 a392096..8fbd162 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 @@ -1,10 +1,13 @@ package pw.binom.agentik.storage.inmemory -import pw.binom.agentik.storageBundle.StorageBundle +import pw.binom.agentik.messageLog.MessageStore +import pw.binom.agentik.messageStore.ConversationStore +import pw.binom.agentik.messageStore.ReflectionStore +import pw.binom.agentik.workingMemory.WorkingMemoryStore import kotlin.time.Clock /** - * Фабрика готового [StorageBundle] на базе in-memory имплов. + * Фабрика готового набора in-memory store'ов. * * Удобно для: * - тестов (быстрая инициализация, не нужен JDBC driver); @@ -12,13 +15,22 @@ import kotlin.time.Clock * - dry-run / preview, где SQLite не нужен. * * Семантика полностью совпадает с SQLite-имплами (`:storage-sqlite`). + * + * Возвращает четыре store'а по отдельности — caller пробрасывает их + * туда, где нужны (раньше был `StorageBundle`, от него отказались). */ object InMemoryStorage { - fun create(clock: Clock = Clock.System): StorageBundle = StorageBundle( + data class Bundle( + val conversationStore: ConversationStore, + val messageStore: MessageStore, + val workingMemoryStore: WorkingMemoryStore, + val reflectionStore: ReflectionStore, + ) + + fun create(clock: Clock = Clock.System): Bundle = Bundle( conversationStore = InMemoryConversationStore(clock), 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 deleted file mode 100644 index de9c764..0000000 --- a/storage-inmemory/src/commonTest/kotlin/pw/binom/agentik/storage/inmemory/InMemoryEventStoreTest.kt +++ /dev/null @@ -1,165 +0,0 @@ -package pw.binom.agentik.storage.inmemory - -import kotlinx.coroutines.coroutineScope -import kotlinx.coroutines.test.runTest -import pw.binom.agentik.messageStore.events.EventRecord -import pw.binom.agentik.messageStore.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-ksqlite/build.gradle.kts b/storage-ksqlite/build.gradle.kts index 38c6db9..8caf6ef 100644 --- a/storage-ksqlite/build.gradle.kts +++ b/storage-ksqlite/build.gradle.kts @@ -25,7 +25,6 @@ kotlin { commonMain.dependencies { api(project(":message-store-api")) api(project(":message-log-api")) - api(project(":storage-bundle")) api(project(":working-memory-api")) // ksqlite ещё не опубликован в Maven Central — только в локальном diff --git a/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteEventStore.kt b/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteEventStore.kt deleted file mode 100644 index f91043f..0000000 --- a/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteEventStore.kt +++ /dev/null @@ -1,139 +0,0 @@ -package pw.binom.agentik.storage.ksqlite - -import pw.binom.agentik.messageStore.events.EventRecord -import pw.binom.agentik.messageStore.events.EventStore -import pw.binom.agentik.messageStore.events.EventType -import pw.binom.db.ksqlite.SQLiteConnection -import pw.binom.db.ksqlite.SQLitePreparedStatement -import pw.binom.db.ksqlite.SQLiteResultSet -import kotlin.time.Instant -import kotlinx.coroutines.CoroutineDispatcher -import kotlinx.coroutines.Dispatchers -import kotlinx.coroutines.sync.Mutex -import kotlinx.coroutines.sync.withLock -import kotlinx.coroutines.withContext - -/** - * KMP-реализация [EventStore] поверх ksqlite (https://github.com/caffeine-mgn/ksqlite). - * - * Схема таблицы `agent_event` повторяет [pw.binom.agentik.storage.sqlite.SqliteEventStore], - * чтобы данные были совместимы между двумя backend'ами (можно мигрировать через дамп SQL). - * - * **Идемпотентность**: append использует `INSERT OR REPLACE` (id — PRIMARY KEY). - * Для retry с тем же id и тем же payload это no-op; для retry с тем же id но - * другим payload — replace. EventStore contract обещает idempotent no-op; это - * ослабление для SQLite (см. [pw.binom.agentik.storage.sqlite.SqliteEventStore]). - * - * **Thread-safety**: ksqlite API синхронный; один [Mutex] сериализует операции - * внутри одного store. Между разными store'ами (если делят [SQLiteConnection]) - * SQLite сам сериализует через внутренний lock. - * - * **KMP coverage**: работает на JVM (через JNI к .so), linuxX64/mingwX64 - * (static C amalgamation), Android Native (нужны NDK headers для включения target'а). - * - * @param connection открытое соединение с БД. Caller владеет lifecycle — - * должен закрыть после [close] EventStore. - */ -class KsqliteEventStore( - private val connection: SQLiteConnection, -) : EventStore { - - private val mutex = Mutex() - - /** - * Диспатчер для блокирующего SQLite I/O. На JVM и native у [Dispatchers.IO] - * разная видимость (на native internal), поэтому для общности используем - * [Dispatchers.Default] — на JVM это ~64 worker thread'а, на native — пул - * для cinterop-blocking вызовов. Mutex сериализует операции внутри store - * и так, так что contention минимальный. - */ - private companion object { - private val IO_DISPATCHER: CoroutineDispatcher = Dispatchers.Default - } - - override suspend fun append(record: EventRecord): Unit = withContext(IO_DISPATCHER) { - mutex.withLock { - connection.prepare( - "INSERT OR REPLACE INTO agent_event " + - "(id, conversation_id, created_at, type, payload) VALUES (?, ?, ?, ?, ?)" - ).use { stmt -> - stmt.bindText(1, record.id) - val convId = record.conversationId - if (convId != null) { - stmt.bindText(2, convId) - } else { - stmt.bindNull(2) - } - stmt.bindLong(3, record.createdAt.toEpochMilliseconds()) - stmt.bindText(4, record.type.name) - stmt.bindText(5, record.payload) - stmt.executeUpdate() - } - } - } - - override suspend fun query( - conversationId: String?, - afterId: String?, - limit: Int, - ): List = withContext(IO_DISPATCHER) { - mutex.withLock { - val sql = buildString { - append("SELECT id, conversation_id, created_at, type, payload FROM agent_event WHERE 1=1") - if (conversationId != null) append(" AND conversation_id = ?") - if (afterId != null) append(" AND id > ?") - append(" ORDER BY created_at ASC, id ASC LIMIT ?") - } - connection.prepare(sql).use { stmt -> - var idx = 1 - if (conversationId != null) stmt.bindText(idx++, conversationId) - if (afterId != null) stmt.bindText(idx++, afterId) - stmt.bindLong(idx, limit.toLong()) - collectQuery(stmt) - } - } - } - - private fun collectQuery(stmt: SQLitePreparedStatement): List { - val result = mutableListOf() - stmt.executeQuery().use { rs: SQLiteResultSet -> - while (rs.next()) { - result.add( - EventRecord( - id = rs.getText(0)!!, - conversationId = rs.getText(1), - createdAt = Instant.fromEpochMilliseconds(rs.getLong(2)!!), - type = runCatching { EventType.valueOf(rs.getText(3)!!) } - .getOrDefault(EventType.AGENT_CREATED), - payload = rs.getText(4)!!, - ) - ) - } - } - return result - } - - override suspend fun pruneOlderThan(olderThan: Instant): Int = withContext(IO_DISPATCHER) { - mutex.withLock { - connection.prepare("DELETE FROM agent_event WHERE created_at < ?").use { stmt -> - stmt.bindLong(1, olderThan.toEpochMilliseconds()) - stmt.executeUpdate() - } - } - } - - override suspend fun count(): Int = withContext(IO_DISPATCHER) { - mutex.withLock { - // Не нашёл queryForInt в API, делаем через prepare. - connection.prepare("SELECT COUNT(*) FROM agent_event").use { stmt -> - stmt.executeQuery().use { rs -> - if (rs.next()) (rs.getLong(0) ?: 0L).toInt() else 0 - } - } - } - } - - override fun close() { - // Connection lifecycle — на caller'е (фабрика KsqliteStores). - } -} diff --git a/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteStores.kt b/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteStores.kt index 799970d..2903a7f 100644 --- a/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteStores.kt +++ b/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteStores.kt @@ -3,16 +3,14 @@ package pw.binom.agentik.storage.ksqlite import pw.binom.agentik.messageStore.ConversationStore import pw.binom.agentik.messageLog.MessageStore import pw.binom.agentik.messageStore.ReflectionStore -import pw.binom.agentik.storageBundle.StorageBundle import pw.binom.agentik.workingMemory.WorkingMemoryStore -import pw.binom.agentik.messageStore.events.EventStore import pw.binom.db.ksqlite.SQLiteConnection /** - * Фабрика всех 5 store'ов из :storage-core поверх ksqlite. + * Фабрика 4 store'ов поверх ksqlite. * * Lifecycle: открывает [SQLiteConnection], гарантирует наличие таблиц - * (CREATE TABLE IF NOT EXISTS), возвращает bundle из 5 store'ов. Caller + * (CREATE TABLE IF NOT EXISTS), возвращает bundle из 4 store'ов. Caller * ДОЛЖЕН вызвать [close] при завершении. * * @param path путь к .db файлу, либо URI для in-memory/shared-cache. @@ -23,24 +21,14 @@ class KsqliteStores private constructor( val messages: MessageStore, val workingMemory: WorkingMemoryStore, val reflections: ReflectionStore, - val events: EventStore, ) : AutoCloseable { - fun asBundle(): StorageBundle = StorageBundle( - conversationStore = conversations, - messageStore = messages, - workingMemoryStore = workingMemory, - reflectionStore = reflections, - eventStore = events, - ) - override fun close() { // Закрытие в правильном порядке: зависимые → владелец connection. conversations.close() messages.close() workingMemory.close() reflections.close() - events.close() connection.close() } @@ -88,16 +76,6 @@ class KsqliteStores 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); - - 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 TEXT 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); """ fun open(path: String): KsqliteStores { @@ -121,7 +99,6 @@ class KsqliteStores private constructor( messages = messages, workingMemory = working, reflections = KsqliteReflectionStore(conn), - events = KsqliteEventStore(conn), ) } } diff --git a/storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteEventStoreTest.kt b/storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteEventStoreTest.kt deleted file mode 100644 index 8fa73d2..0000000 --- a/storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteEventStoreTest.kt +++ /dev/null @@ -1,150 +0,0 @@ -package pw.binom.agentik.storage.ksqlite - -import kotlinx.coroutines.test.runTest -import pw.binom.agentik.messageStore.events.EventRecord -import pw.binom.agentik.messageStore.events.EventType -import kotlin.test.AfterTest -import kotlin.test.BeforeTest -import kotlin.test.Test -import kotlin.test.assertEquals -import kotlin.test.assertNull -import kotlin.test.assertTrue -import kotlin.time.Instant - -class KsqliteEventStoreTest { - - private lateinit var stores: KsqliteStores - - @BeforeTest - fun setup() { - stores = KsqliteStores.inMemory("test-${kotlin.random.Random.nextLong()}") - } - - @AfterTest - fun tearDown() { - stores.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 testAppendThenQueryReturnsRecord() = runTest { - val r = rec("ev-1", ts = 1000) - stores.events.append(r) - assertEquals(listOf(r), stores.events.query()) - } - - @Test - fun testQueryWithConversationIdFilters() = runTest { - stores.events.append(rec("ev-1", ts = 1000, conv = "c-1")) - stores.events.append(rec("ev-2", ts = 2000, conv = "c-2")) - stores.events.append(rec("ev-3", ts = 3000, conv = "c-1")) - - val c1 = stores.events.query(conversationId = "c-1") - assertEquals(listOf("ev-1", "ev-3"), c1.map { it.id }) - } - - @Test - fun testQueryWithAfterIdReturnsEventsStrictlyAfterCursor() = runTest { - stores.events.append(rec("ev-1", ts = 1000)) - stores.events.append(rec("ev-2", ts = 2000)) - stores.events.append(rec("ev-3", ts = 3000)) - - // id-based cursor: id > 'ev-1' returns ev-2, ev-3 - assertEquals(listOf("ev-2", "ev-3"), stores.events.query(afterId = "ev-1").map { it.id }) - } - - @Test - fun testAppendIsIdempotentOnIdWithInsertOrReplace() = 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) - stores.events.append(r1) - stores.events.append(r2) - // INSERT OR REPLACE means replace wins. Documented in KsqliteEventStore KDoc. - val result = stores.events.query() - assertEquals(1, result.size) - assertEquals(EventType.CONVERSATION_TOOL_CALL, result[0].type) - } - - @Test - fun testRecordsOrderedByCreatedAtThenId() = runTest { - stores.events.append(rec("ev-c", ts = 2000)) - stores.events.append(rec("ev-a", ts = 1000)) - stores.events.append(rec("ev-d", ts = 1000)) - stores.events.append(rec("ev-b", ts = 1500)) - - assertEquals(listOf("ev-a", "ev-d", "ev-b", "ev-c"), - stores.events.query().map { it.id }) - } - - @Test - fun testQueryWithLimitCaps() = runTest { - repeat(10) { i -> stores.events.append(rec("ev-$i", ts = (i * 100).toLong())) } - val first3 = stores.events.query(limit = 3) - assertEquals(3, first3.size) - assertEquals(listOf("ev-0", "ev-1", "ev-2"), first3.map { it.id }) - } - - @Test - fun testPruneOlderThanRemovesOld() = runTest { - stores.events.append(rec("ev-1", ts = 1000)) - stores.events.append(rec("ev-2", ts = 2000)) - stores.events.append(rec("ev-3", ts = 3000)) - - val removed = stores.events.pruneOlderThan(Instant.fromEpochMilliseconds(2500)) - assertEquals(2, removed) - assertEquals(listOf("ev-3"), stores.events.query().map { it.id }) - } - - @Test - fun testCountReturnsTotal() = runTest { - assertEquals(0, stores.events.count()) - stores.events.append(rec("ev-1", ts = 1000)) - stores.events.append(rec("ev-2", ts = 2000)) - assertEquals(2, stores.events.count()) - } - - @Test - fun testUnknownEventTypeLoadedAsFallback() = runTest { - // Forward-compat: write record with unknown type, read via store. - stores.connection.exec( - "INSERT INTO agent_event (id, conversation_id, created_at, type, payload) VALUES " + - "('ev-future', NULL, 5000, 'SOME_FUTURE_TYPE', '{}')" - ) - val result = stores.events.query() - assertEquals(1, result.size) - assertEquals("ev-future", result[0].id) - assertEquals(EventType.AGENT_CREATED, result[0].type) - } - - @Test - fun testNullConversationIdStoredAndRetrieved() = runTest { - stores.events.append(rec("ev-no-conv", ts = 1000, conv = null)) - val result = stores.events.query() - assertEquals(1, result.size) - assertNull(result[0].conversationId) - } - - @Test - fun testSchemaCreatedOnFirstOpen() = runTest { - val count = stores.connection.prepare( - "SELECT COUNT(*) FROM sqlite_master WHERE type='table' AND name='agent_event'" - ).use { stmt -> - stmt.executeQuery().use { rs -> - if (rs.next()) (rs.getLong(0) ?: 0L).toInt() else 0 - } - } - assertEquals(1, count) - } -} diff --git a/storage-sqlite/build.gradle.kts b/storage-sqlite/build.gradle.kts index 7fe0c1e..3db81be 100644 --- a/storage-sqlite/build.gradle.kts +++ b/storage-sqlite/build.gradle.kts @@ -17,7 +17,6 @@ kotlin { commonMain.dependencies { api(project(":message-store-api")) api(project(":message-log-api")) - api(project(":storage-bundle")) api(project(":working-memory-api")) api(libs.sqldelight.runtime) api(libs.sqldelight.coroutines) 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 deleted file mode 100644 index 631a851..0000000 --- a/storage-sqlite/src/jvmMain/kotlin/pw/binom/agentik/storage/sqlite/SqliteEventStore.kt +++ /dev/null @@ -1,86 +0,0 @@ -package pw.binom.agentik.storage.sqlite - -import kotlin.time.Instant -import mu.KotlinLogging -import pw.binom.agentik.messageStore.events.EventRecord -import pw.binom.agentik.messageStore.events.EventStore -import pw.binom.agentik.messageStore.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 f39b54a..5fe4b2b 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 @@ -6,14 +6,11 @@ import app.cash.sqldelight.driver.jdbc.sqlite.JdbcSqliteDriver import pw.binom.agentik.messageStore.ConversationStore import pw.binom.agentik.messageLog.MessageStore import pw.binom.agentik.messageStore.ReflectionStore -import pw.binom.agentik.storageBundle.StorageBundle -import pw.binom.agentik.messageStore.events.EventStore import pw.binom.agentik.storage.sqlite.SqliteReflectionStore import pw.binom.agentik.workingMemory.WorkingMemoryStore /** - * Корневой объект SQLite-слоя: держит [SqlDriver] и пять store'ов (включая - * [EventStore] — replay-after-disconnect). + * Корневой объект SQLite-слоя: держит [SqlDriver] и четыре store'а. */ class SqliteStores private constructor( val driver: SqlDriver, @@ -21,24 +18,9 @@ class SqliteStores private constructor( val messages: MessageStore, val workingMemory: WorkingMemoryStore, val reflections: ReflectionStore, - val events: EventStore, ) : AutoCloseable { - /** - * Удобный агрегатор — превращает SqliteStores в [StorageBundle] для передачи - * в агенты, которые работают через общий контракт хранения (ChatAgent после - * commit 6 в :standalone принимает [StorageBundle], не [SqliteStores]). - */ - fun asBundle(): StorageBundle = StorageBundle( - conversationStore = conversations, - messageStore = messages, - workingMemoryStore = workingMemory, - reflectionStore = reflections, - eventStore = events, - ) - override fun close() { - events.close() conversations.close() messages.close() workingMemory.close() @@ -59,7 +41,6 @@ class SqliteStores private constructor( messages = SqliteMessageStore(db), workingMemory = SqliteWorkingMemoryStore(db), reflections = SqliteReflectionStore(db), - events = SqliteEventStore(db), ) } @@ -74,7 +55,6 @@ class SqliteStores private constructor( messages = SqliteMessageStore(db), workingMemory = SqliteWorkingMemoryStore(db), reflections = SqliteReflectionStore(db), - events = SqliteEventStore(db), ) } @@ -124,18 +104,6 @@ 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 deleted file mode 100644 index 6e0a867..0000000 --- a/storage-sqlite/src/jvmMain/sqldelight/pw/binom/agentik/storage/sqlite/EventStore.sq +++ /dev/null @@ -1,49 +0,0 @@ --- 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 deleted file mode 100644 index cec7b39..0000000 --- a/storage-sqlite/src/jvmTest/kotlin/pw/binom/agentik/storage/sqlite/SqliteEventStoreTest.kt +++ /dev/null @@ -1,152 +0,0 @@ -package pw.binom.agentik.storage.sqlite - -import kotlinx.coroutines.test.runTest -import pw.binom.agentik.messageStore.events.EventRecord -import pw.binom.agentik.messageStore.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()) - } -} diff --git a/working-memory-api/build.gradle.kts b/working-memory-api/build.gradle.kts index 2c180a1..db65ee3 100644 --- a/working-memory-api/build.gradle.kts +++ b/working-memory-api/build.gradle.kts @@ -18,7 +18,6 @@ kotlin { sourceSets { commonMain.dependencies { - api(project(":message-store-api")) api(project(":message-log-api")) api(libs.kotlinx.coroutines.core) api(libs.kotlinx.serialization.core)