storage-sqlite: выделить SQLDelight + SQLite-импл в отдельный модуль
Перенесён SQLDelight (4 .sq файла, конфигурация databases { AgentikDatabase })
и 5 SQLite-импл классов (SqliteStores, SqliteConversationStore, SqliteMessageStore,
SqliteWorkingMemoryStore, SqliteReflectionStore) из :standalone в новый JVM-only
модуль :storage-sqlite под пакетом pw.binom.agentik.storage.sqlite.
Изменения:
- Новый :storage-sqlite модуль с sqldelight-плагином + sqlite JDBC driver
- Все .sq файлы и Kotlin-классы переехали с переименованием пакета
- ReflectionStore.kt в :standalone (только SQLite-импл) удалён — функционал
живёт в :storage-sqlite/SqliteReflectionStore.kt
- :standalone/build.gradle.kts: убран sqldelight-плагин и конфигурация,
добавлена зависимость :storage-sqlite
- Все импорты в :standalone (8 main + 7 test) перенаправлены на новый пакет
- ReflectionStoreTest.kt переехал в :storage-sqlite/jvmTest (тестирует
internal fun encode/decodeStringArray в :storage-sqlite)
Совместимость:
- SqliteStores доступен по новому пути pw.binom.agentik.storage.sqlite.SqliteStores
- Старые импорты в тестах обновлены (минимум diff — 1 строка на файл)
- В commit 6 ChatAgent переключится на StorageBundle API; SqliteStores
станет деталью реализации :standalone
Тесты: 299/299 green. Fatjar standalone-all.jar 240 MB.
Преимущества:
- :standalone больше не зависит от SQLDelight плагина (легче поддерживать)
- :storage-sqlite может быть заменён/расширен (например, :storage-android)
- Тесты storage-слоя сгруппированы по модулю реализации
This commit is contained in:
@@ -0,0 +1,44 @@
|
||||
plugins {
|
||||
alias(libs.plugins.kotlin.multiplatform)
|
||||
alias(libs.plugins.kotlin.serialization)
|
||||
alias(libs.plugins.sqldelight)
|
||||
}
|
||||
|
||||
kotlin {
|
||||
jvmToolchain(21)
|
||||
|
||||
// JVM-only: SQLDelight сейчас не имеет KMP-таргетов за пределами JVM/Android.
|
||||
// Если в будущем понадобится Native (iOS/macOS), придётся либо подключать
|
||||
// platform-native SQLDelight driver (https://github.com/cashapp/sqldelight
|
||||
// multiplatform module), либо использовать :storage-inmemory там.
|
||||
jvm()
|
||||
|
||||
sourceSets {
|
||||
commonMain.dependencies {
|
||||
api(project(":storage-core"))
|
||||
api(libs.sqldelight.runtime)
|
||||
api(libs.sqldelight.coroutines)
|
||||
}
|
||||
jvmMain.dependencies {
|
||||
implementation(libs.sqldelight.sqlite.driver)
|
||||
implementation(libs.kotlin.logging)
|
||||
}
|
||||
commonTest.dependencies {
|
||||
implementation(kotlin("test"))
|
||||
implementation(libs.kotlinx.coroutines.test)
|
||||
}
|
||||
jvmTest.dependencies {
|
||||
// SQLDelight JDBC driver — для тестов, которые поднимают in-memory DB
|
||||
implementation(libs.sqldelight.sqlite.driver)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
sqldelight {
|
||||
databases {
|
||||
create("AgentikDatabase") {
|
||||
packageName.set("pw.binom.agentik.storage.sqlite")
|
||||
srcDirs.setFrom("src/jvmMain/sqldelight")
|
||||
}
|
||||
}
|
||||
}
|
||||
+80
@@ -0,0 +1,80 @@
|
||||
package pw.binom.agentik.storage.sqlite
|
||||
|
||||
import kotlin.time.Instant
|
||||
import pw.binom.agentik.storage.ConversationRecord
|
||||
import pw.binom.agentik.storage.ConversationStore
|
||||
|
||||
/**
|
||||
* SQLite-реализация [ConversationStore].
|
||||
*/
|
||||
class SqliteConversationStore(private val db: AgentikDatabase) : ConversationStore {
|
||||
|
||||
private val q get() = db.conversationQueries
|
||||
|
||||
override suspend fun upsert(record: ConversationRecord) {
|
||||
// SQLite-конфликт по PRIMARY KEY → сначала пробуем insert, при ошибке → update
|
||||
val existing = q.getById(record.id).executeAsOneOrNull()
|
||||
if (existing == null) {
|
||||
q.insert(
|
||||
id = record.id,
|
||||
title = record.title,
|
||||
is_temporal = if (record.isTemporal) 1L else 0L,
|
||||
created_at = record.createdAt.toEpochMilliseconds(),
|
||||
updated_at = record.updatedAt.toEpochMilliseconds(),
|
||||
)
|
||||
} else {
|
||||
q.update(
|
||||
title = record.title,
|
||||
is_temporal = if (record.isTemporal) 1L else 0L,
|
||||
updated_at = record.updatedAt.toEpochMilliseconds(),
|
||||
id = record.id,
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
override suspend fun get(id: String): ConversationRecord? {
|
||||
val row = q.getById(id).executeAsOneOrNull() ?: return null
|
||||
return row.toRecord()
|
||||
}
|
||||
|
||||
override suspend fun delete(id: String): Boolean {
|
||||
// Проверяем существование ДО удаления — иначе пустой `deleteById` вернёт
|
||||
// «успех», и get() == null будет true, хотя диалога и не было.
|
||||
val existed = q.getById(id).executeAsOneOrNull() != null
|
||||
if (!existed) return false
|
||||
db.transaction {
|
||||
db.messageQueries.deleteByConversation(id)
|
||||
db.workingMemoryQueries.clearByConversation(id)
|
||||
q.deleteById(id)
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
override suspend fun list(offset: Int, limit: Int): List<ConversationRecord> =
|
||||
q.list(limit = limit.toLong(), offset = offset.toLong())
|
||||
.executeAsList()
|
||||
.map { it.toRecord() }
|
||||
|
||||
override suspend fun rename(id: String, title: String?): Instant? {
|
||||
val now = Instant.fromEpochMilliseconds(System.currentTimeMillis())
|
||||
q.rename(title = title, updated_at = now.toEpochMilliseconds(), id = id)
|
||||
val ts = q.getUpdatedAt(id).executeAsOneOrNull() ?: return null
|
||||
return Instant.fromEpochMilliseconds(ts)
|
||||
}
|
||||
|
||||
override suspend fun touch(id: String, now: Instant) {
|
||||
q.touch(updated_at = now.toEpochMilliseconds(), id = id)
|
||||
}
|
||||
|
||||
override fun close() {
|
||||
// driver закрывается во внешнем SqliteStores
|
||||
}
|
||||
}
|
||||
|
||||
private fun Conversation.toRecord(): ConversationRecord = ConversationRecord(
|
||||
id = id,
|
||||
title = title,
|
||||
isTemporal = is_temporal != 0L,
|
||||
createdAt = Instant.fromEpochMilliseconds(created_at),
|
||||
updatedAt = Instant.fromEpochMilliseconds(updated_at),
|
||||
)
|
||||
+171
@@ -0,0 +1,171 @@
|
||||
package pw.binom.agentik.storage.sqlite
|
||||
|
||||
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 kotlin.time.Instant
|
||||
|
||||
/**
|
||||
* SQLite-реализация [MessageStore] (append-only audit log).
|
||||
*
|
||||
* `payloadJson` хранит JSON-сериализованные kind-specific поля. Для
|
||||
* `user`/`assistant` это `List<Content>` (см. [encodeBodyPayload]).
|
||||
* Для `tool_call`/`tool_result` payload хранит JSON-объект
|
||||
* (см. [CallPayload]/[ResultPayload]).
|
||||
*
|
||||
* SQLDelight сохраняет snake_case в сгенерированной data class (`Message`),
|
||||
* поэтому обращаемся через `conversation_id`, `payload_json`, `created_at`.
|
||||
*/
|
||||
class SqliteMessageStore(private val db: AgentikDatabase) : MessageStore {
|
||||
|
||||
private val q get() = db.messageQueries
|
||||
|
||||
override suspend fun append(record: MessageRecord) {
|
||||
val (kind, payload) = encodeRecord(record)
|
||||
q.insert(
|
||||
conversation_id = record.conversationId,
|
||||
kind = kind,
|
||||
payload_json = payload,
|
||||
created_at = record.createdAt.toEpochMilliseconds(),
|
||||
id = record.id,
|
||||
)
|
||||
}
|
||||
|
||||
override suspend fun list(
|
||||
conversationId: String,
|
||||
after: Instant,
|
||||
offset: Int,
|
||||
limit: Int,
|
||||
): List<MessageRecord> = q.listAfter(
|
||||
conversation_id = conversationId,
|
||||
created_at = after.toEpochMilliseconds(),
|
||||
limit = limit.toLong(),
|
||||
offset = offset.toLong(),
|
||||
).executeAsList().map { it.toRecord() }
|
||||
|
||||
override suspend fun listAll(conversationId: String): List<MessageRecord> =
|
||||
q.listByConversationAll(conversation_id = conversationId).executeAsList().map { it.toRecord() }
|
||||
|
||||
override suspend fun tokenStats(conversationId: String): TokenStats {
|
||||
// Загружаем все assistant-сообщения и считаем локально. Для >10K ходов
|
||||
// это можно оптимизировать SQL aggregation с JSON_EXTRACT, но пока
|
||||
// узких мест нет — реальные диалоги редко длиннее нескольких сотен ходов.
|
||||
val messages = q.listByConversationAll(conversation_id = conversationId).executeAsList()
|
||||
var turns = 0
|
||||
var inputTotal = 0L
|
||||
var outputTotal = 0L
|
||||
for (m in messages) {
|
||||
if (m.kind != "assistant") continue
|
||||
val decoded = try {
|
||||
decodeBodyPayload(m.payload_json)
|
||||
} catch (_: Throwable) {
|
||||
// skip malformed legacy payloads
|
||||
continue
|
||||
}
|
||||
val tokens = decoded.tokens ?: continue
|
||||
turns++
|
||||
inputTotal += tokens.input
|
||||
outputTotal += tokens.output
|
||||
}
|
||||
return TokenStats(turns = turns, inputTokens = inputTotal, outputTokens = outputTotal)
|
||||
}
|
||||
|
||||
override fun close() {}
|
||||
}
|
||||
|
||||
private fun encodeRecord(record: MessageRecord): Pair<String, String> = 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")
|
||||
}
|
||||
|
||||
@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?)
|
||||
|
||||
private fun Message.toRecord(): MessageRecord {
|
||||
val id = id
|
||||
val convId = conversation_id
|
||||
val createdAt = Instant.fromEpochMilliseconds(created_at)
|
||||
return when (kind) {
|
||||
"user" -> {
|
||||
val decoded = decodeBodyPayload(payload_json)
|
||||
MessageRecord.UserMessage(
|
||||
id = id,
|
||||
conversationId = convId,
|
||||
content = decoded.content,
|
||||
createdAt = createdAt,
|
||||
context = decoded.context,
|
||||
)
|
||||
}
|
||||
"assistant" -> {
|
||||
val decoded = decodeBodyPayload(payload_json)
|
||||
MessageRecord.AssistantMessage(
|
||||
id = id,
|
||||
conversationId = convId,
|
||||
content = decoded.content,
|
||||
createdAt = createdAt,
|
||||
tokens = decoded.tokens,
|
||||
)
|
||||
}
|
||||
"tool_call" -> {
|
||||
val p = Json.decodeFromString(CallPayload.serializer(), payload_json)
|
||||
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_json)
|
||||
MessageRecord.ToolResult(
|
||||
id = id,
|
||||
conversationId = convId,
|
||||
toolCallId = p.toolCallId,
|
||||
result = p.result,
|
||||
createdAt = createdAt,
|
||||
)
|
||||
}
|
||||
"error" -> {
|
||||
val p = Json.decodeFromString(ErrorPayload.serializer(), payload_json)
|
||||
MessageRecord.Error(
|
||||
id = id,
|
||||
conversationId = convId,
|
||||
message = p.message,
|
||||
code = p.code,
|
||||
createdAt = createdAt,
|
||||
)
|
||||
}
|
||||
else -> error("Unknown message kind in audit log: $kind")
|
||||
}
|
||||
}
|
||||
+144
@@ -0,0 +1,144 @@
|
||||
package pw.binom.agentik.storage.sqlite
|
||||
|
||||
import kotlin.time.Clock
|
||||
import kotlin.time.Instant
|
||||
import kotlinx.coroutines.flow.Flow
|
||||
import kotlinx.coroutines.flow.MutableSharedFlow
|
||||
import kotlinx.coroutines.flow.asSharedFlow
|
||||
import mu.KotlinLogging
|
||||
import pw.binom.agentik.storage.Reflection
|
||||
import pw.binom.agentik.storage.ReflectionEvent
|
||||
import pw.binom.agentik.storage.ReflectionStore
|
||||
import pw.binom.agentik.storage.sqlite.AgentikDatabase
|
||||
import pw.binom.agentik.storage.sqlite.Reflection as DbReflection
|
||||
import pw.binom.agentik.storage.sqlite.ReflectionQueries
|
||||
|
||||
private val log = KotlinLogging.logger {}
|
||||
|
||||
/**
|
||||
* SQLDelight-реализация [ReflectionStore].
|
||||
*
|
||||
* `weakSpots` хранится как JSON-массив (маленький массив строк, ~ десяток);
|
||||
* парсим минимальным ручным сканером, чтобы не тащить kotlinx-serialization
|
||||
* в этот модуль (он уже есть в commonMain через standalone — но не хочется
|
||||
* парсить JSON ради одного массива строк).
|
||||
*
|
||||
* В commit 3 переедет в `:storage-sqlite`. Сейчас остаётся в `:standalone`,
|
||||
* потому что зависит от SQLDelight'ной `AgentikDatabase`, которую мы ещё
|
||||
* не отвязали от `:standalone` (`generateMainDatabase` нужно подключить к
|
||||
* новому модулю).
|
||||
*/
|
||||
class SqliteReflectionStore(
|
||||
private val db: AgentikDatabase,
|
||||
private val clock: Clock = Clock.System,
|
||||
) : ReflectionStore {
|
||||
private val queries: ReflectionQueries get() = db.reflectionQueries
|
||||
private val ev = MutableSharedFlow<ReflectionEvent>(extraBufferCapacity = 16)
|
||||
|
||||
override suspend fun insert(reflection: Reflection) {
|
||||
log.debug { "insert reflection id=${reflection.id} conv=${reflection.conversationId} score=${reflection.score}" }
|
||||
queries.insert(
|
||||
id = reflection.id,
|
||||
conversation_id = reflection.conversationId,
|
||||
created_at = reflection.createdAt.toEpochMilliseconds(),
|
||||
turns_analyzed = reflection.turnsAnalyzed.toLong(),
|
||||
score = reflection.score.toLong(),
|
||||
summary = reflection.summary,
|
||||
weak_spots_json = encodeStringArray(reflection.weakSpots),
|
||||
)
|
||||
ev.tryEmit(ReflectionEvent.Created(reflection))
|
||||
}
|
||||
|
||||
override suspend fun get(id: String): Reflection? =
|
||||
queries.getById(id).executeAsOneOrNull()?.toDomain()
|
||||
|
||||
override suspend fun listRecent(limit: Int): List<Reflection> =
|
||||
queries.listRecent(limit.toLong()).executeAsList().map { it.toDomain() }
|
||||
|
||||
override suspend fun listForConversation(conversationId: String, limit: Int): List<Reflection> =
|
||||
queries.listForConversation(conversationId, limit.toLong()).executeAsList().map { it.toDomain() }
|
||||
|
||||
override suspend fun deleteOlderThan(cutoff: Instant) {
|
||||
val n = queries.deleteOlderThan(cutoff.toEpochMilliseconds()).value
|
||||
if (n > 0) log.info { "deleted $n reflections older than $cutoff" }
|
||||
}
|
||||
|
||||
override suspend fun count(): Int = queries.count().executeAsOne().toInt()
|
||||
|
||||
override fun events(): Flow<ReflectionEvent> = ev.asSharedFlow()
|
||||
|
||||
override fun close() {
|
||||
// no-op: lifecycle owned by AgentikDatabase / SqliteStores
|
||||
}
|
||||
|
||||
private fun DbReflection.toDomain(): Reflection = Reflection(
|
||||
id = id,
|
||||
conversationId = conversation_id,
|
||||
createdAt = Instant.fromEpochMilliseconds(created_at),
|
||||
turnsAnalyzed = turns_analyzed.toInt(),
|
||||
score = score.toInt(),
|
||||
summary = summary,
|
||||
weakSpots = decodeStringArray(weak_spots_json),
|
||||
)
|
||||
}
|
||||
|
||||
/**
|
||||
* Простой JSON-encode списка строк как `"[\"a\",\"b\"]"`.
|
||||
*
|
||||
* Hand-rolled чтобы не тащить kotlinx-serialization ради одного типа; строки
|
||||
* экранируются только от `"` и `\` (этого достаточно для текста заметок).
|
||||
*/
|
||||
internal fun encodeStringArray(items: List<String>): 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(']')
|
||||
}
|
||||
|
||||
/**
|
||||
* Минимальный JSON-парсер для массива строк: ожидает форму `[...]`,
|
||||
* внутри — строки в `"..."` с `\"`/`\\` escape. Любые невалидные символы
|
||||
* возвращаются как пустой список — лучше пустая рефлексия, чем exception.
|
||||
*/
|
||||
internal fun decodeStringArray(raw: String): List<String> {
|
||||
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<String>()
|
||||
var i = 0
|
||||
while (i < inner.length) {
|
||||
while (i < inner.length && inner[i].isWhitespace() || inner[i] == ',') i++
|
||||
if (i >= inner.length) break
|
||||
if (inner[i] != '"') return emptyList()
|
||||
i++
|
||||
val sb = StringBuilder()
|
||||
while (i < inner.length && inner[i] != '"') {
|
||||
if (inner[i] == '\\' && i + 1 < inner.length) {
|
||||
when (inner[i + 1]) {
|
||||
'"' -> sb.append('"')
|
||||
'\\' -> sb.append('\\')
|
||||
else -> sb.append(inner[i + 1])
|
||||
}
|
||||
i += 2
|
||||
} else {
|
||||
sb.append(inner[i])
|
||||
i++
|
||||
}
|
||||
}
|
||||
if (i >= inner.length) return emptyList() // unterminated
|
||||
i++ // closing "
|
||||
out.add(sb.toString())
|
||||
}
|
||||
return out
|
||||
}
|
||||
@@ -0,0 +1,114 @@
|
||||
package pw.binom.agentik.storage.sqlite
|
||||
|
||||
import app.cash.sqldelight.db.QueryResult
|
||||
import app.cash.sqldelight.db.SqlDriver
|
||||
import app.cash.sqldelight.driver.jdbc.sqlite.JdbcSqliteDriver
|
||||
import pw.binom.agentik.storage.ConversationStore
|
||||
import pw.binom.agentik.storage.MessageStore
|
||||
import pw.binom.agentik.storage.ReflectionStore
|
||||
import pw.binom.agentik.storage.sqlite.SqliteReflectionStore
|
||||
import pw.binom.agentik.storage.WorkingMemoryStore
|
||||
|
||||
/**
|
||||
* Корневой объект SQLite-слоя: держит [SqlDriver] и три [WorkingMemoryStore]/[MessageStore]/[ConversationStore].
|
||||
* Закрывается вместе с приложением.
|
||||
*/
|
||||
class SqliteStores private constructor(
|
||||
val driver: SqlDriver,
|
||||
val conversations: ConversationStore,
|
||||
val messages: MessageStore,
|
||||
val workingMemory: WorkingMemoryStore,
|
||||
val reflections: ReflectionStore,
|
||||
) : AutoCloseable {
|
||||
|
||||
override fun close() {
|
||||
conversations.close()
|
||||
messages.close()
|
||||
workingMemory.close()
|
||||
reflections.close()
|
||||
driver.close()
|
||||
}
|
||||
|
||||
companion object {
|
||||
|
||||
/** Открыть/создать БД по пути `dbPath` (например, `"./agentik.db"` или абсолютный путь). */
|
||||
fun open(dbPath: String): SqliteStores {
|
||||
val driver = JdbcSqliteDriver("jdbc:sqlite:$dbPath")
|
||||
createSchema(driver)
|
||||
val db = AgentikDatabase(driver)
|
||||
return SqliteStores(
|
||||
driver = driver,
|
||||
conversations = SqliteConversationStore(db),
|
||||
messages = SqliteMessageStore(db),
|
||||
workingMemory = SqliteWorkingMemoryStore(db),
|
||||
reflections = SqliteReflectionStore(db),
|
||||
)
|
||||
}
|
||||
|
||||
/** Открыть/создать БД в памяти (для тестов). */
|
||||
fun inMemory(): SqliteStores {
|
||||
val driver = JdbcSqliteDriver(JdbcSqliteDriver.IN_MEMORY)
|
||||
createSchema(driver)
|
||||
val db = AgentikDatabase(driver)
|
||||
return SqliteStores(
|
||||
driver = driver,
|
||||
conversations = SqliteConversationStore(db),
|
||||
messages = SqliteMessageStore(db),
|
||||
workingMemory = SqliteWorkingMemoryStore(db),
|
||||
reflections = SqliteReflectionStore(db),
|
||||
)
|
||||
}
|
||||
|
||||
private fun createSchema(driver: SqlDriver) {
|
||||
// Если таблица `conversation` уже есть — БД уже инициализирована.
|
||||
// Всё равно прогоняем additive-миграции (см. runMigrations), потому что
|
||||
// в новых версиях могли появиться таблицы, которых нет в этой БД.
|
||||
val existing = driver.executeQuery(
|
||||
identifier = null,
|
||||
sql = "SELECT name FROM sqlite_master WHERE type='table' AND name='conversation'",
|
||||
mapper = { cursor ->
|
||||
QueryResult.Value(
|
||||
if (cursor.next().value) cursor.getString(0) else null,
|
||||
)
|
||||
},
|
||||
parameters = 0,
|
||||
).value
|
||||
if (existing == null) {
|
||||
// Свежая БД — пусть SQLDelight создаст всё сразу.
|
||||
AgentikDatabase.Schema.create(driver)
|
||||
return
|
||||
}
|
||||
// БД уже есть — прогоняем миграции для новых таблиц.
|
||||
runMigrations(driver)
|
||||
}
|
||||
|
||||
/**
|
||||
* Additive-миграции для таблиц, добавленных после первого релиза.
|
||||
* Каждая миграция — `CREATE TABLE IF NOT EXISTS ...`, идемпотентна.
|
||||
* Без Schema.version (SQLDelight migrations .sqm файлов) — для простой
|
||||
* additive-семантики это OK: удалять/переименовывать таблицы мы
|
||||
* всё равно пока не планируем.
|
||||
*/
|
||||
private fun runMigrations(driver: SqlDriver) {
|
||||
val migrations = listOf(
|
||||
// v2: self-reflection (см. Reflection.sq)
|
||||
"""
|
||||
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);
|
||||
""".trimIndent(),
|
||||
)
|
||||
for (sql in migrations) {
|
||||
driver.execute(null, sql, 0)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
+93
@@ -0,0 +1,93 @@
|
||||
package pw.binom.agentik.storage.sqlite
|
||||
|
||||
import kotlinx.serialization.json.Json
|
||||
import pw.binom.agentik.storage.WorkingMemoryEntry
|
||||
import pw.binom.agentik.storage.WorkingMemoryRow
|
||||
import pw.binom.agentik.storage.WorkingMemoryStore
|
||||
import kotlin.time.Instant
|
||||
|
||||
/**
|
||||
* SQLite-реализация [WorkingMemoryStore] (мутируемый LLM-контекст).
|
||||
*
|
||||
* Строки упорядочены по `order_idx ASC`. Compact — атомарная операция:
|
||||
* DELETE строк `>= dropFromOrderIdx` (без summary в v1).
|
||||
*/
|
||||
class SqliteWorkingMemoryStore(private val db: AgentikDatabase) : WorkingMemoryStore {
|
||||
|
||||
private val q get() = db.workingMemoryQueries
|
||||
private val json = Json { ignoreUnknownKeys = true }
|
||||
|
||||
override suspend fun append(conversationId: String, entry: WorkingMemoryEntry, now: Instant) {
|
||||
val newIdx = (q.maxOrderIdx(conversationId).executeAsOne()) + 1
|
||||
q.insert(
|
||||
id = newId(),
|
||||
conversation_id = conversationId,
|
||||
order_idx = newIdx,
|
||||
source_message_id = entry.sourceMessageId,
|
||||
kind = entryKind(entry),
|
||||
payload_json = json.encodeToString(WorkingMemoryEntry.serializer(), entry),
|
||||
created_at = now.toEpochMilliseconds(),
|
||||
)
|
||||
}
|
||||
|
||||
override suspend fun list(conversationId: String): List<WorkingMemoryRow> =
|
||||
q.listByConversation(conversationId).executeAsList().map { it.toRow() }
|
||||
|
||||
override suspend fun clear(conversationId: String) {
|
||||
q.clearByConversation(conversationId)
|
||||
}
|
||||
|
||||
override suspend fun compact(
|
||||
dropFromOrderIdx: Long,
|
||||
conversationId: String,
|
||||
summaryText: String?,
|
||||
): Long {
|
||||
var newMax = 0L
|
||||
val nowMs = System.currentTimeMillis()
|
||||
val summaryId = newId()
|
||||
db.transaction {
|
||||
q.compactDelete(conversation_id = conversationId, order_idx = dropFromOrderIdx)
|
||||
if (!summaryText.isNullOrBlank()) {
|
||||
val afterDelete = q.maxOrderIdx(conversationId).executeAsOne()
|
||||
val newIdx = afterDelete + 1
|
||||
q.compactInsert(
|
||||
id = summaryId,
|
||||
conversation_id = conversationId,
|
||||
order_idx = newIdx,
|
||||
payload_json = json.encodeToString(
|
||||
WorkingMemoryEntry.serializer(),
|
||||
WorkingMemoryEntry.Summary(text = summaryText),
|
||||
),
|
||||
created_at = nowMs,
|
||||
)
|
||||
newMax = newIdx
|
||||
} else {
|
||||
newMax = q.maxOrderIdx(conversationId).executeAsOne()
|
||||
}
|
||||
}
|
||||
return newMax
|
||||
}
|
||||
|
||||
override fun close() {}
|
||||
}
|
||||
|
||||
private fun entryKind(e: WorkingMemoryEntry): String = when (e) {
|
||||
is WorkingMemoryEntry.System -> "system"
|
||||
is WorkingMemoryEntry.User -> "user"
|
||||
is WorkingMemoryEntry.Assistant -> "assistant"
|
||||
is WorkingMemoryEntry.Summary -> "summary"
|
||||
}
|
||||
|
||||
private fun Working_memory.toRow(): WorkingMemoryRow {
|
||||
val entry = Json.decodeFromString(WorkingMemoryEntry.serializer(), payload_json)
|
||||
return WorkingMemoryRow(
|
||||
id = id,
|
||||
conversationId = conversation_id,
|
||||
orderIdx = order_idx,
|
||||
sourceMessageId = source_message_id,
|
||||
entry = entry,
|
||||
createdAt = Instant.fromEpochMilliseconds(created_at),
|
||||
)
|
||||
}
|
||||
|
||||
private fun newId(): String = pw.binom.agentik.storage.Ids.new("wm")
|
||||
@@ -0,0 +1,42 @@
|
||||
CREATE TABLE 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 idx_conv_updated ON conversation(updated_at DESC);
|
||||
|
||||
insert:
|
||||
INSERT INTO conversation (id, title, is_temporal, created_at, updated_at)
|
||||
VALUES (?, ?, ?, ?, ?);
|
||||
|
||||
update:
|
||||
UPDATE conversation SET title = ?, is_temporal = ?, updated_at = ?
|
||||
WHERE id = ?;
|
||||
|
||||
getById:
|
||||
SELECT * FROM conversation WHERE id = ?;
|
||||
|
||||
list:
|
||||
SELECT * FROM conversation
|
||||
WHERE is_temporal = 0
|
||||
ORDER BY updated_at DESC
|
||||
LIMIT :limit OFFSET :offset;
|
||||
|
||||
listAll:
|
||||
SELECT * FROM conversation
|
||||
ORDER BY updated_at DESC;
|
||||
|
||||
deleteById:
|
||||
DELETE FROM conversation WHERE id = ?;
|
||||
|
||||
rename:
|
||||
UPDATE conversation SET title = ?, updated_at = ? WHERE id = ?;
|
||||
|
||||
getUpdatedAt:
|
||||
SELECT updated_at FROM conversation WHERE id = ?;
|
||||
|
||||
touch:
|
||||
UPDATE conversation SET updated_at = ? WHERE id = ?;
|
||||
@@ -0,0 +1,33 @@
|
||||
CREATE TABLE 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 idx_msg_conv ON message(conversation_id, created_at);
|
||||
|
||||
insert:
|
||||
INSERT INTO message (id, conversation_id, kind, payload_json, created_at)
|
||||
VALUES (?, ?, ?, ?, ?);
|
||||
|
||||
listByConversation:
|
||||
SELECT * FROM message
|
||||
WHERE conversation_id = ?
|
||||
ORDER BY created_at ASC, id ASC
|
||||
LIMIT :limit OFFSET :offset;
|
||||
|
||||
listByConversationAll:
|
||||
SELECT * FROM message
|
||||
WHERE conversation_id = ?
|
||||
ORDER BY created_at ASC, id ASC;
|
||||
|
||||
listAfter:
|
||||
SELECT * FROM message
|
||||
WHERE conversation_id = ? AND created_at > ?
|
||||
ORDER BY created_at ASC, id ASC
|
||||
LIMIT :limit OFFSET :offset;
|
||||
|
||||
deleteByConversation:
|
||||
DELETE FROM message WHERE conversation_id = ?;
|
||||
@@ -0,0 +1,43 @@
|
||||
-- Self-reflection log: что агент "думает" о качестве своих последних ответов.
|
||||
-- Каждая запись — это one-shot LLM-размышление после N ходов (см. AGENTIK_REFLECTION_INTERVAL).
|
||||
-- Используется в buildSystemPrompt как "слабые места" (top-N последних).
|
||||
|
||||
CREATE TABLE reflection (
|
||||
id TEXT NOT NULL PRIMARY KEY,
|
||||
conversation_id TEXT, -- может быть NULL для глобальных рефлексий
|
||||
created_at INTEGER NOT NULL, -- epoch millis
|
||||
turns_analyzed INTEGER NOT NULL, -- сколько ходов оценили
|
||||
score INTEGER NOT NULL, -- 1..5 (самооценка качества)
|
||||
summary TEXT NOT NULL, -- краткое резюме (markdown)
|
||||
weak_spots_json TEXT NOT NULL DEFAULT '[]' -- JSON-массив строк: ["медленно отвечаю на X", ...]
|
||||
);
|
||||
|
||||
CREATE INDEX idx_reflection_created ON reflection(created_at DESC);
|
||||
CREATE INDEX idx_reflection_conv ON reflection(conversation_id, created_at DESC);
|
||||
|
||||
insert:
|
||||
INSERT INTO reflection (id, conversation_id, created_at, turns_analyzed, score, summary, weak_spots_json)
|
||||
VALUES (?, ?, ?, ?, ?, ?, ?);
|
||||
|
||||
getById:
|
||||
SELECT * FROM reflection WHERE id = ?;
|
||||
|
||||
-- Top-N последних рефлексий для всего агента (для system prompt).
|
||||
listRecent:
|
||||
SELECT * FROM reflection
|
||||
ORDER BY created_at DESC
|
||||
LIMIT :limit;
|
||||
|
||||
-- Top-N рефлексий для конкретного диалога.
|
||||
listForConversation:
|
||||
SELECT * FROM reflection
|
||||
WHERE conversation_id = :convId
|
||||
ORDER BY created_at DESC
|
||||
LIMIT :limit;
|
||||
|
||||
deleteOlderThan:
|
||||
DELETE FROM reflection
|
||||
WHERE created_at < :cutoffEpochMillis;
|
||||
|
||||
count:
|
||||
SELECT COUNT(*) FROM reflection;
|
||||
+38
@@ -0,0 +1,38 @@
|
||||
CREATE TABLE 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 idx_wm_unique ON working_memory(conversation_id, order_idx);
|
||||
CREATE INDEX idx_wm_conv ON working_memory(conversation_id, order_idx);
|
||||
|
||||
insert:
|
||||
INSERT INTO working_memory (id, conversation_id, order_idx, source_message_id, kind, payload_json, created_at)
|
||||
VALUES (?, ?, ?, ?, ?, ?, ?);
|
||||
|
||||
listByConversation:
|
||||
SELECT * FROM working_memory
|
||||
WHERE conversation_id = ?
|
||||
ORDER BY order_idx ASC;
|
||||
|
||||
clearByConversation:
|
||||
DELETE FROM working_memory WHERE conversation_id = ?;
|
||||
|
||||
maxOrderIdx:
|
||||
SELECT COALESCE(MAX(order_idx), 0) FROM working_memory WHERE conversation_id = ?;
|
||||
|
||||
-- Atomic compact: delete rows >= dropFromOrderIdx, then insert summary at
|
||||
-- the next order_idx (max(remaining) + 1). Caller supplies summaryId, summaryText,
|
||||
-- and current epoch millis for created_at.
|
||||
compactDelete:
|
||||
DELETE FROM working_memory
|
||||
WHERE conversation_id = ? AND order_idx >= ?;
|
||||
|
||||
compactInsert:
|
||||
INSERT INTO working_memory (id, conversation_id, order_idx, source_message_id, kind, payload_json, created_at)
|
||||
VALUES (?, ?, ?, NULL, 'summary', ?, ?);
|
||||
+124
@@ -0,0 +1,124 @@
|
||||
package pw.binom.agentik.storage.sqlite
|
||||
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertNull
|
||||
import kotlin.test.assertNotNull
|
||||
import kotlin.time.Instant
|
||||
import pw.binom.agentik.storage.Reflection
|
||||
import pw.binom.agentik.storage.ReflectionEvent
|
||||
|
||||
class ReflectionStoreTest {
|
||||
|
||||
@Test
|
||||
fun `insert and get roundtrip preserves all fields`() {
|
||||
val stores = SqliteStores.inMemory()
|
||||
try {
|
||||
val store = stores.reflections
|
||||
val r = sample(id = "r1", createdAt = Instant.parse("2026-09-15T12:00:00Z"))
|
||||
kotlinx.coroutines.runBlocking { store.insert(r) }
|
||||
|
||||
val loaded = kotlinx.coroutines.runBlocking { store.get("r1") }
|
||||
assertNotNull(loaded)
|
||||
assertEquals("r1", loaded.id)
|
||||
assertEquals("conv-1", loaded.conversationId)
|
||||
assertEquals(Instant.parse("2026-09-15T12:00:00Z"), loaded.createdAt)
|
||||
assertEquals(5, loaded.turnsAnalyzed)
|
||||
assertEquals(4, loaded.score)
|
||||
assertEquals("норм", loaded.summary)
|
||||
assertEquals(listOf("медленно", "путаю"), loaded.weakSpots)
|
||||
} finally { stores.close() }
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `listRecent returns newest first with limit`() {
|
||||
val stores = SqliteStores.inMemory()
|
||||
try {
|
||||
val at = Instant.parse("2026-09-15T12:00:00Z")
|
||||
repeat(5) { i ->
|
||||
kotlinx.coroutines.runBlocking {
|
||||
stores.reflections.insert(sample(id = "r$i", createdAt = at.plus(kotlin.time.Duration.parse("${i}s"))))
|
||||
}
|
||||
}
|
||||
val top3 = kotlinx.coroutines.runBlocking { stores.reflections.listRecent(limit = 3) }
|
||||
assertEquals(3, top3.size)
|
||||
assertEquals(listOf("r4", "r3", "r2"), top3.map { it.id })
|
||||
} finally { stores.close() }
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `listForConversation filters by convId`() {
|
||||
val stores = SqliteStores.inMemory()
|
||||
try {
|
||||
val at = Instant.parse("2026-09-15T12:00:00Z")
|
||||
kotlinx.coroutines.runBlocking {
|
||||
stores.reflections.insert(sample(id = "a", createdAt = at, convId = "conv-A"))
|
||||
stores.reflections.insert(sample(id = "b", createdAt = at, convId = "conv-B"))
|
||||
stores.reflections.insert(sample(id = "c", createdAt = at, convId = "conv-A"))
|
||||
}
|
||||
val onlyA = kotlinx.coroutines.runBlocking { stores.reflections.listForConversation("conv-A", limit = 10) }
|
||||
assertEquals(2, onlyA.size)
|
||||
assertEquals(setOf("a", "c"), onlyA.map { it.id }.toSet())
|
||||
} finally { stores.close() }
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `deleteOlderThan removes only old records`() {
|
||||
val stores = SqliteStores.inMemory()
|
||||
try {
|
||||
val t0 = Instant.parse("2026-09-15T12:00:00Z")
|
||||
kotlinx.coroutines.runBlocking {
|
||||
stores.reflections.insert(sample(id = "old", createdAt = t0))
|
||||
stores.reflections.insert(sample(id = "new", createdAt = t0.plus(kotlin.time.Duration.parse("1h"))))
|
||||
}
|
||||
val cutoff = t0.plus(kotlin.time.Duration.parse("1m"))
|
||||
kotlinx.coroutines.runBlocking { stores.reflections.deleteOlderThan(cutoff) }
|
||||
// "old" до cutoff — удаляется; "new" после cutoff — остаётся.
|
||||
assertNull(kotlinx.coroutines.runBlocking { stores.reflections.get("old") })
|
||||
assertNotNull(kotlinx.coroutines.runBlocking { stores.reflections.get("new") })
|
||||
} finally { stores.close() }
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `count returns total records`() {
|
||||
val stores = SqliteStores.inMemory()
|
||||
try {
|
||||
assertEquals(0, kotlinx.coroutines.runBlocking { stores.reflections.count() })
|
||||
repeat(3) { i ->
|
||||
kotlinx.coroutines.runBlocking {
|
||||
stores.reflections.insert(sample(id = "r$i", createdAt = Instant.parse("2026-09-15T12:00:00Z")))
|
||||
}
|
||||
}
|
||||
assertEquals(3, kotlinx.coroutines.runBlocking { stores.reflections.count() })
|
||||
} finally { stores.close() }
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `encodeStringArray handles quotes and backslashes`() {
|
||||
val encoded = encodeStringArray(listOf("a\"b", "c\\d", "plain"))
|
||||
val decoded = decodeStringArray(encoded)
|
||||
assertEquals(listOf("a\"b", "c\\d", "plain"), decoded)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `decodeStringArray handles empty and invalid`() {
|
||||
assertEquals(emptyList<String>(), decodeStringArray("[]"))
|
||||
assertEquals(emptyList<String>(), decodeStringArray(""))
|
||||
assertEquals(emptyList<String>(), decodeStringArray("not json"))
|
||||
assertEquals(emptyList<String>(), decodeStringArray("[unterminated"))
|
||||
}
|
||||
|
||||
private fun sample(
|
||||
id: String,
|
||||
createdAt: Instant,
|
||||
convId: String = "conv-1",
|
||||
): Reflection = Reflection(
|
||||
id = id,
|
||||
conversationId = convId,
|
||||
createdAt = createdAt,
|
||||
turnsAnalyzed = 5,
|
||||
score = 4,
|
||||
summary = "норм",
|
||||
weakSpots = listOf("медленно", "путаю"),
|
||||
)
|
||||
}
|
||||
Reference in New Issue
Block a user