diff --git a/context-ksqlite/build.gradle.kts b/context-ksqlite/build.gradle.kts new file mode 100644 index 0000000..65c62f0 --- /dev/null +++ b/context-ksqlite/build.gradle.kts @@ -0,0 +1,40 @@ +plugins { + alias(libs.plugins.kotlin.multiplatform) + alias(libs.plugins.kotlin.serialization) +} + +// KMP-реализация :context-api (ContextStore) поверх ksqlite. +// Минимальная — только таблица `working_memory` + 2 индекса по ней. +// Остальные таблицы (`conversation`, `message`, `reflection`) живут в +// других ksqlite-модулях; этот модуль не претендует на полную схему +// агента. +// +// Цели сборки — jvm() + linuxX64() + mingwX64(); Apple targets auto-disabled +// на Linux (см. KDoc :storage-ksqlite). + +kotlin { + jvmToolchain(21) + + jvm() + linuxX64() + mingwX64() + + sourceSets { + commonMain.dependencies { + implementation("pw.binom.db:ksqlite:0.1.1-SNAPSHOT") + implementation(libs.kotlinx.serialization.json) + + api(project(":context-api")) + api(project(":message-store-api")) + // :context-api ссылается на Content / MessageContext из + // :message-log-api (старый canonical). Транзитивно через api, + // но фиксируем явно чтобы тестовый код видел Content без + // обхода через :context-api. + api(project(":message-log-api")) + } + commonTest.dependencies { + implementation(kotlin("test")) + implementation(libs.kotlinx.coroutines.test) + } + } +} diff --git a/context-ksqlite/src/commonMain/kotlin/pw/binom/agentik/context/ksqlite/KsqliteContextStore.kt b/context-ksqlite/src/commonMain/kotlin/pw/binom/agentik/context/ksqlite/KsqliteContextStore.kt new file mode 100644 index 0000000..a49302f --- /dev/null +++ b/context-ksqlite/src/commonMain/kotlin/pw/binom/agentik/context/ksqlite/KsqliteContextStore.kt @@ -0,0 +1,200 @@ +package pw.binom.agentik.context.ksqlite + +import kotlinx.serialization.json.Json +import pw.binom.agentik.context.ContextStore +import pw.binom.agentik.context.WorkingMemoryEntry +import pw.binom.agentik.context.WorkingMemoryRow +import pw.binom.agentik.messageStore.Ids +import pw.binom.db.ksqlite.SQLiteConnection +import pw.binom.db.ksqlite.SQLitePreparedStatement +import kotlin.time.Clock +import kotlin.time.Instant +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.sync.Mutex +import kotlinx.coroutines.sync.withLock +import kotlinx.coroutines.withContext + +/** + * ksqlite-реализация [ContextStore] (таблица `working_memory`). + * + * Структура — копия [pw.binom.agentik.storage.ksqlite.KsqliteWorkingMemoryStore] + * из `:storage-ksqlite`, но: + * - лежит в собственном модуле `:context-ksqlite`; + * - реализует переименованный [ContextStore] (раньше был `WorkingMemoryStore`, + * теперь главный класс — `ContextStore`); сами типы строк + * [WorkingMemoryEntry] / [WorkingMemoryRow] не переименовывались. + * + * ВНИМАНИЕ: `:storage-ksqlite/KsqliteWorkingMemoryStore.kt` остаётся на диске — + * это копия, не замена. Не удалять старый файл; миграция consumers'ов — отдельно. + */ +class KsqliteContextStore( + private val connection: SQLiteConnection, +) : ContextStore { + + private val mutex = Mutex() + private val json = Json { ignoreUnknownKeys = true } + + // pre-prepare (см. KsqliteMessageStore KDoc — почему это критично против + // SIGSEGV в StmtHolder.finalize на закрытой connection). + private val insertStmt: SQLitePreparedStatement = connection.prepare( + """ + INSERT INTO ${Schema.TABLE_WORKING_MEMORY} + (${Schema.COL_ID}, ${Schema.COL_CONVERSATION_ID}, ${Schema.COL_ORDER_IDX}, + ${Schema.COL_SOURCE_MESSAGE_ID}, ${Schema.COL_KIND}, + ${Schema.COL_PAYLOAD_JSON}, ${Schema.COL_CREATED_AT}) + VALUES (?, ?, ?, ?, ?, ?, ?) + """.trimIndent() + ) + private val listStmt: SQLitePreparedStatement = connection.prepare( + """ + SELECT ${Schema.COL_ID}, ${Schema.COL_CONVERSATION_ID}, ${Schema.COL_ORDER_IDX}, + ${Schema.COL_SOURCE_MESSAGE_ID}, ${Schema.COL_PAYLOAD_JSON}, ${Schema.COL_CREATED_AT} + FROM ${Schema.TABLE_WORKING_MEMORY} + WHERE ${Schema.COL_CONVERSATION_ID} = ? + ORDER BY ${Schema.COL_ORDER_IDX} ASC + """.trimIndent() + ) + private val clearStmt: SQLitePreparedStatement = connection.prepare( + "DELETE FROM ${Schema.TABLE_WORKING_MEMORY} WHERE ${Schema.COL_CONVERSATION_ID} = ?" + ) + private val maxOrderIdxStmt: SQLitePreparedStatement = connection.prepare( + """ + SELECT COALESCE(MAX(${Schema.COL_ORDER_IDX}), 0) + FROM ${Schema.TABLE_WORKING_MEMORY} + WHERE ${Schema.COL_CONVERSATION_ID} = ? + """.trimIndent() + ) + private val dropFromIdxStmt: SQLitePreparedStatement = connection.prepare( + """ + DELETE FROM ${Schema.TABLE_WORKING_MEMORY} + WHERE ${Schema.COL_CONVERSATION_ID} = ? AND ${Schema.COL_ORDER_IDX} >= ? + """.trimIndent() + ) + private val insertSummaryStmt: SQLitePreparedStatement = connection.prepare( + """ + INSERT INTO ${Schema.TABLE_WORKING_MEMORY} + (${Schema.COL_ID}, ${Schema.COL_CONVERSATION_ID}, ${Schema.COL_ORDER_IDX}, + ${Schema.COL_SOURCE_MESSAGE_ID}, ${Schema.COL_KIND}, + ${Schema.COL_PAYLOAD_JSON}, ${Schema.COL_CREATED_AT}) + VALUES (?, ?, ?, NULL, ?, ?, ?) + """.trimIndent() + ) + + override suspend fun append(conversationId: String, entry: WorkingMemoryEntry, now: Instant): Unit = withContext(Dispatchers.Default) { + mutex.withLock { + val newIdx = maxOrderIdx(conversationId) + 1 + insertStmt.reset() + insertStmt.clearBindings() + insertStmt.bindText(1, Ids.new("wm")) + insertStmt.bindText(2, conversationId) + insertStmt.bindLong(3, newIdx) + val srcId = entry.sourceMessageId + if (srcId != null) insertStmt.bindText(4, srcId) else insertStmt.bindNull(4) + insertStmt.bindText(5, entryKind(entry)) + insertStmt.bindText(6, json.encodeToString(WorkingMemoryEntry.serializer(), entry)) + insertStmt.bindLong(7, now.toEpochMilliseconds()) + insertStmt.executeUpdate() + } + } + + override suspend fun list(conversationId: String): List = withContext(Dispatchers.Default) { + mutex.withLock { + listStmt.reset() + listStmt.clearBindings() + listStmt.bindText(1, conversationId) + val out = mutableListOf() + listStmt.executeQuery().use { rs -> + while (rs.next()) { + out.add( + WorkingMemoryRow( + id = rs.getText(0)!!, + conversationId = rs.getText(1)!!, + orderIdx = rs.getLong(2)!!, + sourceMessageId = rs.getText(3), + entry = Json.decodeFromString(WorkingMemoryEntry.serializer(), rs.getText(4)!!), + createdAt = Instant.fromEpochMilliseconds(rs.getLong(5)!!), + ) + ) + } + } + out + } + } + + override suspend fun clear(conversationId: String): Unit = withContext(Dispatchers.Default) { + mutex.withLock { + clearStmt.reset() + clearStmt.clearBindings() + clearStmt.bindText(1, conversationId) + clearStmt.executeUpdate() + } + } + + override suspend fun compact( + dropFromOrderIdx: Long, + conversationId: String, + summaryText: String?, + ): Long = withContext(Dispatchers.Default) { + mutex.withLock { + var newMax = 0L + val nowMs = Clock.System.now().toEpochMilliseconds() + val summaryId = Ids.new("wm") + connection.exec("BEGIN") + try { + dropFromIdxStmt.reset() + dropFromIdxStmt.clearBindings() + dropFromIdxStmt.bindText(1, conversationId) + dropFromIdxStmt.bindLong(2, dropFromOrderIdx) + dropFromIdxStmt.executeUpdate() + + if (!summaryText.isNullOrBlank()) { + val afterDelete = maxOrderIdx(conversationId) + val newIdx = afterDelete + 1 + insertSummaryStmt.reset() + insertSummaryStmt.clearBindings() + insertSummaryStmt.bindText(1, summaryId) + insertSummaryStmt.bindText(2, conversationId) + insertSummaryStmt.bindLong(3, newIdx) + insertSummaryStmt.bindText(4, "summary") + insertSummaryStmt.bindText(5, json.encodeToString(WorkingMemoryEntry.serializer(), WorkingMemoryEntry.Summary(text = summaryText))) + insertSummaryStmt.bindLong(6, nowMs) + insertSummaryStmt.executeUpdate() + newMax = newIdx + } else { + newMax = maxOrderIdx(conversationId) + } + connection.exec("COMMIT") + } catch (t: Throwable) { + runCatching { connection.exec("ROLLBACK") } + throw t + } + newMax + } + } + + override fun close() { + insertStmt.close() + listStmt.close() + clearStmt.close() + maxOrderIdxStmt.close() + dropFromIdxStmt.close() + insertSummaryStmt.close() + } + + private fun maxOrderIdx(conversationId: String): Long { + maxOrderIdxStmt.reset() + maxOrderIdxStmt.clearBindings() + maxOrderIdxStmt.bindText(1, conversationId) + maxOrderIdxStmt.executeQuery().use { rs -> + if (rs.next()) return rs.getLong(0) ?: 0L + } + return 0L + } + + private fun entryKind(e: WorkingMemoryEntry): String = when (e) { + is WorkingMemoryEntry.User -> "user" + is WorkingMemoryEntry.Assistant -> "assistant" + is WorkingMemoryEntry.ToolExchange -> "tool_exchange" + is WorkingMemoryEntry.Summary -> "summary" + } +} diff --git a/context-ksqlite/src/commonMain/kotlin/pw/binom/agentik/context/ksqlite/Schema.kt b/context-ksqlite/src/commonMain/kotlin/pw/binom/agentik/context/ksqlite/Schema.kt new file mode 100644 index 0000000..abf7caf --- /dev/null +++ b/context-ksqlite/src/commonMain/kotlin/pw/binom/agentik/context/ksqlite/Schema.kt @@ -0,0 +1,98 @@ +package pw.binom.agentik.context.ksqlite + +import pw.binom.db.ksqlite.SQLiteConnection + +/** + * Имена таблиц/колонок/индексов для ksqlite-бэкенда `:context-api`. + * + * Минимум — только то, что относится к `working_memory` (реализация + * [KsqliteContextStore]). Остальные таблицы агента (`conversation`, + * `message`, `reflection`) живут в других ksqlite-модулях. + * + * Все DDL/DML в этом модуле должны ссылаться на эти константы — никаких + * хардкоженных литералов в `prepare("SELECT ... FROM foo ...")` в store'е. + */ +internal object Schema { + + /** Версия схемы модуля. Увеличивать при ЛЮБОМ изменении DDL. */ + const val CURRENT_VERSION: Int = 1 + + // ───── Таблица ───── + const val TABLE_WORKING_MEMORY = "working_memory" + + // ───── Колонки ───── + const val COL_ID = "id" + const val COL_CONVERSATION_ID = "conversation_id" + const val COL_ORDER_IDX = "order_idx" + const val COL_SOURCE_MESSAGE_ID = "source_message_id" + const val COL_KIND = "kind" + const val COL_PAYLOAD_JSON = "payload_json" + const val COL_CREATED_AT = "created_at" + + // ───── Индексы ───── + const val IDX_WM_UNIQUE = "idx_wm_unique" + const val IDX_WM_CONV = "idx_wm_conv" + + private val v1Ddl = """ + CREATE TABLE IF NOT EXISTS $TABLE_WORKING_MEMORY ( + $COL_ID TEXT NOT NULL PRIMARY KEY, + $COL_CONVERSATION_ID TEXT NOT NULL, + $COL_ORDER_IDX INTEGER NOT NULL, + $COL_SOURCE_MESSAGE_ID TEXT, + $COL_KIND TEXT NOT NULL, + $COL_PAYLOAD_JSON TEXT NOT NULL, + $COL_CREATED_AT INTEGER NOT NULL + ); + """.trimIndent() + + private val v1IndexesDdl = """ + CREATE UNIQUE INDEX IF NOT EXISTS $IDX_WM_UNIQUE + ON $TABLE_WORKING_MEMORY($COL_CONVERSATION_ID, $COL_ORDER_IDX); + CREATE INDEX IF NOT EXISTS $IDX_WM_CONV + ON $TABLE_WORKING_MEMORY($COL_CONVERSATION_ID, $COL_ORDER_IDX); + """.trimIndent() + + /** + * Прогоняет миграцию схемы до [CURRENT_VERSION] на пустой или существующей БД. + * + * Версия хранится в `PRAGMA user_version` (стандартный SQLite-механизм, + * 32-bit int в заголовке БД — без своей таблицы). Каждая миграция — + * блок DDL под номером `fromV+1`, выполняется в транзакции. Если миграция + * упадёт посередине — `ROLLBACK` оставит БД на предыдущей версии. + * + * Идемпотентен: повторный вызов на уже мигрированной БД — no-op. + */ + fun migrate(conn: SQLiteConnection) { + val current = readUserVersion(conn) + if (current >= CURRENT_VERSION) return + + conn.exec("BEGIN") + try { + if (current < 1) { + conn.exec(v1Ddl) + conn.exec(v1IndexesDdl) + } + // future: if (current < 2) { conn.exec(v2Ddl) } + writeUserVersion(conn, CURRENT_VERSION) + conn.exec("COMMIT") + } catch (t: Throwable) { + runCatching { conn.exec("ROLLBACK") } + throw t + } + } + + private fun readUserVersion(conn: SQLiteConnection): Int { + conn.prepare("PRAGMA user_version").use { stmt -> + stmt.executeQuery().use { rs -> + if (rs.next()) return rs.getLong(0)?.toInt() ?: 0 + } + } + return 0 + } + + private fun writeUserVersion(conn: SQLiteConnection, version: Int) { + // SQLite PRAGMA с literal-аргументом нельзя параметризовать через `?`, + // поэтому собираем SQL строкой (значение контролируемое, не user input). + conn.exec("PRAGMA user_version = $version") + } +} diff --git a/context-ksqlite/src/commonTest/kotlin/pw/binom/agentik/context/ksqlite/KsqliteContextStoreTest.kt b/context-ksqlite/src/commonTest/kotlin/pw/binom/agentik/context/ksqlite/KsqliteContextStoreTest.kt new file mode 100644 index 0000000..d39c749 --- /dev/null +++ b/context-ksqlite/src/commonTest/kotlin/pw/binom/agentik/context/ksqlite/KsqliteContextStoreTest.kt @@ -0,0 +1,107 @@ +package pw.binom.agentik.context.ksqlite + +import kotlinx.coroutines.test.runTest +import pw.binom.agentik.context.WorkingMemoryEntry +import pw.binom.agentik.messageLog.Content +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.assertTrue +import kotlin.time.Instant + +/** + * Тесты для [KsqliteContextStore] — точная копия + * `KsqliteWorkingMemoryStoreTest` из `:storage-ksqlite`, с переименованием + * типов (`WorkingMemoryStore` → `ContextStore`) и обновлённым пакетом для + * `Content` (`pw.binom.agentik.journal` — новый canonical, но структура + * та же). + * + * Тестовая фикстура: in-memory SQLiteConnection, [Schema.migrate] в @BeforeTest, + * `KsqliteContextStore(conn)` + ручной close в @AfterTest. Никакой внешней + * зависимости от `KsqliteStores` из `:storage-ksqlite` — этот модуль + * автономный. + */ +class KsqliteContextStoreTest { + + private lateinit var conn: SQLiteConnection + private lateinit var store: KsqliteContextStore + + @BeforeTest + fun setup() { + conn = SQLiteConnection.memory("ctx-${kotlin.random.Random.nextLong()}") + Schema.migrate(conn) + store = KsqliteContextStore(conn) + } + + @AfterTest + fun tearDown() { + store.close() + conn.close() + } + + private fun userMsg(content: String, srcId: String = "m-${content.hashCode()}"): WorkingMemoryEntry.User = + WorkingMemoryEntry.User(sourceMessageId = srcId, content = listOf(Content.Text(content))) + + private fun asstMsg(content: String): WorkingMemoryEntry.Assistant = + WorkingMemoryEntry.Assistant(sourceMessageId = "m-${content.hashCode()}", content = listOf(Content.Text(content))) + + @Test + fun testAppendAndListReturnsInOrder() = runTest { + val t = Instant.parse("2026-09-15T10:00:00Z") + store.append("c1", userMsg("first"), t) + store.append("c1", asstMsg("reply"), t) + val list = store.list("c1") + assertEquals(2, list.size) + assertEquals(1L, list[0].orderIdx) + assertEquals(2L, list[1].orderIdx) + } + + @Test + fun testListIsolatesConversations() = runTest { + val t = Instant.parse("2026-09-15T10:00:00Z") + store.append("c1", userMsg("c1-msg"), t) + store.append("c2", userMsg("c2-msg"), t) + assertEquals(1, store.list("c1").size) + assertEquals(1, store.list("c2").size) + } + + @Test + fun testClearRemovesAllForConversation() = runTest { + val t = Instant.parse("2026-09-15T10:00:00Z") + store.append("c1", userMsg("a"), t) + store.append("c1", userMsg("b"), t) + store.clear("c1") + assertEquals(emptyList(), store.list("c1")) + } + + @Test + fun testCompactDeletesAndInsertsSummary() = runTest { + val t = Instant.parse("2026-09-15T10:00:00Z") + store.append("c1", userMsg("a"), t) + store.append("c1", userMsg("b"), t) + store.append("c1", userMsg("c"), t) + // dropFromOrderIdx=2: удаляет idx=2 и idx=3 (b и c), остаётся idx=1 (a). + // Summary встаёт на idx=2 (= max(remaining)+1). Возвращает newMax=2. + val newMax = store.compact(dropFromOrderIdx = 2, conversationId = "c1", summaryText = "summary") + assertEquals(2L, newMax) + val remaining = store.list("c1") + assertEquals(2, remaining.size) + assertEquals(1L, remaining[0].orderIdx) + assertEquals(2L, remaining[1].orderIdx) + assertTrue(remaining[1].entry is WorkingMemoryEntry.Summary) + } + + @Test + fun testCompactWithoutSummaryKeepsTailBelow() = runTest { + val t = Instant.parse("2026-09-15T10:00:00Z") + store.append("c1", userMsg("a"), t) + // dropFromOrderIdx=2: удаляет idx >= 2, остаётся idx=1. + val newMax = store.compact(dropFromOrderIdx = 2, conversationId = "c1", summaryText = null) + assertEquals(1L, newMax) + val remaining = store.list("c1") + assertEquals(1, remaining.size) + assertEquals(1L, remaining[0].orderIdx) + } +} diff --git a/settings.gradle.kts b/settings.gradle.kts index 55cc669..1f49731 100644 --- a/settings.gradle.kts +++ b/settings.gradle.kts @@ -95,6 +95,11 @@ include(":storage-inmemory") // остальные store'ы мигрируют после стабилизации ksqlite и реального // использования на Android-агенте. include(":storage-ksqlite") +// ksqlite-реализация :context-api (ContextStore / working_memory table). +// Минимальный модуль: только таблица `working_memory` + 2 индекса. +// Параллельно существует :storage-ksqlite/KsqliteWorkingMemoryStore.kt — +// миграция consumers'ов по чуть-чуть, отдельно. +include(":context-ksqlite") // KMP-реализация EventStore через ksqlite (https://github.com/caffeine-mgn/ksqlite). // Цель: проверить что pure-Kotlin SQLite с sqlite-vec заменяет SQLDelight+JVector // на KMP-таргетах (JVM + linuxX64 + mingwX64). Apple targets auto-disabled на