storage-core: новый KMP-модуль с интерфейсами хранилища
Выносим интерфейсы и data-классы истории диалога (MessageStore / WorkingMemoryStore / ConversationStore / ReflectionStore + соответствующие sealed-иерархии MessageRecord / WorkingMemoryEntry / Content / ConversationRecord / Reflection + payload-утилиты) из :standalone в отдельный KMP-модуль :storage-core (pw.binom.agentik.storage). Цель — подготовка к Android-портированию и подключению альтернативных реализаций хранилища без затягивания всей :standalone. Дальше (commit 2/3) — :storage-inmemory и :storage-sqlite как самостоятельные модули, плюс :storage-android (deferred). Изменения: - Новый :storage-core (KMP, commonMain only, jvm + native таргеты) — 12 файлов - StorageBundle агрегатор (conversationStore + messageStore + workingMemoryStore + reflectionStore; SkillStore живёт в :skills и подключается отдельно) - 11 файлов импортов в :standalone переключены на новый пакет - SqliteReflectionStore оставлен в :standalone до commit 3 (зависит от SQLDelight AgentikDatabase, которую ещё не отвязали от :standalone) - 4 теста перенесены в :standalone/.../storage/ с обновлённым пакетом - PayloadTest переехал в :storage-core/commonTest (тестирует чистые типы) Tests: 264/264 green (179 :standalone + 6 :storage-core + прочие JVM-модули)
This commit is contained in:
@@ -0,0 +1,29 @@
|
||||
package pw.binom.agentik.storage
|
||||
|
||||
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()
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,14 @@
|
||||
package pw.binom.agentik.storage
|
||||
|
||||
import kotlin.time.Instant
|
||||
|
||||
/**
|
||||
* Snapshot диалога. В таблице `conversation` хранится как есть.
|
||||
*/
|
||||
data class ConversationRecord(
|
||||
val id: String,
|
||||
val title: String?,
|
||||
val isTemporal: Boolean,
|
||||
val createdAt: Instant,
|
||||
val updatedAt: Instant,
|
||||
)
|
||||
@@ -0,0 +1,27 @@
|
||||
package pw.binom.agentik.storage
|
||||
|
||||
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.storage
|
||||
|
||||
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")
|
||||
}
|
||||
@@ -0,0 +1,40 @@
|
||||
package pw.binom.agentik.storage
|
||||
|
||||
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,
|
||||
)
|
||||
@@ -0,0 +1,144 @@
|
||||
package pw.binom.agentik.storage
|
||||
|
||||
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.storage
|
||||
|
||||
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.storage
|
||||
|
||||
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)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,51 @@
|
||||
package pw.binom.agentik.storage
|
||||
|
||||
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
|
||||
}
|
||||
@@ -0,0 +1,29 @@
|
||||
package pw.binom.agentik.storage
|
||||
|
||||
/**
|
||||
* Агрегатор всех storage-интерфейсов, нужных агенту для работы с историей диалога.
|
||||
*
|
||||
* Реализации:
|
||||
* - [SqliteStorageBundle] — JVM-only, основная. SQLDelight + SQLite.
|
||||
* - [InMemoryStorageBundle] — Map-based, для тестов.
|
||||
* - `:storage-android` — Android-специфичная (deferred, Android-сборка в планах).
|
||||
*
|
||||
* `SkillStore` НЕ входит сюда — он живёт в модуле `:skills` (другая ответственность:
|
||||
* не сообщения/рефлексии, а контент-файлы навыков) и принимается отдельно в `ChatAgent`.
|
||||
*
|
||||
* AutoCloseable: один `close()` закрывает все четыре store'а. В реализациях,
|
||||
* которые не владеют ресурсами (in-memory), close — no-op.
|
||||
*/
|
||||
data class StorageBundle(
|
||||
val conversationStore: ConversationStore,
|
||||
val messageStore: MessageStore,
|
||||
val workingMemoryStore: WorkingMemoryStore,
|
||||
val reflectionStore: ReflectionStore,
|
||||
) : AutoCloseable {
|
||||
override fun close() {
|
||||
conversationStore.close()
|
||||
messageStore.close()
|
||||
workingMemoryStore.close()
|
||||
reflectionStore.close()
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,64 @@
|
||||
package pw.binom.agentik.storage
|
||||
|
||||
import kotlinx.serialization.SerialName
|
||||
import kotlinx.serialization.Serializable
|
||||
|
||||
/**
|
||||
* Запись в working memory диалога: ровно то, что агент сейчас видит в
|
||||
* LLM-контексте. Упорядочено по `order_idx` (заполняется в store при append).
|
||||
*
|
||||
* Sealed-иерархия: для v1 — `System` (синтетический system-prompt),
|
||||
* `User`/`Assistant` (реплики с ссылкой на audit log через [sourceMessageId]).
|
||||
* Суммаризация (для v2) добавит подтип `Summary`.
|
||||
*/
|
||||
@Serializable
|
||||
sealed interface WorkingMemoryEntry {
|
||||
|
||||
/** Ссылка на исходное сообщение в audit log (`message.id`). `null` для синтетических строк. */
|
||||
val sourceMessageId: String?
|
||||
|
||||
/** Синтетический system-prompt, добавляется при создании диалога. */
|
||||
@Serializable
|
||||
@SerialName("system")
|
||||
data class System(val text: String) : WorkingMemoryEntry {
|
||||
override val sourceMessageId: String? = null
|
||||
}
|
||||
|
||||
/** Реплика пользователя. */
|
||||
@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
|
||||
|
||||
/**
|
||||
* Синтетический блок: суммаризация старых ходов, сгенерированная при
|
||||
* 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
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,58 @@
|
||||
package pw.binom.agentik.storage
|
||||
|
||||
import kotlin.time.Instant
|
||||
|
||||
/**
|
||||
* Одна строка `working_memory` таблицы (внутреннее представление store).
|
||||
*
|
||||
* Используется для тестов и для перестроения [WorkingMemoryEntry] из row.
|
||||
* Агент не должен с этим типом работать напрямую — он работает с
|
||||
* [WorkingMemoryEntry] через [WorkingMemoryStore].
|
||||
*/
|
||||
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 WorkingMemoryStore : 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,81 @@
|
||||
package pw.binom.agentik.storage
|
||||
|
||||
import kotlinx.serialization.json.JsonPrimitive
|
||||
import kotlinx.serialization.json.buildJsonObject
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertNull
|
||||
import kotlin.test.assertTrue
|
||||
|
||||
/**
|
||||
* Тесты payload-формата: контекст инициации хода персистится через
|
||||
* `payload_json` SQLite, старый plain-array формат читается без потерь.
|
||||
*/
|
||||
class PayloadTest {
|
||||
|
||||
@Test
|
||||
fun `user context roundtrips through MessageBodyPayload`() {
|
||||
val ctx = MessageContext(
|
||||
origin = MessageOrigin.USER,
|
||||
description = "irc PRIVMSG",
|
||||
sourceId = "irc:agentik",
|
||||
)
|
||||
val encoded = encodeBodyPayload(listOf(Content.Text("hi")), ctx)
|
||||
// новый формат: wrapper-объект с полем context
|
||||
assertTrue(encoded.startsWith("{"), "expected wrapped object, got: $encoded")
|
||||
assertTrue(encoded.contains("\"context\""), "expected context field, got: $encoded")
|
||||
|
||||
val decoded = decodeBodyPayload(encoded)
|
||||
assertEquals(listOf(Content.Text("hi")), decoded.content)
|
||||
assertEquals(ctx, decoded.context)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `event context roundtrips with metadata`() {
|
||||
val ctx = MessageContext(
|
||||
origin = MessageOrigin.EVENT,
|
||||
description = "scheduled cron morning-briefing",
|
||||
sourceId = "cron-42",
|
||||
metadata = buildJsonObject {
|
||||
put("scheduledAt", JsonPrimitive("2026-09-14T08:00:00Z"))
|
||||
put("rule", JsonPrimitive("0 8 * * *"))
|
||||
},
|
||||
)
|
||||
val encoded = encodeBodyPayload(listOf(Content.Text("wake up")), ctx)
|
||||
val decoded = decodeBodyPayload(encoded)
|
||||
assertEquals(MessageOrigin.EVENT, decoded.context?.origin)
|
||||
assertEquals("scheduled cron morning-briefing", decoded.context?.description)
|
||||
assertEquals("cron-42", decoded.context?.sourceId)
|
||||
assertEquals("0 8 * * *", (decoded.context?.metadata?.let { (it as kotlinx.serialization.json.JsonObject)["rule"] } as? JsonPrimitive)?.content)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `null context produces wrapper without context field`() {
|
||||
val encoded = encodeBodyPayload(listOf(Content.Text("hello")))
|
||||
assertTrue(encoded.contains("\"content\""), "expected content field, got: $encoded")
|
||||
val decoded = decodeBodyPayload(encoded)
|
||||
assertEquals(listOf(Content.Text("hello")), decoded.content)
|
||||
assertNull(decoded.context)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `legacy plain-array payload still decodes (backward compat)`() {
|
||||
val legacy = "[" +
|
||||
"""{"type":"text","body":"old message"}""" +
|
||||
"]"
|
||||
val decoded = decodeBodyPayload(legacy)
|
||||
assertEquals(listOf(Content.Text("old message")), decoded.content)
|
||||
assertNull(decoded.context)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `system context with description only roundtrips`() {
|
||||
val ctx = MessageContext(origin = MessageOrigin.SYSTEM, description = "agent startup greeting")
|
||||
val encoded = encodeBodyPayload(listOf(Content.Text("boot")), ctx)
|
||||
val decoded = decodeBodyPayload(encoded)
|
||||
assertEquals(MessageOrigin.SYSTEM, decoded.context?.origin)
|
||||
assertEquals("agent startup greeting", decoded.context?.description)
|
||||
assertNull(decoded.context?.sourceId)
|
||||
assertNull(decoded.context?.metadata)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user