Коммит заменяет SQLDelight на ksqlite и добавляет journal/context/outbox
ci / JVM build + tests (push) Failing after 1m24s

This commit is contained in:
2026-09-21 01:12:31 +03:00
parent bd65c29b48
commit acc7237e51
61 changed files with 1623 additions and 1750 deletions
+11 -8
View File
@@ -23,18 +23,21 @@ kotlin {
sourceSets {
commonMain.dependencies {
api(project(":message-store-api"))
api(project(":message-log-api"))
api(project(":working-memory-api"))
// :message-store-api / :working-memory-api / :message-log-api объявлены
// как api-зависимости здесь, в commonMain — без этого commonMain
// не скомпилируется (KsqliteConversationStore, KsqliteMessageStore
// и т.д. используют их типы в commonMain). У них самих есть Apple
// targets (macosX64/Arm64, iosX64/Arm64/SimulatorArm64, linuxArm64),
// так что KMP-метаданные корректно резолвятся для всех таргетов.
// ksqlite ещё не опубликован в Maven Central — только в локальном
// caffeine-репо. Версия пока 0.1.0-SNAPSHOT (CI fallback из README).
implementation("pw.binom.db:ksqlite:0.1.0")
implementation("pw.binom.db:ksqlite:0.1.1-SNAPSHOT")
implementation(libs.kotlinx.serialization.json)
}
jvmMain.dependencies {
// Логирование — JVM-only, для native targets в kotlin-logging нет KMP-артефакта.
implementation("io.github.microutils:kotlin-logging-jvm:3.0.5")
api(project(":message-store-api"))
api(project(":message-log-api"))
api(project(":working-memory-api"))
}
commonTest.dependencies {
implementation(kotlin("test"))
@@ -1,5 +1,6 @@
package pw.binom.agentik.storage.ksqlite
import kotlin.time.Clock
import kotlin.time.Instant
import pw.binom.agentik.messageStore.ConversationRecord
import pw.binom.agentik.messageStore.ConversationStore
@@ -10,8 +11,9 @@ import kotlinx.coroutines.sync.withLock
import kotlinx.coroutines.withContext
/**
* ksqlite-реализация [ConversationStore]. Схема таблицы `conversation` повторяет
* [pw.binom.agentik.storage.sqlite.SqliteConversationStore] для совместимости данных.
* ksqlite-реализация [ConversationStore]. Схема таблицы `conversation` живёт
* в [Schema] (миграция через PRAGMA user_version) — этот класс только
* готовит и выполняет SQL, ссылаясь на `Schema.COL_*` / `Schema.TABLE_*`.
*/
class KsqliteConversationStore(
private val connection: SQLiteConnection,
@@ -21,50 +23,100 @@ class KsqliteConversationStore(
private val mutex = Mutex()
// pre-prepare всех statement'ов — аналогично KsqliteMessageStore (см.
// KDoc там — почему GC-finalize на StmtHolder'е роняет JVM, если stmt
// живёт после закрытия connection).
private val existsStmt = connection.prepare(
"SELECT 1 FROM ${Schema.TABLE_CONVERSATION} WHERE ${Schema.COL_ID} = ?"
)
private val updateStmt = connection.prepare(
"""
UPDATE ${Schema.TABLE_CONVERSATION}
SET ${Schema.COL_TITLE} = ?, ${Schema.COL_IS_TEMPORAL} = ?, ${Schema.COL_UPDATED_AT} = ?
WHERE ${Schema.COL_ID} = ?
""".trimIndent()
)
private val insertStmt = connection.prepare(
"""
INSERT INTO ${Schema.TABLE_CONVERSATION}
(${Schema.COL_ID}, ${Schema.COL_TITLE}, ${Schema.COL_IS_TEMPORAL},
${Schema.COL_CREATED_AT}, ${Schema.COL_UPDATED_AT})
VALUES (?, ?, ?, ?, ?)
""".trimIndent()
)
private val getStmt = connection.prepare(
"""
SELECT ${Schema.COL_ID}, ${Schema.COL_TITLE}, ${Schema.COL_IS_TEMPORAL},
${Schema.COL_CREATED_AT}, ${Schema.COL_UPDATED_AT}
FROM ${Schema.TABLE_CONVERSATION}
WHERE ${Schema.COL_ID} = ?
""".trimIndent()
)
private val deleteStmt = connection.prepare(
"DELETE FROM ${Schema.TABLE_CONVERSATION} WHERE ${Schema.COL_ID} = ?"
)
private val listStmt = connection.prepare(
"""
SELECT ${Schema.COL_ID}, ${Schema.COL_TITLE}, ${Schema.COL_IS_TEMPORAL},
${Schema.COL_CREATED_AT}, ${Schema.COL_UPDATED_AT}
FROM ${Schema.TABLE_CONVERSATION}
WHERE ${Schema.COL_IS_TEMPORAL} = 0
ORDER BY ${Schema.COL_UPDATED_AT} DESC
LIMIT ? OFFSET ?
""".trimIndent()
)
private val renameStmt = connection.prepare(
"""
UPDATE ${Schema.TABLE_CONVERSATION}
SET ${Schema.COL_TITLE} = ?, ${Schema.COL_UPDATED_AT} = ?
WHERE ${Schema.COL_ID} = ?
""".trimIndent()
)
private val renameUpdatedAtStmt = connection.prepare(
"SELECT ${Schema.COL_UPDATED_AT} FROM ${Schema.TABLE_CONVERSATION} WHERE ${Schema.COL_ID} = ?"
)
private val touchStmt = connection.prepare(
"""
UPDATE ${Schema.TABLE_CONVERSATION}
SET ${Schema.COL_UPDATED_AT} = ?
WHERE ${Schema.COL_ID} = ?
""".trimIndent()
)
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()
}
}
}
val exists = execExists(record.id)
if (exists) {
updateStmt.reset()
updateStmt.clearBindings()
val t = record.title
if (t != null) updateStmt.bindText(1, t) else updateStmt.bindNull(1)
updateStmt.bindInt(2, if (record.isTemporal) 1 else 0)
updateStmt.bindLong(3, record.updatedAt.toEpochMilliseconds())
updateStmt.bindText(4, record.id)
updateStmt.executeUpdate()
} else {
insertStmt.reset()
insertStmt.clearBindings()
insertStmt.bindText(1, record.id)
val title = record.title
if (title != null) insertStmt.bindText(2, title) else insertStmt.bindNull(2)
insertStmt.bindInt(3, if (record.isTemporal) 1 else 0)
insertStmt.bindLong(4, record.createdAt.toEpochMilliseconds())
insertStmt.bindLong(5, record.updatedAt.toEpochMilliseconds())
insertStmt.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
}
}
getStmt.reset()
getStmt.clearBindings()
getStmt.bindText(1, id)
getStmt.executeQuery().use { rs ->
if (rs.next()) rs.toRecord() else null
}
}
}
@@ -72,68 +124,77 @@ class KsqliteConversationStore(
mutex.withLock {
// Проверяем существование через raw query, НЕ через get() — get() тоже
// берёт mutex (не реентрант), что привело бы к deadlock.
connection.prepare("SELECT 1 FROM conversation WHERE id = ?").use { check ->
check.bindText(1, id)
check.executeQuery().use { rs -> if (!rs.next()) return@withContext false }
}
if (!execExists(id)) return@withContext false
messageStore?.clear(id)
workingMemoryStore?.clear(id)
connection.prepare("DELETE FROM conversation WHERE id = ?").use { stmt ->
stmt.bindText(1, id)
stmt.executeUpdate()
}
deleteStmt.reset()
deleteStmt.clearBindings()
deleteStmt.bindText(1, id)
deleteStmt.executeUpdate()
true
}
}
override suspend fun list(offset: Int, limit: Int): List<ConversationRecord> = 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<ConversationRecord>()
stmt.executeQuery().use { rs ->
while (rs.next()) result.add(rs.toRecord())
}
result
listStmt.reset()
listStmt.clearBindings()
listStmt.bindLong(1, limit.toLong())
listStmt.bindLong(2, offset.toLong())
val result = mutableListOf<ConversationRecord>()
listStmt.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()
}
// raw query instead of get() (deadlock — get() also takes mutex)
connection.prepare("SELECT updated_at FROM conversation WHERE id = ?").use { stmt ->
stmt.bindText(1, id)
stmt.executeQuery().use { rs ->
if (rs.next()) Instant.fromEpochMilliseconds(rs.getLong(0)!!) else null
}
val nowMs = Clock.System.now().toEpochMilliseconds()
renameStmt.reset()
renameStmt.clearBindings()
if (title != null) renameStmt.bindText(1, title) else renameStmt.bindNull(1)
renameStmt.bindLong(2, nowMs)
renameStmt.bindText(3, id)
renameStmt.executeUpdate()
renameUpdatedAtStmt.reset()
renameUpdatedAtStmt.clearBindings()
renameUpdatedAtStmt.bindText(1, id)
renameUpdatedAtStmt.executeQuery().use { rs ->
if (rs.next()) Instant.fromEpochMilliseconds(rs.getLong(0)!!) else null
}
}
}
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()
}
touchStmt.reset()
touchStmt.clearBindings()
touchStmt.bindLong(1, now.toEpochMilliseconds())
touchStmt.bindText(2, id)
touchStmt.executeUpdate()
}
}
override fun close() {
// Connection lifecycle — на caller'е (фабрика KsqliteStores).
existsStmt.close()
updateStmt.close()
insertStmt.close()
getStmt.close()
deleteStmt.close()
listStmt.close()
renameStmt.close()
renameUpdatedAtStmt.close()
touchStmt.close()
}
private fun execExists(id: String): Boolean {
existsStmt.reset()
existsStmt.clearBindings()
existsStmt.bindText(1, id)
existsStmt.executeQuery().use { rs -> return rs.next() }
}
private fun pw.binom.db.ksqlite.SQLiteResultSet.toRecord(): ConversationRecord = ConversationRecord(
@@ -3,11 +3,8 @@ package pw.binom.agentik.storage.ksqlite
import kotlinx.serialization.json.Json
import pw.binom.agentik.messageLog.MessageRecord
import pw.binom.agentik.messageLog.MutableMessageStore
import pw.binom.agentik.messageLog.TokenStats
import pw.binom.agentik.messageLog.decodeBodyPayload
import pw.binom.agentik.messageLog.encodeBodyPayload
import pw.binom.db.ksqlite.SQLiteConnection
import pw.binom.db.ksqlite.SQLiteResultSet
import pw.binom.db.ksqlite.SQLitePreparedStatement
import kotlin.time.Instant
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.sync.Mutex
@@ -15,33 +12,58 @@ import kotlinx.coroutines.sync.withLock
import kotlinx.coroutines.withContext
/**
* ksqlite-реализация [MessageStore]. Схема таблицы `message` повторяет
* [pw.binom.agentik.storage.sqlite.SqliteMessageStore].
* ksqlite-реализация [MutableMessageStore] (append-only audit log).
*
* encoding helpers (`encodeRecord`/`toMessageRecord`/`CallPayload`/...)
* переиспользуются из `:storage-sqlite` чтобы избежать дрейфа между
* двумя backend'ами.
* 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] рядом.
*/
class KsqliteMessageStore(
class KsqliteMessageStore internal constructor(
private val connection: SQLiteConnection,
) : MutableMessageStore {
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 {
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()
}
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()
}
}
@@ -52,63 +74,32 @@ class KsqliteMessageStore(
limit: Int,
): List<MessageRecord> = 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<MessageRecord>()
stmt.executeQuery().use { rs ->
while (rs.next()) out.add(rs.toMessageRecord(json))
}
out
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<MessageRecord>()
listStmt.executeQuery().use { rs ->
while (rs.next()) out.add(rs.toMessageRecord(json))
}
out
}
}
override suspend fun listAll(conversationId: String): List<MessageRecord> = 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<MessageRecord>()
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()
}
clearStmt.reset()
clearStmt.clearBindings()
clearStmt.bindText(1, conversationId)
clearStmt.executeUpdate()
}
}
override fun close() {
insertStmt.close()
listStmt.close()
clearStmt.close()
}
}
@@ -4,6 +4,7 @@ import pw.binom.agentik.messageStore.Reflection
import pw.binom.agentik.messageStore.ReflectionEvent
import pw.binom.agentik.messageStore.ReflectionStore
import pw.binom.db.ksqlite.SQLiteConnection
import pw.binom.db.ksqlite.SQLitePreparedStatement
import pw.binom.db.ksqlite.SQLiteResultSet
import kotlin.time.Clock
import kotlin.time.Instant
@@ -14,16 +15,7 @@ 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,
@@ -32,91 +24,138 @@ class KsqliteReflectionStore(
private val mutex = Mutex()
private val ev = MutableSharedFlow<ReflectionEvent>(extraBufferCapacity = 16)
// pre-prepare (см. KsqliteMessageStore KDoc — почему это критично против
// SIGSEGV в StmtHolder.finalize на закрытой connection).
private val insertStmt = connection.prepare(
"""
INSERT INTO ${Schema.TABLE_REFLECTION}
(${Schema.COL_ID}, ${Schema.COL_CONVERSATION_ID}, ${Schema.COL_CREATED_AT},
${Schema.COL_TURNS_ANALYZED}, ${Schema.COL_SCORE},
${Schema.COL_SUMMARY}, ${Schema.COL_WEAK_SPOTS_JSON})
VALUES (?, ?, ?, ?, ?, ?, ?)
""".trimIndent()
)
private val getStmt = connection.prepare(
"""
SELECT ${Schema.COL_ID}, ${Schema.COL_CONVERSATION_ID}, ${Schema.COL_CREATED_AT},
${Schema.COL_TURNS_ANALYZED}, ${Schema.COL_SCORE},
${Schema.COL_SUMMARY}, ${Schema.COL_WEAK_SPOTS_JSON}
FROM ${Schema.TABLE_REFLECTION}
WHERE ${Schema.COL_ID} = ?
""".trimIndent()
)
private val listRecentStmt = connection.prepare(
"""
SELECT ${Schema.COL_ID}, ${Schema.COL_CONVERSATION_ID}, ${Schema.COL_CREATED_AT},
${Schema.COL_TURNS_ANALYZED}, ${Schema.COL_SCORE},
${Schema.COL_SUMMARY}, ${Schema.COL_WEAK_SPOTS_JSON}
FROM ${Schema.TABLE_REFLECTION}
ORDER BY ${Schema.COL_CREATED_AT} DESC
LIMIT ?
""".trimIndent()
)
private val listForConvStmt = connection.prepare(
"""
SELECT ${Schema.COL_ID}, ${Schema.COL_CONVERSATION_ID}, ${Schema.COL_CREATED_AT},
${Schema.COL_TURNS_ANALYZED}, ${Schema.COL_SCORE},
${Schema.COL_SUMMARY}, ${Schema.COL_WEAK_SPOTS_JSON}
FROM ${Schema.TABLE_REFLECTION}
WHERE ${Schema.COL_CONVERSATION_ID} = ?
ORDER BY ${Schema.COL_CREATED_AT} DESC
LIMIT ?
""".trimIndent()
)
private val deleteOlderThanStmt = connection.prepare(
"DELETE FROM ${Schema.TABLE_REFLECTION} WHERE ${Schema.COL_CREATED_AT} < ?"
)
private val countStmt = connection.prepare(
"SELECT COUNT(*) FROM ${Schema.TABLE_REFLECTION}"
)
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()
}
insertStmt.reset()
insertStmt.clearBindings()
insertStmt.bindText(1, reflection.id)
val convId = reflection.conversationId
if (convId != null) insertStmt.bindText(2, convId) else insertStmt.bindNull(2)
insertStmt.bindLong(3, reflection.createdAt.toEpochMilliseconds())
insertStmt.bindLong(4, reflection.turnsAnalyzed.toLong())
insertStmt.bindLong(5, reflection.score.toLong())
insertStmt.bindText(6, reflection.summary)
insertStmt.bindText(7, encodeStringArray(reflection.weakSpots))
insertStmt.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
}
}
getStmt.reset()
getStmt.clearBindings()
getStmt.bindText(1, id)
getStmt.executeQuery().use { rs ->
if (rs.next()) rs.toDomain() else null
}
}
}
override suspend fun listRecent(limit: Int): List<Reflection> = 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<Reflection>()
stmt.executeQuery().use { rs ->
while (rs.next()) out.add(rs.toDomain())
}
out
}
listRecentStmt.reset()
listRecentStmt.clearBindings()
listRecentStmt.bindLong(1, limit.toLong())
val out = mutableListOf<Reflection>()
listRecentStmt.executeQuery().use { rs ->
while (rs.next()) out.add(rs.toDomain())
}
out
}
}
override suspend fun listForConversation(conversationId: String, limit: Int): List<Reflection> = 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<Reflection>()
stmt.executeQuery().use { rs ->
while (rs.next()) out.add(rs.toDomain())
}
out
}
listForConvStmt.reset()
listForConvStmt.clearBindings()
listForConvStmt.bindText(1, conversationId)
listForConvStmt.bindLong(2, limit.toLong())
val out = mutableListOf<Reflection>()
listForConvStmt.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" }
}
deleteOlderThanStmt.reset()
deleteOlderThanStmt.clearBindings()
deleteOlderThanStmt.bindLong(1, cutoff.toEpochMilliseconds())
deleteOlderThanStmt.executeUpdate()
}
}
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
}
countStmt.reset()
countStmt.clearBindings()
countStmt.executeQuery().use { rs ->
if (rs.next()) (rs.getLong(0) ?: 0L).toInt() else 0
}
}
}
override fun events(): Flow<ReflectionEvent> = ev.asSharedFlow()
override fun close() {}
override fun close() {
insertStmt.close()
getStmt.close()
listRecentStmt.close()
listForConvStmt.close()
deleteOlderThanStmt.close()
countStmt.close()
}
private fun SQLiteResultSet.toDomain(): Reflection = Reflection(
id = getText(0)!!,
@@ -9,13 +9,13 @@ import pw.binom.db.ksqlite.SQLiteConnection
/**
* Фабрика 4 store'ов поверх ksqlite.
*
* Lifecycle: открывает [SQLiteConnection], гарантирует наличие таблиц
* (CREATE TABLE IF NOT EXISTS), возвращает bundle из 4 store'ов. Caller
* ДОЛЖЕН вызвать [close] при завершении.
* Lifecycle: открывает [SQLiteConnection], прогоняет [Schema.migrate] (создаёт
* таблицы/индексы если их нет, догоняет версию схемы до [Schema.CURRENT_VERSION]),
* возвращает bundle из 4 store'ов. Caller ДОЛЖЕН вызвать [close] при завершении.
*
* @param path путь к .db файлу, либо URI для in-memory/shared-cache.
*/
class KsqliteStores private constructor(
class KsqliteStores internal constructor(
val connection: SQLiteConnection,
val conversations: ConversationStore,
val messages: MutableMessageStore,
@@ -34,59 +34,15 @@ class KsqliteStores private constructor(
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);
"""
fun open(path: String): KsqliteStores {
val conn = SQLiteConnection.open(path)
conn.exec(SCHEMA)
Schema.migrate(conn)
return assemble(conn)
}
fun inMemory(name: String = "agentik-test"): KsqliteStores {
val conn = SQLiteConnection.memory(name)
conn.exec(SCHEMA)
Schema.migrate(conn)
return assemble(conn)
}
@@ -6,6 +6,8 @@ import pw.binom.agentik.workingMemory.WorkingMemoryRow
import pw.binom.agentik.workingMemory.WorkingMemoryStore
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
@@ -19,59 +21,99 @@ class KsqliteWorkingMemoryStore(
private val mutex = Mutex()
private val json = Json { ignoreUnknownKeys = true }
// pre-prepare (см. KsqliteMessageStore KDoc — почему это критично против
// SIGSEGV в StmtHolder.finalize на закрытой connection).
private val insertStmt = 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 = 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 = connection.prepare(
"DELETE FROM ${Schema.TABLE_WORKING_MEMORY} WHERE ${Schema.COL_CONVERSATION_ID} = ?"
)
private val maxOrderIdxStmt = connection.prepare(
"""
SELECT COALESCE(MAX(${Schema.COL_ORDER_IDX}), 0)
FROM ${Schema.TABLE_WORKING_MEMORY}
WHERE ${Schema.COL_CONVERSATION_ID} = ?
""".trimIndent()
)
private val dropFromIdxStmt = connection.prepare(
"""
DELETE FROM ${Schema.TABLE_WORKING_MEMORY}
WHERE ${Schema.COL_CONVERSATION_ID} = ? AND ${Schema.COL_ORDER_IDX} >= ?
""".trimIndent()
)
private val insertSummaryStmt = 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
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()
}
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 {
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<WorkingMemoryRow>()
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)!!),
)
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
}
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()
}
clearStmt.reset()
clearStmt.clearBindings()
clearStmt.bindText(1, conversationId)
clearStmt.executeUpdate()
}
}
@@ -82,51 +124,56 @@ class KsqliteWorkingMemoryStore(
): Long = withContext(Dispatchers.Default) {
mutex.withLock {
var newMax = 0L
val nowMs = System.currentTimeMillis()
val nowMs = Clock.System.now().toEpochMilliseconds()
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()
}
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
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()
}
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) {
connection.exec("ROLLBACK")
runCatching { connection.exec("ROLLBACK") }
throw t
}
newMax
}
}
override fun close() {}
override fun close() {
insertStmt.close()
listStmt.close()
clearStmt.close()
maxOrderIdxStmt.close()
dropFromIdxStmt.close()
insertSummaryStmt.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
}
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
}
@@ -0,0 +1,188 @@
package pw.binom.agentik.storage.ksqlite
import pw.binom.db.ksqlite.SQLiteConnection
/**
* Имена таблиц/колонок/индексов для ksqlite-бэкенда agentik'а.
*
* Все DDL/DML через ksqlite должны ссылаться на эти константы — никаких
* хардкоженных литералов в `prepare("SELECT ... FROM foo ...")` в каждом
* store'е. Это (1) даёт единую точку правды при будущих миграциях и
* (2) делает rename'ы безопасными (компилятор поймает все использования).
*/
internal object Schema {
/** Версия схемы. Увеличивать при ЛЮБОМ изменении DDL. */
const val CURRENT_VERSION: Int = 1
// ───── Таблицы ─────
const val TABLE_CONVERSATION = "conversation"
const val TABLE_MESSAGE = "message"
const val TABLE_WORKING_MEMORY = "working_memory"
const val TABLE_REFLECTION = "reflection"
const val TABLE_MIGRATION = "migration" // (резерв на будущее, сейчас версия в user_version)
// ───── Колонки conversation ─────
const val COL_ID = "id"
const val COL_TITLE = "title"
const val COL_IS_TEMPORAL = "is_temporal"
const val COL_CREATED_AT = "created_at"
const val COL_UPDATED_AT = "updated_at"
// ───── Колонки message ─────
const val COL_CONVERSATION_ID = "conversation_id"
const val COL_KIND = "kind"
const val COL_PAYLOAD_JSON = "payload_json"
// ───── Колонки working_memory ─────
const val COL_ORDER_IDX = "order_idx"
const val COL_SOURCE_MESSAGE_ID = "source_message_id"
// ───── Колонки reflection ─────
const val COL_TURNS_ANALYZED = "turns_analyzed"
const val COL_SCORE = "score"
const val COL_SUMMARY = "summary"
const val COL_WEAK_SPOTS_JSON = "weak_spots_json"
// ───── Индексы ─────
const val IDX_CONV_UPDATED = "idx_conv_updated"
const val IDX_MSG_CONV = "idx_msg_conv"
const val IDX_WM_UNIQUE = "idx_wm_unique"
const val IDX_WM_CONV = "idx_wm_conv"
const val IDX_REFLECTION_CREATED = "idx_reflection_created"
const val IDX_REFLECTION_CONV = "idx_reflection_conv"
/**
* Прогоняет миграцию схемы до [CURRENT_VERSION] на пустой или существующей БД.
*
* Версия хранится в `PRAGMA user_version` (стандартный SQLite-механизм,
* 32-bit int в заголовке БД — без своей таблицы). Каждая миграция —
* блок DDL+данных под номером `fromV+1`, выполняется в транзакции.
*
* Гарантии:
* - идемпотентность: повторный вызов после достижения текущей версии —
* no-op (`PRAGMA user_version` совпадает с CURRENT_VERSION);
* - атомарность: каждая миграция в BEGIN/COMMIT — упал посреди →
* PRAGMA остаётся на предыдущей версии, БД консистентна;
* - CREATE TABLE/INDEX через `IF NOT EXISTS` — безопасно на partially-
* migrated БД (если ручной ROLLBACK оставил схему в полупосаженном виде).
*/
fun migrate(conn: SQLiteConnection) {
val current = readUserVersion(conn)
if (current >= CURRENT_VERSION) return
// Миграции строго последовательны — каждая стартует с (current) и
// выставляет user_version = current+1 в конце (внутри транзакции).
if (current < 1) {
conn.exec("BEGIN")
try {
conn.exec(v1ConversationDdl)
conn.exec(v1MessageDdl)
conn.exec(v1WorkingMemoryDdl)
conn.exec(v1ReflectionDdl)
conn.exec(v1IndexesDdl)
writeUserVersion(conn, 1)
conn.exec("COMMIT")
} catch (t: Throwable) {
runCatching { conn.exec("ROLLBACK") }
throw t
}
}
// sanity check — после всех миграций обязаны достичь CURRENT_VERSION
check(readUserVersion(conn) == CURRENT_VERSION) {
"Schema migration failed to reach version $CURRENT_VERSION"
}
}
private fun readUserVersion(conn: SQLiteConnection): Int {
// user_version живёт в sqlite_master-метаданных; PRAGMA возвращает его
// как row column "user_version". Используем прямой SELECT к internal
// pragma function через ksqlite (prepared + executeQuery).
var version = 0
conn.prepare("PRAGMA user_version").use { stmt ->
stmt.executeQuery().use { rs ->
if (rs.next()) {
version = (rs.getLong(0) ?: 0L).toInt()
}
}
}
return version
}
private fun writeUserVersion(conn: SQLiteConnection, version: Int) {
// SQLite PRAGMA с literal-аргументом нельзя параметризовать через `?`,
// поэтому собираем SQL строкой (значение контролируемое, не user input).
conn.exec("PRAGMA user_version = $version")
}
// ───── DDL миграций ─────
// Ниже идут блоки по одной миграции. Конкатенация в [v1*] — потому что
// v1 — начальная схема (нет pre-existing DB с user_version=0 в проде,
// но мы поддерживаем эту ветку на случай dev-БД под `agentik-dev.db`).
private val v1ConversationDdl = """
CREATE TABLE IF NOT EXISTS $TABLE_CONVERSATION (
$COL_ID TEXT NOT NULL PRIMARY KEY,
$COL_TITLE TEXT,
$COL_IS_TEMPORAL INTEGER NOT NULL DEFAULT 0,
$COL_CREATED_AT INTEGER NOT NULL,
$COL_UPDATED_AT INTEGER NOT NULL
);
"""
private val v1MessageDdl = """
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
);
"""
private val v1WorkingMemoryDdl = """
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
);
"""
private val v1ReflectionDdl = """
CREATE TABLE IF NOT EXISTS $TABLE_REFLECTION (
$COL_ID TEXT NOT NULL PRIMARY KEY,
$COL_CONVERSATION_ID TEXT,
$COL_CREATED_AT INTEGER NOT NULL,
$COL_TURNS_ANALYZED INTEGER NOT NULL,
$COL_SCORE INTEGER NOT NULL,
$COL_SUMMARY TEXT NOT NULL,
$COL_WEAK_SPOTS_JSON TEXT NOT NULL DEFAULT '[]'
);
"""
private val v1IndexesDdl = """
CREATE INDEX IF NOT EXISTS $IDX_CONV_UPDATED
ON $TABLE_CONVERSATION($COL_UPDATED_AT DESC);
-- Главный hot-path индекс для list/conversation: фильтр по conv +
-- сортировка по created_at (используется list(), cascade-delete, etc.)
CREATE INDEX IF NOT EXISTS $IDX_MSG_CONV
ON $TABLE_MESSAGE($COL_CONVERSATION_ID, $COL_CREATED_AT);
-- Working memory: гарантия уникального order_idx внутри conv'а
-- (порядок имеет значение — compaction полагается на монотонность).
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);
CREATE INDEX IF NOT EXISTS $IDX_REFLECTION_CREATED
ON $TABLE_REFLECTION($COL_CREATED_AT DESC);
CREATE INDEX IF NOT EXISTS $IDX_REFLECTION_CONV
ON $TABLE_REFLECTION($COL_CONVERSATION_ID, $COL_CREATED_AT DESC);
"""
}
@@ -1,5 +1,6 @@
package pw.binom.agentik.storage.ksqlite
import kotlinx.coroutines.flow.toList
import kotlinx.coroutines.test.runTest
import pw.binom.agentik.messageLog.Content
import pw.binom.agentik.messageLog.MessageRecord
@@ -8,8 +9,6 @@ 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 {
@@ -41,7 +40,7 @@ class KsqliteMessageStoreTest {
createdAt = Instant.parse("2026-09-15T10:01:00Z"),
context = null,
))
val list = stores.messages.listAll("conv1")
val list = stores.messages.listFlow("conv1", Instant.DISTANT_PAST).toList()
assertEquals(1, list.size)
val msg = list[0]
assertEquals("m1", msg.id)
@@ -57,10 +56,10 @@ class KsqliteMessageStoreTest {
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)
val list = stores.messages.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
@@ -78,20 +77,20 @@ class KsqliteMessageStoreTest {
}
@Test
fun testListAllReturnsAllInOrder() = runTest {
fun testListFlowReturnsAllInOrder() = 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 })
assertEquals(listOf("m1", "m2", "m3"), stores.messages.listFlow("conv1", Instant.DISTANT_PAST).toList().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)
fun testClearRemovesByConversation() = runTest {
stores.messages.append(MessageRecord.UserMessage("m1", "conv1", listOf(Content.Text("a")), Instant.parse("2026-09-15T10:00:00Z"), null))
stores.messages.append(MessageRecord.UserMessage("m2", "conv2", listOf(Content.Text("b")), Instant.parse("2026-09-15T10:00:00Z"), null))
stores.conversations.delete("conv1")
assertEquals(emptyList(), stores.messages.listFlow("conv1", Instant.DISTANT_PAST).toList())
assertEquals(1, stores.messages.listFlow("conv2", Instant.DISTANT_PAST).toList().size)
}
}
@@ -0,0 +1,173 @@
package pw.binom.agentik.storage.ksqlite
import kotlinx.coroutines.test.runTest
import pw.binom.db.ksqlite.SQLiteConnection
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertTrue
/**
* Тесты на Schema.migrate():
* - fresh DB → создаются все 4 таблицы + индексы + user_version = CURRENT_VERSION;
* - уже мигрированная БД → migrate() идемпотентен (no-op, не падает на
* повторных CREATE);
* - DB, открытая напрямую через SQLiteConnection (минуя KsqliteStores),
* migrate() приводит её в боевое состояние.
*
* Также проверяем что наличие индекса idx_msg_conv (conversation_id +
* created_at) — обязательный hot-path для list()/cascade-delete.
*/
class SchemaMigrationTest {
@Test
fun `fresh DB gets all tables indexes and CURRENT_VERSION`() = runTest {
val conn = SQLiteConnection.memory("mig-fresh-${kotlin.random.Random.nextLong()}")
try {
Schema.migrate(conn)
for (table in listOf(
Schema.TABLE_CONVERSATION,
Schema.TABLE_MESSAGE,
Schema.TABLE_WORKING_MEMORY,
Schema.TABLE_REFLECTION,
)) {
assertTrue(tableExists(conn, table), "table '$table' should exist after migrate()")
}
for (index in listOf(
Schema.IDX_CONV_UPDATED,
Schema.IDX_MSG_CONV,
Schema.IDX_WM_UNIQUE,
Schema.IDX_WM_CONV,
Schema.IDX_REFLECTION_CREATED,
Schema.IDX_REFLECTION_CONV,
)) {
assertTrue(indexExists(conn, index), "index '$index' should exist after migrate()")
}
assertEquals(Schema.CURRENT_VERSION, readUserVersion(conn))
} finally {
conn.close()
}
}
@Test
fun `migrate is idempotent on already-migrated DB`() = runTest {
val conn = SQLiteConnection.memory("mig-idem-${kotlin.random.Random.nextLong()}")
try {
Schema.migrate(conn)
val versionAfterFirst = readUserVersion(conn)
// повторный вызов не должен ни упасть, ни изменить версию, ни
// пересоздать таблицы/индексы (CREATE IF NOT EXISTS — no-op)
Schema.migrate(conn)
assertEquals(versionAfterFirst, readUserVersion(conn))
} finally {
conn.close()
}
}
@Test
fun `raw SQLiteConnection plus migrate gives working bundle`() = runTest {
// Имитируем сценарий: существующая БД без schema, открываем через
// ksqlite и прогоняем migrate руками (тот же путь, что в
// KsqliteStores.open, но без зависимости от фабрики).
val conn = SQLiteConnection.memory("mig-bundle-${kotlin.random.Random.nextLong()}")
Schema.migrate(conn)
// Сборка bundle через internal-конструктор — KsqliteStores primary
// constructor internal, тест в том же модуле и может его звать.
val stores = KsqliteStores(
connection = conn,
conversations = KsqliteConversationStore(conn),
messages = KsqliteMessageStore(conn),
workingMemory = KsqliteWorkingMemoryStore(conn),
reflections = KsqliteReflectionStore(conn),
)
try {
// bundle работает end-to-end — conversation upsert + message append +
// list. Никаких "no such table" или подобного.
stores.conversations.upsert(
pw.binom.agentik.messageStore.ConversationRecord(
id = "c1", title = "t", isTemporal = false,
createdAt = kotlin.time.Instant.parse("2026-09-15T10:00:00Z"),
updatedAt = kotlin.time.Instant.parse("2026-09-15T10:00:00Z"),
)
)
stores.messages.append(
pw.binom.agentik.messageLog.MessageRecord.UserMessage(
id = "m1", conversationId = "c1",
content = listOf(pw.binom.agentik.messageLog.Content.Text("hi")),
createdAt = kotlin.time.Instant.parse("2026-09-15T10:00:01Z"),
)
)
val got = stores.messages.list("c1", kotlin.time.Instant.DISTANT_PAST, offset = 0, limit = 10)
assertEquals(1, got.size)
assertEquals("m1", got[0].id)
} finally {
// Закрываем store'ы → они закроют свои pre-prepared statements
// (StmtHolder.finalize увидит isOpen == false и не полезет в
// нативный sqlite3_finalize с уже-разрушенным db mutex).
stores.close()
}
}
@Test
fun `idx_msg_conv covers conversation_id and created_at columns`() = runTest {
// Проверяем что индекс действительно покрывает обе колонки — без
// этого list()/cascade-delete будут делать full-scan по message.
val conn = SQLiteConnection.memory("mig-idx-${kotlin.random.Random.nextLong()}")
try {
Schema.migrate(conn)
val cols = indexColumns(conn, Schema.IDX_MSG_CONV)
assertEquals(listOf(Schema.COL_CONVERSATION_ID, Schema.COL_CREATED_AT), cols)
} finally {
conn.close()
}
}
private fun tableExists(conn: SQLiteConnection, name: String): Boolean {
conn.prepare(
"SELECT 1 FROM sqlite_master WHERE type = 'table' AND name = ?"
).use { stmt ->
stmt.bindText(1, name)
stmt.executeQuery().use { rs -> return rs.next() }
}
}
private fun indexExists(conn: SQLiteConnection, name: String): Boolean {
conn.prepare(
"SELECT 1 FROM sqlite_master WHERE type = 'index' AND name = ?"
).use { stmt ->
stmt.bindText(1, name)
stmt.executeQuery().use { rs -> return rs.next() }
}
}
private fun indexColumns(conn: SQLiteConnection, indexName: String): List<String> {
// PRAGMA index_info возвращает одну строку на колонку индекса
// (seqno, cid, name). Параметризовать через `?` нельзя — собираем
// строку (name — контролируемая константа, не user input).
val cols = mutableListOf<String>()
conn.prepare("PRAGMA index_info($indexName)").use { stmt ->
stmt.executeQuery().use { rs ->
while (rs.next()) {
rs.getText(2)?.let(cols::add)
}
}
}
return cols
}
private fun readUserVersion(conn: SQLiteConnection): Int {
var version = 0
conn.prepare("PRAGMA user_version").use { stmt ->
stmt.executeQuery().use { rs ->
if (rs.next()) {
version = (rs.getLong(0) ?: 0L).toInt()
}
}
}
return version
}
}