Коммит заменяет 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
+1 -1
View File
@@ -77,7 +77,7 @@ budget exhaustion, registry filter, parallel dispatch.
## Чего здесь НЕТ
- Никакого конкретного LLM. Dispatcher вызывает tools, не LLM.
- Никакого persistent storage. Опирается на контракт `WorkingMemoryStore`
- Никакого persistent storage. Опирается на контракт `ContextStore`
(см. `:storage-core`).
## Текущий статус
+37
View File
@@ -0,0 +1,37 @@
plugins {
alias(libs.plugins.kotlin.multiplatform)
alias(libs.plugins.kotlin.serialization)
}
// Public API для mutable working-memory — runtime context агента (compaction,
// order_idx, summary entries). Зависит от :message-log-api для Ids (wm- префикс).
//
// НЕ нужен тонким клиентам — только серверному рантайму (`:standalone`, `:agentik-cli`,
// будущий `:android-agent` core).
kotlin {
jvmToolchain(21)
jvm()
macosX64()
macosArm64()
iosX64()
iosArm64()
iosSimulatorArm64()
linuxX64()
linuxArm64()
mingwX64()
sourceSets {
commonMain.dependencies {
api(project(":message-log-api"))
api(libs.kotlinx.coroutines.core)
api(libs.kotlinx.serialization.core)
api(libs.kotlinx.serialization.json)
}
commonTest.dependencies {
implementation(kotlin("test"))
implementation(libs.kotlinx.coroutines.test)
}
}
}
@@ -0,0 +1,58 @@
package pw.binom.agentik.context
import kotlin.time.Instant
/**
* Одна строка `working_memory` таблицы (внутреннее представление store).
*
* Используется для тестов и для перестроения [WorkingMemoryEntry] из row.
* Агент не должен с этим типом работать напрямую — он работает с
* [WorkingMemoryEntry] через [ContextStore].
*/
data class WorkingMemoryRow(
val id: String,
val conversationId: String,
val orderIdx: Long,
val sourceMessageId: String?,
val entry: WorkingMemoryEntry,
val createdAt: Instant,
)
/**
* Мутируемое представление LLM-контекста диалога (`working_memory` table).
*
* Аудит-лог — [MessageStore], неизменный; здесь — ровно то, что агент сейчас
* «видит»: системный промпт + реплики + (опционально) суммаризации.
* Строки упорядочены по `order_idx` ASC; `source_message_id` NULL указывает
* на синтетические строки (System, а в v2 — Summary).
*
* Суммаризация / чистка — один атомарный вызов [compact].
*/
interface ContextStore : AutoCloseable {
/** Добавить запись в конец working memory (новый максимальный `order_idx`). */
suspend fun append(conversationId: String, entry: WorkingMemoryEntry, now: Instant)
/** Все строки working memory диалога в порядке отправки. */
suspend fun list(conversationId: String): List<WorkingMemoryRow>
/** Очистить working memory диалога (используется при reset/rebuild). */
suspend fun clear(conversationId: String)
/**
* Атомарная суммаризация: удаляет все строки с `order_idx` в диапазоне
* `[dropFromOrderIdx, +∞)`. Если [summaryText] непустое — вместо удалённых
* строк вставляется одна синтетическая [WorkingMemoryEntry.Summary]
* с этим текстом и `order_idx = max(old order_idx after delete) + 1`
* (т.е. summary становится хвостом working memory).
*
* Если [summaryText] == null — работает как «отрезать хвост» (v1 поведение).
*
* Возвращает новый максимальный `order_idx` после операции.
*/
suspend fun compact(
dropFromOrderIdx: Long,
conversationId: String,
summaryText: String? = null,
): Long
}
@@ -0,0 +1,86 @@
package pw.binom.agentik.context
import kotlinx.serialization.SerialName
import kotlinx.serialization.Serializable
import pw.binom.agentik.messageLog.Content
import pw.binom.agentik.messageLog.MessageContext
/**
* Запись в working memory диалога: ровно то, что агент сейчас видит в
* LLM-контексте. Упорядочено по `order_idx` (заполняется в store при append).
*
* Sealed-иерархия: `User`/`Assistant` (реплики с ссылкой на audit log
* через [sourceMessageId]), `ToolExchange` (синтетическая запись об одном
* tool-вызове + его результате — для replay в LiteMessage(TOOL, ToolResult)
* при пересоздании LiteConv), `Summary` (суммаризация при compaction).
*/
@Serializable
sealed interface WorkingMemoryEntry {
/** Ссылка на исходное сообщение в audit log (`message.id`). `null` для синтетических строк. */
val sourceMessageId: String?
/** Реплика пользователя. */
@Serializable
@SerialName("user")
data class User(
override val sourceMessageId: String,
val content: List<Content>,
/**
* Контекст инициации хода. Применяется при сборке `LiteConversation`:
* если `origin != USER`, текст префиксуется `[origin] description (sourceId=…)`,
* чтобы модель видела, что её разбудил не пользователь.
* `null` = обычное user-сообщение.
*/
val context: MessageContext? = null,
) : WorkingMemoryEntry
/** Реплика ассистента. */
@Serializable
@SerialName("assistant")
data class Assistant(
override val sourceMessageId: String,
val content: List<Content>,
) : WorkingMemoryEntry
/**
* Синтетический блок: один tool-вызов + его результат. Синтетический — потому
* что в audit log это две отдельные записи (`MessageRecord.ToolCall` +
* `MessageRecord.ToolResult`), а в working_memory мы храним одной строкой
* для удобства replay'а.
*
* При создании новой LiteConv каждая такая запись превращается в
* `LiteMessage(TOOL, [ToolResult(callId, name, response)])` — LiteRT-LM
* матчит по `name`, `callId` берётся из [sourceMessageId] (= id исходного
* [MessageRecord.ToolCall]). Если [wasCancelled] = true, [resultText]
* содержит маркер `[cancelled by user]` — модель видит честную причину
* отсутствия результата.
*
* [sourceMessageId] = id исходного [MessageRecord.ToolCall] (для трассировки
* в audit log).
*/
@Serializable
@SerialName("tool_exchange")
data class ToolExchange(
override val sourceMessageId: String,
val toolName: String,
val toolArgsJson: String,
val resultText: String,
val wasCancelled: Boolean = false,
) : WorkingMemoryEntry
/**
* Синтетический блок: суммаризация старых ходов, сгенерированная при
* compaction'е working memory. Не имеет ссылки на конкретное сообщение
* в audit log — это наша собственная интерпретация контекста.
*/
@Serializable
@SerialName("summary")
data class Summary(
val text: String,
/** Ходы, которые были свёрнуты в этот summary (диапазон order_idx в виде меты). */
val coversUpToOrderIdx: Long? = null,
) : WorkingMemoryEntry {
override val sourceMessageId: String? = null
}
}
+1 -1
View File
@@ -118,7 +118,7 @@ Main.kt
- `compact(dropFromOrderIdx, conversationId)` — v1: DELETE rows ≥ order_idx,
summarization-вставка отложена (нужен дизайн-проработка).
**`MessageStore`** — append-only аудит. На каждый ход дописываются
**`JournalStore`** — append-only аудит. На каждый ход дописываются
`UserMessage`, `AssistantMessage`, `ToolCall`, `ToolResult`, `Error`. Никаких
update/delete кроме каскада из `ConversationStore.delete`.
+2 -2
View File
@@ -175,7 +175,7 @@ fun main() {
**Ошибки хода персистятся.** Если ход провалился (LLM/движок недоступны — например, HTTP 400 от endpoint'а), `ChatConversation.failTurn` пишет терминальную запись `MessageRecord.Error` в audit и эмитит `Event.Error` + `Event.End`. Благодаря audit-записи ошибка видна не только подписчику live-SSE, но и клиенту, который делает backfill через `getMessages` (polling/переподключение): в истории будет `Message.Error(id, message, code?)`, а для этого user-сообщения не будет `AssistantMessage`. При ошибке стрима живой `LiteConversation` сбрасывается — следующий `send` пересоберёт его из `working_memory`. В working_memory `Error` не пишется (модель не должна видеть ошибки прошлых ходов).
### `MessageStore`
### `JournalStore`
```kotlin
suspend fun append(record: MessageRecord)
@@ -183,7 +183,7 @@ suspend fun list(conversationId: String, after: Instant, offset: Int, limit: Int
suspend fun listAll(conversationId: String): List<MessageRecord>
```
### `WorkingMemoryStore`
### `ContextStore`
```kotlin
suspend fun append(conversationId: String, entry: WorkingMemoryEntry, now: Instant)
+3 -3
View File
@@ -20,7 +20,7 @@ bounded tail** для live-SSE и недавнего replay. Полный audit
## Где используется
- `:standalone` ChatAgent — append через `MutableEventStore` (заменяет
- `:standalone` ChatAgent — append через `MutableOutboxStore` (заменяет
текущий `agentEvents: MutableSharedFlow` + `persistAgentEvent`).
- `:server` Routes.kt — `/events/all` SSE endpoint читает через
`EventStore.events(after)`.
@@ -60,7 +60,7 @@ eventStore.events(after = client.lastSeen).collect { apply(it) }
## API
### `EventStore` (read-only, для consumer'ов)
### `OutboxStore` (read-only, для consumer'ов)
```kotlin
interface EventStore : AutoCloseable {
@@ -123,7 +123,7 @@ append по TTL/cap. Producer не должен полагаться на то,
## Текущее состояние
- ✅ Interface дизайн (`EventStore` + `MutableEventStore`)
- ✅ Interface дизайн (`OutboxStore` + `MutableOutboxStore`)
- ✅ KMP build (jvm + linuxX64 + mingwX64)
- ⏳ Нет implementations (next: `InMemoryEventStore` для тестов)
- ⏳ Не интегрирован в `:standalone`/`:server`
+29
View File
@@ -0,0 +1,29 @@
plugins {
alias(libs.plugins.kotlin.multiplatform)
alias(libs.plugins.kotlin.serialization)
}
kotlin {
jvmToolchain(21)
jvm()
macosX64()
macosArm64()
iosX64()
iosArm64()
iosSimulatorArm64()
linuxX64()
linuxArm64()
mingwX64()
sourceSets {
commonMain.dependencies {
api(libs.kotlinx.coroutines.core)
api(libs.kotlinx.serialization.core)
api(libs.kotlinx.serialization.json)
}
commonTest.dependencies {
implementation(kotlin("test"))
implementation(libs.kotlinx.coroutines.test)
}
}
}
@@ -0,0 +1,24 @@
package pw.binom.agentik.journal
import kotlinx.serialization.SerialName
import kotlinx.serialization.Serializable
/**
* Часть контента сообщения на уровне хранилища. Намеренно НЕ зависит от
* `pw.binom.agentik.proto.Content` — маппинг `:proto.Content ↔ Content` живёт
* в `Mapping.kt` storage impl'ов.
*/
@Serializable
sealed interface Content {
@Serializable
@SerialName("text")
data class Text(val body: String) : Content
@Serializable
@SerialName("image")
data class Image(val data: ByteArray, val mime: String) : Content {
override fun equals(other: Any?): Boolean =
this === other || (other is Image && mime == other.mime && data.contentEquals(other.data))
override fun hashCode(): Int = 31 * mime.hashCode() + data.contentHashCode()
}
}
@@ -0,0 +1,41 @@
package pw.binom.agentik.journal
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.flow
import kotlin.time.Instant
/**
* Append-only audit log сообщений — read-only представление.
*
* Producer-операция [MutableJournalStore.append] находится на
* [MutableJournalStore] — этот интерфейс только для чтения, чтобы
* consumer'ы физически не могли писать в audit log.
*
* Никаких обновлений, никакого удаления (кроме каскадного вместе
* с ConversationStore.delete).
*/
interface JournalStore : AutoCloseable {
suspend fun list(conversationId: String, after: Instant, offset: Int, limit: Int): List<MessageRecord>
/**
* Cold-flow paging через [list]. Default-реализация делает N+1 round-trip
* (по странице через `list()` пока не получит короткую страницу). Для
* in-memory backend'ов это OK; remote/SQLite impl'ы могут override'нуть
* на `Channel` / cursor-батчинг, чтобы избежать per-page round-trip.
*/
fun listFlow(conversationId: String, after: Instant, pageSize: Int = PAGE_SIZE): Flow<MessageRecord> = flow {
var offset = 0
while (true) {
val page = list(conversationId, after, offset, pageSize)
if (page.isEmpty()) return@flow
for (rec in page) emit(rec)
if (page.size < pageSize) return@flow
offset += page.size
}
}
companion object {
const val PAGE_SIZE = 100
}
}
@@ -0,0 +1,12 @@
package pw.binom.agentik.journal
import kotlinx.serialization.Serializable
import kotlinx.serialization.json.JsonElement
@Serializable
data class MessageContext(
val origin: MessageOrigin,
val sourceId: String? = null,
val description: String? = null,
val metadata: JsonElement? = null,
)
@@ -0,0 +1,20 @@
package pw.binom.agentik.journal
import kotlinx.serialization.SerialName
import kotlinx.serialization.Serializable
/**
* Контекст инициации хода (кто/что и почему). Дубликат типа из `:proto` —
* живёт здесь чтобы не тащить `:proto` в слой хранения данных.
*/
@Serializable
enum class MessageOrigin {
@SerialName("user")
USER,
@SerialName("system")
SYSTEM,
@SerialName("event")
EVENT,
}
@@ -0,0 +1,71 @@
package pw.binom.agentik.journal
import kotlinx.serialization.SerialName
import kotlinx.serialization.Serializable
import kotlin.time.Instant
/**
* Запись в таблице `message` (append-only audit).
*/
@Serializable
sealed interface MessageRecord {
val id: String
val conversationId: String
val createdAt: Instant
@Serializable
sealed interface Body : MessageRecord {
val content: List<Content>
}
@Serializable
@SerialName("user")
data class UserMessage(
override val id: String,
override val conversationId: String,
override val content: List<Content>,
override val createdAt: Instant,
val context: MessageContext? = null,
) : Body
@Serializable
@SerialName("assistant")
data class AssistantMessage(
override val id: String,
override val conversationId: String,
override val content: List<Content>,
override val createdAt: Instant,
val tokens: TurnTokens? = null,
) : Body
@Serializable
@SerialName("tool_call")
data class ToolCall(
override val id: String,
override val conversationId: String,
val toolName: String,
val toolTitle: String?,
val toolArgsJson: String,
override val createdAt: Instant,
) : MessageRecord
@Serializable
@SerialName("tool_result")
data class ToolResult(
override val id: String,
override val conversationId: String,
val toolCallId: String,
val result: String?,
override val createdAt: Instant,
) : MessageRecord
@Serializable
@SerialName("error")
data class Error(
override val id: String,
override val conversationId: String,
val message: String,
val code: String?,
override val createdAt: Instant,
) : MessageRecord
}
@@ -0,0 +1,17 @@
package pw.binom.agentik.journal
/**
* Mutable вариант [JournalStore] — добавляет producer-операцию [append].
*
* Этот интерфейс предназначен **только для producer'ов** (ChatAgent,
* ConversationLoop, ToolDispatcher, sub-agents, A2A-bridge).
* Consumer'ы (DebugRoutes, admin dashboards, parent agents) должны
* принимать **read-only** [JournalStore] — тогда невозможно случайно
* записать в audit log из observer'а.
*
* **Append семантика**: см. KDoc [MessageStore.append][JournalStore] —
* на этом интерфейсе (не дублируем).
*/
interface MutableJournalStore : JournalStore {
suspend fun append(record: MessageRecord)
}
@@ -0,0 +1,47 @@
package pw.binom.agentik.journal
import kotlinx.serialization.SerialName
import kotlinx.serialization.Serializable
import kotlinx.serialization.builtins.ListSerializer
import kotlinx.serialization.json.Json
private val bodyJson = Json {
ignoreUnknownKeys = true
encodeDefaults = true
explicitNulls = false
}
@Serializable
data class MessageBodyPayload(
val content: List<Content>,
@SerialName("context")
val context: MessageContext? = null,
val tokens: TurnTokens? = null,
)
fun encodeBodyPayload(
content: List<Content>,
context: MessageContext? = null,
tokens: TurnTokens? = null,
): String = bodyJson.encodeToString(
MessageBodyPayload.serializer(),
MessageBodyPayload(content = content, context = context, tokens = tokens),
)
fun decodeBodyPayload(json: String): BodyDecoded = readPayload(json)
data class BodyDecoded(
val content: List<Content>,
val context: MessageContext?,
val tokens: TurnTokens? = null,
)
private fun readPayload(json: String): BodyDecoded {
return try {
val p = bodyJson.decodeFromString(MessageBodyPayload.serializer(), json)
BodyDecoded(p.content, p.context, p.tokens)
} catch (e: kotlinx.serialization.SerializationException) {
val arr = bodyJson.decodeFromString(ListSerializer(Content.serializer()), json)
BodyDecoded(arr, null, null)
}
}
@@ -0,0 +1,18 @@
package pw.binom.agentik.journal
import kotlinx.serialization.Serializable
/**
* Token usage одного assistant turn'а.
*/
@Serializable
data class TurnTokens(
val input: Int,
val output: Int,
) {
val total: Int get() = input + output
init {
require(input >= 0) { "input tokens must be non-negative, got $input" }
require(output >= 0) { "output tokens must be non-negative, got $output" }
}
}
+6
View File
@@ -6,7 +6,13 @@ plugins {
kotlin {
jvmToolchain(21)
jvm()
macosX64()
macosArm64()
iosX64()
iosArm64()
iosSimulatorArm64()
linuxX64()
linuxArm64()
mingwX64()
sourceSets {
@@ -1,6 +0,0 @@
package pw.binom.agentik.messageLog
sealed interface MessageEvent {
val conversationId: String
data class Appended(override val conversationId: String, val record: MessageRecord) : MessageEvent
}
@@ -1,6 +1,7 @@
package pw.binom.agentik.messageLog
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.flow
import kotlin.time.Instant
/**
@@ -17,14 +18,24 @@ interface MessageStore : AutoCloseable {
suspend fun list(conversationId: String, after: Instant, offset: Int, limit: Int): List<MessageRecord>
suspend fun listAll(conversationId: String): List<MessageRecord>
/**
* Subscribe на события append'ов. Default — `emptyFlow()` для store'ов
* без live-уведомлений (например, SQLite без триггеров); impl'ы с
* in-memory notification (`:storage-inmemory`) override'ят.
* Cold-flow paging через [list]. Default-реализация делает N+1 round-trip
* (по странице через `list()` пока не получит короткую страницу). Для
* in-memory backend'ов это OK; remote/SQLite impl'ы могут override'нуть
* на `Channel` / cursor-батчинг, чтобы избежать per-page round-trip.
*/
fun events(): Flow<MessageEvent> = kotlinx.coroutines.flow.emptyFlow()
fun listFlow(conversationId: String, after: Instant, pageSize: Int = PAGE_SIZE): Flow<MessageRecord> = flow {
var offset = 0
while (true) {
val page = list(conversationId, after, offset, pageSize)
if (page.isEmpty()) return@flow
for (rec in page) emit(rec)
if (page.size < pageSize) return@flow
offset += page.size
}
}
suspend fun tokenStats(conversationId: String): TokenStats
}
companion object {
const val PAGE_SIZE = 100
}
}
@@ -1,9 +0,0 @@
package pw.binom.agentik.messageLog
data class TokenStats(
val turns: Int,
val inputTokens: Long,
val outputTokens: Long,
) {
val totalTokens: Long get() = inputTokens + outputTokens
}
+6
View File
@@ -16,7 +16,13 @@ kotlin {
jvmToolchain(21)
jvm()
macosX64()
macosArm64()
iosX64()
iosArm64()
iosSimulatorArm64()
linuxX64()
linuxArm64()
mingwX64()
sourceSets {
+28
View File
@@ -0,0 +1,28 @@
plugins {
alias(libs.plugins.kotlin.multiplatform)
}
kotlin {
jvmToolchain(21)
jvm()
macosX64()
macosArm64()
iosX64()
iosArm64()
iosSimulatorArm64()
linuxX64()
linuxArm64()
mingwX64()
sourceSets {
commonMain.dependencies {
api(project(":proto"))
api(libs.kotlinx.coroutines.core)
}
commonTest.dependencies {
implementation(kotlin("test"))
implementation(libs.kotlinx.coroutines.test)
}
}
}
@@ -0,0 +1,52 @@
package pw.binom.agentik.outbox
import pw.binom.agentik.proto.CommonEvent
/**
* Mutable вариант [OutboxStore] — добавляет producer-операцию [append].
*
* Этот интерфейс предназначен **только для producer'ов** (ChatAgent,
* sub-agents, A2A-bridge). Consumer'ы (server SSE endpoints, admin
* dashboards, parent agents) должны принимать **read-only** [OutboxStore]
* — тогда невозможно случайно писать в store из observer'а.
*
* Типичное использование:
* ```
* // Producer
* class ChatAgent(private val events: MutableEventStore) {
* suspend fun doSomething() {
* events.append(CommonEvent.Agent(date = now, event = AgentEvent.Created(...)))
* }
* }
*
* // Consumer
* class EventStreamEndpoint(private val events: EventStore) {
* fun stream() = events.events(after = null)
* // Ошибка компиляции если раскомментировать:
* // events.append(...) // ← нельзя, MutableEventStore нет в типе
* }
* ```
*
* **Append НЕ идемпотентен**: [CommonEvent] не имеет уникального id,
* поэтому retry с тем же logical event (например, после network failure
* между producer и store) приведёт к дубликату в tail'е. Это OK для
* use case'a bounded-tail — клиент, делающий catchup через [events](after),
* получит свой диапазон ровно один раз при подключении, а последующие
* retry producer'а просто насытят tail повторами, не задевая уже
* обработанные. Для гарантированной exactly-once — dedup через
* [message-store] (там есть монотонный `id`).
*
* **Silently evicted**: implementation может выкинуть этот event сразу
* после append (TTL/cap) без уведомления producer'а. Producer **не
* должен** полагаться на то, что event дойдёт до клиента, если он
* вне retention window.
*/
interface MutableOutboxStore : OutboxStore {
/**
* Положить event в log.
*
* - **Не идемпотентно** — см. KDoc интерфейса.
* - **Suspend** для KMP I/O impl'ов (SQLite через JNI).
*/
suspend fun append(event: CommonEvent)
}
@@ -0,0 +1,145 @@
package pw.binom.agentik.outbox
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.filter
import kotlinx.coroutines.flow.filterIsInstance
import kotlin.time.Instant
import pw.binom.agentik.proto.CommonEvent
/**
* Bounded-tail event log с автоматическим управлением TTL.
*
* **Архитектура двухуровневого хранилища событий**:
* 1. **Этот store** = короткий bounded tail (live SSE + недавний replay).
* События автоматически эвиктятся по TTL/cap (implementation-defined).
* 2. **Message store (`:message-store-api`)** = полный audit log, никогда не
* эвиктится. Source of truth для всего прошлого.
*
* **Паттерн reconnect** (caller'ы):
* ```
* val earliest = store.earliestEventDate()
* if (client.lastSeen < earliest) {
* // gap обнаружен — идём в message store за прошлым
* val gap = messageStore.query(after = client.lastSeen, before = earliest)
* applyAll(gap)
* }
* store.events(after = client.lastSeen).collect { apply(it) }
* ```
*
* **Нет delete/cleanup методов** — TTL/cap eviction полностью на стороне
* implementation. Это:
* - Убирает single source of truth дублирование (caller не может забыть cleanup).
* - Позволяет impl выбирать retention strategy (TTL, size cap, sliding window).
* - Сохраняет контракт clean: интерфейс только о put/get.
*
* **Read-only**: этот интерфейс предоставляет только read-операции.
* Для записи см. [MutableOutboxStore].
*
* **Подписки нереентрантные**: каждый вызов [events] создаёт **новую
* подписку** (cold Flow). Один [events] НЕ видит события, добавленные до
* его вызова, если [after] == null. Если нужен catchup — передавайте
* `after = lastSeenDate` явно.
*
* **Multi-consumer**: разные [events] подписки видят одно и то же live
* tail. Каждая подписка — независимая projection.
*/
interface OutboxStore : AutoCloseable {
/**
* Subscribe на events.
*
* **`after == null`** → только **live** (события с момента вызова
* `events()`). Каждое новое событие от любого producer'а немедленно
* появится в Flow. Буфер replay не отдаётся.
*
* **`after != null`** → сначала **catchup**: эмитт все буферизованные
* события с `date > after`, порядок `date ASC` (ties по `id ASC`).
* Затем **live** (как null-case).
*
* Cold Flow: каждый вызов — новая подписка. Вызов **после** append'а
* не увидит этот конкретный event (если `after == null`); для catchup
* передавайте явный `after`.
*
* ВАЖНО: `Flow` НЕ бросает ошибку при потере сети между producer и
* store — такие события просто не дойдут до этого Flow. Для гарантии
* полноты клиент обязан cross-check с [earliestEventDate] и fallback
* в message store при gap'е (см. KDoc интерфейса).
*/
fun events(after: Instant?): Flow<CommonEvent>
/**
* Subscribe на **только conversation events** (т.е. [CommonEvent.Conversation]).
*
* - [conversationId] == null → события **всех** диалогов.
* - [conversationId] != null → события **только этого** диалога.
*
* Семантика `after` идентична [events] (catchup + live).
* Возвращаемый тип — конкретный subtype [CommonEvent.Conversation].
*/
/**
* **Default implementation** (читает все events + фильтрует).
*
* Простая реализация через [events] + filterIsInstance. Реализации
* могут override'нуть для эффективности (например, добавить SQL
* `WHERE conversation_id = ?` чтобы не тянуть всё в память), но
* контракт корректен и без override.
*/
fun conversationEvents(after: Instant?, conversationId: String? = null): Flow<CommonEvent.Conversation> =
events(after)
.filterIsInstance<CommonEvent.Conversation>()
.let { filtered ->
if (conversationId == null) filtered
else filtered.filter { it.conversationId == conversationId }
}
/**
* Subscribe на **только agent events** ([CommonEvent.Agent] —
* создание/удаление/переименование диалога).
*
* Семантика `after` идентична [events] (catchup + live).
* Возвращаемый тип — конкретный subtype [CommonEvent.Agent].
*
* Полезно для admin-дашборда, который хочет видеть только lifecycle
* диалогов без деталей ходов.
*/
/**
* **Default implementation** (читает все events + фильтрует по типу).
*
* Простая реализация через [events] + filterIsInstance. Реализации
* могут override'нуть для эффективности (например, читать только agent
* row'ы из БД), но контракт корректен и без override.
*/
fun agentEvents(after: Instant?): Flow<CommonEvent.Agent> =
events(after).filterIsInstance<CommonEvent.Agent>()
/**
* Date **стартовой точки** буфера.
*
* - Если буфер не пуст → `date` самого старого буферизованного event'а.
* - Если буфер пуст → текущее время (`Clock.System.now()` на момент вызова).
*
* **Семантика "now если пусто"** важна: позволяет клиенту безопасно
* подписаться на [events](after = earliest) сразу — он получит только
* новые live event'ы, без ложного catchup. Если бы возвращалось
* `Instant.DISTANT_PAST` или `null` (с проверкой), клиент мог бы
* ошибочно подписаться на несуществующий catchup и зависнуть в ожидании.
*
* **Используется клиентом для gap detection**:
* - `lastSeen < earliest` → есть дыра в покрытии, нужен fallback
* в message store за диапазоном `[lastSeen, earliest)`.
* - `lastSeen >= earliest` → всё доступно через [events](after),
* fallback не нужен.
* - `lastSeen == earliest` → OK, первый live event будет > earliest.
*
* **Edge case**: клиент, подключившийся до того как store увидел хоть
* один event, получает `earliest ≈ now`. Его `lastSeen` будет < earliest
* — адаптируется в первом же poll'е и пойдёт через fallback если
* сообщения audit log существуют (для consistency с прошлым).
*
* Suspend потому что в persistent impl'ах требует SQL query (`MIN(date)`
* или `Clock.now()` для пустого буфера).
*/
suspend fun earliestEventDate(): Instant
override fun close()
}
+1 -1
View File
@@ -114,7 +114,7 @@ sealed interface Event {
- Никакого HTTP/SSE/JSON. Это контракт. Сериализация живёт в `:server`
и `:client`.
- Никакого хранения. Реализации `MessageStore` живут в `:storage-*`.
- Никакого хранения. Реализации `JournalStore` живут в `:storage-*`.
- Никакой логики прерывания / инструментов / LLM-вызовов. Это всё
внутри `:standalone` (ChatAgent) и выше.
+10 -6
View File
@@ -68,6 +68,15 @@ include(":memory-vector")
// использоваться в Android-сборке (JVector/SQLite не подходят для ART out-of-box).
include(":message-store-api")
include(":message-log-api")
// Новые API-модули трёх сущностей (canonical имена):
// - :journal-api — append-only audit log (бывший :message-log-api)
// - :context-api — то что видит LLM (бывший :working-memory-api)
// - :outbox-api — bounded-tail event stream (бывший :event-store)
// Старые модули :message-log-api / :working-memory-api / :event-store
// остаются на диске — миграция consumers'ов по чуть-чуть, отдельно.
include(":journal-api")
include(":context-api")
include(":outbox-api")
// Bounded-tail event log с auto-TTL. Двухуровневое хранилище: этот модуль —
// короткий live tail + recent replay; полный audit log живёт в :message-store-api
// (там — MessageStore + ConversationStore). EventStore сам управляет eviction,
@@ -79,12 +88,7 @@ include(":event-store-in-memory")
include(":working-memory-api")
include(":storage-inmemory")
// SQLDelight-реализация store'ов из :message-store-api и :working-memory-api. JVM-only
// native драйверов для KMP вне JVM пока не публикует). Содержит 4 .sq-файла +
// 5 классов: SqliteStores, SqliteConversationStore, SqliteMessageStore,
// SqliteWorkingMemoryStore, SqliteReflectionStore. Бэкенд для прод-запуска
// :standalone (путь к .db файлу в AGENTIK_DB).
include(":storage-sqlite")
// KMP-реализация EventStore через ksqlite (https://github.com/caffeine-mgn/ksqlite).
// KMP-реализация EventStore поверх ksqlite (https://github.com/caffeine-mgn/ksqlite).
// Цель: проверить что pure-Kotlin SQLite с sqlite-vec заменяет SQLDelight+JVector
// на KMP-таргетах (JVM + linuxX64 + mingwX64). Apple targets auto-disabled на
// Linux — собираются локально на macOS. Пока покрывает только EventStore;
+8 -1
View File
@@ -22,6 +22,13 @@ val skipVectorMemory: Boolean =
kotlin {
jvmToolchain(21)
// Отключаем auto-propagation Apple-таргетов от KMP-зависимостей (`:proto`,
// `:server`, `:agent-toolsets` имеют Apple-варианты). `:standalone` JVM-only,
// Apple таргеты не нужны и ломают резолюцию `:message-store-api` в appleMain
// (он KMP jvm+linuxX64+mingwX64 без Apple).
// @OptIn(org.jetbrains.kotlin.gradle.ExperimentalKotlinGradlePluginApi::class)
// applyDefaultHierarchyTemplate { }
jvm {
binaries {
executable {
@@ -61,7 +68,7 @@ kotlin {
}
implementation(project(":message-store-api"))
implementation(project(":working-memory-api"))
implementation(project(":storage-sqlite"))
implementation(project(":storage-ksqlite"))
// Новый единый канал событий агента — заменил старые
// `agentEvents: MutableSharedFlow<AgentEvent>` и per-conv `ConversationEvents._flow`.
@@ -10,7 +10,6 @@ import kotlinx.serialization.json.buildJsonObject
import kotlinx.serialization.json.put
import pw.binom.agentik.memory.ConversationTurn
import pw.binom.agentik.messageStore.ReflectionStore
import pw.binom.agentik.messageLog.MessageStore
import pw.binom.agentik.proto.Agent
import pw.binom.agentik.standalone.agent.ChatConversation
import pw.binom.agentik.llm.tools.LlmReflector
@@ -27,7 +26,10 @@ import pw.binom.agentik.workingMemory.WorkingMemoryStore
* - `POST /debug/skill-mine?conversationId=...` — прогон [SkillMiner] прямо сейчас
* - `POST /debug/curate` — прогон [Curator.runPass] прямо сейчас
* - `POST /debug/compact?conversationId=...` — принудительный compaction
* - `GET /debug/tokens?conversationId=...` — token-статистика диалога из БД
*
* Token-stats эндпоинт убран вместе с `MessageStore.tokenStats()` (см.
* :message-log-api/MessageStore.kt). Сейчас token accounting доступен
* только через assistant-сообщения с `TurnTokens` (см. MessageRecord).
*
* Каждый возвращает JSON с результатом (что сохранил / нашёл / сжал), чтобы в
* тестах было видно не только "триггер сработал", а что именно LLM намайнила.
@@ -35,7 +37,6 @@ import pw.binom.agentik.workingMemory.WorkingMemoryStore
*/
internal fun Route.debugRoutes(
agent: Agent,
messageStore: MessageStore,
workingMemoryStore: WorkingMemoryStore,
reflectionStore: ReflectionStore,
reflector: LlmReflector?,
@@ -110,20 +111,6 @@ internal fun Route.debugRoutes(
val ok = chatConv.forceCompactNow()
call.respondText("""{"compacted":$ok}""", contentType = ContentType.Application.Json)
}
get("/debug/tokens") {
val convId = call.parameters["conversationId"]
?: return@get call.respondText("conversationId required", status = HttpStatusCode.BadRequest)
val stats = messageStore.tokenStats(convId)
val json = buildJsonObject {
put("conversationId", convId)
put("turns", stats.turns.toString())
put("inputTokens", stats.inputTokens.toString())
put("outputTokens", stats.outputTokens.toString())
put("totalTokens", (stats.inputTokens + stats.outputTokens).toString())
}.toString()
call.respondText(json, contentType = ContentType.Application.Json)
}
}
/**
@@ -29,7 +29,7 @@ import pw.binom.agentik.standalone.config.AppConfig.MemoryBackend
import pw.binom.agentik.standalone.llm.LlmBackend
import pw.binom.agentik.standalone.llm.ModelDownloader
import pw.binom.agentik.mcp.bridge.McpRegistry
import pw.binom.agentik.storage.sqlite.SqliteStores
import pw.binom.agentik.storage.ksqlite.KsqliteStores
import java.io.File
import pw.binom.agentik.llm.tools.SkillMiner
/**
@@ -198,7 +198,7 @@ private fun runServer() {
}
val llm = config.llm.createLlm()
val sqliteStores = SqliteStores.open(dbPath = config.agent.dbPath)
val sqliteStores = KsqliteStores.open(path = config.agent.dbPath)
val mcpRegistry = McpRegistry.fromConfig(config.mcp)
// Хранилище скилов: если skillsDir задан, читаем каталог + создаём
@@ -368,7 +368,6 @@ private fun runServer() {
if (config.debug.endpoints) {
debugRoutes(
agent = agent,
messageStore = sqliteStores.messages,
workingMemoryStore = sqliteStores.workingMemory,
reflectionStore = sqliteStores.reflections,
reflector = reflector,
@@ -405,24 +404,13 @@ private fun runServer() {
if (config.debug.endpoints) {
println(" debug endpoints: enabled (/debug/reflect, /debug/skill-mine, /debug/curate, /debug/compact, /debug/tokens)")
}
// Token stats по существующим диалогам (агрегат на старте — каждая запись
// парсится из payload_json, ну >100 turns и БД приличная — но в рамках
// стартапа это терпимо).
// Раньше здесь был агрегат tokenStats() по всем conv'ам при старте. Метод
// убран из :message-log-api (MessageStore стал чисто read-only list+listFlow);
// см. agentik :message-log-api/MessageStore.kt. Token accounting теперь
// доступен через assistant-сообщения с TurnTokens (см. MessageRecord.AssistantMessage).
val existingConvs = kotlinx.coroutines.runBlocking { sqliteStores.conversations.list(offset = 0, limit = 1000) }
if (existingConvs.isNotEmpty()) {
var totalTurns = 0
var totalIn = 0L
var totalOut = 0L
for (c in existingConvs) {
if (c.isTemporal) continue
val s = kotlinx.coroutines.runBlocking { sqliteStores.messages.tokenStats(c.id) }
totalTurns += s.turns
totalIn += s.inputTokens
totalOut += s.outputTokens
}
if (totalTurns > 0) {
println(" tokens: ${existingConvs.size} convs, $totalTurns turns, in=${totalIn}, out=${totalOut}, total=${totalIn + totalOut}")
}
println(" conversations: ${existingConvs.size} (active)")
}
Runtime.getRuntime().addShutdownHook(Thread {
agent.close()
@@ -196,8 +196,23 @@ class ConversationLoop(
}
override suspend fun interrupt() {
if (activeTurn?.isActive != true) {
log.info { "interrupt() no-op: no active turn for $id" }
// Всегда ставим флаг — даже если activeTurn ещё не стартовал.
// runTurn проверяет interrupted.get() при входе (short-circuit) и в
// каждой итерации цикла + finally. Если turn запустится ПОСЛЕ нашего
// interrupt() — он увидит флаг на entry и сразу завершится без
// реального LLM-вызова. Если turn уже идёт — Interrupted + End придут
// в finally.
//
// Раньше здесь был early-return при `!activeTurn?.isActive`, но это
// давало race с точки зрения тестов: send() может завершиться
// быстрее (например, на ksqlite-бэкенде, где messageStore.append
// практически мгновенный), и interrupt(), вызванный после
// delay(200) от launch send(), видел completed Job → no-op →
// Interrupted event не эмитится.
val wasActive = activeTurn?.isActive
if (wasActive != null && !wasActive) {
// Предыдущий turn уже завершился — interrupt() действительно no-op.
log.info { "interrupt() no-op: previous turn already completed for $id" }
return
}
interrupted.set(true)
@@ -16,7 +16,7 @@ import pw.binom.agentik.standalone.llm.LlmBackend
import pw.binom.agentik.standalone.llm.LlmConfig
import pw.binom.agentik.messageLog.MessageRecord
import pw.binom.agentik.workingMemory.WorkingMemoryEntry
import pw.binom.agentik.storage.sqlite.SqliteStores
import pw.binom.agentik.storage.ksqlite.KsqliteStores
import pw.binom.litert.LiteContentPart
import pw.binom.litert.LiteConversation
import pw.binom.litert.LiteConversationConfig
@@ -41,12 +41,12 @@ import pw.binom.agentik.toolsets.NamedTool
class ChatAgentTest {
private lateinit var sqliteStores: pw.binom.agentik.storage.sqlite.SqliteStores
private lateinit var sqliteStores: pw.binom.agentik.storage.ksqlite.KsqliteStores
private lateinit var fakeLlm: FakeLiteLlm
@BeforeTest
fun setup() {
sqliteStores = SqliteStores.inMemory()
sqliteStores = KsqliteStores.inMemory("chat-${kotlin.random.Random.nextLong()}")
fakeLlm = FakeLiteLlm()
}
@@ -56,7 +56,7 @@ class ChatAgentTest {
}
private fun newAgent(
sqliteStores: pw.binom.agentik.storage.sqlite.SqliteStores = this.sqliteStores,
sqliteStores: pw.binom.agentik.storage.ksqlite.KsqliteStores = this.sqliteStores,
llm: LiteLlm = this.fakeLlm,
tools: List<NamedTool> = emptyList(),
skills: SkillCatalog = SkillCatalog.EMPTY,
@@ -173,7 +173,7 @@ class ChatAgentTest {
assertTrue(agent.deleteConversation(id))
assertNull(agent.getConversation(id))
assertNull(sqliteStores.conversations.get(id))
assertEquals(emptyList(), sqliteStores.messages.listAll(id))
assertEquals(emptyList(), sqliteStores.messages.listFlow(id, Instant.DISTANT_PAST).toList())
}
@Test
@@ -191,7 +191,7 @@ class ChatAgentTest {
conv.send(listOf(Content.Text("hi")))
// user message записан в audit + working memory
val msgs = sqliteStores.messages.listAll(conv.id)
val msgs = sqliteStores.messages.listFlow(conv.id, Instant.DISTANT_PAST).toList()
assertEquals(2, msgs.size)
assertEquals("hi", (msgs[0] as pw.binom.agentik.messageLog.MessageRecord.UserMessage).content.let {
(it[0] as pw.binom.agentik.messageLog.Content.Text).body
@@ -328,7 +328,7 @@ class ChatAgentTest {
assertTrue(events.any { it is ProtoEvent.Error && it.message == "boom from llm" }, "events=$events")
assertTrue(events.any { it is ProtoEvent.End }, "events=$events")
val msgs = sqliteStores.messages.listAll(conv.id)
val msgs = sqliteStores.messages.listFlow(conv.id, Instant.DISTANT_PAST).toList()
assertEquals(2, msgs.size)
assertIs<pw.binom.agentik.messageLog.MessageRecord.UserMessage>(msgs[0])
val err = assertIs<pw.binom.agentik.messageLog.MessageRecord.Error>(msgs[1])
@@ -373,7 +373,7 @@ class ChatAgentTest {
eventsJob.cancel()
// audit: только user (assistant не успел сгенериться)
val msgs = sqliteStores.messages.listAll(conv.id)
val msgs = sqliteStores.messages.listFlow(conv.id, Instant.DISTANT_PAST).toList()
assertEquals(1, msgs.size)
assertIs<pw.binom.agentik.messageLog.MessageRecord.UserMessage>(msgs[0])
@@ -430,7 +430,7 @@ class ChatAgentTest {
eventsJob.cancel()
// audit: user + toolcall + toolresult (tool выполнился), assistant может быть
val msgs = sqliteStores.messages.listAll(conv.id)
val msgs = sqliteStores.messages.listFlow(conv.id, Instant.DISTANT_PAST).toList()
val toolResult = msgs.filterIsInstance<pw.binom.agentik.messageLog.MessageRecord.ToolResult>().firstOrNull()
assertNotNull(toolResult, "tool result должен быть в audit — tool выполнился нормально")
val toolResultResult = toolResult!!.result!!
@@ -455,7 +455,7 @@ class ChatAgentTest {
// Поднимаем file-backed БД, создаём temp-беседу
sqliteStores.close()
val dbPath = (System.getProperty("java.io.tmpdir") + "/agentik-test-${System.nanoTime()}.db")
sqliteStores = SqliteStores.open(dbPath)
sqliteStores = KsqliteStores.open(dbPath)
val agent1 = ChatAgent(
id = "agentik",
conversationStore = sqliteStores.conversations,
@@ -475,7 +475,7 @@ class ChatAgentTest {
// Переоткрываем БД — temp-беседа не должна пережить рестарт
sqliteStores.close()
sqliteStores = SqliteStores.open(dbPath)
sqliteStores = KsqliteStores.open(dbPath)
val agent2 = ChatAgent(
id = "agentik",
conversationStore = sqliteStores.conversations,
@@ -497,7 +497,7 @@ class ChatAgentTest {
fun `non-temp conversation persists across agent instances`() = runTest {
sqliteStores.close()
val dbPath = (System.getProperty("java.io.tmpdir") + "/agentik-test-${System.nanoTime()}.db")
sqliteStores = SqliteStores.open(dbPath)
sqliteStores = KsqliteStores.open(dbPath)
val agent1 = ChatAgent(
id = "agentik",
conversationStore = sqliteStores.conversations,
@@ -515,7 +515,7 @@ class ChatAgentTest {
val id = conv.id
sqliteStores.close()
sqliteStores = SqliteStores.open(dbPath)
sqliteStores = KsqliteStores.open(dbPath)
val agent2 = ChatAgent(
id = "agentik",
conversationStore = sqliteStores.conversations,
@@ -2,7 +2,7 @@ package pw.binom.agentik.standalone.agent
import kotlinx.coroutines.runBlocking
import pw.binom.agentik.standalone.llm.LlmConfig
import pw.binom.agentik.storage.sqlite.SqliteStores
import pw.binom.agentik.storage.ksqlite.KsqliteStores
import pw.binom.agentik.toolsets.ToolsetContribution
import pw.binom.litert.LiteLlm
import pw.binom.litert.LiteTool
@@ -39,8 +39,8 @@ class ChatAgentToolsetsTest {
private fun newAgent(
toolsets: List<ToolsetContribution> = emptyList(),
): Pair<ChatAgent, pw.binom.agentik.storage.sqlite.SqliteStores> {
val sqliteStores = SqliteStores.inMemory()
): Pair<ChatAgent, pw.binom.agentik.storage.ksqlite.KsqliteStores> {
val sqliteStores = KsqliteStores.inMemory("toolsets-${kotlin.random.Random.nextLong()}")
val agent = ChatAgent(
id = "test-agent",
conversationStore = sqliteStores.conversations,
@@ -11,7 +11,7 @@ import pw.binom.agentik.memory.md.KeywordMdReviewer
import pw.binom.agentik.standalone.llm.LlmBackend
import pw.binom.agentik.standalone.llm.LlmConfig
import pw.binom.agentik.standalone.llm.OpenAiConfig
import pw.binom.agentik.storage.sqlite.SqliteStores
import pw.binom.agentik.storage.ksqlite.KsqliteStores
import pw.binom.agentik.proto.Content as ProtoContent
import pw.binom.litert.LiteConversation
import pw.binom.litert.LiteConversationConfig
@@ -40,12 +40,12 @@ import pw.binom.agentik.llm.tools.SummaryTurn
*/
class CompactionTest {
private lateinit var sqliteStores: pw.binom.agentik.storage.sqlite.SqliteStores
private lateinit var sqliteStores: pw.binom.agentik.storage.ksqlite.KsqliteStores
private lateinit var fakeLlm: FakeLiteLlm
@BeforeTest
fun setup() {
sqliteStores = SqliteStores.inMemory()
sqliteStores = KsqliteStores.inMemory("compact-${kotlin.random.Random.nextLong()}")
fakeLlm = FakeLiteLlm()
}
@@ -30,7 +30,7 @@ import pw.binom.agentik.standalone.agent.memory.MemoryToolsFactory
import pw.binom.agentik.standalone.llm.LlmBackend
import pw.binom.agentik.standalone.llm.LlmConfig
import pw.binom.agentik.standalone.llm.OpenAiConfig
import pw.binom.agentik.storage.sqlite.SqliteStores
import pw.binom.agentik.storage.ksqlite.KsqliteStores
import pw.binom.agentik.llm.tools.ContextCompactor
import pw.binom.agentik.llm.tools.SummaryTurn
import kotlin.test.AfterTest
@@ -49,13 +49,13 @@ import kotlin.time.Instant
*/
class MemoryWiringTest {
private lateinit var sqliteStores: pw.binom.agentik.storage.sqlite.SqliteStores
private lateinit var sqliteStores: pw.binom.agentik.storage.ksqlite.KsqliteStores
private lateinit var fakeLlm: FakeLiteLlm
private lateinit var root: Path
@BeforeTest
fun setup() {
sqliteStores = SqliteStores.inMemory()
sqliteStores = KsqliteStores.inMemory("memwire-${kotlin.random.Random.nextLong()}")
fakeLlm = FakeLiteLlm()
root = Path(SystemTemporaryDirectory.toString(), "agentik-mem-${java.util.UUID.randomUUID()}")
SystemFileSystem.createDirectories(root, mustCreate = true)
@@ -6,8 +6,9 @@ import pw.binom.agentik.messageLog.MessageRecord
import pw.binom.agentik.messageLog.Content
import pw.binom.agentik.workingMemory.WorkingMemoryEntry
import kotlinx.coroutines.flow.toList
import kotlinx.coroutines.test.runTest
import pw.binom.agentik.storage.sqlite.SqliteStores
import pw.binom.agentik.storage.ksqlite.KsqliteStores
import kotlin.test.AfterTest
import kotlin.test.BeforeTest
import kotlin.test.Test
@@ -21,11 +22,11 @@ import kotlin.time.Instant
class PersistenceTest {
private lateinit var stores: SqliteStores
private lateinit var stores: KsqliteStores
@BeforeTest
fun setup() {
stores = SqliteStores.inMemory()
stores = KsqliteStores.inMemory("persist-${kotlin.random.Random.nextLong()}")
}
@AfterTest
@@ -109,13 +110,13 @@ class PersistenceTest {
),
now = t0,
)
assertEquals(1, stores.messages.listAll("c1").size)
assertEquals(1, stores.messages.listFlow("c1", Instant.DISTANT_PAST).toList().size)
assertEquals(1, stores.workingMemory.list("c1").size)
val removed = stores.conversations.delete("c1")
assertTrue(removed)
assertNull(stores.conversations.get("c1"))
assertEquals(emptyList(), stores.messages.listAll("c1"))
assertEquals(emptyList(), stores.messages.listFlow("c1", Instant.DISTANT_PAST).toList())
assertEquals(emptyList(), stores.workingMemory.list("c1"))
}
@@ -129,7 +130,7 @@ class PersistenceTest {
MessageRecord.AssistantMessage("m2", "c1", listOf(Content.Text("yo")), t0),
)
val all = stores.messages.listAll("c1")
val all = stores.messages.listFlow("c1", Instant.DISTANT_PAST).toList()
assertEquals(2, all.size)
assertEquals("m1", all[0].id)
assertEquals("m2", all[1].id)
@@ -284,7 +285,7 @@ class PersistenceTest {
createdAt = Instant.fromEpochMilliseconds(1_700_000_001_000),
),
)
val all = stores.messages.listAll("c1")
val all = stores.messages.listFlow("c1", Instant.DISTANT_PAST).toList()
assertEquals(2, all.size)
val first = assertIs<MessageRecord.Error>(all[0])
assertEquals("e1", first.id)
@@ -305,7 +306,7 @@ class PersistenceTest {
createdAt = t0,
),
)
val all = stores.messages.listAll("c1")
val all = stores.messages.listFlow("c1", Instant.DISTANT_PAST).toList()
val image = (all[0] as MessageRecord.UserMessage).content[0] as Content.Image
assertEquals("image/png", image.mime)
assertTrue(bytes.contentEquals(image.data))
@@ -328,7 +329,7 @@ class PersistenceTest {
context = ctx,
),
)
val all = stores.messages.listAll("c1")
val all = stores.messages.listFlow("c1", Instant.DISTANT_PAST).toList()
assertEquals(1, all.size)
val user = assertIs<MessageRecord.UserMessage>(all[0])
assertEquals(ctx, user.context)
@@ -345,7 +346,7 @@ class PersistenceTest {
createdAt = t0,
),
)
val all = stores.messages.listAll("c1")
val all = stores.messages.listFlow("c1", Instant.DISTANT_PAST).toList()
val user = assertIs<MessageRecord.UserMessage>(all[0])
assertNull(user.context)
}
@@ -1,48 +0,0 @@
package pw.binom.agentik.standalone.persistence
import pw.binom.agentik.messageStore.Reflection
import pw.binom.agentik.messageStore.ConversationRecord
import pw.binom.agentik.messageLog.MessageRecord
import pw.binom.agentik.messageLog.Content
import pw.binom.agentik.workingMemory.WorkingMemoryEntry
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertNotNull
import kotlin.test.assertTrue
import pw.binom.agentik.storage.sqlite.SqliteStores
class SqliteStoresMigrationTest {
@Test
fun `reflection table is created on fresh inMemory database`() {
val stores = SqliteStores.inMemory()
try {
// Если таблицы нет — insert упадёт. Проверяем insert round-trip.
val r = Reflection(
id = "m1",
conversationId = null,
createdAt = kotlin.time.Instant.parse("2026-09-15T12:00:00Z"),
turnsAnalyzed = 3,
score = 5,
summary = "fresh",
weakSpots = listOf("none"),
)
kotlinx.coroutines.runBlocking { stores.reflections.insert(r) }
val loaded = kotlinx.coroutines.runBlocking { stores.reflections.get("m1") }
assertNotNull(loaded)
assertEquals("fresh", loaded.summary)
} finally { stores.close() }
}
@Test
fun `all stores exposed`() {
val stores = SqliteStores.inMemory()
try {
assertNotNull(stores.conversations)
assertNotNull(stores.messages)
assertNotNull(stores.workingMemory)
assertNotNull(stores.reflections)
assertTrue(stores.driver.toString().isNotBlank())
} finally { stores.close() }
}
}
@@ -1,137 +0,0 @@
package pw.binom.agentik.standalone.persistence
import pw.binom.agentik.messageLog.TurnTokens
import pw.binom.agentik.messageLog.MessageBodyPayload
import pw.binom.agentik.messageLog.encodeBodyPayload
import pw.binom.agentik.messageLog.decodeBodyPayload
import pw.binom.agentik.messageLog.MessageContext
import pw.binom.agentik.messageLog.MessageOrigin
import pw.binom.agentik.messageStore.ConversationRecord
import pw.binom.agentik.messageLog.MessageRecord
import pw.binom.agentik.messageLog.Content
import pw.binom.agentik.workingMemory.WorkingMemoryEntry
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertNull
import kotlin.time.Instant
import pw.binom.agentik.storage.sqlite.SqliteStores
class TokenStatsTest {
@Test
fun `tokenStats sums across assistant messages of same conversation`() {
val stores = SqliteStores.inMemory()
try {
kotlinx.coroutines.runBlocking {
stores.conversations.upsert(
ConversationRecord(
id = "c1",
title = null,
isTemporal = false,
createdAt = Instant.parse("2026-09-15T12:00:00Z"),
updatedAt = Instant.parse("2026-09-15T12:00:00Z"),
)
)}
kotlinx.coroutines.runBlocking {
stores.messages.append(
assistant("a1", "c1", 100, 50, Instant.parse("2026-09-15T12:01:00Z"))
)
stores.messages.append(
assistant("a2", "c1", 200, 80, Instant.parse("2026-09-15T12:02:00Z"))
)
// User без tokens — не должны считаться.
stores.messages.append(
MessageRecord.UserMessage(
id = "u1",
conversationId = "c1",
content = listOf(Content.Text("hi")),
createdAt = Instant.parse("2026-09-15T12:00:30Z"),
)
)
}
val stats = kotlinx.coroutines.runBlocking { stores.messages.tokenStats("c1") }
assertEquals(2, stats.turns)
assertEquals(300L, stats.inputTokens)
assertEquals(130L, stats.outputTokens)
assertEquals(430L, stats.totalTokens)
} finally { stores.close() }
}
@Test
fun `tokenStats skips legacy assistant messages without tokens field`() {
val stores = SqliteStores.inMemory()
try {
// Вставляем запись со СТАРЫМ payload-форматом (plain array через
// прямой SQL апдейт, минуя типизированный encoder).
val id = "legacy"
val msgId = "m1"
stores.driver.execute(null, "INSERT INTO conversation(id, title, is_temporal, created_at, updated_at) VALUES('c1', NULL, 0, 0, 0)", 0)
stores.driver.execute(
null,
"INSERT INTO message(id, conversation_id, kind, payload_json, created_at) VALUES(?, ?, ?, ?, ?)",
5,
) {
bindString(0, msgId)
bindString(1, "c1")
bindString(2, "assistant")
bindString(3, """[{"kind":"text","body":"legacy"}]""")
bindLong(4, Instant.parse("2026-09-15T12:00:00Z").toEpochMilliseconds())
}
val stats = kotlinx.coroutines.runBlocking { stores.messages.tokenStats("c1") }
assertEquals(0, stats.turns, "legacy без tokens не должен считаться")
assertEquals(0L, stats.inputTokens)
assertEquals(0L, stats.outputTokens)
} finally { stores.close() }
}
@Test
fun `TurnTokens rejects negative values`() {
val tokens = TurnTokens(input = 100, output = 50)
assertEquals(100, tokens.input)
assertEquals(150, tokens.total)
// Sanity для init{}
try {
TurnTokens(input = -1, output = 50)
error("should have thrown")
} catch (_: IllegalArgumentException) {}
}
@Test
fun `backwards compatible encode-decode round trip preserves both context and tokens`() {
val original = MessageBodyPayload(
content = listOf(Content.Text("hi")),
context = MessageContext(
origin = MessageOrigin.USER,
description = "x",
),
tokens = TurnTokens(input = 100, output = 50),
)
val json = kotlinx.serialization.json.Json.encodeToString(MessageBodyPayload.serializer(), original)
val decoded = kotlinx.serialization.json.Json.decodeFromString(MessageBodyPayload.serializer(), json)
assertEquals(100, decoded.tokens?.input)
assertEquals(50, decoded.tokens?.output)
assertEquals(1, decoded.content.size)
}
@Test
fun `decodeBodyPayload round-trip preserves tokens`() {
val tokens = TurnTokens(input = 250, output = 80)
val encoded = encodeBodyPayload(
content = listOf(Content.Text("ok")),
tokens = tokens,
)
val decoded = decodeBodyPayload(encoded)
assertEquals(250, decoded.tokens?.input)
assertEquals(80, decoded.tokens?.output)
}
private fun assistant(
msgId: String, convId: String, input: Int, output: Int, at: Instant,
) = MessageRecord.AssistantMessage(
id = msgId,
conversationId = convId,
content = listOf(Content.Text("reply")),
createdAt = at,
tokens = TurnTokens(input = input, output = output),
)
}
@@ -1,31 +1,24 @@
package pw.binom.agentik.storage.inmemory
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.MutableSharedFlow
import kotlinx.coroutines.flow.asSharedFlow
import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock
import pw.binom.agentik.messageLog.MessageEvent
import pw.binom.agentik.messageLog.MessageRecord
import pw.binom.agentik.messageLog.MutableMessageStore
import pw.binom.agentik.messageLog.TokenStats
import kotlin.time.Instant
/**
* Thread-safe append-only лог сообщений в памяти.
*
* Хранит все записи в одном `List<MessageRecord>`, индексированном
* `conversationId`. `tokenStats` обходит только assistant-записи и
* суммирует `TurnTokens`.
* `conversationId`. Подходит для unit-тестов и ephemeral runtime (например,
* CLI-сессии или in-memory demo), где не нужен долгоживущий persistence.
*
* В отличие от SQLite-импла, не нуждается в SQLDelight и работает в
* любом KMP-таргете (включая iOS/native, где SQLite через Android driver
* недоступен).
* В отличие от ksqlite-импла, не требует native SQLite — работает в любом
* KMP-таргете (включая iOS/native).
*/
class InMemoryMessageStore : MutableMessageStore {
private val byConv: MutableMap<String, MutableList<MessageRecord>> = mutableMapOf()
private val events = MutableSharedFlow<MessageEvent>(extraBufferCapacity = 64)
private val mutex = Mutex()
override suspend fun append(record: MessageRecord) {
@@ -33,7 +26,6 @@ class InMemoryMessageStore : MutableMessageStore {
val list = byConv.getOrPut(record.conversationId) { mutableListOf() }
list.add(record)
}
events.tryEmit(MessageEvent.Appended(record.conversationId, record))
}
override suspend fun list(conversationId: String, after: Instant, offset: Int, limit: Int): List<MessageRecord> {
@@ -45,30 +37,6 @@ class InMemoryMessageStore : MutableMessageStore {
}
}
override suspend fun listAll(conversationId: String): List<MessageRecord> {
mutex.withLock {
return byConv[conversationId].orEmpty().sortedBy { it.createdAt }.toList()
}
}
override fun events(): Flow<MessageEvent> = events.asSharedFlow()
override suspend fun tokenStats(conversationId: String): TokenStats {
mutex.withLock {
var input = 0L
var output = 0L
var turns = 0
for (rec in byConv[conversationId].orEmpty()) {
if (rec !is MessageRecord.AssistantMessage) continue
val t = rec.tokens ?: continue
input += t.input
output += t.output
turns += 1
}
return TokenStats(turns = turns, inputTokens = input, outputTokens = output)
}
}
override fun close() {
// no-op
}
@@ -2,28 +2,26 @@ package pw.binom.agentik.storage.inmemory
import pw.binom.agentik.messageLog.Content
import pw.binom.agentik.messageLog.MessageRecord
import pw.binom.agentik.messageLog.TurnTokens
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertNull
import kotlin.test.assertTrue
import kotlin.time.Instant
import kotlinx.coroutines.async
import kotlinx.coroutines.flow.first
import kotlinx.coroutines.yield
import kotlinx.coroutines.flow.toList
import kotlinx.coroutines.test.runTest
class InMemoryMessageStoreTest {
@Test
fun `append and listAll returns inserted records in createdAt order`() = runTest {
fun `append and listFlow returns inserted records in createdAt order`() = runTest {
val store = InMemoryMessageStore()
val t0 = Instant.parse("2026-09-15T10:00:00Z")
val u = MessageRecord.UserMessage("u1", "c1", listOf(Content.Text("hi")), t0)
val a = MessageRecord.AssistantMessage("a1", "c1", listOf(Content.Text("hello")), t0.plus(kotlin.time.Duration.parse("PT1S")), null)
store.append(u)
store.append(a)
val all = store.listAll("c1")
val all = store.listFlow("c1", Instant.DISTANT_PAST).toList()
assertEquals(listOf("u1", "a1"), all.map { it.id })
}
@@ -51,82 +49,19 @@ class InMemoryMessageStoreTest {
}
@Test
fun `list returns empty for unknown conversation`() = runTest {
fun `listFlow returns empty for unknown conversation`() = runTest {
val store = InMemoryMessageStore()
assertEquals(emptyList(), store.listAll("none"))
assertTrue(store.listFlow("none", Instant.DISTANT_PAST).toList().isEmpty())
}
@Test
fun `tokenStats sums across assistant messages with tokens`() = runTest {
val store = InMemoryMessageStore()
val t0 = Instant.parse("2026-09-15T10:00:00Z")
store.append(
MessageRecord.AssistantMessage(
"a1", "c1", listOf(Content.Text("r")), t0.plus(kotlin.time.Duration.parse("PT1S")),
TurnTokens(input = 100, output = 50),
)
)
store.append(
MessageRecord.AssistantMessage(
"a2", "c1", listOf(Content.Text("r")), t0.plus(kotlin.time.Duration.parse("PT2S")),
TurnTokens(input = 200, output = 80),
)
)
// user без tokens
store.append(
MessageRecord.UserMessage("u1", "c1", listOf(Content.Text("hi")), t0)
)
val stats = store.tokenStats("c1")
assertEquals(2, stats.turns)
assertEquals(300L, stats.inputTokens)
assertEquals(130L, stats.outputTokens)
assertEquals(430L, stats.totalTokens)
}
@Test
fun `tokenStats skips assistant messages without tokens`() = runTest {
val store = InMemoryMessageStore()
val t0 = Instant.parse("2026-09-15T10:00:00Z")
store.append(
MessageRecord.AssistantMessage("a1", "c1", listOf(Content.Text("r")), t0, tokens = null)
)
val stats = store.tokenStats("c1")
assertEquals(0, stats.turns)
assertEquals(0L, stats.inputTokens)
}
@Test
fun `tokenStats returns zeros for unknown conversation`() = runTest {
val store = InMemoryMessageStore()
val stats = store.tokenStats("none")
assertEquals(0, stats.turns)
assertEquals(0L, stats.inputTokens)
}
@Test
fun `events flow emits Appended on append`() = runTest {
val store = InMemoryMessageStore()
val t0 = Instant.parse("2026-09-15T10:00:00Z")
// SharedFlow не реплеит — запускаем коллектор ДО append, чтобы не потерять эвент.
val events = store.events()
val deferred = async { events.first() }
yield() // даём коллектору подписаться ДО append — иначе SharedFlow без replay потеряет эвент
store.append(MessageRecord.UserMessage("u1", "c1", listOf(Content.Text("hi")), t0))
val ev = deferred.await()
assertTrue(ev is pw.binom.agentik.messageLog.MessageEvent.Appended)
val appended = ev as pw.binom.agentik.messageLog.MessageEvent.Appended
assertEquals("c1", appended.conversationId)
assertEquals("u1", appended.record.id)
}
@Test
fun `isolates conversations - list returns only requested conv`() = runTest {
fun `isolates conversations - listFlow returns only requested conv`() = runTest {
val store = InMemoryMessageStore()
val t0 = Instant.parse("2026-09-15T10:00:00Z")
store.append(MessageRecord.UserMessage("u1", "c1", listOf(Content.Text("hi")), t0))
store.append(MessageRecord.UserMessage("u2", "c2", listOf(Content.Text("hello")), t0))
assertEquals(listOf("u1"), store.listAll("c1").map { it.id })
assertEquals(listOf("u2"), store.listAll("c2").map { it.id })
assertEquals(listOf("u1"), store.listFlow("c1", Instant.DISTANT_PAST).toList().map { it.id })
assertEquals(listOf("u2"), store.listFlow("c2", Instant.DISTANT_PAST).toList().map { it.id })
}
@Test
@@ -148,7 +83,7 @@ class InMemoryMessageStoreTest {
context = null,
)
store.append(msg)
val got = store.listAll("c1").first() as MessageRecord.UserMessage
val got = store.listFlow("c1", Instant.DISTANT_PAST).first() as MessageRecord.UserMessage
assertNull(got.context)
}
}
+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
}
}
-78
View File
@@ -1,78 +0,0 @@
# `:storage-sqlite` — SQLite реализация `:storage-core` (JVM-only)
## Что что это
Production persistence для `:standalone` на [SQLDelight](https://cashapp.github.io/sqldelight/):
- **messages** — append-only журнал с `conversation_id`, `created_at`.
- **working_memory** — rolling buffer последних 100 entries, типы
в JSON (`UserMessage / AssistantMessage / ToolExchange / SystemPrompt`).
- **conversations** — метаданные (id, title, model, timestamps).
- **reflections** — произвольные заметки ("I notice you often
prefer short replies").
- Промпт хранителя (`@mem0`) индексирован отдельно для быстрого
доступа.
Решает: стабильная, локальная, нулевая-настройка БД. Подходит и для
desktop-продакшена, и для Android, и для тестов (через Testcontainers).
## Где используется
- `:standalone` подключает по умолчанию (`storage.db` = путь из
`AGENTIK_DB_PATH`).
## Как подключить
```kotlin
jvmMain.dependencies {
implementation("pw.binom.agentik:storage-sqlite:0.1.0")
implementation("pw.binom.agentik:storage-core:0.1.0")
}
val storage = SqliteStorageSystem.open(Path("agentik.db"))
val messages: MessageStore = storage.messages
```
## Версии
`gradle/libs.versions.toml` → `[versions] agentik-storage-sqlite`.
Зависит от `app.cash.sqldelight:sqlite-driver:2.1.0` (через
`gradle/libs.versions.toml`).
## Тесты
```
./gradlew :storage-sqlite:jvmTest
```
Покрывают: миграции (через `migrations/` каталог и SQLDelight
`*.sqm`), round-trip, race-conditions (concurrent append), paged
flow.
## Что в схеме (упрощённо)
```sql
CREATE TABLE messages (
id TEXT PRIMARY KEY,
conversation_id TEXT NOT NULL,
created_at TEXT NOT NULL, -- ISO Instant
kind TEXT NOT NULL, -- 'user', 'assistant', 'tool_call', 'tool_result'
body_json TEXT NOT NULL
);
CREATE INDEX idx_messages_conv_time ON messages(conversation_id, created_at);
CREATE TABLE working_memory (
conversation_id TEXT NOT NULL,
entry_id TEXT PRIMARY KEY,
created_at TEXT NOT NULL,
kind TEXT NOT NULL,
body_json TEXT NOT NULL
);
```
Полная схема + миграции — в `src/jvmMain/sqldelight/`.
## Текущий статус
Используется продакшеном. Миграции 1.0+.
-46
View File
@@ -1,46 +0,0 @@
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(":message-store-api"))
api(project(":message-log-api"))
api(project(":working-memory-api"))
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")
}
}
}
@@ -1,80 +0,0 @@
package pw.binom.agentik.storage.sqlite
import kotlin.time.Instant
import pw.binom.agentik.messageStore.ConversationRecord
import pw.binom.agentik.messageStore.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),
)
@@ -1,169 +0,0 @@
package pw.binom.agentik.storage.sqlite
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 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) : MutableMessageStore {
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),
)
}
@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")
}
}
@@ -1,144 +0,0 @@
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.messageStore.Reflection
import pw.binom.agentik.messageStore.ReflectionEvent
import pw.binom.agentik.messageStore.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
}
@@ -1,113 +0,0 @@
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.messageStore.ConversationStore
import pw.binom.agentik.messageLog.MutableMessageStore
import pw.binom.agentik.messageStore.ReflectionStore
import pw.binom.agentik.storage.sqlite.SqliteReflectionStore
import pw.binom.agentik.workingMemory.WorkingMemoryStore
/**
* Корневой объект SQLite-слоя: держит [SqlDriver] и четыре store'а.
*/
class SqliteStores private constructor(
val driver: SqlDriver,
val conversations: ConversationStore,
val messages: MutableMessageStore,
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)
}
}
}
}
@@ -1,93 +0,0 @@
package pw.binom.agentik.storage.sqlite
import kotlinx.serialization.json.Json
import pw.binom.agentik.workingMemory.WorkingMemoryEntry
import pw.binom.agentik.workingMemory.WorkingMemoryRow
import pw.binom.agentik.workingMemory.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.User -> "user"
is WorkingMemoryEntry.Assistant -> "assistant"
is WorkingMemoryEntry.ToolExchange -> "tool_exchange"
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.messageStore.Ids.new("wm")
@@ -1,42 +0,0 @@
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 = ?;
@@ -1,33 +0,0 @@
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 = ?;
@@ -1,43 +0,0 @@
-- 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;
@@ -1,38 +0,0 @@
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', ?, ?);
@@ -1,124 +0,0 @@
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.messageStore.Reflection
import pw.binom.agentik.messageStore.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("медленно", "путаю"),
)
}
+6
View File
@@ -13,7 +13,13 @@ kotlin {
jvmToolchain(21)
jvm()
macosX64()
macosArm64()
iosX64()
iosArm64()
iosSimulatorArm64()
linuxX64()
linuxArm64()
mingwX64()
sourceSets {