diff --git a/journal-ksqlite/build.gradle.kts b/journal-ksqlite/build.gradle.kts new file mode 100644 index 0000000..8d10f09 --- /dev/null +++ b/journal-ksqlite/build.gradle.kts @@ -0,0 +1,33 @@ +plugins { + alias(libs.plugins.kotlin.multiplatform) + alias(libs.plugins.kotlin.serialization) +} + +// KMP-реализация :journal-api (JournalStore / MutableJournalStore) поверх ksqlite. +// Минимальная — только таблица `message` для append-only audit log'а. +// ConversationStore / ReflectionStore / WorkingMemoryStore живут в своих +// собственных 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(":journal-api")) + } + commonTest.dependencies { + implementation(kotlin("test")) + implementation(libs.kotlinx.coroutines.test) + } + } +} diff --git a/journal-ksqlite/src/commonMain/kotlin/pw/binom/agentik/journal/ksqlite/KsqliteJournalStore.kt b/journal-ksqlite/src/commonMain/kotlin/pw/binom/agentik/journal/ksqlite/KsqliteJournalStore.kt new file mode 100644 index 0000000..db4d957 --- /dev/null +++ b/journal-ksqlite/src/commonMain/kotlin/pw/binom/agentik/journal/ksqlite/KsqliteJournalStore.kt @@ -0,0 +1,117 @@ +package pw.binom.agentik.journal.ksqlite + +import kotlinx.serialization.json.Json +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 +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.sync.Mutex +import kotlinx.coroutines.sync.withLock +import kotlinx.coroutines.withContext + +/** + * ksqlite-реализация [MutableJournalStore] (append-only audit log). + * + * Структура — копия [pw.binom.agentik.storage.ksqlite.KsqliteMessageStore] + * из `:storage-ksqlite`, но: + * - лежит в собственном модуле `:journal-ksqlite`; + * - реализует переименованный [MutableJournalStore] (раньше был + * `MutableMessageStore`, теперь главный класс — `JournalStore` / + * `MutableJournalStore`); сам тип записи [MessageRecord] не + * переименовывался. + * + * Prepared statements (insert / list / clear) препарируются один раз в + * конструкторе и закрываются в [close]. Без этого GC финалайзеры каждого + * StmtHolder'а пытаются `sqlite3_finalize` stmt, чей parent connection уже + * закрыт → SIGSEGV в `pthread_mutex_lock` (см. [pw.binom.db.ksqlite.StmtHolder]). + * + * `payloadJson` хранит JSON-сериализованные kind-specific поля. encoding + * helpers (`encodeRecord` / `toMessageRecord` / `CallPayload` / ...) лежат + * в [MessageCodecs.kt] рядом. + * + * ВНИМАНИЕ: `:storage-ksqlite/KsqliteMessageStore.kt` остаётся на диске — + * это копия, не замена. Не удалять старый файл; миграция consumers'ов — + * отдельно. + */ +class KsqliteJournalStore internal constructor( + private val connection: SQLiteConnection, +) : MutableJournalStore { + + private val mutex = Mutex() + private val json = Json { ignoreUnknownKeys = true } + + private val insertStmt: SQLitePreparedStatement = connection.prepare( + """ + INSERT INTO ${Schema.TABLE_MESSAGE} + (${Schema.COL_ID}, ${Schema.COL_CONVERSATION_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_KIND}, + ${Schema.COL_PAYLOAD_JSON}, ${Schema.COL_CREATED_AT} + FROM ${Schema.TABLE_MESSAGE} + WHERE ${Schema.COL_CONVERSATION_ID} = ? + AND ${Schema.COL_CREATED_AT} > ? + ORDER BY ${Schema.COL_CREATED_AT} ASC, ${Schema.COL_ID} ASC + LIMIT ? OFFSET ? + """.trimIndent() + ) + private val clearStmt: SQLitePreparedStatement = connection.prepare( + "DELETE FROM ${Schema.TABLE_MESSAGE} WHERE ${Schema.COL_CONVERSATION_ID} = ?" + ) + + override suspend fun append(record: MessageRecord): Unit = withContext(Dispatchers.Default) { + val (kind, payload) = encodeRecord(record) + mutex.withLock { + insertStmt.reset() + insertStmt.clearBindings() + insertStmt.bindText(1, record.id) + insertStmt.bindText(2, record.conversationId) + insertStmt.bindText(3, kind) + insertStmt.bindText(4, payload) + insertStmt.bindLong(5, record.createdAt.toEpochMilliseconds()) + insertStmt.executeUpdate() + } + } + + override suspend fun list( + conversationId: String, + after: Instant, + offset: Int, + limit: Int, + ): List = withContext(Dispatchers.Default) { + mutex.withLock { + listStmt.reset() + listStmt.clearBindings() + listStmt.bindText(1, conversationId) + listStmt.bindLong(2, after.toEpochMilliseconds()) + listStmt.bindLong(3, limit.toLong()) + listStmt.bindLong(4, offset.toLong()) + val out = mutableListOf() + listStmt.executeQuery().use { rs -> + while (rs.next()) out.add(rs.toMessageRecord(json)) + } + out + } + } + + internal suspend fun clear(conversationId: String): Unit = withContext(Dispatchers.Default) { + mutex.withLock { + clearStmt.reset() + clearStmt.clearBindings() + clearStmt.bindText(1, conversationId) + clearStmt.executeUpdate() + } + } + + override fun close() { + insertStmt.close() + listStmt.close() + clearStmt.close() + } +} diff --git a/journal-ksqlite/src/commonMain/kotlin/pw/binom/agentik/journal/ksqlite/MessageCodecs.kt b/journal-ksqlite/src/commonMain/kotlin/pw/binom/agentik/journal/ksqlite/MessageCodecs.kt new file mode 100644 index 0000000..babc296 --- /dev/null +++ b/journal-ksqlite/src/commonMain/kotlin/pw/binom/agentik/journal/ksqlite/MessageCodecs.kt @@ -0,0 +1,79 @@ +package pw.binom.agentik.journal.ksqlite + +import kotlinx.serialization.json.Json +import pw.binom.agentik.journal.MessageRecord +import pw.binom.agentik.journal.decodeBodyPayload +import pw.binom.agentik.journal.encodeBodyPayload +import pw.binom.db.ksqlite.SQLiteResultSet +import kotlin.time.Instant + +/** + * Кодирование [MessageRecord] → пара (kind, payloadJson) для SQLite. + * + * Копия `MessageCodecs.kt` из `:storage-ksqlite` — `internal` 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?) diff --git a/journal-ksqlite/src/commonMain/kotlin/pw/binom/agentik/journal/ksqlite/Schema.kt b/journal-ksqlite/src/commonMain/kotlin/pw/binom/agentik/journal/ksqlite/Schema.kt new file mode 100644 index 0000000..4b143d6 --- /dev/null +++ b/journal-ksqlite/src/commonMain/kotlin/pw/binom/agentik/journal/ksqlite/Schema.kt @@ -0,0 +1,91 @@ +package pw.binom.agentik.journal.ksqlite + +import pw.binom.db.ksqlite.SQLiteConnection + +/** + * Имена таблиц/колонок/индексов для ksqlite-бэкенда `:journal-api`. + * + * Минимум — только то, что относится к `message` (append-only audit log). + * Остальные таблицы агента (`conversation`, `working_memory`, `reflection`) + * живут в других ksqlite-модулях. + * + * Все DDL/DML в этом модуле должны ссылаться на эти константы — никаких + * хардкоженных литералов в `prepare("SELECT ... FROM foo ...")` в store'е. + */ +internal object Schema { + + /** Версия схемы модуля. Увеличивать при ЛЮБОМ изменении DDL. */ + const val CURRENT_VERSION: Int = 1 + + // ───── Таблица ───── + const val TABLE_MESSAGE = "message" + + // ───── Колонки ───── + const val COL_ID = "id" + const val COL_CONVERSATION_ID = "conversation_id" + const val COL_KIND = "kind" + const val COL_PAYLOAD_JSON = "payload_json" + const val COL_CREATED_AT = "created_at" + + // ───── Индексы ───── + const val IDX_MSG_CONV = "idx_msg_conv" + + private val v1Ddl = """ + CREATE TABLE IF NOT EXISTS $TABLE_MESSAGE ( + $COL_ID TEXT NOT NULL PRIMARY KEY, + $COL_CONVERSATION_ID TEXT NOT NULL, + $COL_KIND TEXT NOT NULL, + $COL_PAYLOAD_JSON TEXT NOT NULL, + $COL_CREATED_AT INTEGER NOT NULL + ); + """.trimIndent() + + private val v1IndexesDdl = """ + CREATE INDEX IF NOT EXISTS $IDX_MSG_CONV + ON $TABLE_MESSAGE($COL_CONVERSATION_ID, $COL_CREATED_AT); + """.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/journal-ksqlite/src/commonTest/kotlin/pw/binom/agentik/journal/ksqlite/KsqliteJournalStoreTest.kt b/journal-ksqlite/src/commonTest/kotlin/pw/binom/agentik/journal/ksqlite/KsqliteJournalStoreTest.kt new file mode 100644 index 0000000..583d0be --- /dev/null +++ b/journal-ksqlite/src/commonTest/kotlin/pw/binom/agentik/journal/ksqlite/KsqliteJournalStoreTest.kt @@ -0,0 +1,106 @@ +package pw.binom.agentik.journal.ksqlite + +import kotlinx.coroutines.flow.toList +import kotlinx.coroutines.test.runTest +import pw.binom.agentik.journal.Content +import pw.binom.agentik.journal.MessageRecord +import pw.binom.agentik.journal.TurnTokens +import pw.binom.db.ksqlite.SQLiteConnection +import kotlin.test.AfterTest +import kotlin.test.BeforeTest +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.time.Instant + +/** + * Тесты для [KsqliteJournalStore] — точная копия + * `KsqliteMessageStoreTest` из `:storage-ksqlite`, с переименованием типов + * (`MessageStore` → `JournalStore`) и автономной фикстурой (in-memory + * SQLiteConnection + Schema.migrate). + * + * Тест `testClearRemovesByConversation` из оригинала использовал + * `stores.conversations.delete(...)` (cascade через `KsqliteStores`) — здесь + * он заменён на прямой вызов `store.clear(...)`, потому что `:journal-ksqlite` + * автономен и не знает про ConversationStore. + */ +class KsqliteJournalStoreTest { + + private lateinit var conn: SQLiteConnection + private lateinit var store: KsqliteJournalStore + + @BeforeTest + fun setup() { + conn = SQLiteConnection.memory("journal-${kotlin.random.Random.nextLong()}") + Schema.migrate(conn) + store = KsqliteJournalStore(conn) + } + + @AfterTest + fun tearDown() { + store.close() + 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 = TurnTokens(input = 50, output = 30), + )) + val list = store.listFlow("conv1", Instant.DISTANT_PAST).toList() + assertEquals(1, list.size) + val msg = list[0] as MessageRecord.AssistantMessage + assertEquals(TurnTokens(input = 50, output = 30), msg.tokens) + } + + @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 testListFlowReturnsAllInOrder() = 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 { + store.append(MessageRecord.UserMessage("m1", "conv1", listOf(Content.Text("a")), Instant.parse("2026-09-15T10:00:00Z"), null)) + store.append(MessageRecord.UserMessage("m2", "conv2", listOf(Content.Text("b")), Instant.parse("2026-09-15T10:00:00Z"), null)) + store.clear("conv1") + assertEquals(emptyList(), store.listFlow("conv1", Instant.DISTANT_PAST).toList()) + assertEquals(1, store.listFlow("conv2", Instant.DISTANT_PAST).toList().size) + } +} diff --git a/settings.gradle.kts b/settings.gradle.kts index 1f49731..afa6354 100644 --- a/settings.gradle.kts +++ b/settings.gradle.kts @@ -100,6 +100,12 @@ include(":storage-ksqlite") // Параллельно существует :storage-ksqlite/KsqliteWorkingMemoryStore.kt — // миграция consumers'ов по чуть-чуть, отдельно. include(":context-ksqlite") +// ksqlite-реализация :journal-api (JournalStore / message table). +// Минимальный модуль: только таблица `message` + 1 индекс +// `(conversation_id, created_at)`. Параллельно существует +// :storage-ksqlite/KsqliteMessageStore.kt — миграция consumers'ов +// по чуть-чуть, отдельно. +include(":journal-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 на