feat(context-ksqlite): add :context-ksqlite module with KsqliteContextStore implementation
ci / JVM build + tests (push) Failing after 1m6s

- Introduced a new `:context-ksqlite` module implementing `:context-api` with Ksqlite backend.
- Added a minimal schema (working_memory table + 2 indexes) for runtime context store: lightweight, excludes conversation/message/reflection tables, which are handled in separate modules.
- Initial version supports JVM, Linux, and Windows builds; migrates schema using SQLite PRAGMA.
- Includes fully autonomous in-memory schema migration and tests to validate operations like append, list, clear, and compact.
- Partial duplication of `:storage-ksqlite/KsqliteWorkingMemoryStore`; consumers will migrate incrementally.
This commit is contained in:
2026-09-21 01:18:32 +03:00
parent acc7237e51
commit d6afc05c20
5 changed files with 450 additions and 0 deletions
+40
View File
@@ -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)
}
}
}
@@ -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<WorkingMemoryRow> = withContext(Dispatchers.Default) {
mutex.withLock {
listStmt.reset()
listStmt.clearBindings()
listStmt.bindText(1, conversationId)
val out = mutableListOf<WorkingMemoryRow>()
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"
}
}
@@ -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")
}
}
@@ -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)
}
}
+5
View File
@@ -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 на