From 2d6cf89c52866b78e13115d95c2aa53dd708085a Mon Sep 17 00:00:00 2001 From: subochev Date: Mon, 21 Sep 2026 02:57:45 +0300 Subject: [PATCH] remove deprecated `EventStore` and `MessageStore` implementations, along with related in-memory and SQLite code --- agent-toolsets/build.gradle.kts | 3 +- client/build.gradle.kts | 2 +- .../pw/binom/agentik/client/HttpEventStore.kt | 10 +- context-api/build.gradle.kts | 7 +- .../agentik/context/WorkingMemoryEntry.kt | 4 +- context-ksqlite/build.gradle.kts | 2 +- event-store-in-memory/build.gradle.kts | 29 --- .../eventStore/inmemory/InMemoryEventStore.kt | 160 ------------ .../inmemory/InMemoryEventStoreTest.kt | 228 ------------------ .../pw/binom/agentik/eventStore/EventStore.kt | 16 -- .../agentik/eventStore/MutableEventStore.kt | 15 -- llm-tools/build.gradle.kts | 2 +- .../pw/binom/agentik/messageLog/Content.kt | 13 - .../agentik/messageLog/MessageContext.kt | 13 - .../binom/agentik/messageLog/MessageOrigin.kt | 13 - .../binom/agentik/messageLog/MessageRecord.kt | 13 - .../binom/agentik/messageLog/MessageStore.kt | 17 -- .../agentik/messageLog/MutableMessageStore.kt | 15 -- .../pw/binom/agentik/messageLog/Payload.kt | 13 - .../pw/binom/agentik/messageLog/TurnTokens.kt | 13 - message-log-ksqlite/build.gradle.kts | 33 --- .../messageLog/ksqlite/KsqliteMessageStore.kt | 141 ----------- .../messageLog/ksqlite/MessageCodecs.kt | 79 ------ .../ksqlite/KsqliteMessageStoreTest.kt | 122 ---------- settings.gradle.kts | 10 +- standalone/build.gradle.kts | 2 +- .../agentik/standalone/agent/ChatAgent.kt | 2 +- storage-inmemory/build.gradle.kts | 4 +- .../storage/inmemory/InMemoryMessageStore.kt | 6 +- .../storage/inmemory/InMemoryStorage.kt | 6 +- storage-ksqlite/build.gradle.kts | 4 +- .../storage/ksqlite/KsqliteMessageStore.kt | 6 +- .../agentik/storage/ksqlite/KsqliteStores.kt | 8 +- working-memory-api/build.gradle.kts | 2 +- .../workingMemory/WorkingMemoryEntry.kt | 16 -- .../workingMemory/WorkingMemoryStore.kt | 18 -- 36 files changed, 41 insertions(+), 1006 deletions(-) delete mode 100644 event-store-in-memory/build.gradle.kts delete mode 100644 event-store-in-memory/src/commonMain/kotlin/pw/binom/agentik/eventStore/inmemory/InMemoryEventStore.kt delete mode 100644 event-store-in-memory/src/commonTest/kotlin/pw/binom/agentik/eventStore/inmemory/InMemoryEventStoreTest.kt delete mode 100644 event-store/src/commonMain/kotlin/pw/binom/agentik/eventStore/EventStore.kt delete mode 100644 event-store/src/commonMain/kotlin/pw/binom/agentik/eventStore/MutableEventStore.kt delete mode 100644 message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/Content.kt delete mode 100644 message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/MessageContext.kt delete mode 100644 message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/MessageOrigin.kt delete mode 100644 message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/MessageRecord.kt delete mode 100644 message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/MessageStore.kt delete mode 100644 message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/MutableMessageStore.kt delete mode 100644 message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/Payload.kt delete mode 100644 message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/TurnTokens.kt delete mode 100644 message-log-ksqlite/build.gradle.kts delete mode 100644 message-log-ksqlite/src/commonMain/kotlin/pw/binom/agentik/messageLog/ksqlite/KsqliteMessageStore.kt delete mode 100644 message-log-ksqlite/src/commonMain/kotlin/pw/binom/agentik/messageLog/ksqlite/MessageCodecs.kt delete mode 100644 message-log-ksqlite/src/commonTest/kotlin/pw/binom/agentik/messageLog/ksqlite/KsqliteMessageStoreTest.kt delete mode 100644 working-memory-api/src/commonMain/kotlin/pw/binom/agentik/workingMemory/WorkingMemoryEntry.kt delete mode 100644 working-memory-api/src/commonMain/kotlin/pw/binom/agentik/workingMemory/WorkingMemoryStore.kt diff --git a/agent-toolsets/build.gradle.kts b/agent-toolsets/build.gradle.kts index 3366eb7..8932690 100644 --- a/agent-toolsets/build.gradle.kts +++ b/agent-toolsets/build.gradle.kts @@ -22,8 +22,7 @@ kotlin { sourceSets { commonMain.dependencies { api(project(":message-store-api")) - api(project(":message-log-api")) - api(project(":working-memory-api")) + api(project(":context-api")) // litert-kmp: LiteTool интерфейс (sync describe/invoke) api(libs.litert.api) diff --git a/client/build.gradle.kts b/client/build.gradle.kts index 00d451d..36c5389 100644 --- a/client/build.gradle.kts +++ b/client/build.gradle.kts @@ -21,7 +21,7 @@ kotlin { sourceSets { commonMain.dependencies { api(project(":proto")) - api(project(":event-store")) + api(project(":outbox-api")) api(project(":journal-api")) api(libs.ktor.client.core) diff --git a/client/src/commonMain/kotlin/pw/binom/agentik/client/HttpEventStore.kt b/client/src/commonMain/kotlin/pw/binom/agentik/client/HttpEventStore.kt index 2bd3719..a94c8f7 100644 --- a/client/src/commonMain/kotlin/pw/binom/agentik/client/HttpEventStore.kt +++ b/client/src/commonMain/kotlin/pw/binom/agentik/client/HttpEventStore.kt @@ -6,7 +6,7 @@ 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.outbox.OutboxStore import pw.binom.agentik.outbox.AgentEvent import pw.binom.agentik.outbox.CommonEvent import pw.binom.agentik.outbox.Event @@ -14,7 +14,7 @@ import kotlin.time.Clock import kotlin.time.Instant /** - * HTTP-реализация [EventStore] (= [pw.binom.agentik.outbox.OutboxStore]), + * HTTP-реализация [OutboxStore] (= [pw.binom.agentik.outbox.OutboxStore]), * ходящая в `:server`-фасад. * * **Endpoint-раскладка** (новый дизайн — storage handles на [Agent]): @@ -26,12 +26,12 @@ import kotlin.time.Instant * - [conversationEvents] с `conversationId != null` → `GET /conversations/{id}/events` * * Для [conversationEvents] с `conversationId == null` (события всех диалогов) - * fallback на default [EventStore.conversationEvents] — общий поток + * fallback на default [OutboxStore.conversationEvents] — общий поток * `/outbox/events` + filter. Это редкий кейс (admin-дашборды), и * оптимизировать его отдельно нерационально. * * [earliestEventDate] не имеет своего endpoint'а; возвращает `Clock.System.now()` - * (см. KDoc [EventStore.earliestEventDate] — для пустого буфера это и есть + * (см. KDoc [OutboxStore.earliestEventDate] — для пустого буфера это и есть * контрактное значение). Клиент, который полагался на gap detection через * message store, продолжит работать — просто fallback никогда не сработает. * @@ -44,7 +44,7 @@ import kotlin.time.Instant internal class HttpEventStore( private val httpClient: HttpClient, private val baseUrl: String, -) : EventStore { +) : OutboxStore { private val agentUrl: String = baseUrl.trimEnd('/') diff --git a/context-api/build.gradle.kts b/context-api/build.gradle.kts index fcf64a1..d7b3afe 100644 --- a/context-api/build.gradle.kts +++ b/context-api/build.gradle.kts @@ -3,8 +3,9 @@ plugins { alias(libs.plugins.kotlin.serialization) } -// Public API для mutable working-memory — runtime context агента (compaction, -// order_idx, summary entries). Зависит от :message-log-api для Ids (wm- префикс). +// Public API для runtime context агента (compaction, order_idx, summary entries). +// Зависит от :journal-api для типов `Content` / `MessageContext` (audit-log +// payload'ы, которые рабочая память ссылает). // // НЕ нужен тонким клиентам — только серверному рантайму (`:standalone`, `:agentik-cli`, // будущий `:android-agent` core). @@ -24,7 +25,7 @@ kotlin { sourceSets { commonMain.dependencies { - api(project(":message-log-api")) + api(project(":journal-api")) api(libs.kotlinx.coroutines.core) api(libs.kotlinx.serialization.core) api(libs.kotlinx.serialization.json) diff --git a/context-api/src/commonMain/kotlin/pw/binom/agentik/context/WorkingMemoryEntry.kt b/context-api/src/commonMain/kotlin/pw/binom/agentik/context/WorkingMemoryEntry.kt index 503957e..4b66b5c 100644 --- a/context-api/src/commonMain/kotlin/pw/binom/agentik/context/WorkingMemoryEntry.kt +++ b/context-api/src/commonMain/kotlin/pw/binom/agentik/context/WorkingMemoryEntry.kt @@ -2,8 +2,8 @@ package pw.binom.agentik.context import kotlinx.serialization.SerialName import kotlinx.serialization.Serializable -import pw.binom.agentik.messageLog.Content -import pw.binom.agentik.messageLog.MessageContext +import pw.binom.agentik.journal.Content +import pw.binom.agentik.journal.MessageContext /** * Запись в working memory диалога: ровно то, что агент сейчас видит в diff --git a/context-ksqlite/build.gradle.kts b/context-ksqlite/build.gradle.kts index 65c62f0..3a64810 100644 --- a/context-ksqlite/build.gradle.kts +++ b/context-ksqlite/build.gradle.kts @@ -30,7 +30,7 @@ kotlin { // :message-log-api (старый canonical). Транзитивно через api, // но фиксируем явно чтобы тестовый код видел Content без // обхода через :context-api. - api(project(":message-log-api")) +// api(project(":message-log-api")) } commonTest.dependencies { implementation(kotlin("test")) diff --git a/event-store-in-memory/build.gradle.kts b/event-store-in-memory/build.gradle.kts deleted file mode 100644 index 2c75078..0000000 --- a/event-store-in-memory/build.gradle.kts +++ /dev/null @@ -1,29 +0,0 @@ -plugins { - alias(libs.plugins.kotlin.multiplatform) -} - -// KMP-реализация [MutableEventStore] на `ConcurrentLinkedDeque` — для тестов, -// dev-режима и embedded-сценариев (Android core, CLI). TTL и size-cap eviction -// вызываются на каждом [append], в одном проходе с amortized O(1) для стабильного -// размера буфера. -// -// Зависимости: только `:event-store` (api → `:proto` транзитивно). -// Никакого I/O — pure in-memory. - -kotlin { - jvmToolchain(21) - - jvm() - linuxX64() - mingwX64() - - sourceSets { - commonMain.dependencies { - api(project(":event-store")) - } - commonTest.dependencies { - implementation(kotlin("test")) - implementation(libs.kotlinx.coroutines.test) - } - } -} diff --git a/event-store-in-memory/src/commonMain/kotlin/pw/binom/agentik/eventStore/inmemory/InMemoryEventStore.kt b/event-store-in-memory/src/commonMain/kotlin/pw/binom/agentik/eventStore/inmemory/InMemoryEventStore.kt deleted file mode 100644 index 2c88009..0000000 --- a/event-store-in-memory/src/commonMain/kotlin/pw/binom/agentik/eventStore/inmemory/InMemoryEventStore.kt +++ /dev/null @@ -1,160 +0,0 @@ -package pw.binom.agentik.eventStore.inmemory - -import kotlin.time.Clock -import kotlin.time.Duration -import kotlin.time.Instant -import kotlinx.coroutines.channels.BufferOverflow -import kotlinx.coroutines.coroutineScope -import kotlinx.coroutines.flow.Flow -import kotlinx.coroutines.flow.MutableSharedFlow -import kotlinx.coroutines.flow.asSharedFlow -import kotlinx.coroutines.flow.channelFlow -import kotlinx.coroutines.flow.flow -import kotlinx.coroutines.sync.Mutex -import kotlinx.coroutines.sync.withLock -import pw.binom.agentik.eventStore.MutableEventStore -import pw.binom.agentik.proto.CommonEvent - -/** - * In-memory реализация [MutableEventStore] на `ArrayDeque` + [Mutex]. - * - * **Retention policy** — оба параметра **nullable** без default'ов - * (контракт: caller явно решает что ему нужно, не получает "удобные дефолты"): - * - [maxMessages] `null` → неограниченно по количеству. - * - [ttl] `null` → нет time-based eviction (храним вечно, **пока maxMessages тоже null**). - * - **Оба `null` → вечное хранилище.** - * - Любой non-null → соответствующая граница применяется **на каждом - * [append]** (amortized O(1) при стабильном размере буфера). - * - * **Concurrency**: [Mutex] защищает append/evict от concurrent writer'ов; - * reader'ы [events] не блокируются — снимают snapshot под lock'ом, дальше - * итерируют без него. Snapshot под `mutex.withLock` даёт weakly-consistent - * точку обзора: append'ы, попавшие в окно между snapshot и live-collect, - * обрабатываются через **monotonic sequence boundary** (см. [events] KDoc). - * - * **Live tail**: [MutableSharedFlow] с DROP_OLDEST policy. Producer никогда - * не блокируется — если буфер live-flow переполнен (4096 подписчиков - * медленных), старые события дропаются без уведомления. Это OK: каждый - * subscriber видит **свой** late tail, а за полным покрытием — fallback - * в `:message-store-api`. - * - * **Threading model**: append происходит из любого dispatcher'а; eviction - * — best-effort, синхронный, в том же вызове append (это нормально - * для in-memory, добавляет O(evicted) работы). - */ -class InMemoryEventStore( - private val maxMessages: Int?, - private val ttl: Duration?, - private val clock: Clock = Clock.System, -) : MutableEventStore { - - private val mutex = Mutex() - private val buffer = ArrayDeque() - private val liveFlow = MutableSharedFlow( - replay = 0, - extraBufferCapacity = LIVE_BUFFER_CAPACITY, - onBufferOverflow = BufferOverflow.DROP_OLDEST, - ) - - init { - // Аргументы — НЕ optional default'ы; explicit null = "не применяется". - // Если caller передал отрицательный max — это ошибка конфигурации, - // пробрасываем сразу при инициализации. - require(maxMessages == null || maxMessages > 0) { - "maxMessages must be > 0 or null, got $maxMessages" - } - } - - override suspend fun append(event: CommonEvent) { - mutex.withLock { - buffer.addLast(event) - } - liveFlow.tryEmit(event) - evictExpired() - evictOverCapacity() - } - - /** - * Удалить с головы все event'ы старше [ttl]. Amortized O(evicted). - * Если [ttl] null — no-op. - */ - private suspend fun evictExpired() { - val ttlValue = ttl ?: return - val cutoff = clock.now() - ttlValue - mutex.withLock { - while (true) { - val head = buffer.firstOrNull() ?: return@withLock - if (head.date >= cutoff) return@withLock - buffer.removeFirst() - } - } - } - - /** - * Удалить с головы пока размер > [maxMessages]. Amortized O(evicted). - * Если [maxMessages] null — no-op. - */ - private suspend fun evictOverCapacity() { - val cap = maxMessages ?: return - mutex.withLock { - while (buffer.size > cap) { - if (buffer.isEmpty()) return@withLock - buffer.removeFirst() - } - } - } - - override fun events(after: Instant?): Flow = flow { - // Replay buffer — snapshot под mutex'ом, дальше iterate без lock'а. - // Append'ы в окне между snapshot и live-collect компенсируются - // через monotonic sequence boundary: append нумерует события - // последовательно, live-collect фильтрует по last-seen-seq. - val snapshot: List = mutex.withLock { - if (after == null) { - buffer.toList() - } else { - buffer.filter { it.date > after } - } - } - snapshot.forEach { emit(it) } - // Live tail — `coroutineScope` гарантирует proper cleanup: когда - // collector отменяется (take(N)), scope отменяется, liveFlow.collect - // выходит чисто. Без этого — runTest видит "uncompleted coroutine" - // и валит тест с UncompletedCoroutinesError. - coroutineScope { - liveFlow.collect { emit(it) } - } - } - - override suspend fun earliestEventDate(): Instant { - val earliest = mutex.withLock { buffer.firstOrNull()?.date } - // Не nullable: для пустого буфера возвращаем "сейчас" — это позволяет - // клиенту безопасно подписаться на `events(after = earliest)`. - return earliest ?: clock.now() - } - - /** - * **Test-only helper** — снимок буфера в текущий момент. - * - * `internal` потому что production код не должен ходить напрямую в буфер - * (для этого есть `events(after)`). Доступно только из `commonTest`. - * - * Returns: иммутабельный snapshot (копия). Под `mutex.withLock` — - * consistency на момент снятия; concurrent append'ы могут расширить - * буфер сразу после, но для single-threaded тестов OK. - */ - internal suspend fun snapshot(): List = mutex.withLock { buffer.toList() } - - override fun close() { - // mutex не закрываем (kotlinx Mutex не AutoCloseable; для in-memory - // store GC соберёт всё при выходе ссылки). buffer чистим. - buffer.clear() - } - - private companion object { - // Live-flow capacity — generous default. Если реально 4096 подписчиков - // отстают настолько что переполняют буфер, проблема upstream, не здесь. - private const val LIVE_BUFFER_CAPACITY = 4096 - } -} - diff --git a/event-store-in-memory/src/commonTest/kotlin/pw/binom/agentik/eventStore/inmemory/InMemoryEventStoreTest.kt b/event-store-in-memory/src/commonTest/kotlin/pw/binom/agentik/eventStore/inmemory/InMemoryEventStoreTest.kt deleted file mode 100644 index c746481..0000000 --- a/event-store-in-memory/src/commonTest/kotlin/pw/binom/agentik/eventStore/inmemory/InMemoryEventStoreTest.kt +++ /dev/null @@ -1,228 +0,0 @@ -package pw.binom.agentik.eventStore.inmemory - -import kotlin.test.Test -import kotlin.test.assertEquals -import kotlin.test.assertTrue -import kotlin.time.Clock -import kotlin.time.Duration -import kotlin.time.Instant -import kotlinx.coroutines.CompletableDeferred -import kotlinx.coroutines.delay -import kotlinx.coroutines.launch -import kotlinx.coroutines.runBlocking -// Импортируем напрямую из :outbox-api — typealias'ы в :proto для -// CommonEvent/AgentEvent/Event НЕ поддерживают nested-class access -// (`CommonEvent.Agent` через alias даёт "Unresolved qualified name"). -import pw.binom.agentik.outbox.AgentEvent -import pw.binom.agentik.outbox.CommonEvent -import pw.binom.agentik.outbox.Event - -class InMemoryEventStoreTest { - - private class FixedClock(private var nowMs: Long = 1_000_000_000L) : Clock { - fun advance(delta: Duration) { nowMs += delta.inWholeMilliseconds } - override fun now(): Instant = Instant.fromEpochMilliseconds(nowMs) - } - - private fun evtAt(clock: Clock, body: String): CommonEvent = - CommonEvent.Agent(date = clock.now(), event = AgentEvent.Created(date = clock.now(), conversationId = body)) - - @Test - fun `append stores all events when both limits are null store-forever`() = runBlocking { - val store = InMemoryEventStore(maxMessages = null, ttl = null) - repeat(100) { i -> - store.append(CommonEvent.Agent( - date = Instant.fromEpochSeconds(i.toLong()), - event = AgentEvent.Created(date = Instant.fromEpochSeconds(i.toLong()), conversationId = "c-$i"), - )) - } - assertEquals(100, store.snapshot().size) - } - - @Test - fun `maxMessages cap evicts oldest when exceeded`() = runBlocking { - val store = InMemoryEventStore(maxMessages = 3, ttl = null) - for (i in 1..5) { - store.append(CommonEvent.Agent( - date = Instant.fromEpochSeconds(i.toLong()), - event = AgentEvent.Created(date = Instant.fromEpochSeconds(i.toLong()), conversationId = "c-$i"), - )) - } - val ids = store.snapshot().map { ((it as CommonEvent.Agent).event as AgentEvent.Created).conversationId } - assertEquals(listOf("c-3", "c-4", "c-5"), ids) - } - - @Test - fun `ttl evicts events older than threshold`() = runBlocking { - val clock = FixedClock() - val store = InMemoryEventStore(maxMessages = null, ttl = 100.milliseconds, clock = clock) - - store.append(evtAt(clock, "old")) - clock.advance(50.milliseconds) - store.append(evtAt(clock, "middle")) - clock.advance(70.milliseconds) - store.append(evtAt(clock, "fresh")) - - val ids = store.snapshot().map { ((it as CommonEvent.Agent).event as AgentEvent.Created).conversationId } - assertEquals(listOf("middle", "fresh"), ids) - } - - @Test - fun `both maxMessages and ttl apply together`() = runBlocking { - val clock = FixedClock() - // ttl=100ms so b at t=20 (deadline=120) survives when c is appended at t=80. - // Cap=2 evicts oldest. Result: [b, c]. - val store = InMemoryEventStore(maxMessages = 2, ttl = 100.milliseconds, clock = clock) - - store.append(evtAt(clock, "a")) - clock.advance(20.milliseconds) - store.append(evtAt(clock, "b")) - clock.advance(60.milliseconds) - store.append(evtAt(clock, "c")) - - val ids = store.snapshot().map { ((it as CommonEvent.Agent).event as AgentEvent.Created).conversationId } - assertEquals(listOf("b", "c"), ids) - } - - @Test - fun `events with null after replays buffer then collects live`() = runBlocking { - val store = InMemoryEventStore(maxMessages = null, ttl = null) - store.append(evtAt(Clock.System, "e1")) - store.append(evtAt(Clock.System, "e2")) - - val collected = mutableListOf() - val done = CompletableDeferred() - val job = launch { - store.events(after = null).collect { e -> - collected.add(e) - if (collected.size >= 3) done.complete(Unit) - } - } - delay(20) - store.append(evtAt(Clock.System, "e3")) - done.await() - job.cancel() - assertEquals(3, collected.size) - } - - @Test - fun `events with after catches up then continues with live`() = runBlocking { - val store = InMemoryEventStore(maxMessages = null, ttl = null) - val t0 = Instant.fromEpochSeconds(0) - val t1 = Instant.fromEpochSeconds(10) - val t2 = Instant.fromEpochSeconds(20) - - store.append(CommonEvent.Agent(date = t0, event = AgentEvent.Created(date = t0, conversationId = "e1"))) - store.append(CommonEvent.Agent(date = t1, event = AgentEvent.Created(date = t1, conversationId = "e2"))) - store.append(CommonEvent.Agent(date = t2, event = AgentEvent.Created(date = t2, conversationId = "e3"))) - - val collected = mutableListOf() - val done = CompletableDeferred() - val job = launch { - store.events(after = t0).collect { e -> - collected.add(e) - if (collected.size >= 3) done.complete(Unit) - } - } - delay(20) - store.append(CommonEvent.Agent( - date = Instant.fromEpochSeconds(30), - event = AgentEvent.Created(date = Instant.fromEpochSeconds(30), conversationId = "e4"), - )) - done.await() - job.cancel() - val ids = collected.map { ((it as CommonEvent.Agent).event as AgentEvent.Created).conversationId } - assertEquals(listOf("e2", "e3", "e4"), ids) - } - - @Test - fun `earliestEventDate returns oldest buffered date`() = runBlocking { - val clock = FixedClock() - val store = InMemoryEventStore(maxMessages = null, ttl = null, clock = clock) - store.append(evtAt(clock, "e1")) - clock.advance(100.milliseconds) - store.append(evtAt(clock, "e2")) - - assertEquals(Instant.fromEpochMilliseconds(1_000_000_000L), store.earliestEventDate()) - } - - @Test - fun `earliestEventDate returns current time when buffer is empty`() = runBlocking { - val clock = FixedClock(nowMs = 5_000_000_000L) - val store = InMemoryEventStore(maxMessages = null, ttl = null, clock = clock) - assertEquals(Instant.fromEpochMilliseconds(5_000_000_000L), store.earliestEventDate()) - } - - @Test - fun `conversationEvents default impl filters to conversation variant`() = runBlocking { - val store = InMemoryEventStore(maxMessages = null, ttl = null) - val now = Instant.fromEpochSeconds(0) - store.append(CommonEvent.Agent( - date = now, - event = AgentEvent.Created(date = now, conversationId = "agent-event"), - )) - store.append(CommonEvent.Conversation( - date = now, - conversationId = "c-1", - event = Event.AppendText(date = now, body = "hi"), - )) - - // Snapshot-based test of the default impl (uses events() + filterIsInstance). - // We test the post-condition directly: there should be exactly 1 - // conversation event. - val all = store.snapshot() - assertEquals(2, all.size) - assertEquals(1, all.count { it is CommonEvent.Conversation }) - assertEquals(1, all.count { it is CommonEvent.Agent }) - } - - @Test - fun `conversationEvents with conversationId filters to that conversation`() = runBlocking { - val store = InMemoryEventStore(maxMessages = null, ttl = null) - val now = Instant.fromEpochSeconds(0) - store.append(CommonEvent.Conversation(now, "c-1", Event.AppendText(now, "a"))) - store.append(CommonEvent.Conversation(now, "c-2", Event.AppendText(now, "b"))) - store.append(CommonEvent.Conversation(now, "c-1", Event.AppendText(now, "c"))) - - // Test the filter logic by manually filtering snapshot. - val c1 = store.snapshot() - .filterIsInstance() - .filter { it.conversationId == "c-1" } - assertEquals(2, c1.size) - assertTrue(c1.all { it.conversationId == "c-1" }) - } - - @Test - fun `agentEvents default impl filters to agent variant`() = runBlocking { - val store = InMemoryEventStore(maxMessages = null, ttl = null) - val now = Instant.fromEpochSeconds(0) - store.append(CommonEvent.Agent( - date = now, - event = AgentEvent.Created(date = now, conversationId = "created"), - )) - store.append(CommonEvent.Conversation(now, "c-1", Event.AppendText(now, "hi"))) - - val all = store.snapshot() - val agents = all.filterIsInstance() - assertEquals(1, agents.size) - val created = agents[0].event as AgentEvent.Created - assertEquals("created", created.conversationId) - } - - @Test - fun `close clears buffer`() = runBlocking { - val store = InMemoryEventStore(maxMessages = null, ttl = null) - store.append(evtAt(Clock.System, "e1")) - store.close() - assertEquals(emptyList(), store.snapshot()) - } - - @Test - fun `negative maxMessages throws at construction`() { - kotlin.runCatching { InMemoryEventStore(maxMessages = -1, ttl = null) } - .onFailure { /* expected */ } - .onSuccess { kotlin.test.fail("should have thrown") } - } -} - -private val Int.milliseconds: Duration get() = Duration.parse("${this}ms") diff --git a/event-store/src/commonMain/kotlin/pw/binom/agentik/eventStore/EventStore.kt b/event-store/src/commonMain/kotlin/pw/binom/agentik/eventStore/EventStore.kt deleted file mode 100644 index 0213343..0000000 --- a/event-store/src/commonMain/kotlin/pw/binom/agentik/eventStore/EventStore.kt +++ /dev/null @@ -1,16 +0,0 @@ -package pw.binom.agentik.eventStore - -/** - * @Deprecated - * Перенесено в `:outbox-api`. Используй `pw.binom.agentik.outbox.OutboxStore`. - * - * Backward-compat typealias. Существующие импорты `pw.binom.agentik.eventStore.EventStore` - * продолжают работать; новый код в `:server` использует `:outbox-api` напрямую. - * - * Удалить когда все импорты будут на `:outbox-api`. - */ -@Deprecated( - message = "Перенесено в :outbox-api. Используй pw.binom.agentik.outbox.OutboxStore.", - replaceWith = ReplaceWith("OutboxStore", "pw.binom.agentik.outbox.OutboxStore"), -) -typealias EventStore = pw.binom.agentik.outbox.OutboxStore diff --git a/event-store/src/commonMain/kotlin/pw/binom/agentik/eventStore/MutableEventStore.kt b/event-store/src/commonMain/kotlin/pw/binom/agentik/eventStore/MutableEventStore.kt deleted file mode 100644 index b32dd37..0000000 --- a/event-store/src/commonMain/kotlin/pw/binom/agentik/eventStore/MutableEventStore.kt +++ /dev/null @@ -1,15 +0,0 @@ -package pw.binom.agentik.eventStore - -/** - * @Deprecated - * Перенесено в `:outbox-api`. Используй `pw.binom.agentik.outbox.MutableOutboxStore`. - * - * Backward-compat typealias. `MutableEventStore` == `MutableOutboxStore` (один тип). - * - * Удалить когда все импорты будут на `:outbox-api`. - */ -@Deprecated( - message = "Перенесено в :outbox-api. Используй pw.binom.agentik.outbox.MutableOutboxStore.", - replaceWith = ReplaceWith("MutableOutboxStore", "pw.binom.agentik.outbox.MutableOutboxStore"), -) -typealias MutableEventStore = pw.binom.agentik.outbox.MutableOutboxStore diff --git a/llm-tools/build.gradle.kts b/llm-tools/build.gradle.kts index 5f0ae86..bdcb198 100644 --- a/llm-tools/build.gradle.kts +++ b/llm-tools/build.gradle.kts @@ -20,7 +20,7 @@ kotlin { commonMain.dependencies { api(project(":memory-api")) api(project(":message-store-api")) - api(project(":working-memory-api")) + api(project(":context-api")) api(project(":skills")) api(libs.litert.api) implementation(libs.kotlinx.coroutines.core) diff --git a/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/Content.kt b/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/Content.kt deleted file mode 100644 index b907735..0000000 --- a/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/Content.kt +++ /dev/null @@ -1,13 +0,0 @@ -package pw.binom.agentik.messageLog - -/** - * @Deprecated - * Перенесено в `:journal-api`. Используй `pw.binom.agentik.journal.Content`. - * - * Backward-compat typealias. - */ -@Deprecated( - message = "Перенесено в :journal-api. Используй pw.binom.agentik.journal.Content.", - replaceWith = ReplaceWith("Content", "pw.binom.agentik.journal.Content"), -) -typealias Content = pw.binom.agentik.journal.Content diff --git a/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/MessageContext.kt b/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/MessageContext.kt deleted file mode 100644 index 2b4e5b5..0000000 --- a/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/MessageContext.kt +++ /dev/null @@ -1,13 +0,0 @@ -package pw.binom.agentik.messageLog - -/** - * @Deprecated - * Перенесено в `:journal-api`. Используй `pw.binom.agentik.journal.MessageContext`. - * - * Backward-compat typealias. - */ -@Deprecated( - message = "Перенесено в :journal-api. Используй pw.binom.agentik.journal.MessageContext.", - replaceWith = ReplaceWith("MessageContext", "pw.binom.agentik.journal.MessageContext"), -) -typealias MessageContext = pw.binom.agentik.journal.MessageContext diff --git a/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/MessageOrigin.kt b/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/MessageOrigin.kt deleted file mode 100644 index 41ced42..0000000 --- a/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/MessageOrigin.kt +++ /dev/null @@ -1,13 +0,0 @@ -package pw.binom.agentik.messageLog - -/** - * @Deprecated - * Перенесено в `:journal-api`. Используй `pw.binom.agentik.journal.MessageOrigin`. - * - * Backward-compat typealias. - */ -@Deprecated( - message = "Перенесено в :journal-api. Используй pw.binom.agentik.journal.MessageOrigin.", - replaceWith = ReplaceWith("MessageOrigin", "pw.binom.agentik.journal.MessageOrigin"), -) -typealias MessageOrigin = pw.binom.agentik.journal.MessageOrigin diff --git a/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/MessageRecord.kt b/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/MessageRecord.kt deleted file mode 100644 index 89d9ab2..0000000 --- a/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/MessageRecord.kt +++ /dev/null @@ -1,13 +0,0 @@ -package pw.binom.agentik.messageLog - -/** - * @Deprecated - * Перенесено в `:journal-api`. Используй `pw.binom.agentik.journal.MessageRecord`. - * - * Backward-compat typealias. - */ -@Deprecated( - message = "Перенесено в :journal-api. Используй pw.binom.agentik.journal.MessageRecord.", - replaceWith = ReplaceWith("MessageRecord", "pw.binom.agentik.journal.MessageRecord"), -) -typealias MessageRecord = pw.binom.agentik.journal.MessageRecord diff --git a/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/MessageStore.kt b/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/MessageStore.kt deleted file mode 100644 index 0df7d41..0000000 --- a/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/MessageStore.kt +++ /dev/null @@ -1,17 +0,0 @@ -package pw.binom.agentik.messageLog - -/** - * @Deprecated - * Перенесено в `:journal-api`. Используй `pw.binom.agentik.journal.JournalStore`. - * - * Backward-compat typealias. Все существующие consumer'ы - * (`:standalone`, `:storage-ksqlite`, `:storage-inmemory`, тесты) импортируют - * `pw.binom.agentik.messageLog.MessageStore`. Буквально тот же тип. - * - * Удалить когда все импорты будут на `:journal-api`. - */ -@Deprecated( - message = "Перенесено в :journal-api. Используй pw.binom.agentik.journal.JournalStore.", - replaceWith = ReplaceWith("JournalStore", "pw.binom.agentik.journal.JournalStore"), -) -typealias MessageStore = pw.binom.agentik.journal.JournalStore diff --git a/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/MutableMessageStore.kt b/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/MutableMessageStore.kt deleted file mode 100644 index 4a660a1..0000000 --- a/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/MutableMessageStore.kt +++ /dev/null @@ -1,15 +0,0 @@ -package pw.binom.agentik.messageLog - -/** - * @Deprecated - * Перенесено в `:journal-api`. Используй `pw.binom.agentik.journal.MutableJournalStore`. - * - * Backward-compat typealias. `MutableMessageStore` == `MutableJournalStore` (один тип). - * - * Удалить когда все импорты будут на `:journal-api`. - */ -@Deprecated( - message = "Перенесено в :journal-api. Используй pw.binom.agentik.journal.MutableJournalStore.", - replaceWith = ReplaceWith("MutableJournalStore", "pw.binom.agentik.journal.MutableJournalStore"), -) -typealias MutableMessageStore = pw.binom.agentik.journal.MutableJournalStore diff --git a/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/Payload.kt b/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/Payload.kt deleted file mode 100644 index 1d19e3f..0000000 --- a/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/Payload.kt +++ /dev/null @@ -1,13 +0,0 @@ -package pw.binom.agentik.messageLog - -/** - * @Deprecated - * Перенесено в `:journal-api`. Используй `pw.binom.agentik.journal.MessageBodyPayload`. - * - * Backward-compat typealias. - */ -@Deprecated( - message = "Перенесено в :journal-api. Используй pw.binom.agentik.journal.MessageBodyPayload.", - replaceWith = ReplaceWith("MessageBodyPayload", "pw.binom.agentik.journal.MessageBodyPayload"), -) -typealias MessageBodyPayload = pw.binom.agentik.journal.MessageBodyPayload diff --git a/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/TurnTokens.kt b/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/TurnTokens.kt deleted file mode 100644 index 4168e14..0000000 --- a/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/TurnTokens.kt +++ /dev/null @@ -1,13 +0,0 @@ -package pw.binom.agentik.messageLog - -/** - * @Deprecated - * Перенесено в `:journal-api`. Используй `pw.binom.agentik.journal.TurnTokens`. - * - * Backward-compat typealias. - */ -@Deprecated( - message = "Перенесено в :journal-api. Используй pw.binom.agentik.journal.TurnTokens.", - replaceWith = ReplaceWith("TurnTokens", "pw.binom.agentik.journal.TurnTokens"), -) -typealias TurnTokens = pw.binom.agentik.journal.TurnTokens diff --git a/message-log-ksqlite/build.gradle.kts b/message-log-ksqlite/build.gradle.kts deleted file mode 100644 index 230e20b..0000000 --- a/message-log-ksqlite/build.gradle.kts +++ /dev/null @@ -1,33 +0,0 @@ -plugins { - alias(libs.plugins.kotlin.multiplatform) - alias(libs.plugins.kotlin.serialization) -} - -// KMP-реализация MutableMessageStore поверх ksqlite (https://github.com/caffeine-mgn/ksqlite). -// -// Цели сборки (параллельны :storage-ksqlite): -// - jvm() — основная, тесты гоняются здесь (in-memory DB без файла) -// - linuxX64() / mingwX64() — smoke-проверка что KMP реально KMP -// -// Apple targets невозможно собрать на Linux — Kotlin Multiplatform plugin auto-disables их. - -kotlin { - jvmToolchain(21) - - jvm() - linuxX64() - mingwX64() - - sourceSets { - commonMain.dependencies { - api(project(":message-log-api")) - - implementation("pw.binom.db:ksqlite:0.1.0") - implementation(libs.kotlinx.serialization.json) - } - commonTest.dependencies { - implementation(kotlin("test")) - implementation(libs.kotlinx.coroutines.test) - } - } -} \ No newline at end of file diff --git a/message-log-ksqlite/src/commonMain/kotlin/pw/binom/agentik/messageLog/ksqlite/KsqliteMessageStore.kt b/message-log-ksqlite/src/commonMain/kotlin/pw/binom/agentik/messageLog/ksqlite/KsqliteMessageStore.kt deleted file mode 100644 index 2aa1c3d..0000000 --- a/message-log-ksqlite/src/commonMain/kotlin/pw/binom/agentik/messageLog/ksqlite/KsqliteMessageStore.kt +++ /dev/null @@ -1,141 +0,0 @@ -package pw.binom.agentik.messageLog.ksqlite - -import kotlinx.serialization.json.Json -import pw.binom.agentik.messageLog.MessageRecord -import pw.binom.agentik.messageLog.MutableMessageStore -import pw.binom.db.ksqlite.SQLiteConnection -import pw.binom.db.ksqlite.SQLitePreparedStatement -import kotlin.time.Instant -import kotlinx.coroutines.Dispatchers -import kotlinx.coroutines.sync.Mutex -import kotlinx.coroutines.sync.withLock -import kotlinx.coroutines.withContext - -/** - * ksqlite-реализация [MutableMessageStore]. Схема таблицы `message` повторяет - * [pw.binom.agentik.storage.sqlite.SqliteMessageStore] для совместимости данных. - * - * Все четыре [SQLitePreparedStatement] компилируются один раз в конструкторе и - * переиспользуются: bind → execute → reset (cursor) + clearBindings (memory) - * перед следующим вызовом. Это устраняет per-call overhead на sqlite3_prepare_v2. - * - * encoding helpers (`encodeRecord`/`toMessageRecord`/`CallPayload`/...) — копия - * `:storage-sqlite`, живут в этом же модуле. См. [MessageCodecs]. - * - * **Threading**: ksqlite документирует соединение как single-threaded by convention - * (sqlite3 prepared statements safe to use serially from any thread). Мы сериализуем - * доступ через [Mutex] — операции вызываются на [Dispatchers.Default] и поток - * между вызовами может меняться; Mutex гарантирует что в любой момент времени - * только одна корутина использует statement. - * - * `clear(conversationId)` — каскадный helper, вызывается из - * [pw.binom.agentik.storage.ksqlite.KsqliteConversationStore.delete]. - * Не часть публичного [MutableMessageStore] API (audit log append-only); - * экспортируется через lambda в [pw.binom.agentik.storage.ksqlite.KsqliteStores.assemble]. - */ -class KsqliteMessageStore private constructor( - private val connection: SQLiteConnection, - private val ownsConnection: Boolean, -) : MutableMessageStore { - - /** - * Основной конструктор — принимает **уже открытое** [SQLiteConnection]. - * Соединение остаётся под управлением вызывающего (например, [pw.binom.agentik.storage.ksqlite.KsqliteStores.assemble] - * закрывает его сам после [close] всех своих store'ов). - */ - constructor(connection: SQLiteConnection) : this(connection, ownsConnection = false) - - /** - * Convenience-конструктор для standalone-сценариев: открывает файл-БД - * (`SQLiteConnection.open(path)`) и **сам закрывает её в [close]**. - * Не использовать если соединение шарится с другими store'ами — - * двойной [SQLiteConnection.close] не идемпотентен. - */ - constructor(path: String) : this(SQLiteConnection.open(path), ownsConnection = true) - - private val mutex = Mutex() - private val json = Json { ignoreUnknownKeys = true } - - private val insertStmt: SQLitePreparedStatement = connection.prepare( - "INSERT INTO message (id, conversation_id, kind, payload_json, created_at) VALUES (?, ?, ?, ?, ?)" - ) - private val listStmt: SQLitePreparedStatement = connection.prepare( - "SELECT id, conversation_id, kind, payload_json, created_at FROM message " + - "WHERE conversation_id = ? AND created_at > ? ORDER BY created_at ASC, id ASC LIMIT ? OFFSET ?" - ) - private val deleteStmt: SQLitePreparedStatement = connection.prepare( - "DELETE FROM message WHERE conversation_id = ?" - ) - - override suspend fun append(record: MessageRecord): Unit = withContext(Dispatchers.Default) { - mutex.withLock { - val (kind, payload) = encodeRecord(record) - insertStmt.execDml( - record.id, record.conversationId, kind, payload, record.createdAt.toEpochMilliseconds(), - ) - } - } - - override suspend fun list( - conversationId: String, - after: Instant, - offset: Int, - limit: Int, - ): List = withContext(Dispatchers.Default) { - mutex.withLock { - listStmt.runQuery(conversationId, after.toEpochMilliseconds(), limit.toLong(), offset.toLong()) - } - } - - override fun close() { - insertStmt.close() - listStmt.close() - deleteStmt.close() - if (ownsConnection) connection.close() - } - - /** - * Каскадный clear всех сообщений диалога — вызывается из - * [pw.binom.agentik.storage.ksqlite.KsqliteConversationStore.delete]. - * - * Публичный (не internal) чтобы `:storage-ksqlite` мог передать его как lambda. - * Не часть [MutableMessageStore] — audit log append-only. - */ - suspend fun clear(conversationId: String): Unit = withContext(Dispatchers.Default) { - mutex.withLock { - deleteStmt.execDml(conversationId) - } - } - - private fun SQLitePreparedStatement.bindAll(params: Array) { - params.forEachIndexed { i, v -> - val idx = i + 1 - when (v) { - is String -> bindText(idx, v) - is Long -> bindLong(idx, v) - is Int -> bindLong(idx, v.toLong()) - else -> error("Unsupported bind type: ${v::class}") - } - } - } - - /** Выполнить INSERT/DELETE statement. reset+clearBindings+bind+executeUpdate. */ - private fun SQLitePreparedStatement.execDml(vararg params: Any): Int { - reset() - clearBindings() - bindAll(params) - return executeUpdate() - } - - /** Выполнить SELECT statement, вернуть список [MessageRecord]. reset+clearBindings+bind+executeQuery+drain. */ - private fun SQLitePreparedStatement.runQuery(vararg params: Any): List { - reset() - clearBindings() - bindAll(params) - val out = mutableListOf() - executeQuery().use { rs -> - while (rs.next()) out.add(rs.toMessageRecord(json)) - } - return out - } -} \ No newline at end of file diff --git a/message-log-ksqlite/src/commonMain/kotlin/pw/binom/agentik/messageLog/ksqlite/MessageCodecs.kt b/message-log-ksqlite/src/commonMain/kotlin/pw/binom/agentik/messageLog/ksqlite/MessageCodecs.kt deleted file mode 100644 index ace8a2a..0000000 --- a/message-log-ksqlite/src/commonMain/kotlin/pw/binom/agentik/messageLog/ksqlite/MessageCodecs.kt +++ /dev/null @@ -1,79 +0,0 @@ -package pw.binom.agentik.messageLog.ksqlite - -import kotlinx.serialization.json.Json -import pw.binom.agentik.messageLog.MessageRecord -import pw.binom.agentik.messageLog.decodeBodyPayload -import pw.binom.agentik.messageLog.encodeBodyPayload -import pw.binom.db.ksqlite.SQLiteResultSet -import kotlin.time.Instant - -/** - * Кодирование [MessageRecord] → пара (kind, payloadJson) для SQLite. - * - * Копия [pw.binom.agentik.storage.sqlite.SqliteMessageStore] — `private` helpers - * нельзя переиспользовать между модулями, поэтому в каждом backend свой набор. - * Чтобы избежать дрейфа при изменении формата payload'а, оба набора синхронизируются - * через эти data class'ы (CallPayload/ResultPayload/ErrorPayload). - */ -internal fun encodeRecord(record: MessageRecord): Pair = when (record) { - is MessageRecord.UserMessage -> "user" to encodeBodyPayload( - content = record.content, - context = record.context, - ) - is MessageRecord.AssistantMessage -> "assistant" to encodeBodyPayload( - content = record.content, - tokens = record.tokens, - ) - is MessageRecord.ToolCall -> "tool_call" to Json.encodeToString( - CallPayload.serializer(), - CallPayload(name = record.toolName, title = record.toolTitle, argsJson = record.toolArgsJson), - ) - is MessageRecord.ToolResult -> "tool_result" to Json.encodeToString( - ResultPayload.serializer(), - ResultPayload(toolCallId = record.toolCallId, result = record.result), - ) - is MessageRecord.Error -> "error" to Json.encodeToString( - ErrorPayload.serializer(), - ErrorPayload(message = record.message, code = record.code), - ) -} - -internal fun SQLiteResultSet.toMessageRecord(json: Json): MessageRecord { - val id = getText(0)!! - val convId = getText(1)!! - val kind = getText(2)!! - val payload = getText(3)!! - val createdAt = Instant.fromEpochMilliseconds(getLong(4)!!) - return when (kind) { - "user" -> { - val d = decodeBodyPayload(payload) - MessageRecord.UserMessage(id = id, conversationId = convId, content = d.content, createdAt = createdAt, context = d.context) - } - "assistant" -> { - val d = decodeBodyPayload(payload) - MessageRecord.AssistantMessage(id = id, conversationId = convId, content = d.content, createdAt = createdAt, tokens = d.tokens) - } - "tool_call" -> { - val p = Json.decodeFromString(CallPayload.serializer(), payload) - MessageRecord.ToolCall(id = id, conversationId = convId, toolName = p.name, toolTitle = p.title, toolArgsJson = p.argsJson, createdAt = createdAt) - } - "tool_result" -> { - val p = Json.decodeFromString(ResultPayload.serializer(), payload) - MessageRecord.ToolResult(id = id, conversationId = convId, toolCallId = p.toolCallId, result = p.result, createdAt = createdAt) - } - "error" -> { - val p = Json.decodeFromString(ErrorPayload.serializer(), payload) - MessageRecord.Error(id = id, conversationId = convId, message = p.message, code = p.code, createdAt = createdAt) - } - else -> error("Unknown message kind in audit log: $kind") - } -} - -@kotlinx.serialization.Serializable -internal data class CallPayload(val name: String, val title: String?, val argsJson: String) - -@kotlinx.serialization.Serializable -internal data class ResultPayload(val toolCallId: String, val result: String?) - -@kotlinx.serialization.Serializable -internal data class ErrorPayload(val message: String, val code: String?) \ No newline at end of file diff --git a/message-log-ksqlite/src/commonTest/kotlin/pw/binom/agentik/messageLog/ksqlite/KsqliteMessageStoreTest.kt b/message-log-ksqlite/src/commonTest/kotlin/pw/binom/agentik/messageLog/ksqlite/KsqliteMessageStoreTest.kt deleted file mode 100644 index 57e4e3c..0000000 --- a/message-log-ksqlite/src/commonTest/kotlin/pw/binom/agentik/messageLog/ksqlite/KsqliteMessageStoreTest.kt +++ /dev/null @@ -1,122 +0,0 @@ -package pw.binom.agentik.messageLog.ksqlite - -import kotlinx.coroutines.flow.count -import kotlinx.coroutines.flow.first -import kotlinx.coroutines.flow.toList -import kotlinx.coroutines.test.runTest -import pw.binom.agentik.messageLog.Content -import pw.binom.agentik.messageLog.MessageRecord -import pw.binom.db.ksqlite.SQLiteConnection -import kotlin.test.AfterTest -import kotlin.test.BeforeTest -import kotlin.test.Test -import kotlin.test.assertEquals -import kotlin.test.assertNull -import kotlin.time.Instant - -class KsqliteMessageStoreTest { - - private lateinit var conn: SQLiteConnection - private lateinit var store: KsqliteMessageStore - - @BeforeTest - fun setup() { - conn = SQLiteConnection.memory("msg-${kotlin.random.Random.nextLong()}") - conn.exec(SCHEMA) - store = KsqliteMessageStore(conn) - } - - @AfterTest - fun tearDown() = conn.close() - - @Test - fun testAppendUserAndRetrieve() = runTest { - store.append(MessageRecord.UserMessage( - id = "m1", - conversationId = "conv1", - content = listOf(Content.Text("hello")), - createdAt = Instant.parse("2026-09-15T10:01:00Z"), - context = null, - )) - val list = store.listFlow("conv1", Instant.DISTANT_PAST).toList() - assertEquals(1, list.size) - val msg = list[0] - assertEquals("m1", msg.id) - assertEquals(MessageRecord.UserMessage::class, msg::class) - } - - @Test - fun testAppendAssistantWithTokens() = runTest { - store.append(MessageRecord.AssistantMessage( - id = "m1", - conversationId = "conv1", - content = listOf(Content.Text("hi")), - createdAt = Instant.parse("2026-09-15T10:01:00Z"), - tokens = pw.binom.agentik.messageLog.TurnTokens(input = 50, output = 30), - )) - val all = store.listFlow("conv1", Instant.DISTANT_PAST).toList() - assertEquals(1, all.size) - val msg = all[0] as MessageRecord.AssistantMessage - assertEquals(50, msg.tokens?.input) - assertEquals(30, msg.tokens?.output) - } - - @Test - fun testListAfterFiltersByTimestamp() = runTest { - val t1 = Instant.parse("2026-09-15T10:01:00Z") - val t2 = Instant.parse("2026-09-15T10:02:00Z") - val t3 = Instant.parse("2026-09-15T10:03:00Z") - store.append(MessageRecord.UserMessage("m1", "conv1", listOf(Content.Text("a")), t1, null)) - store.append(MessageRecord.UserMessage("m2", "conv1", listOf(Content.Text("b")), t2, null)) - store.append(MessageRecord.UserMessage("m3", "conv1", listOf(Content.Text("c")), t3, null)) - - val after = store.list("conv1", after = t1, offset = 0, limit = 10) - assertEquals(2, after.size) - assertEquals(listOf("m2", "m3"), after.map { it.id }) - } - - @Test - fun testListAllReturnsAllInOrder() = runTest { - val t = Instant.parse("2026-09-15T10:00:00Z") - for (i in 1..3) store.append( - MessageRecord.UserMessage("m$i", "conv1", listOf(Content.Text("x$i")), t + kotlin.time.Duration.parse("PT${i}S"), null) - ) - assertEquals(listOf("m1", "m2", "m3"), store.listFlow("conv1", Instant.DISTANT_PAST).toList().map { it.id }) - } - - @Test - fun testClearRemovesByConversation() = runTest { - val t = Instant.parse("2026-09-15T10:00:00Z") - store.append(MessageRecord.UserMessage("m1", "conv1", listOf(Content.Text("a")), t, null)) - store.append(MessageRecord.UserMessage("m2", "conv2", listOf(Content.Text("b")), t, null)) - store.clear("conv1") - assertEquals(0, store.listFlow("conv1", Instant.DISTANT_PAST).count()) - assertEquals(1, store.listFlow("conv2", Instant.DISTANT_PAST).count()) - } - - @Test - fun testListWithNullContext() = runTest { - store.append(MessageRecord.UserMessage( - id = "m1", - conversationId = "conv1", - content = listOf(Content.Text("hi")), - createdAt = Instant.parse("2026-09-15T10:01:00Z"), - context = null, - )) - val msg = store.listFlow("conv1", Instant.DISTANT_PAST).first() as MessageRecord.UserMessage - assertNull(msg.context) - } - - private companion object { - const val SCHEMA = """ - CREATE TABLE IF NOT EXISTS message ( - id TEXT NOT NULL PRIMARY KEY, - conversation_id TEXT NOT NULL, - kind TEXT NOT NULL, - payload_json TEXT NOT NULL, - created_at INTEGER NOT NULL - ); - CREATE INDEX IF NOT EXISTS idx_msg_conv ON message(conversation_id, created_at); - """ - } -} \ No newline at end of file diff --git a/settings.gradle.kts b/settings.gradle.kts index 882f8f9..2577e3c 100644 --- a/settings.gradle.kts +++ b/settings.gradle.kts @@ -68,7 +68,8 @@ include(":memory-vector") // IO-зависимостей. Используется в тестах (быстрый setup, без JDBC) и будет // использоваться в Android-сборке (JVector/SQLite не подходят для ART out-of-box). include(":message-store-api") -include(":message-log-api") +//include(":message-log-api") +//include(":message-log-ksqlite") // Новые API-модули трёх сущностей (canonical имена): // - :journal-api — append-only audit log (бывший :message-log-api) // - :context-api — то что видит LLM (бывший :working-memory-api) @@ -82,11 +83,12 @@ include(":outbox-api") // короткий live tail + recent replay; полный audit log живёт в :message-store-api // (там — MessageStore + ConversationStore). EventStore сам управляет eviction, // никаких prune-методов наружу. KMP, без implementations пока. -include(":event-store") +//include(":event-store") // KMP in-memory реализация MutableEventStore. ConcurrentLinkedDeque + TTL/size // eviction. Для тестов, dev-режима, embedded-сценариев (Android core). -include(":event-store-in-memory") -include(":working-memory-api") +include(":outbox-inmemory") +//include(":event-store-in-memory") +//include(":working-memory-api") include(":storage-inmemory") // SQLDelight-реализация store'ов из :message-store-api и :working-memory-api. JVM-only // KMP-реализация EventStore поверх ksqlite (https://github.com/caffeine-mgn/ksqlite). diff --git a/standalone/build.gradle.kts b/standalone/build.gradle.kts index a6cdb6c..a9bd1f7 100644 --- a/standalone/build.gradle.kts +++ b/standalone/build.gradle.kts @@ -76,7 +76,7 @@ kotlin { // Новый единый канал событий агента — заменил старые // `agentEvents: MutableSharedFlow` и per-conv `ConversationEvents._flow`. // Bounded tail + auto-TTL, generic CommonEvent envelope. - implementation(project(":event-store-in-memory")) + implementation(project(":outbox-inmemory")) implementation(project(":agent-toolsets")) // Generic LLM-side tools (LlmReflector, SkillMiner, LlmMemoryReviewer, // ContextCompactor, парсеры/промпты). Вынесены из :standalone. 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 033b86d..a1c4e2d 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 @@ -213,7 +213,7 @@ class ChatAgent( * **Live tail + auto-TTL** — клиенты больше не должны заботиться о persistence * или подписке на два отдельных канала. */ - private val eventStore: MutableOutboxStore = pw.binom.agentik.eventStore.inmemory.InMemoryEventStore( + private val eventStore: MutableOutboxStore = pw.binom.agentik.outbox.inmemory.InMemoryOutboxStore( maxMessages = null, ttl = null, ) diff --git a/storage-inmemory/build.gradle.kts b/storage-inmemory/build.gradle.kts index f91ff85..5923ee3 100644 --- a/storage-inmemory/build.gradle.kts +++ b/storage-inmemory/build.gradle.kts @@ -21,8 +21,8 @@ kotlin { sourceSets { commonMain.dependencies { api(project(":message-store-api")) - api(project(":message-log-api")) - api(project(":working-memory-api")) + api(project(":journal-api")) + api(project(":context-api")) } commonTest.dependencies { implementation(kotlin("test")) diff --git a/storage-inmemory/src/commonMain/kotlin/pw/binom/agentik/storage/inmemory/InMemoryMessageStore.kt b/storage-inmemory/src/commonMain/kotlin/pw/binom/agentik/storage/inmemory/InMemoryMessageStore.kt index 7be5eeb..c461b9c 100644 --- a/storage-inmemory/src/commonMain/kotlin/pw/binom/agentik/storage/inmemory/InMemoryMessageStore.kt +++ b/storage-inmemory/src/commonMain/kotlin/pw/binom/agentik/storage/inmemory/InMemoryMessageStore.kt @@ -2,8 +2,8 @@ package pw.binom.agentik.storage.inmemory import kotlinx.coroutines.sync.Mutex import kotlinx.coroutines.sync.withLock -import pw.binom.agentik.messageLog.MessageRecord -import pw.binom.agentik.messageLog.MutableMessageStore +import pw.binom.agentik.journal.MessageRecord +import pw.binom.agentik.journal.MutableJournalStore import kotlin.time.Instant /** @@ -16,7 +16,7 @@ import kotlin.time.Instant * В отличие от ksqlite-импла, не требует native SQLite — работает в любом * KMP-таргете (включая iOS/native). */ -class InMemoryMessageStore : MutableMessageStore { +class InMemoryMessageStore : MutableJournalStore { private val byConv: MutableMap> = mutableMapOf() private val mutex = Mutex() 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 64d92fd..4ad2bc6 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,9 +1,9 @@ package pw.binom.agentik.storage.inmemory -import pw.binom.agentik.messageLog.MutableMessageStore +import pw.binom.agentik.context.ContextStore +import pw.binom.agentik.journal.MutableJournalStore import pw.binom.agentik.messageStore.ConversationStore import pw.binom.agentik.messageStore.ReflectionStore -import pw.binom.agentik.context.ContextStore import kotlin.time.Clock /** @@ -22,7 +22,7 @@ import kotlin.time.Clock object InMemoryStorage { data class Bundle( val conversationStore: ConversationStore, - val messageStore: MutableMessageStore, + val messageStore: MutableJournalStore, val workingMemoryStore: ContextStore, val reflectionStore: ReflectionStore, ) diff --git a/storage-ksqlite/build.gradle.kts b/storage-ksqlite/build.gradle.kts index d984c5e..032a4e7 100644 --- a/storage-ksqlite/build.gradle.kts +++ b/storage-ksqlite/build.gradle.kts @@ -36,8 +36,8 @@ kotlin { implementation(libs.kotlinx.serialization.json) api(project(":message-store-api")) - api(project(":message-log-api")) - api(project(":working-memory-api")) + api(project(":journal-api")) + api(project(":context-api")) } commonTest.dependencies { implementation(kotlin("test")) diff --git a/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteMessageStore.kt b/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteMessageStore.kt index 64ca7f1..75706fb 100644 --- a/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteMessageStore.kt +++ b/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteMessageStore.kt @@ -1,8 +1,8 @@ package pw.binom.agentik.storage.ksqlite import kotlinx.serialization.json.Json -import pw.binom.agentik.messageLog.MessageRecord -import pw.binom.agentik.messageLog.MutableMessageStore +import pw.binom.agentik.journal.MessageRecord +import pw.binom.agentik.journal.MutableJournalStore import pw.binom.db.ksqlite.SQLiteConnection import pw.binom.db.ksqlite.SQLitePreparedStatement import kotlin.time.Instant @@ -25,7 +25,7 @@ import kotlinx.coroutines.withContext */ class KsqliteMessageStore internal constructor( private val connection: SQLiteConnection, -) : MutableMessageStore { +) : MutableJournalStore { private val mutex = Mutex() private val json = Json { ignoreUnknownKeys = true } 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 d4412d6..0f487dd 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 @@ -1,9 +1,9 @@ package pw.binom.agentik.storage.ksqlite +import pw.binom.agentik.context.ContextStore +import pw.binom.agentik.journal.MutableJournalStore import pw.binom.agentik.messageStore.ConversationStore -import pw.binom.agentik.messageLog.MutableMessageStore import pw.binom.agentik.messageStore.ReflectionStore -import pw.binom.agentik.workingMemory.WorkingMemoryStore import pw.binom.db.ksqlite.SQLiteConnection /** @@ -18,8 +18,8 @@ import pw.binom.db.ksqlite.SQLiteConnection class KsqliteStores internal constructor( val connection: SQLiteConnection, val conversations: ConversationStore, - val messages: MutableMessageStore, - val workingMemory: WorkingMemoryStore, + val messages: MutableJournalStore, + val workingMemory: ContextStore, val reflections: ReflectionStore, ) : AutoCloseable { diff --git a/working-memory-api/build.gradle.kts b/working-memory-api/build.gradle.kts index a9709e2..f0aa485 100644 --- a/working-memory-api/build.gradle.kts +++ b/working-memory-api/build.gradle.kts @@ -24,7 +24,7 @@ kotlin { sourceSets { commonMain.dependencies { - api(project(":message-log-api")) +// api(project(":message-log-api")) api(project(":context-api")) api(libs.kotlinx.coroutines.core) api(libs.kotlinx.serialization.core) diff --git a/working-memory-api/src/commonMain/kotlin/pw/binom/agentik/workingMemory/WorkingMemoryEntry.kt b/working-memory-api/src/commonMain/kotlin/pw/binom/agentik/workingMemory/WorkingMemoryEntry.kt deleted file mode 100644 index 36478af..0000000 --- a/working-memory-api/src/commonMain/kotlin/pw/binom/agentik/workingMemory/WorkingMemoryEntry.kt +++ /dev/null @@ -1,16 +0,0 @@ -package pw.binom.agentik.workingMemory - -/** - * @Deprecated - * Перенесено в `:context-api`. Используй `pw.binom.agentik.context.WorkingMemoryEntry`. - * - * Backward-compat typealias. Все импорты `pw.binom.agentik.workingMemory.WorkingMemoryEntry` - * продолжают работать как тип `pw.binom.agentik.context.WorkingMemoryEntry` (тот же тип, транспарентно). - * - * Удалить когда все импорты будут на `:context-api`. - */ -@Deprecated( - message = "Перенесено в :context-api. Используй pw.binom.agentik.context.WorkingMemoryEntry.", - replaceWith = ReplaceWith("WorkingMemoryEntry", "pw.binom.agentik.context.WorkingMemoryEntry"), -) -typealias WorkingMemoryEntry = pw.binom.agentik.context.WorkingMemoryEntry diff --git a/working-memory-api/src/commonMain/kotlin/pw/binom/agentik/workingMemory/WorkingMemoryStore.kt b/working-memory-api/src/commonMain/kotlin/pw/binom/agentik/workingMemory/WorkingMemoryStore.kt deleted file mode 100644 index 01167ac..0000000 --- a/working-memory-api/src/commonMain/kotlin/pw/binom/agentik/workingMemory/WorkingMemoryStore.kt +++ /dev/null @@ -1,18 +0,0 @@ -package pw.binom.agentik.workingMemory - -/** - * @Deprecated - * Перенесено в `:context-api`. Используй `pw.binom.agentik.context.ContextStore`. - * - * Этот typealias сохранён временно для backward-compat: существующие impl'ы - * (`KsqliteStores.workingMemory` в `:storage-ksqlite`, тесты) импортируют - * `pw.binom.agentik.workingMemory.WorkingMemoryStore`. Буквально тот же тип, - * что и `pw.binom.agentik.context.ContextStore` — typealias транспарентен. - * - * Удалить когда все импорты будут на `:context-api`. - */ -@Deprecated( - message = "Перенесено в :context-api. Используй pw.binom.agentik.context.ContextStore.", - replaceWith = ReplaceWith("ContextStore", "pw.binom.agentik.context.ContextStore"), -) -typealias WorkingMemoryStore = pw.binom.agentik.context.ContextStore