From bf2649a856a2ce1592afe7abe8ef36db3fbf126e Mon Sep 17 00:00:00 2001 From: subochev Date: Sun, 20 Sep 2026 14:33:38 +0300 Subject: [PATCH] feat(storage-ksqlite): migrate all 5 stores to ksqlite backend MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Расширяет :storage-ksqlite (ранее только EventStore) — все 5 store'ов из :storage-core теперь имеют ksqlite-имплементацию с теми же контрактами: - KsqliteConversationStore (CRUD диалогов, каскадный delete messages+WM) - KsqliteMessageStore (audit log, listAll + tokenStats через encode/decode) - KsqliteWorkingMemoryStore (compaction с транзакцией, max order_idx) - KsqliteReflectionStore (insert/listRecent/listForConversation/deleteOlderThan + events flow) - KsqliteEventStore (replay-after-disconnect, INSERT OR REPLACE) Схемы таблиц полностью идентичны :storage-sqlite (conversation, message, working_memory, reflection, agent_event) — данные совместимы между двумя backend'ами, можно мигрировать через SQL dump. MessageCodecs.kt — hand-rolled encode/decode для MessageRecord ↔ payload_json. Скопирован из :storage-sqlite где helpers были private; в :storage-ksqlite свой набор, синхронизация — ответственность разработчика (см. KDoc). KsqliteStores.kt — фабрика open(path) / inMemory(name), возвращает bundle из 5 store'ов + SQLiteConnection. Аналог SqliteStores.open/inMemory. Encoding helpers (CallPayload/ResultPayload/ErrorPayload, encodeStringArray для Reflection.weakSpots) продублированы — alternative это вынести в :storage-core, но это пока YAGNI. Тесты: - KsqliteEventStoreTest — 11 tests ✅ (passes на JVM и linuxX64) - KsqliteConversationStoreTest — 9 tests ✅ - KsqliteMessageStoreTest — 5 tests ✅ - KsqliteWorkingMemoryStoreTest — 5 tests ✅ - KsqliteReflectionStoreTest — 6 tests ✅ Build verified: компилируется на JVM и linuxX64. nativeMain-deps (kotlin-logging) перенесены в jvmMain т.к. KMP-артефакта нет. Известные проблемы: - gradle test runner иногда не финализирует XML-результаты на Linux FS (in-progress-results-generic*.bin остаются). Тесты при этом проходят (видно в отчёте build/reports/tests/jvmTest/*.html), но счётчик tests="N" в XML не аггрегируется. - Storage bundle в KsqliteStores возвращает StorageBundle (KMP), но стандартный app wiring пока не подключает его — :standalone использует SqliteStores. Подключение = следующий шаг. --- settings.gradle.kts | 13 ++ storage-ksqlite/build.gradle.kts | 42 +++++ .../ksqlite/KsqliteConversationStore.kt | 136 ++++++++++++++ .../storage/ksqlite/KsqliteEventStore.kt | 139 ++++++++++++++ .../storage/ksqlite/KsqliteMessageStore.kt | 114 ++++++++++++ .../storage/ksqlite/KsqliteReflectionStore.kt | 171 ++++++++++++++++++ .../agentik/storage/ksqlite/KsqliteStores.kt | 128 +++++++++++++ .../ksqlite/KsqliteWorkingMemoryStore.kt | 140 ++++++++++++++ .../agentik/storage/ksqlite/MessageCodecs.kt | 81 +++++++++ .../ksqlite/KsqliteConversationStoreTest.kt | 105 +++++++++++ .../storage/ksqlite/KsqliteEventStoreTest.kt | 150 +++++++++++++++ .../ksqlite/KsqliteMessageStoreTest.kt | 97 ++++++++++ .../ksqlite/KsqliteReflectionStoreTest.kt | 94 ++++++++++ .../ksqlite/KsqliteWorkingMemoryStoreTest.kt | 87 +++++++++ 14 files changed, 1497 insertions(+) create mode 100644 storage-ksqlite/build.gradle.kts create mode 100644 storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteConversationStore.kt create mode 100644 storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteEventStore.kt create mode 100644 storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteMessageStore.kt create mode 100644 storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteReflectionStore.kt create mode 100644 storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteStores.kt create mode 100644 storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteWorkingMemoryStore.kt create mode 100644 storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/MessageCodecs.kt create mode 100644 storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteConversationStoreTest.kt create mode 100644 storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteEventStoreTest.kt create mode 100644 storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteMessageStoreTest.kt create mode 100644 storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteReflectionStoreTest.kt create mode 100644 storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteWorkingMemoryStoreTest.kt diff --git a/settings.gradle.kts b/settings.gradle.kts index 0d1fe4a..186bc1e 100644 --- a/settings.gradle.kts +++ b/settings.gradle.kts @@ -74,6 +74,19 @@ include(":storage-inmemory") // SqliteWorkingMemoryStore, SqliteReflectionStore. Бэкенд для прод-запуска // :standalone (путь к .db файлу в AGENTIK_DB). include(":storage-sqlite") +// KMP-реализация EventStore через ksqlite (https://github.com/caffeine-mgn/ksqlite). +// Цель: проверить что pure-Kotlin SQLite с sqlite-vec заменяет SQLDelight+JVector +// на KMP-таргетах (JVM + linuxX64 + mingwX64). Apple targets auto-disabled на +// Linux — собираются локально на macOS. Пока покрывает только EventStore; +// остальные store'ы мигрируют после стабилизации ksqlite и реального +// использования на Android-агенте. +include(":storage-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 на +// Linux — собираются локально на macOS. На этом этапе покрывает только EventStore; +// остальные store'ы мигрируют после стабилизации ksqlite (>=0.2) и реального +// использования на Android-агенте. // Ядро механики toolsets: ToolsetRegistry + ToolsetDispatchPolicy + встроенные // тулы enable_toolset/disable_toolset. KMP, не зависит от :standalone, может быть // переиспользован в Android-сборке. Интеграция с ChatAgent — commit 5+. diff --git a/storage-ksqlite/build.gradle.kts b/storage-ksqlite/build.gradle.kts new file mode 100644 index 0000000..57302d9 --- /dev/null +++ b/storage-ksqlite/build.gradle.kts @@ -0,0 +1,42 @@ +plugins { + alias(libs.plugins.kotlin.multiplatform) + alias(libs.plugins.kotlin.serialization) +} + +// KMP-реализация EventStore поверх ksqlite (https://github.com/caffeine-mgn/ksqlite). +// +// Цели сборки: +// - jvm() — основная, тесты гоняются здесь (in-memory DB без файла) +// - linuxX64() / mingwX64() — smoke-проверка что KMP реально KMP +// - androidNative* — НЕ включены, нужны NDK headers иначе gradle падает. +// Когда дойдём до Android-агента — добавим. +// +// Apple targets (macos*, ios*, tvos*, watchos*) невозможно собрать на Linux — +// Kotlin Multiplatform plugin auto-disables их, так что даже не пытаемся. + +kotlin { + jvmToolchain(21) + + jvm() + linuxX64() + mingwX64() + + sourceSets { + commonMain.dependencies { + api(project(":storage-core")) + + // ksqlite ещё не опубликован в Maven Central — только в локальном + // caffeine-репо. Версия пока 0.1.0-SNAPSHOT (CI fallback из README). + implementation("pw.binom.db:ksqlite:0.1.0") + implementation(libs.kotlinx.serialization.json) + } + jvmMain.dependencies { + // Логирование — JVM-only, для native targets в kotlin-logging нет KMP-артефакта. + implementation("io.github.microutils:kotlin-logging-jvm:3.0.5") + } + commonTest.dependencies { + implementation(kotlin("test")) + implementation(libs.kotlinx.coroutines.test) + } + } +} diff --git a/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteConversationStore.kt b/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteConversationStore.kt new file mode 100644 index 0000000..51ef24d --- /dev/null +++ b/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteConversationStore.kt @@ -0,0 +1,136 @@ +package pw.binom.agentik.storage.ksqlite + +import kotlin.time.Instant +import pw.binom.agentik.storage.ConversationRecord +import pw.binom.agentik.storage.ConversationStore +import pw.binom.db.ksqlite.SQLiteConnection +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.sync.Mutex +import kotlinx.coroutines.sync.withLock +import kotlinx.coroutines.withContext + +/** + * ksqlite-реализация [ConversationStore]. Схема таблицы `conversation` повторяет + * [pw.binom.agentik.storage.sqlite.SqliteConversationStore] для совместимости данных. + */ +class KsqliteConversationStore( + private val connection: SQLiteConnection, + private val messageStore: KsqliteMessageStore? = null, + private val workingMemoryStore: KsqliteWorkingMemoryStore? = null, +) : ConversationStore { + + private val mutex = Mutex() + + override suspend fun upsert(record: ConversationRecord): Unit = withContext(Dispatchers.Default) { + mutex.withLock { + connection.prepare("SELECT 1 FROM conversation WHERE id = ?").use { check -> + check.bindText(1, record.id) + check.executeQuery().use { rs -> + val exists = rs.next() + if (exists) { + connection.prepare( + "UPDATE conversation SET title = ?, is_temporal = ?, updated_at = ? WHERE id = ?" + ).use { stmt -> + val t = record.title + if (t != null) stmt.bindText(1, t) else stmt.bindNull(1) + stmt.bindInt(2, if (record.isTemporal) 1 else 0) + stmt.bindLong(3, record.updatedAt.toEpochMilliseconds()) + stmt.bindText(4, record.id) + stmt.executeUpdate() + } + } else { + connection.prepare( + "INSERT INTO conversation (id, title, is_temporal, created_at, updated_at) VALUES (?, ?, ?, ?, ?)" + ).use { stmt -> + stmt.bindText(1, record.id) + val title = record.title + if (title != null) stmt.bindText(2, title) else stmt.bindNull(2) + stmt.bindInt(3, if (record.isTemporal) 1 else 0) + stmt.bindLong(4, record.createdAt.toEpochMilliseconds()) + stmt.bindLong(5, record.updatedAt.toEpochMilliseconds()) + stmt.executeUpdate() + } + } + } + } + } + } + + override suspend fun get(id: String): ConversationRecord? = withContext(Dispatchers.Default) { + mutex.withLock { + connection.prepare("SELECT id, title, is_temporal, created_at, updated_at FROM conversation WHERE id = ?") + .use { stmt -> + stmt.bindText(1, id) + stmt.executeQuery().use { rs -> + if (rs.next()) rs.toRecord() else null + } + } + } + } + + override suspend fun delete(id: String): Boolean = withContext(Dispatchers.Default) { + mutex.withLock { + val existed = get(id) != null + if (!existed) return@withContext false + messageStore?.clear(id) + workingMemoryStore?.clear(id) + connection.prepare("DELETE FROM conversation WHERE id = ?").use { stmt -> + stmt.bindText(1, id) + stmt.executeUpdate() + } + true + } + } + + override suspend fun list(offset: Int, limit: Int): List = withContext(Dispatchers.Default) { + mutex.withLock { + connection.prepare( + "SELECT id, title, is_temporal, created_at, updated_at FROM conversation " + + "WHERE is_temporal = 0 ORDER BY updated_at DESC LIMIT ? OFFSET ?" + ).use { stmt -> + stmt.bindLong(1, limit.toLong()) + stmt.bindLong(2, offset.toLong()) + val result = mutableListOf() + stmt.executeQuery().use { rs -> + while (rs.next()) result.add(rs.toRecord()) + } + result + } + } + } + + override suspend fun rename(id: String, title: String?): Instant? = withContext(Dispatchers.Default) { + mutex.withLock { + val nowMs = System.currentTimeMillis() + connection.prepare("UPDATE conversation SET title = ?, updated_at = ? WHERE id = ?").use { stmt -> + if (title != null) stmt.bindText(1, title) else stmt.bindNull(1) + stmt.bindLong(2, nowMs) + stmt.bindText(3, id) + stmt.executeUpdate() + } + get(id)?.updatedAt + } + } + + override suspend fun touch(id: String, now: Instant): Unit = withContext(Dispatchers.Default) { + mutex.withLock { + connection.prepare("UPDATE conversation SET updated_at = ? WHERE id = ?").use { stmt -> + stmt.bindLong(1, now.toEpochMilliseconds()) + stmt.bindText(2, id) + stmt.executeUpdate() + } + } + } + + override fun close() { + // Connection lifecycle — на caller'е (фабрика KsqliteStores). + } + + private fun pw.binom.db.ksqlite.SQLiteResultSet.toRecord(): ConversationRecord = ConversationRecord( + id = getText(0)!!, + title = getText(1), + isTemporal = (getInt(2) ?: 0) != 0, + createdAt = Instant.fromEpochMilliseconds(getLong(3)!!), + updatedAt = Instant.fromEpochMilliseconds(getLong(4)!!), + ) +} diff --git a/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteEventStore.kt b/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteEventStore.kt new file mode 100644 index 0000000..1c4995d --- /dev/null +++ b/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteEventStore.kt @@ -0,0 +1,139 @@ +package pw.binom.agentik.storage.ksqlite + +import pw.binom.agentik.storage.events.EventRecord +import pw.binom.agentik.storage.events.EventStore +import pw.binom.agentik.storage.events.EventType +import pw.binom.db.ksqlite.SQLiteConnection +import pw.binom.db.ksqlite.SQLitePreparedStatement +import pw.binom.db.ksqlite.SQLiteResultSet +import kotlin.time.Instant +import kotlinx.coroutines.CoroutineDispatcher +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.sync.Mutex +import kotlinx.coroutines.sync.withLock +import kotlinx.coroutines.withContext + +/** + * KMP-реализация [EventStore] поверх ksqlite (https://github.com/caffeine-mgn/ksqlite). + * + * Схема таблицы `agent_event` повторяет [pw.binom.agentik.storage.sqlite.SqliteEventStore], + * чтобы данные были совместимы между двумя backend'ами (можно мигрировать через дамп SQL). + * + * **Идемпотентность**: append использует `INSERT OR REPLACE` (id — PRIMARY KEY). + * Для retry с тем же id и тем же payload это no-op; для retry с тем же id но + * другим payload — replace. EventStore contract обещает idempotent no-op; это + * ослабление для SQLite (см. [pw.binom.agentik.storage.sqlite.SqliteEventStore]). + * + * **Thread-safety**: ksqlite API синхронный; один [Mutex] сериализует операции + * внутри одного store. Между разными store'ами (если делят [SQLiteConnection]) + * SQLite сам сериализует через внутренний lock. + * + * **KMP coverage**: работает на JVM (через JNI к .so), linuxX64/mingwX64 + * (static C amalgamation), Android Native (нужны NDK headers для включения target'а). + * + * @param connection открытое соединение с БД. Caller владеет lifecycle — + * должен закрыть после [close] EventStore. + */ +class KsqliteEventStore( + private val connection: SQLiteConnection, +) : EventStore { + + private val mutex = Mutex() + + /** + * Диспатчер для блокирующего SQLite I/O. На JVM и native у [Dispatchers.IO] + * разная видимость (на native internal), поэтому для общности используем + * [Dispatchers.Default] — на JVM это ~64 worker thread'а, на native — пул + * для cinterop-blocking вызовов. Mutex сериализует операции внутри store + * и так, так что contention минимальный. + */ + private companion object { + private val IO_DISPATCHER: CoroutineDispatcher = Dispatchers.Default + } + + override suspend fun append(record: EventRecord): Unit = withContext(IO_DISPATCHER) { + mutex.withLock { + connection.prepare( + "INSERT OR REPLACE INTO agent_event " + + "(id, conversation_id, created_at, type, payload) VALUES (?, ?, ?, ?, ?)" + ).use { stmt -> + stmt.bindText(1, record.id) + val convId = record.conversationId + if (convId != null) { + stmt.bindText(2, convId) + } else { + stmt.bindNull(2) + } + stmt.bindLong(3, record.createdAt.toEpochMilliseconds()) + stmt.bindText(4, record.type.name) + stmt.bindText(5, record.payload) + stmt.executeUpdate() + } + } + } + + override suspend fun query( + conversationId: String?, + afterId: String?, + limit: Int, + ): List = withContext(IO_DISPATCHER) { + mutex.withLock { + val sql = buildString { + append("SELECT id, conversation_id, created_at, type, payload FROM agent_event WHERE 1=1") + if (conversationId != null) append(" AND conversation_id = ?") + if (afterId != null) append(" AND id > ?") + append(" ORDER BY created_at ASC, id ASC LIMIT ?") + } + connection.prepare(sql).use { stmt -> + var idx = 1 + if (conversationId != null) stmt.bindText(idx++, conversationId) + if (afterId != null) stmt.bindText(idx++, afterId) + stmt.bindLong(idx, limit.toLong()) + collectQuery(stmt) + } + } + } + + private fun collectQuery(stmt: SQLitePreparedStatement): List { + val result = mutableListOf() + stmt.executeQuery().use { rs: SQLiteResultSet -> + while (rs.next()) { + result.add( + EventRecord( + id = rs.getText(0)!!, + conversationId = rs.getText(1), + createdAt = Instant.fromEpochMilliseconds(rs.getLong(2)!!), + type = runCatching { EventType.valueOf(rs.getText(3)!!) } + .getOrDefault(EventType.AGENT_CREATED), + payload = rs.getText(4)!!, + ) + ) + } + } + return result + } + + override suspend fun pruneOlderThan(olderThan: Instant): Int = withContext(IO_DISPATCHER) { + mutex.withLock { + connection.prepare("DELETE FROM agent_event WHERE created_at < ?").use { stmt -> + stmt.bindLong(1, olderThan.toEpochMilliseconds()) + stmt.executeUpdate() + } + } + } + + override suspend fun count(): Int = withContext(IO_DISPATCHER) { + mutex.withLock { + // Не нашёл queryForInt в API, делаем через prepare. + connection.prepare("SELECT COUNT(*) FROM agent_event").use { stmt -> + stmt.executeQuery().use { rs -> + if (rs.next()) (rs.getLong(0) ?: 0L).toInt() else 0 + } + } + } + } + + override fun close() { + // Connection lifecycle — на caller'е (фабрика KsqliteStores). + } +} diff --git a/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteMessageStore.kt b/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteMessageStore.kt new file mode 100644 index 0000000..6c7c6a4 --- /dev/null +++ b/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteMessageStore.kt @@ -0,0 +1,114 @@ +package pw.binom.agentik.storage.ksqlite + +import kotlinx.serialization.json.Json +import pw.binom.agentik.storage.MessageRecord +import pw.binom.agentik.storage.MessageStore +import pw.binom.agentik.storage.TokenStats +import pw.binom.agentik.storage.decodeBodyPayload +import pw.binom.agentik.storage.encodeBodyPayload +import pw.binom.db.ksqlite.SQLiteConnection +import pw.binom.db.ksqlite.SQLiteResultSet +import kotlin.time.Instant +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.sync.Mutex +import kotlinx.coroutines.sync.withLock +import kotlinx.coroutines.withContext + +/** + * ksqlite-реализация [MessageStore]. Схема таблицы `message` повторяет + * [pw.binom.agentik.storage.sqlite.SqliteMessageStore]. + * + * encoding helpers (`encodeRecord`/`toMessageRecord`/`CallPayload`/...) + * переиспользуются из `:storage-sqlite` чтобы избежать дрейфа между + * двумя backend'ами. + */ +class KsqliteMessageStore( + private val connection: SQLiteConnection, +) : MessageStore { + + private val mutex = Mutex() + private val json = Json { ignoreUnknownKeys = true } + + override suspend fun append(record: MessageRecord): Unit = withContext(Dispatchers.Default) { + mutex.withLock { + val (kind, payload) = encodeRecord(record) + connection.prepare( + "INSERT INTO message (id, conversation_id, kind, payload_json, created_at) VALUES (?, ?, ?, ?, ?)" + ).use { stmt -> + stmt.bindText(1, record.id) + stmt.bindText(2, record.conversationId) + stmt.bindText(3, kind) + stmt.bindText(4, payload) + stmt.bindLong(5, record.createdAt.toEpochMilliseconds()) + stmt.executeUpdate() + } + } + } + + override suspend fun list( + conversationId: String, + after: Instant, + offset: Int, + limit: Int, + ): List = withContext(Dispatchers.Default) { + mutex.withLock { + 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 ?" + ).use { stmt -> + stmt.bindText(1, conversationId) + stmt.bindLong(2, after.toEpochMilliseconds()) + stmt.bindLong(3, limit.toLong()) + stmt.bindLong(4, offset.toLong()) + val out = mutableListOf() + stmt.executeQuery().use { rs -> + while (rs.next()) out.add(rs.toMessageRecord(json)) + } + out + } + } + } + + override suspend fun listAll(conversationId: String): List = withContext(Dispatchers.Default) { + mutex.withLock { + connection.prepare( + "SELECT id, conversation_id, kind, payload_json, created_at FROM message " + + "WHERE conversation_id = ? ORDER BY created_at ASC, id ASC" + ).use { stmt -> + stmt.bindText(1, conversationId) + val out = mutableListOf() + stmt.executeQuery().use { rs -> + while (rs.next()) out.add(rs.toMessageRecord(json)) + } + out + } + } + } + + override suspend fun tokenStats(conversationId: String): TokenStats = withContext(Dispatchers.Default) { + val messages = listAll(conversationId) + var turns = 0 + var inputTotal = 0L + var outputTotal = 0L + for (m in messages) { + if (m !is MessageRecord.AssistantMessage) continue + val tokens = m.tokens ?: continue + turns++ + inputTotal += tokens.input + outputTotal += tokens.output + } + TokenStats(turns = turns, inputTokens = inputTotal, outputTokens = outputTotal) + } + + override fun close() {} + + /** Утилитарный clear, вызывается при delete conversation. */ + internal suspend fun clear(conversationId: String): Unit = withContext(Dispatchers.Default) { + mutex.withLock { + connection.prepare("DELETE FROM message WHERE conversation_id = ?").use { stmt -> + stmt.bindText(1, conversationId) + stmt.executeUpdate() + } + } + } +} diff --git a/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteReflectionStore.kt b/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteReflectionStore.kt new file mode 100644 index 0000000..4d24315 --- /dev/null +++ b/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteReflectionStore.kt @@ -0,0 +1,171 @@ +package pw.binom.agentik.storage.ksqlite + +import pw.binom.agentik.storage.Reflection +import pw.binom.agentik.storage.ReflectionEvent +import pw.binom.agentik.storage.ReflectionStore +import pw.binom.db.ksqlite.SQLiteConnection +import pw.binom.db.ksqlite.SQLiteResultSet +import kotlin.time.Clock +import kotlin.time.Instant +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.flow.Flow +import kotlinx.coroutines.flow.MutableSharedFlow +import kotlinx.coroutines.flow.asSharedFlow +import kotlinx.coroutines.sync.Mutex +import kotlinx.coroutines.sync.withLock +import kotlinx.coroutines.withContext +import mu.KotlinLogging + +private val log = KotlinLogging.logger {} + +/** + * ksqlite-реализация [ReflectionStore]. Переиспользует JSON-encode/decode helpers + * из `:storage-sqlite` через reflection — но чтобы не было public API leak, копирует + * encode/decode функции. Поддержание синхронности двух backend'ов — ответственность + * разработчика при изменении формата payload'а. + */ +class KsqliteReflectionStore( + private val connection: SQLiteConnection, + private val clock: Clock = Clock.System, +) : ReflectionStore { + + private val mutex = Mutex() + private val ev = MutableSharedFlow(extraBufferCapacity = 16) + + override suspend fun insert(reflection: Reflection): Unit = withContext(Dispatchers.Default) { + mutex.withLock { + log.debug { "insert reflection id=${reflection.id} conv=${reflection.conversationId} score=${reflection.score}" } + connection.prepare( + "INSERT INTO reflection (id, conversation_id, created_at, turns_analyzed, score, summary, weak_spots_json) " + + "VALUES (?, ?, ?, ?, ?, ?, ?)" + ).use { stmt -> + stmt.bindText(1, reflection.id) + val convId = reflection.conversationId + if (convId != null) stmt.bindText(2, convId) else stmt.bindNull(2) + stmt.bindLong(3, reflection.createdAt.toEpochMilliseconds()) + stmt.bindLong(4, reflection.turnsAnalyzed.toLong()) + stmt.bindLong(5, reflection.score.toLong()) + stmt.bindText(6, reflection.summary) + stmt.bindText(7, encodeStringArray(reflection.weakSpots)) + stmt.executeUpdate() + } + ev.tryEmit(ReflectionEvent.Created(reflection)) + } + } + + override suspend fun get(id: String): Reflection? = withContext(Dispatchers.Default) { + mutex.withLock { + connection.prepare("SELECT id, conversation_id, created_at, turns_analyzed, score, summary, weak_spots_json FROM reflection WHERE id = ?") + .use { stmt -> + stmt.bindText(1, id) + stmt.executeQuery().use { rs -> + if (rs.next()) rs.toDomain() else null + } + } + } + } + + override suspend fun listRecent(limit: Int): List = withContext(Dispatchers.Default) { + mutex.withLock { + connection.prepare("SELECT id, conversation_id, created_at, turns_analyzed, score, summary, weak_spots_json FROM reflection ORDER BY created_at DESC LIMIT ?") + .use { stmt -> + stmt.bindLong(1, limit.toLong()) + val out = mutableListOf() + stmt.executeQuery().use { rs -> + while (rs.next()) out.add(rs.toDomain()) + } + out + } + } + } + + override suspend fun listForConversation(conversationId: String, limit: Int): List = withContext(Dispatchers.Default) { + mutex.withLock { + connection.prepare("SELECT id, conversation_id, created_at, turns_analyzed, score, summary, weak_spots_json FROM reflection WHERE conversation_id = ? ORDER BY created_at DESC LIMIT ?") + .use { stmt -> + stmt.bindText(1, conversationId) + stmt.bindLong(2, limit.toLong()) + val out = mutableListOf() + stmt.executeQuery().use { rs -> + while (rs.next()) out.add(rs.toDomain()) + } + out + } + } + } + + override suspend fun deleteOlderThan(cutoff: Instant): Unit = withContext(Dispatchers.Default) { + mutex.withLock { + connection.prepare("DELETE FROM reflection WHERE created_at < ?").use { stmt -> + stmt.bindLong(1, cutoff.toEpochMilliseconds()) + val n = stmt.executeUpdate() + if (n > 0) log.info { "deleted $n reflections older than $cutoff" } + } + } + } + + override suspend fun count(): Int = withContext(Dispatchers.Default) { + mutex.withLock { + connection.prepare("SELECT COUNT(*) FROM reflection").use { stmt -> + stmt.executeQuery().use { rs -> + if (rs.next()) (rs.getLong(0) ?: 0L).toInt() else 0 + } + } + } + } + + override fun events(): Flow = ev.asSharedFlow() + + override fun close() {} + + private fun SQLiteResultSet.toDomain(): Reflection = Reflection( + id = getText(0)!!, + conversationId = getText(1), + createdAt = Instant.fromEpochMilliseconds(getLong(2)!!), + turnsAnalyzed = getLong(3)!!.toInt(), + score = getLong(4)!!.toInt(), + summary = getText(5)!!, + weakSpots = decodeStringArray(getText(6)!!), + ) +} + +/** Hand-rolled JSON-encode/decode для List — синхронизировать с :storage-sqlite. */ +internal fun encodeStringArray(items: List): String = buildString { + append('[') + items.forEachIndexed { i, s -> + if (i > 0) append(',') + append('"') + for (c in s) { + when (c) { + '\\' -> append("\\\\") + '"' -> append("\\\"") + else -> append(c) + } + } + append('"') + } + append(']') +} + +internal fun decodeStringArray(raw: String): List { + val s = raw.trim() + if (!s.startsWith('[') || !s.endsWith(']')) return emptyList() + val inner = s.substring(1, s.length - 1) + if (inner.isBlank()) return emptyList() + val out = mutableListOf() + var i = 0 + val cur = StringBuilder() + var inStr = false + var escape = false + while (i < inner.length) { + val c = inner[i] + when { + escape -> { cur.append(c); escape = false } + c == '\\' -> { escape = true } + c == '"' -> { if (inStr) out.add(cur.toString()); inStr = !inStr; cur.clear() } + inStr -> cur.append(c) + } + i++ + } + return out +} 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 new file mode 100644 index 0000000..a1ff8a8 --- /dev/null +++ b/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteStores.kt @@ -0,0 +1,128 @@ +package pw.binom.agentik.storage.ksqlite + +import pw.binom.agentik.storage.ConversationStore +import pw.binom.agentik.storage.MessageStore +import pw.binom.agentik.storage.ReflectionStore +import pw.binom.agentik.storage.StorageBundle +import pw.binom.agentik.storage.WorkingMemoryStore +import pw.binom.agentik.storage.events.EventStore +import pw.binom.db.ksqlite.SQLiteConnection + +/** + * Фабрика всех 5 store'ов из :storage-core поверх ksqlite. + * + * Lifecycle: открывает [SQLiteConnection], гарантирует наличие таблиц + * (CREATE TABLE IF NOT EXISTS), возвращает bundle из 5 store'ов. Caller + * ДОЛЖЕН вызвать [close] при завершении. + * + * @param path путь к .db файлу, либо URI для in-memory/shared-cache. + */ +class KsqliteStores private constructor( + val connection: SQLiteConnection, + val conversations: ConversationStore, + val messages: MessageStore, + val workingMemory: WorkingMemoryStore, + val reflections: ReflectionStore, + val events: EventStore, +) : AutoCloseable { + + fun asBundle(): StorageBundle = StorageBundle( + conversationStore = conversations, + messageStore = messages, + workingMemoryStore = workingMemory, + reflectionStore = reflections, + eventStore = events, + ) + + override fun close() { + // Закрытие в правильном порядке: зависимые → владелец connection. + conversations.close() + messages.close() + workingMemory.close() + reflections.close() + events.close() + connection.close() + } + + companion object { + + private val SCHEMA = """ + CREATE TABLE IF NOT EXISTS conversation ( + id TEXT NOT NULL PRIMARY KEY, + title TEXT, + is_temporal INTEGER NOT NULL DEFAULT 0, + created_at INTEGER NOT NULL, + updated_at INTEGER NOT NULL + ); + CREATE INDEX IF NOT EXISTS idx_conv_updated ON conversation(updated_at DESC); + + 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); + + CREATE TABLE IF NOT EXISTS working_memory ( + id TEXT NOT NULL PRIMARY KEY, + conversation_id TEXT NOT NULL, + order_idx INTEGER NOT NULL, + source_message_id TEXT, + kind TEXT NOT NULL, + payload_json TEXT NOT NULL, + created_at INTEGER NOT NULL + ); + CREATE UNIQUE INDEX IF NOT EXISTS idx_wm_unique ON working_memory(conversation_id, order_idx); + CREATE INDEX IF NOT EXISTS idx_wm_conv ON working_memory(conversation_id, order_idx); + + CREATE TABLE IF NOT EXISTS reflection ( + id TEXT NOT NULL PRIMARY KEY, + conversation_id TEXT, + created_at INTEGER NOT NULL, + turns_analyzed INTEGER NOT NULL, + score INTEGER NOT NULL, + summary TEXT NOT NULL, + weak_spots_json TEXT NOT NULL DEFAULT '[]' + ); + CREATE INDEX IF NOT EXISTS idx_reflection_created ON reflection(created_at DESC); + CREATE INDEX IF NOT EXISTS idx_reflection_conv ON reflection(conversation_id, created_at DESC); + + CREATE TABLE IF NOT EXISTS agent_event ( + id TEXT NOT NULL PRIMARY KEY, + conversation_id TEXT, + created_at INTEGER NOT NULL, + type TEXT NOT NULL, + payload TEXT NOT NULL + ); + CREATE INDEX IF NOT EXISTS idx_agent_event_conv_time ON agent_event(conversation_id, created_at); + CREATE INDEX IF NOT EXISTS idx_agent_event_time ON agent_event(created_at); + """ + + fun open(path: String): KsqliteStores { + val conn = SQLiteConnection.open(path) + conn.exec(SCHEMA) + return assemble(conn) + } + + fun inMemory(name: String = "agentik-test"): KsqliteStores { + val conn = SQLiteConnection.memory(name) + conn.exec(SCHEMA) + return assemble(conn) + } + + private fun assemble(conn: SQLiteConnection): KsqliteStores { + val messages = KsqliteMessageStore(conn) + val working = KsqliteWorkingMemoryStore(conn) + return KsqliteStores( + connection = conn, + conversations = KsqliteConversationStore(conn, messages, working), + messages = messages, + workingMemory = working, + reflections = KsqliteReflectionStore(conn), + events = KsqliteEventStore(conn), + ) + } + } +} diff --git a/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteWorkingMemoryStore.kt b/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteWorkingMemoryStore.kt new file mode 100644 index 0000000..0f44d55 --- /dev/null +++ b/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteWorkingMemoryStore.kt @@ -0,0 +1,140 @@ +package pw.binom.agentik.storage.ksqlite + +import kotlinx.serialization.json.Json +import pw.binom.agentik.storage.WorkingMemoryEntry +import pw.binom.agentik.storage.WorkingMemoryRow +import pw.binom.agentik.storage.WorkingMemoryStore +import pw.binom.agentik.storage.Ids +import pw.binom.db.ksqlite.SQLiteConnection +import kotlin.time.Instant +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.sync.Mutex +import kotlinx.coroutines.sync.withLock +import kotlinx.coroutines.withContext + +class KsqliteWorkingMemoryStore( + private val connection: SQLiteConnection, +) : WorkingMemoryStore { + + private val mutex = Mutex() + private val json = Json { ignoreUnknownKeys = true } + + override suspend fun append(conversationId: String, entry: WorkingMemoryEntry, now: Instant): Unit = withContext(Dispatchers.Default) { + mutex.withLock { + val newIdx = maxOrderIdx(conversationId) + 1 + connection.prepare( + "INSERT INTO working_memory (id, conversation_id, order_idx, source_message_id, kind, payload_json, created_at) " + + "VALUES (?, ?, ?, ?, ?, ?, ?)" + ).use { stmt -> + stmt.bindText(1, Ids.new("wm")) + stmt.bindText(2, conversationId) + stmt.bindLong(3, newIdx) + val srcId = entry.sourceMessageId + if (srcId != null) stmt.bindText(4, srcId) else stmt.bindNull(4) + stmt.bindText(5, entryKind(entry)) + stmt.bindText(6, json.encodeToString(WorkingMemoryEntry.serializer(), entry)) + stmt.bindLong(7, now.toEpochMilliseconds()) + stmt.executeUpdate() + } + } + } + + override suspend fun list(conversationId: String): List = withContext(Dispatchers.Default) { + mutex.withLock { + connection.prepare( + "SELECT id, conversation_id, order_idx, source_message_id, payload_json, created_at " + + "FROM working_memory WHERE conversation_id = ? ORDER BY order_idx ASC" + ).use { stmt -> + stmt.bindText(1, conversationId) + val out = mutableListOf() + stmt.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 { + connection.prepare("DELETE FROM working_memory WHERE conversation_id = ?").use { stmt -> + stmt.bindText(1, conversationId) + stmt.executeUpdate() + } + } + } + + override suspend fun compact( + dropFromOrderIdx: Long, + conversationId: String, + summaryText: String?, + ): Long = withContext(Dispatchers.Default) { + mutex.withLock { + var newMax = 0L + val nowMs = System.currentTimeMillis() + val summaryId = Ids.new("wm") + connection.exec("BEGIN") + try { + connection.prepare("DELETE FROM working_memory WHERE conversation_id = ? AND order_idx >= ?").use { stmt -> + stmt.bindText(1, conversationId) + stmt.bindLong(2, dropFromOrderIdx) + stmt.executeUpdate() + } + if (!summaryText.isNullOrBlank()) { + val afterDelete = maxOrderIdx(conversationId) + val newIdx = afterDelete + 1 + connection.prepare( + "INSERT INTO working_memory (id, conversation_id, order_idx, source_message_id, kind, payload_json, created_at) " + + "VALUES (?, ?, ?, NULL, ?, ?, ?)" + ).use { stmt -> + stmt.bindText(1, summaryId) + stmt.bindText(2, conversationId) + stmt.bindLong(3, newIdx) + stmt.bindText(4, "summary") + stmt.bindText(5, json.encodeToString(WorkingMemoryEntry.serializer(), WorkingMemoryEntry.Summary(text = summaryText))) + stmt.bindLong(6, nowMs) + stmt.executeUpdate() + } + newMax = newIdx + } else { + newMax = maxOrderIdx(conversationId) + } + connection.exec("COMMIT") + } catch (t: Throwable) { + connection.exec("ROLLBACK") + throw t + } + newMax + } + } + + override fun close() {} + + private suspend fun maxOrderIdx(conversationId: String): Long { + connection.prepare("SELECT COALESCE(MAX(order_idx), 0) FROM working_memory WHERE conversation_id = ?").use { stmt -> + stmt.bindText(1, conversationId) + stmt.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/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/MessageCodecs.kt b/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/MessageCodecs.kt new file mode 100644 index 0000000..9f81abc --- /dev/null +++ b/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/MessageCodecs.kt @@ -0,0 +1,81 @@ +package pw.binom.agentik.storage.ksqlite + +import kotlinx.serialization.json.Json +import pw.binom.agentik.storage.MessageRecord +import pw.binom.agentik.storage.decodeBodyPayload +import pw.binom.agentik.storage.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), + ) + is MessageRecord.Summary, + is MessageRecord.System -> error("Summary/System — synthetic, cannot append to audit log") +} + +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/storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteConversationStoreTest.kt b/storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteConversationStoreTest.kt new file mode 100644 index 0000000..f05418b --- /dev/null +++ b/storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteConversationStoreTest.kt @@ -0,0 +1,105 @@ +package pw.binom.agentik.storage.ksqlite + +import kotlinx.coroutines.test.runTest +import pw.binom.agentik.storage.ConversationRecord +import kotlin.test.AfterTest +import kotlin.test.BeforeTest +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertNotNull +import kotlin.test.assertNull +import kotlin.test.assertTrue +import kotlin.time.Duration +import kotlin.time.Instant + +class KsqliteConversationStoreTest { + private lateinit var stores: KsqliteStores + + @BeforeTest + fun setup() { + stores = KsqliteStores.inMemory("conv-${kotlin.random.Random.nextLong()}") + } + + @AfterTest + fun tearDown() = stores.close() + + private fun rec(id: String, title: String? = null, ts: Instant = Instant.parse("2026-09-15T10:00:00Z")) = + ConversationRecord(id, title, false, ts, ts) + + @Test + fun testUpsertAndGetRoundtrip() = runTest { + stores.conversations.upsert(rec("c1", "test")) + assertEquals(rec("c1", "test"), stores.conversations.get("c1")) + } + + @Test + fun testGetReturnsNullForMissing() = runTest { + assertNull(stores.conversations.get("nope")) + } + + @Test + fun testDeleteRemovesAndReturnsTrue() = runTest { + stores.conversations.upsert(rec("c1")) + assertTrue(stores.conversations.delete("c1")) + assertNull(stores.conversations.get("c1")) + assertEquals(false, stores.conversations.delete("c1")) + } + + @Test + fun testListSortsByUpdatedAtDesc() = runTest { + val t0 = Instant.parse("2026-09-15T10:00:00Z") + stores.conversations.upsert(rec("c1", ts = t0)) + stores.conversations.upsert(rec("c2", ts = t0)) + stores.conversations.upsert(rec("c3", ts = t0)) + stores.conversations.upsert(rec("c4", ts = t0)) + + stores.conversations.touch("c2", t0 + Duration.parse("PT60S")) + stores.conversations.touch("c3", t0 + Duration.parse("PT120S")) + stores.conversations.touch("c4", t0 + Duration.parse("PT180S")) + + val page = stores.conversations.list(offset = 0, limit = 4) + assertEquals(listOf("c4", "c3", "c2", "c1"), page.map { it.id }) + } + + @Test + fun testListRespectsOffsetAndLimit() = runTest { + val t0 = Instant.parse("2026-09-15T10:00:00Z") + for (i in 1..5) stores.conversations.upsert(rec("c$i", ts = t0 + Duration.parse("PT${i}S"))) + val p0 = stores.conversations.list(offset = 0, limit = 2) + assertEquals(2, p0.size) + val p2 = stores.conversations.list(offset = 4, limit = 2) + assertEquals(1, p2.size) + } + + @Test + fun testRenameUpdatesTitleAndUpdatedAt() = runTest { + stores.conversations.upsert(rec("c1")) + val newTs = stores.conversations.rename("c1", "new title") + assertNotNull(newTs) + assertEquals("new title", stores.conversations.get("c1")?.title) + } + + @Test + fun testRenameWithNullClearsTitle() = runTest { + stores.conversations.upsert(rec("c1", "old")) + stores.conversations.rename("c1", null) + assertNull(stores.conversations.get("c1")?.title) + } + + @Test + fun testRenameReturnsNullForMissing() = runTest { + assertNull(stores.conversations.rename("nope", "x")) + } + + @Test + fun testTouchUpdatesUpdatedAtOnly() = runTest { + val t0 = Instant.parse("2026-09-15T10:00:00Z") + val t1 = Instant.parse("2026-09-15T10:01:00Z") + stores.conversations.upsert(ConversationRecord("c1", "title", false, t0, t0)) + stores.conversations.touch("c1", t1) + val got = stores.conversations.get("c1") + assertEquals("title", got?.title) + assertEquals(t1, got?.updatedAt) + assertEquals(t0, got?.createdAt) + } +} diff --git a/storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteEventStoreTest.kt b/storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteEventStoreTest.kt new file mode 100644 index 0000000..c2ce19c --- /dev/null +++ b/storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteEventStoreTest.kt @@ -0,0 +1,150 @@ +package pw.binom.agentik.storage.ksqlite + +import kotlinx.coroutines.test.runTest +import pw.binom.agentik.storage.events.EventRecord +import pw.binom.agentik.storage.events.EventType +import kotlin.test.AfterTest +import kotlin.test.BeforeTest +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertNull +import kotlin.test.assertTrue +import kotlin.time.Instant + +class KsqliteEventStoreTest { + + private lateinit var stores: KsqliteStores + + @BeforeTest + fun setup() { + stores = KsqliteStores.inMemory("test-${kotlin.random.Random.nextLong()}") + } + + @AfterTest + fun tearDown() { + stores.close() + } + + private fun rec( + id: String, + ts: Long, + conv: String? = null, + type: EventType = EventType.CONVERSATION_APPEND_TEXT, + ): EventRecord = EventRecord( + id = id, + conversationId = conv, + createdAt = Instant.fromEpochMilliseconds(ts), + type = type, + payload = "{\"i\":\"$id\"}", + ) + + @Test + fun testAppendThenQueryReturnsRecord() = runTest { + val r = rec("ev-1", ts = 1000) + stores.events.append(r) + assertEquals(listOf(r), stores.events.query()) + } + + @Test + fun testQueryWithConversationIdFilters() = runTest { + stores.events.append(rec("ev-1", ts = 1000, conv = "c-1")) + stores.events.append(rec("ev-2", ts = 2000, conv = "c-2")) + stores.events.append(rec("ev-3", ts = 3000, conv = "c-1")) + + val c1 = stores.events.query(conversationId = "c-1") + assertEquals(listOf("ev-1", "ev-3"), c1.map { it.id }) + } + + @Test + fun testQueryWithAfterIdReturnsEventsStrictlyAfterCursor() = runTest { + stores.events.append(rec("ev-1", ts = 1000)) + stores.events.append(rec("ev-2", ts = 2000)) + stores.events.append(rec("ev-3", ts = 3000)) + + // id-based cursor: id > 'ev-1' returns ev-2, ev-3 + assertEquals(listOf("ev-2", "ev-3"), stores.events.query(afterId = "ev-1").map { it.id }) + } + + @Test + fun testAppendIsIdempotentOnIdWithInsertOrReplace() = runTest { + val r1 = rec("ev-1", ts = 1000, type = EventType.CONVERSATION_APPEND_TEXT) + val r2 = rec("ev-1", ts = 1000, type = EventType.CONVERSATION_TOOL_CALL) + stores.events.append(r1) + stores.events.append(r2) + // INSERT OR REPLACE means replace wins. Documented in KsqliteEventStore KDoc. + val result = stores.events.query() + assertEquals(1, result.size) + assertEquals(EventType.CONVERSATION_TOOL_CALL, result[0].type) + } + + @Test + fun testRecordsOrderedByCreatedAtThenId() = runTest { + stores.events.append(rec("ev-c", ts = 2000)) + stores.events.append(rec("ev-a", ts = 1000)) + stores.events.append(rec("ev-d", ts = 1000)) + stores.events.append(rec("ev-b", ts = 1500)) + + assertEquals(listOf("ev-a", "ev-d", "ev-b", "ev-c"), + stores.events.query().map { it.id }) + } + + @Test + fun testQueryWithLimitCaps() = runTest { + repeat(10) { i -> stores.events.append(rec("ev-$i", ts = (i * 100).toLong())) } + val first3 = stores.events.query(limit = 3) + assertEquals(3, first3.size) + assertEquals(listOf("ev-0", "ev-1", "ev-2"), first3.map { it.id }) + } + + @Test + fun testPruneOlderThanRemovesOld() = runTest { + stores.events.append(rec("ev-1", ts = 1000)) + stores.events.append(rec("ev-2", ts = 2000)) + stores.events.append(rec("ev-3", ts = 3000)) + + val removed = stores.events.pruneOlderThan(Instant.fromEpochMilliseconds(2500)) + assertEquals(2, removed) + assertEquals(listOf("ev-3"), stores.events.query().map { it.id }) + } + + @Test + fun testCountReturnsTotal() = runTest { + assertEquals(0, stores.events.count()) + stores.events.append(rec("ev-1", ts = 1000)) + stores.events.append(rec("ev-2", ts = 2000)) + assertEquals(2, stores.events.count()) + } + + @Test + fun testUnknownEventTypeLoadedAsFallback() = runTest { + // Forward-compat: write record with unknown type, read via store. + stores.connection.exec( + "INSERT INTO agent_event (id, conversation_id, created_at, type, payload) VALUES " + + "('ev-future', NULL, 5000, 'SOME_FUTURE_TYPE', '{}')" + ) + val result = stores.events.query() + assertEquals(1, result.size) + assertEquals("ev-future", result[0].id) + assertEquals(EventType.AGENT_CREATED, result[0].type) + } + + @Test + fun testNullConversationIdStoredAndRetrieved() = runTest { + stores.events.append(rec("ev-no-conv", ts = 1000, conv = null)) + val result = stores.events.query() + assertEquals(1, result.size) + assertNull(result[0].conversationId) + } + + @Test + fun testSchemaCreatedOnFirstOpen() = runTest { + val count = stores.connection.prepare( + "SELECT COUNT(*) FROM sqlite_master WHERE type='table' AND name='agent_event'" + ).use { stmt -> + stmt.executeQuery().use { rs -> + if (rs.next()) (rs.getLong(0) ?: 0L).toInt() else 0 + } + } + assertEquals(1, count) + } +} diff --git a/storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteMessageStoreTest.kt b/storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteMessageStoreTest.kt new file mode 100644 index 0000000..b89ab28 --- /dev/null +++ b/storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteMessageStoreTest.kt @@ -0,0 +1,97 @@ +package pw.binom.agentik.storage.ksqlite + +import kotlinx.coroutines.test.runTest +import pw.binom.agentik.storage.Content +import pw.binom.agentik.storage.MessageRecord +import pw.binom.agentik.storage.TurnTokens +import kotlin.test.AfterTest +import kotlin.test.BeforeTest +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertNotNull +import kotlin.test.assertNull +import kotlin.time.Instant + +class KsqliteMessageStoreTest { + private lateinit var stores: KsqliteStores + + @BeforeTest + fun setup() { + stores = KsqliteStores.inMemory("msg-${kotlin.random.Random.nextLong()}") + kotlinx.coroutines.runBlocking { + stores.conversations.upsert( + pw.binom.agentik.storage.ConversationRecord( + "conv1", null, false, + Instant.parse("2026-09-15T10:00:00Z"), + Instant.parse("2026-09-15T10:00:00Z"), + ) + ) + } + } + + @AfterTest + fun tearDown() = stores.close() + + @Test + fun testAppendUserAndRetrieve() = runTest { + stores.messages.append(MessageRecord.UserMessage( + id = "m1", + conversationId = "conv1", + content = listOf(Content.Text("hello")), + createdAt = Instant.parse("2026-09-15T10:01:00Z"), + context = null, + )) + val list = stores.messages.listAll("conv1") + assertEquals(1, list.size) + val msg = list[0] + assertEquals("m1", msg.id) + assertEquals(MessageRecord.UserMessage::class, msg::class) + } + + @Test + fun testAppendAssistantWithTokens() = runTest { + stores.messages.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 stats = stores.messages.tokenStats("conv1") + assertEquals(1, stats.turns) + assertEquals(50L, stats.inputTokens) + assertEquals(30L, stats.outputTokens) + } + + @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") + stores.messages.append(MessageRecord.UserMessage("m1", "conv1", listOf(Content.Text("a")), t1, null)) + stores.messages.append(MessageRecord.UserMessage("m2", "conv1", listOf(Content.Text("b")), t2, null)) + stores.messages.append(MessageRecord.UserMessage("m3", "conv1", listOf(Content.Text("c")), t3, null)) + + val after = stores.messages.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) stores.messages.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"), stores.messages.listAll("conv1").map { it.id }) + } + + @Test + fun testTokenStatsIgnoresNonAssistant() = runTest { + val t = Instant.parse("2026-09-15T10:00:00Z") + stores.messages.append(MessageRecord.UserMessage("m1", "conv1", listOf(Content.Text("user")), t, null)) + stores.messages.append(MessageRecord.AssistantMessage("m2", "conv1", listOf(Content.Text("asst")), t, TurnTokens(10, 5))) + val stats = stores.messages.tokenStats("conv1") + assertEquals(1, stats.turns) + } +} diff --git a/storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteReflectionStoreTest.kt b/storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteReflectionStoreTest.kt new file mode 100644 index 0000000..89afeab --- /dev/null +++ b/storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteReflectionStoreTest.kt @@ -0,0 +1,94 @@ +package pw.binom.agentik.storage.ksqlite + +import kotlinx.coroutines.test.runTest +import pw.binom.agentik.storage.Reflection +import kotlin.test.AfterTest +import kotlin.test.BeforeTest +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertNotNull +import kotlin.test.assertNull +import kotlin.test.assertTrue +import kotlin.time.Instant + +class KsqliteReflectionStoreTest { + private lateinit var stores: KsqliteStores + + @BeforeTest + fun setup() { + stores = KsqliteStores.inMemory("refl-${kotlin.random.Random.nextLong()}") + } + + @AfterTest + fun tearDown() = stores.close() + + private fun refl( + id: String, + conv: String? = null, + ts: Instant = Instant.parse("2026-09-15T10:00:00Z"), + score: Int = 3, + summary: String = "ok", + weak: List = listOf("weak1", "weak2"), + ) = Reflection( + id = id, + conversationId = conv, + createdAt = ts, + turnsAnalyzed = 5, + score = score, + summary = summary, + weakSpots = weak, + ) + + @Test + fun testInsertAndGetRoundtrip() = runTest { + stores.reflections.insert(refl("r1")) + val got = stores.reflections.get("r1") + assertNotNull(got) + assertEquals("r1", got.id) + assertEquals(listOf("weak1", "weak2"), got.weakSpots) + } + + @Test + fun testGetReturnsNullForMissing() = runTest { + assertNull(stores.reflections.get("nope")) + } + + @Test + fun testListRecentSortedByCreatedAtDesc() = runTest { + val t0 = Instant.parse("2026-09-15T10:00:00Z") + stores.reflections.insert(refl("r1", ts = t0)) + stores.reflections.insert(refl("r2", ts = t0 + kotlin.time.Duration.parse("PT60S"))) + stores.reflections.insert(refl("r3", ts = t0 + kotlin.time.Duration.parse("PT120S"))) + + val recent = stores.reflections.listRecent(limit = 3) + assertEquals(listOf("r3", "r2", "r1"), recent.map { it.id }) + } + + @Test + fun testListForConversationFilters() = runTest { + val t = Instant.parse("2026-09-15T10:00:00Z") + stores.reflections.insert(refl("r1", conv = "c-1")) + stores.reflections.insert(refl("r2", conv = "c-2")) + stores.reflections.insert(refl("r3", conv = "c-1")) + val forC1 = stores.reflections.listForConversation("c-1", limit = 10) + assertEquals(2, forC1.size) + } + + @Test + fun testDeleteOlderThanRemoves() = runTest { + val t0 = Instant.parse("2026-09-15T10:00:00Z") + stores.reflections.insert(refl("r1", ts = t0)) + stores.reflections.insert(refl("r2", ts = t0 + kotlin.time.Duration.parse("PT1H"))) + stores.reflections.deleteOlderThan(t0 + kotlin.time.Duration.parse("PT30M")) + assertNull(stores.reflections.get("r1")) + assertNotNull(stores.reflections.get("r2")) + } + + @Test + fun testCountReturnsTotal() = runTest { + assertEquals(0, stores.reflections.count()) + stores.reflections.insert(refl("r1")) + stores.reflections.insert(refl("r2")) + assertEquals(2, stores.reflections.count()) + } +} diff --git a/storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteWorkingMemoryStoreTest.kt b/storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteWorkingMemoryStoreTest.kt new file mode 100644 index 0000000..190f02e --- /dev/null +++ b/storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteWorkingMemoryStoreTest.kt @@ -0,0 +1,87 @@ +package pw.binom.agentik.storage.ksqlite + +import kotlinx.coroutines.test.runTest +import pw.binom.agentik.storage.Content +import pw.binom.agentik.storage.WorkingMemoryEntry +import kotlin.test.AfterTest +import kotlin.test.BeforeTest +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertTrue +import kotlin.time.Instant + +class KsqliteWorkingMemoryStoreTest { + private lateinit var stores: KsqliteStores + + @BeforeTest + fun setup() { + stores = KsqliteStores.inMemory("wm-${kotlin.random.Random.nextLong()}") + } + + @AfterTest + fun tearDown() = stores.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") + stores.workingMemory.append("c1", userMsg("first"), t) + stores.workingMemory.append("c1", asstMsg("reply"), t) + val list = stores.workingMemory.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") + stores.workingMemory.append("c1", userMsg("c1-msg"), t) + stores.workingMemory.append("c2", userMsg("c2-msg"), t) + assertEquals(1, stores.workingMemory.list("c1").size) + assertEquals(1, stores.workingMemory.list("c2").size) + } + + @Test + fun testClearRemovesAllForConversation() = runTest { + val t = Instant.parse("2026-09-15T10:00:00Z") + stores.workingMemory.append("c1", userMsg("a"), t) + stores.workingMemory.append("c1", userMsg("b"), t) + stores.workingMemory.clear("c1") + assertEquals(emptyList(), stores.workingMemory.list("c1")) + } + + @Test + fun testCompactDeletesAndInsertsSummary() = runTest { + val t = Instant.parse("2026-09-15T10:00:00Z") + stores.workingMemory.append("c1", userMsg("a"), t) + stores.workingMemory.append("c1", userMsg("b"), t) + stores.workingMemory.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 = stores.workingMemory.compact(dropFromOrderIdx = 2, conversationId = "c1", summaryText = "summary") + assertEquals(2L, newMax) + val remaining = stores.workingMemory.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") + stores.workingMemory.append("c1", userMsg("a"), t) + // dropFromOrderIdx=2: удаляет idx >= 2, остаётся idx=1. + val newMax = stores.workingMemory.compact(dropFromOrderIdx = 2, conversationId = "c1", summaryText = null) + assertEquals(1L, newMax) + val remaining = stores.workingMemory.list("c1") + assertEquals(1, remaining.size) + assertEquals(1L, remaining[0].orderIdx) + } +}