refactor(storage): split :storage-core into message-store-api + working-memory-api
ci / JVM build + tests (push) Failing after 2m5s
ci / JVM build + tests (push) Failing after 2m5s
Разделяет монолитный :storage-core на 3 модуля с чёткими границами:
:message-store-api — MessageStore, ReflectionStore, EventStore, ConversationStore +
Content, Payload, MessageContext, Ids, MessageEvent
(audit log + event stream)
:working-memory-api — WorkingMemoryStore + WorkingMemoryEntry
(runtime context с compaction)
:storage-bundle — StorageBundle агрегатор, зависит от обоих
(только для server-side runtime)
Пакеты:
pw.binom.agentik.storage.* → УДАЛЕНО
pw.binom.agentik.messageStore.* — append-only API
pw.binom.agentik.messageStore.events.* — EventStore + EventRecord
pw.binom.agentik.workingMemory.* — WM API
pw.binom.agentik.storageBundle.* — aggregator
Зачем:
- Тонкий клиент может подтянуть ТОЛЬКО :message-store-api (~15KB, нет
compaction-логики, нет MessageStore+WorkingMemoryStore cross-deps).
- Android-agent в будущем подключит :message-store-api для audit log,
серверный runtime — :storage-bundle со всем.
- Компиляционные границы защищают от случайной зависимости от WM
в read-only клиентах (раньше один :storage-core не давал такой
гарантии).
Миграция:
- Имплементации (:storage-inmemory, :storage-sqlite, :storage-ksqlite)
обновили package + добавили deps на оба API модуля + :storage-bundle.
- Тесты из :storage-core (PersistenceTest, SqliteStoresMigrationTest,
TokenStatsTest) переехали в :standalone, получили testImplementation
на оба API модуля и импорты новых типов.
- 52 файла в :standalone, :agent-toolsets, :llm-tools, :server, :client,
:agentik-cli обновили FQN.
- :storage-core удалён.
Совместимость схем не меняется — все 5 impl'ов (3 backend × 5 store) хранят
данные в тех же таблицах, миграция между Sqlite и Ksqlite возможна через SQL dump.
Тесты:
standalone 178 ✅
agent-toolsets 36 ✅
storage-inmemory 47 ✅
storage-sqlite 17 ✅ (включая переехавшие persistence/* + tokenStats)
storage-ksqlite 36 ✅
---
Total: 314 tests, 0 failures
This commit is contained in:
@@ -0,0 +1,33 @@
|
||||
plugins {
|
||||
alias(libs.plugins.kotlin.multiplatform)
|
||||
alias(libs.plugins.kotlin.serialization)
|
||||
}
|
||||
|
||||
// Public API для append-only storage агента (audit log, replay-after-disconnect,
|
||||
// conversation metadata). Этот модуль НЕ содержит working memory — он для
|
||||
// тонких клиентов, которые хотят только события/сообщения без тяжёлой
|
||||
// runtime-логики.
|
||||
//
|
||||
// Зависимости:
|
||||
// - kotlinx.coroutines / serialization (KMP).
|
||||
// - БЕЗ :storage-core — мы его полностью выпиливаем, интерфейсы переезжают сюда.
|
||||
|
||||
kotlin {
|
||||
jvmToolchain(21)
|
||||
|
||||
jvm()
|
||||
linuxX64()
|
||||
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,29 @@
|
||||
package pw.binom.agentik.messageStore
|
||||
|
||||
import kotlinx.serialization.SerialName
|
||||
import kotlinx.serialization.Serializable
|
||||
|
||||
/**
|
||||
* Часть контента сообщения на уровне хранилища.
|
||||
*
|
||||
* Намеренно НЕ зависит от [pw.binom.agentik.proto.Content] — маппинг
|
||||
* `:proto.Content ↔ Content` живёт в `Mapping.kt`. Структурно типы
|
||||
* идентичны, но даёт возможность заменить transport-протокол без миграции
|
||||
* таблиц.
|
||||
*
|
||||
* Image сериализуется в JSON через base64 (стандарт для kotlinx-serialization).
|
||||
*/
|
||||
@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()
|
||||
}
|
||||
}
|
||||
+14
@@ -0,0 +1,14 @@
|
||||
package pw.binom.agentik.messageStore
|
||||
|
||||
import kotlin.time.Instant
|
||||
|
||||
/**
|
||||
* Snapshot диалога. В таблице `conversation` хранится как есть.
|
||||
*/
|
||||
data class ConversationRecord(
|
||||
val id: String,
|
||||
val title: String?,
|
||||
val isTemporal: Boolean,
|
||||
val createdAt: Instant,
|
||||
val updatedAt: Instant,
|
||||
)
|
||||
+27
@@ -0,0 +1,27 @@
|
||||
package pw.binom.agentik.messageStore
|
||||
|
||||
import kotlin.time.Instant
|
||||
|
||||
/**
|
||||
* CRUD по таблице `conversation`.
|
||||
*/
|
||||
interface ConversationStore : AutoCloseable {
|
||||
|
||||
/** Создать или обновить snapshot диалога. */
|
||||
suspend fun upsert(record: ConversationRecord)
|
||||
|
||||
/** Диалог по id, или `null`. */
|
||||
suspend fun get(id: String): ConversationRecord?
|
||||
|
||||
/** Удалить диалог (вместе с его сообщениями и working memory). */
|
||||
suspend fun delete(id: String): Boolean
|
||||
|
||||
/** Список диалогов, отсортированный по `updatedAt` DESC. */
|
||||
suspend fun list(offset: Int, limit: Int): List<ConversationRecord>
|
||||
|
||||
/** Переименовать диалог; `null` для сброса заголовка. Возвращает новый `updatedAt` или `null`, если не найден. */
|
||||
suspend fun rename(id: String, title: String?): Instant?
|
||||
|
||||
/** Обновить `updatedAt` диалога (например, после отправки сообщения). */
|
||||
suspend fun touch(id: String, now: Instant)
|
||||
}
|
||||
@@ -0,0 +1,15 @@
|
||||
package pw.binom.agentik.messageStore
|
||||
|
||||
import kotlin.uuid.Uuid
|
||||
|
||||
/**
|
||||
* Генератор id. Использует `kotlin.uuid.Uuid` из stdlib (KMP: jvm + native),
|
||||
* чтобы не зависеть от `java.util.UUID` и подготовить код к linuxX64-сборке.
|
||||
*
|
||||
* Сохраняет формат `<prefix>-xxxxxxxx-xxxx-xxxx-xxxx-xxxxxxxxxxxx` —
|
||||
* `:server` его парсит как opaque string, без знания внутренней структуры.
|
||||
*/
|
||||
object Ids {
|
||||
fun new(prefix: String): String = "$prefix-${Uuid.random()}"
|
||||
fun reflection(): String = new("refl")
|
||||
}
|
||||
+40
@@ -0,0 +1,40 @@
|
||||
package pw.binom.agentik.messageStore
|
||||
|
||||
import kotlinx.serialization.SerialName
|
||||
import kotlinx.serialization.Serializable
|
||||
import kotlinx.serialization.json.JsonElement
|
||||
|
||||
/**
|
||||
* Контекст инициации хода (кто/что и почему).
|
||||
*
|
||||
* Дубликат типа из `:proto` (`pw.binom.agentik.proto.MessageContext`):
|
||||
* живёт в `:storage-core` чтобы не тащить `:proto` в слой хранения данных.
|
||||
* Маппинг между ними — в `MessageRecord.toProto()` / `ProtoMessage.toStorage()`.
|
||||
*
|
||||
* Используется:
|
||||
* - в `MessageRecord.UserMessage.context` — фиксируется в audit log;
|
||||
* - в `WorkingMemoryEntry.User.context` — попадает в LLM-нагрузку
|
||||
* как префикс к user-тексту (для не-USER origin'ов).
|
||||
*
|
||||
* Сериализация в `payload_json` (SQLite) — через kotlinx-serialization,
|
||||
* формат snake_case для enum origin.
|
||||
*/
|
||||
@Serializable
|
||||
enum class MessageOrigin {
|
||||
@SerialName("user")
|
||||
USER,
|
||||
|
||||
@SerialName("system")
|
||||
SYSTEM,
|
||||
|
||||
@SerialName("event")
|
||||
EVENT,
|
||||
}
|
||||
|
||||
@Serializable
|
||||
data class MessageContext(
|
||||
val origin: MessageOrigin,
|
||||
val description: String? = null,
|
||||
val sourceId: String? = null,
|
||||
val metadata: JsonElement? = null,
|
||||
)
|
||||
+144
@@ -0,0 +1,144 @@
|
||||
package pw.binom.agentik.messageStore
|
||||
|
||||
import kotlinx.serialization.SerialName
|
||||
import kotlinx.serialization.Serializable
|
||||
import kotlin.time.Instant
|
||||
|
||||
/**
|
||||
* Token usage одного assistant turn'а: сколько токенов модель обработала
|
||||
* на входе (system + history + tools + user message) и сколько сгенерировала
|
||||
* (assistant text + tool calls + tool results, всё что LiteConversation
|
||||
* добавила к истории за этот turn).
|
||||
*
|
||||
* `input` — снимок [LiteConversation.tokenCount] перед первым send() в turn'е
|
||||
* (после подготовки user-сообщения). `output` — дельта после завершения turn'а
|
||||
* (включая все tool loop итерации).
|
||||
*
|
||||
* Persisted в `message.payload_json` — никаких schema-миграций при добавлении
|
||||
* полей. Optional: `null` для исторических сообщений или для бэкендов, не
|
||||
* отдающих tokenCount (например off-line embedded LLM без контекст-счётчика).
|
||||
*/
|
||||
@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" }
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Запись в таблице `message` (append-only audit) и `working_memory` (mutable view).
|
||||
*
|
||||
* Использует sealed-иерархию: подтипы `User`/`Assistant`/`ToolCall`/`ToolResult`
|
||||
* живут и там, и там. `Summary`/`System` — только в `working_memory`
|
||||
* (синтетические строки, созданные при суммаризации или как system-prompt).
|
||||
*
|
||||
* Все подтипы несут [id] (UUID, стабильный между лайв-стримом Event и историей),
|
||||
* [conversationId] и [createdAt].
|
||||
*
|
||||
* Поля, специфичные для подтипа, сериализуются в JSON в `payload_json`
|
||||
* колонке SQLite — это даёт гибкость без миграций при добавлении полей.
|
||||
*/
|
||||
@Serializable
|
||||
sealed interface MessageRecord {
|
||||
|
||||
val id: String
|
||||
val conversationId: String
|
||||
val createdAt: Instant
|
||||
|
||||
/** Подтип сообщения с телом из [Content]. */
|
||||
@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,
|
||||
/**
|
||||
* Контекст инициации хода: кто/что вызвал этот turn. `null` —
|
||||
* обычное user-сообщение. См. [MessageContext].
|
||||
*
|
||||
* Persisted через `payload_json` SQLite (см. `Payload.kt`).
|
||||
*/
|
||||
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,
|
||||
/**
|
||||
* Token usage этого turn'а: сколько input+output токенов обработала
|
||||
* модель. Заполняется в [ChatConversation.runTurn] через
|
||||
* `LiteConversation.tokenCount()` (до/после send).
|
||||
*/
|
||||
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
|
||||
|
||||
/**
|
||||
* Терминальная запись провалившегося хода. Только audit log
|
||||
* (в working_memory не пишется — модель не должна видеть ошибки прошлых ходов).
|
||||
*/
|
||||
@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
|
||||
|
||||
/** Синтетическое: суммаризация старого контекста. Только в working_memory. */
|
||||
@Serializable
|
||||
@SerialName("summary")
|
||||
data class Summary(
|
||||
override val id: String,
|
||||
override val conversationId: String,
|
||||
val text: String,
|
||||
override val createdAt: Instant,
|
||||
) : MessageRecord
|
||||
|
||||
/** Синтетическое: system-prompt, введённый при создании диалога. Только в working_memory. */
|
||||
@Serializable
|
||||
@SerialName("system")
|
||||
data class System(
|
||||
override val id: String,
|
||||
override val conversationId: String,
|
||||
val text: String,
|
||||
override val createdAt: Instant,
|
||||
) : MessageRecord
|
||||
}
|
||||
@@ -0,0 +1,56 @@
|
||||
package pw.binom.agentik.messageStore
|
||||
|
||||
import kotlinx.coroutines.flow.Flow
|
||||
import kotlin.time.Instant
|
||||
|
||||
/**
|
||||
* Суммарная статистика токенов диалога — aggregate по всем assistant-сообщениям
|
||||
* в audit log. Делит input/output и считает число assistant-ходов.
|
||||
*
|
||||
* Используется:
|
||||
* - В startup banner'е агента (см. `Main.kt` → "conversation stats").
|
||||
* - На HTTP фасаде `/agentik/conversations/{id}/stats` (если будет endpoint).
|
||||
* - В клиентских дашбордах для оценки cost.
|
||||
*/
|
||||
data class TokenStats(
|
||||
val turns: Int,
|
||||
val inputTokens: Long,
|
||||
val outputTokens: Long,
|
||||
) {
|
||||
val totalTokens: Long get() = inputTokens + outputTokens
|
||||
}
|
||||
|
||||
/**
|
||||
* Append-only audit log сообщений (`message` table).
|
||||
*
|
||||
* Только `insert` и чтение. Никаких обновлений, никакого удаления (кроме
|
||||
* каскадного удаления вместе с [ConversationStore.delete]).
|
||||
*/
|
||||
interface MessageStore : AutoCloseable {
|
||||
|
||||
/** Добавить запись в audit log. `conversationId` берётся из [MessageRecord.conversationId]. */
|
||||
suspend fun append(record: MessageRecord)
|
||||
|
||||
/**
|
||||
* Страница audit-сообщений диалога после [after] (UTC), отсортированная
|
||||
* по `createdAt ASC`. Для первоначальной загрузки передай `Instant.DISTANT_PAST`.
|
||||
*/
|
||||
suspend fun list(conversationId: String, after: Instant, offset: Int, limit: Int): List<MessageRecord>
|
||||
|
||||
/** Все сообщения диалога, отсортированные по `createdAt ASC` (для rebuild working memory). */
|
||||
suspend fun listAll(conversationId: String): List<MessageRecord>
|
||||
|
||||
/** Лайв-стрим новых сообщений (для SSE-подписчиков). По умолчанию — пустой. */
|
||||
fun events(): Flow<MessageEvent> = kotlinx.coroutines.flow.emptyFlow()
|
||||
|
||||
/**
|
||||
* Суммарная token-статистика по диалогу: input/output/turns. Один проход
|
||||
* по всем assistant-сообщениям. Дёшево (на практике < 1мс на SQLite).
|
||||
*/
|
||||
suspend fun tokenStats(conversationId: String): TokenStats
|
||||
}
|
||||
|
||||
sealed interface MessageEvent {
|
||||
val conversationId: String
|
||||
data class Appended(override val conversationId: String, val record: MessageRecord) : MessageEvent
|
||||
}
|
||||
@@ -0,0 +1,79 @@
|
||||
package pw.binom.agentik.messageStore
|
||||
|
||||
import kotlinx.serialization.SerialName
|
||||
import kotlinx.serialization.Serializable
|
||||
import kotlinx.serialization.builtins.ListSerializer
|
||||
import kotlinx.serialization.json.Json
|
||||
|
||||
/**
|
||||
* JSON-формат для тел user/assistant сообщений: список [Content], опционально
|
||||
* с [MessageContext] (для user — кто инициировал ход) и [TurnTokens]
|
||||
* (для assistant — сколько токенов стоил этот turn).
|
||||
*
|
||||
* Encoded-формат:
|
||||
* ```
|
||||
* {"content": [ ...Content ], "context": {...MessageContext?}, "tokens": {...TurnTokens?}}
|
||||
* ```
|
||||
*
|
||||
* Backward compat: при чтении старых строк, где payload был просто
|
||||
* `[ ... ]` (без обёртки), парсер падает на wrapper-формат и fallback'ит
|
||||
* к `ListSerializer<Content>` — такие строки возвращаются с `context = null`,
|
||||
* `tokens = null`.
|
||||
*/
|
||||
private val bodyJson = Json {
|
||||
ignoreUnknownKeys = true
|
||||
encodeDefaults = true
|
||||
explicitNulls = false
|
||||
}
|
||||
|
||||
@Serializable
|
||||
data class MessageBodyPayload(
|
||||
val content: List<Content>,
|
||||
@SerialName("context")
|
||||
val context: MessageContext? = null,
|
||||
/**
|
||||
* Token usage для assistant (input + output). `null` для user-сообщений,
|
||||
* для исторических assistant-сообщений без метрики и для бэкендов без
|
||||
* tokenCount() (off-line модели).
|
||||
*/
|
||||
val tokens: TurnTokens? = null,
|
||||
)
|
||||
|
||||
/**
|
||||
* Сериализует тело user (или assistant) сообщения в JSON-строку для
|
||||
* `payload_json` SQLite. Для user может нести [context] — кто инициировал ход;
|
||||
* для assistant может нести [tokens] — token usage этого turn'а.
|
||||
*/
|
||||
fun encodeBodyPayload(
|
||||
content: List<Content>,
|
||||
context: MessageContext? = null,
|
||||
tokens: TurnTokens? = null,
|
||||
): String = bodyJson.encodeToString(
|
||||
MessageBodyPayload.serializer(),
|
||||
MessageBodyPayload(content = content, context = context, tokens = tokens),
|
||||
)
|
||||
|
||||
/**
|
||||
* Десериализует тело сообщения: возвращает тройку `(content, context, tokens)`.
|
||||
* Контекст и токены — null если:
|
||||
* - поля отсутствуют в новом формате;
|
||||
* - payload в старом plain-array формате (миграция не нужна — fallback).
|
||||
*/
|
||||
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) {
|
||||
// Старый формат: голый JSON-массив Content, без обёртки.
|
||||
val arr = bodyJson.decodeFromString(ListSerializer(Content.serializer()), json)
|
||||
BodyDecoded(arr, null, null)
|
||||
}
|
||||
}
|
||||
+51
@@ -0,0 +1,51 @@
|
||||
package pw.binom.agentik.messageStore
|
||||
|
||||
import kotlinx.coroutines.flow.Flow
|
||||
import kotlin.time.Instant
|
||||
|
||||
/**
|
||||
* Self-reflection запись — что агент "думает" о качестве своих последних ходов.
|
||||
*
|
||||
* - [score]: 1..5 (самооценка качества)
|
||||
* - [weakSpots]: конкретные слабые места ("медленно ищу в Y", "путаю A и B")
|
||||
* - [summary]: свободный комментарий в markdown (что заметил, что улучшить)
|
||||
*
|
||||
* Источник: one-shot LLM-размышление после каждых N ходов (см. AGENTIK_REFLECTION_INTERVAL).
|
||||
* Используется в `ChatAgent.buildSystemPrompt` как "слабые места за последнее время".
|
||||
*/
|
||||
data class Reflection(
|
||||
val id: String,
|
||||
val conversationId: String?,
|
||||
val createdAt: Instant,
|
||||
val turnsAnalyzed: Int,
|
||||
val score: Int,
|
||||
val summary: String,
|
||||
val weakSpots: List<String>,
|
||||
)
|
||||
|
||||
/**
|
||||
* Хранилище рефлексий. Backed by SQLDelight `reflection` таблицу (impl в
|
||||
* `:storage-sqlite`) или in-memory (impl в `:storage-inmemory`).
|
||||
*
|
||||
* Рефлексии — append-only: старые записи удаляются [deleteOlderThan] (cleanup)
|
||||
* или архивируются через [Curator]-подобный процесс, но не редактируются.
|
||||
*/
|
||||
interface ReflectionStore : AutoCloseable {
|
||||
suspend fun insert(reflection: Reflection)
|
||||
suspend fun get(id: String): Reflection?
|
||||
/** Самые свежие рефлексии (по всему агенту). */
|
||||
suspend fun listRecent(limit: Int = 10): List<Reflection>
|
||||
/** Рефлексии для конкретного диалога. */
|
||||
suspend fun listForConversation(conversationId: String, limit: Int = 10): List<Reflection>
|
||||
suspend fun deleteOlderThan(cutoff: Instant)
|
||||
suspend fun count(): Int
|
||||
|
||||
/** Стрим новых рефлексий для подписчиков (для UI в будущем). */
|
||||
fun events(): Flow<ReflectionEvent> = kotlinx.coroutines.flow.emptyFlow()
|
||||
|
||||
override fun close()
|
||||
}
|
||||
|
||||
sealed interface ReflectionEvent {
|
||||
data class Created(val reflection: Reflection) : ReflectionEvent
|
||||
}
|
||||
+109
@@ -0,0 +1,109 @@
|
||||
package pw.binom.agentik.messageStore.events
|
||||
|
||||
import kotlinx.serialization.Serializable
|
||||
import kotlin.time.Instant
|
||||
|
||||
/**
|
||||
* Persistent event log для replay после disconnect.
|
||||
*
|
||||
* Зачем: SSE-подписка на `/events` и `/conversations/{id}/events` — cold (no replay).
|
||||
* Если клиент отвалился на час, он пропустил всё. [EventStore] даёт:
|
||||
* - append() — producer (ChatAgent) пишет при каждом event
|
||||
* - query() — consumer (server SSE replay endpoint) читает по cursor
|
||||
* - prune() — maintenance: удалить старые events по TTL
|
||||
*
|
||||
* Не заменяет live-подписку на [MutableSharedFlow] — это для долговременного
|
||||
* хранения, а live-streaming идёт через in-memory channel.
|
||||
*
|
||||
* Платформо-агностичный interface (KMP): impl в `:storage-sqlite` (JVM-only),
|
||||
* `:storage-inmemory` (KMP, для тестов и dev), и в будущем `:storage-sqlite-android`
|
||||
* для Android-агента.
|
||||
*
|
||||
* Payload — opaque JSON string. [storage-core] не должен знать про
|
||||
* kotlinx.serialization или [AgentEvent]/[Conversation.Event] типы (это `:proto`-шный
|
||||
* слой). Конвертация — на стороне producer'а (:standalone ChatAgent).
|
||||
*/
|
||||
interface EventStore : AutoCloseable {
|
||||
/**
|
||||
* Записать event. Идемпотентен по [EventRecord.id] — повторный append с тем же
|
||||
* id это no-op (важно для retry при network failure между producer'ом и БД).
|
||||
*/
|
||||
suspend fun append(record: EventRecord)
|
||||
|
||||
/**
|
||||
* Catchup query для reconnect.
|
||||
*
|
||||
* @param conversationId если `null` — глобальный catchup (для `/events/replay`).
|
||||
* если задан — только этот диалог (для `/conversations/{id}/events/replay`).
|
||||
* @param afterId exclusive cursor: вернуть events СТРОГО после этого id.
|
||||
* Если `null` — с начала.
|
||||
* @param limit max количество records (default 100). Caller делает пагинацию
|
||||
* пока `result.size == limit`.
|
||||
*
|
||||
* Сортировка: по [EventRecord.createdAt] ASC, ties broken по [EventRecord.id] ASC
|
||||
* (т.к. id содержит timestamp-like prefix в нашей схеме, это даёт стабильный порядок).
|
||||
*/
|
||||
suspend fun query(
|
||||
conversationId: String? = null,
|
||||
afterId: String? = null,
|
||||
limit: Int = 100,
|
||||
): List<EventRecord>
|
||||
|
||||
/**
|
||||
* Maintenance: удалить events старше [olderThan]. Возвращает количество удалённых.
|
||||
* Default вызывается из background scope раз в час (TTL = 24h типично).
|
||||
*/
|
||||
suspend fun pruneOlderThan(olderThan: Instant): Int
|
||||
|
||||
/** Сколько events всего хранится (для observability). */
|
||||
suspend fun count(): Int
|
||||
|
||||
override fun close()
|
||||
}
|
||||
|
||||
/**
|
||||
* Платформо-агностичная запись event'а.
|
||||
*
|
||||
* @param id уникальный в пределах EventStore. Convention: `"ev-<uuid>"`.
|
||||
* Используется как cursor для [EventStore.query].
|
||||
* @param conversationId `null` для agent-level events (Created/Deleted/Renamed).
|
||||
* Задан для conversation events.
|
||||
* @param createdAt UTC timestamp. Используется для сортировки в query() и для TTL в prune().
|
||||
* @param type kind of event (для индексирования/фильтрации; payload всё равно opaque).
|
||||
* @param payload opaque JSON string. Producer (:standalone ChatAgent) сериализует
|
||||
* [pw.binom.agentik.proto.AgentEvent] или [pw.binom.agentik.proto.Event]
|
||||
* в JSON перед append. Consumer (:server Routes) парсит обратно.
|
||||
*
|
||||
* Note: payload хранится as String, не ByteArray, чтобы не зависеть от kotlinx
|
||||
* serialization и platform-specific binary encoding в [storage-core].
|
||||
*/
|
||||
@Serializable
|
||||
data class EventRecord(
|
||||
val id: String,
|
||||
val conversationId: String?,
|
||||
val createdAt: Instant,
|
||||
val type: EventType,
|
||||
val payload: String,
|
||||
)
|
||||
|
||||
/**
|
||||
* Категория event'а — для индексирования и для фильтрации в query().
|
||||
*
|
||||
* Naming: AGENT_* — agent-level, CONVERSATION_* — turn-level.
|
||||
*/
|
||||
enum class EventType {
|
||||
AGENT_CREATED,
|
||||
AGENT_DELETED,
|
||||
AGENT_RENAMED,
|
||||
|
||||
CONVERSATION_START_REASONING,
|
||||
CONVERSATION_START_RESPONSE,
|
||||
CONVERSATION_APPEND_TEXT,
|
||||
CONVERSATION_APPEND_IMAGE,
|
||||
CONVERSATION_TOOL_CALL,
|
||||
CONVERSATION_TOOL_RESULT,
|
||||
CONVERSATION_END,
|
||||
CONVERSATION_INTERRUPTED,
|
||||
CONVERSATION_ERROR,
|
||||
// reserved for future — adding new variants doesn't break older consumers
|
||||
}
|
||||
Reference in New Issue
Block a user