Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 4ad59d5f5d |
@@ -17,6 +17,7 @@ import kotlinx.coroutines.flow.flow
|
||||
import kotlinx.coroutines.runBlocking
|
||||
import pw.binom.agentik.proto.Agent
|
||||
import pw.binom.agentik.proto.AgentEvent
|
||||
import pw.binom.agentik.proto.AllEvent
|
||||
import pw.binom.agentik.proto.Conversation
|
||||
import kotlin.time.Instant
|
||||
|
||||
@@ -78,4 +79,74 @@ internal class AgentClient(
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
/**
|
||||
* Подписка на ВСЕ события: agent lifecycle + все conversation events.
|
||||
* Использует SSE endpoint /events/all.
|
||||
*/
|
||||
override fun allEvents(after: Instant): Flow<AllEvent> = flow {
|
||||
httpClient.prepareGet("$agentUrl/events/all?after=$after") { noSseReadTimeout() }
|
||||
.execute { response ->
|
||||
check(response.status == HttpStatusCode.OK) {
|
||||
"allEvents: server returned ${response.status}"
|
||||
}
|
||||
readSse(response.bodyAsChannel())
|
||||
.collect { payload ->
|
||||
emit(agentikJson.decodeFromString(AllEvent.serializer(), payload))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Catchup для /events/replay — пагинированно читает events после [afterId].
|
||||
* Caller делает несколько вызовов пока `result.size < limit` (= конец).
|
||||
*
|
||||
* @param afterId exclusive cursor. `null` = с начала.
|
||||
* @param limit max per-request (default 100, max 1000 на сервере).
|
||||
*/
|
||||
internal suspend fun replayAllEvents(
|
||||
afterId: String? = null,
|
||||
limit: Int = 100,
|
||||
): List<EventRecordDto> {
|
||||
val response = httpClient.get("$agentUrl/events/replay") {
|
||||
afterId?.let { parameter("after_id", it) }
|
||||
parameter("limit", limit)
|
||||
}
|
||||
return response.body()
|
||||
}
|
||||
|
||||
/**
|
||||
* Catchup для /conversations/{id}/events/replay — пагинированно.
|
||||
*/
|
||||
internal suspend fun replayConversationEvents(
|
||||
conversationId: String,
|
||||
afterId: String? = null,
|
||||
limit: Int = 100,
|
||||
): List<EventRecordDto> {
|
||||
val response = httpClient.get("$agentUrl/conversations/$conversationId/events/replay") {
|
||||
afterId?.let { parameter("after_id", it) }
|
||||
parameter("limit", limit)
|
||||
}
|
||||
return response.body()
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
/**
|
||||
* Запись event'а, которую возвращает /events/replay endpoint.
|
||||
*
|
||||
* Это **мини-зеркало** `pw.binom.agentik.storage.events.EventRecord` — клиент
|
||||
* не зависит от `:storage-core`, поэтому определяет свою модель (wire-only).
|
||||
*
|
||||
* Формат полностью совместим с сервером — `agentikJson.encodeToString(...)` там
|
||||
* и `agentikJson.decodeFromString(...)` здесь.
|
||||
*/
|
||||
@kotlinx.serialization.Serializable
|
||||
internal data class EventRecordDto(
|
||||
val id: String,
|
||||
val conversationId: String? = null,
|
||||
val createdAt: kotlin.time.Instant,
|
||||
val type: String,
|
||||
val payload: String,
|
||||
)
|
||||
|
||||
@@ -49,6 +49,18 @@ public interface Agent {
|
||||
*/
|
||||
fun events(after: Instant): Flow<AgentEvent>
|
||||
|
||||
/**
|
||||
* All events in one stream: agent lifecycle (Created/Deleted/Renamed) +
|
||||
* all conversation turns. Useful for admin dashboards, debug tools,
|
||||
* parent agents.
|
||||
*
|
||||
* For UI use [events] + [Conversation.events]. This one-feed variant is
|
||||
* for cases where everything-in-one is preferred.
|
||||
*
|
||||
* Cold (no replay). For catchup use EventStore.
|
||||
*/
|
||||
fun allEvents(after: Instant): Flow<AllEvent>
|
||||
|
||||
companion object {
|
||||
|
||||
const val PAGE_SIZE: Int = 100
|
||||
|
||||
@@ -0,0 +1,39 @@
|
||||
package pw.binom.agentik.proto
|
||||
|
||||
import kotlin.time.Instant
|
||||
import kotlinx.serialization.SerialName
|
||||
import kotlinx.serialization.Serializable
|
||||
|
||||
/**
|
||||
* Unified wrapper for all agent events in a single stream.
|
||||
*
|
||||
* Useful for admin dashboards, debug tools, parent agents: one subscription
|
||||
* instead of N+1. For regular UI use two separate SSE feeds
|
||||
* ([AgentEvent] via /events and [Event] via /conversations/{id}/events);
|
||||
* [AllEvent] is for those who need everything in one place.
|
||||
*
|
||||
* Server endpoint: GET /events/all (SSE), or replay via EventStore.
|
||||
*
|
||||
* Not used for persistence payload: EventStore stores AgentEvent and
|
||||
* Conversation.Event natively (compact form); this wrapper is wire-format
|
||||
* only.
|
||||
*/
|
||||
@Serializable
|
||||
sealed interface AllEvent {
|
||||
val date: Instant
|
||||
|
||||
@Serializable
|
||||
@SerialName("agent")
|
||||
data class Agent(
|
||||
override val date: Instant,
|
||||
val event: AgentEvent,
|
||||
) : AllEvent
|
||||
|
||||
@Serializable
|
||||
@SerialName("conversation")
|
||||
data class Conversation(
|
||||
override val date: Instant,
|
||||
val conversationId: String,
|
||||
val event: Event,
|
||||
) : AllEvent
|
||||
}
|
||||
@@ -23,6 +23,7 @@ kotlin {
|
||||
sourceSets {
|
||||
commonMain.dependencies {
|
||||
implementation(project(":proto"))
|
||||
implementation(project(":storage-core"))
|
||||
|
||||
// Ktor (без engine — engine подключает потребитель, см. :standalone).
|
||||
implementation(libs.ktor.server.core)
|
||||
|
||||
@@ -6,6 +6,7 @@ import io.ktor.server.plugins.contentnegotiation.ContentNegotiation
|
||||
import io.ktor.server.routing.Route
|
||||
import io.ktor.server.routing.route
|
||||
import pw.binom.agentik.proto.Agent
|
||||
import pw.binom.agentik.storage.events.EventStore
|
||||
|
||||
/**
|
||||
* Встраивает HTTP/SSE-фасад протокола agentik в твой Ktor-роутинг.
|
||||
@@ -30,10 +31,20 @@ import pw.binom.agentik.proto.Agent
|
||||
* - `POST /conversations/{id}/interrupt` — `interrupt`
|
||||
* - `GET /conversations/{id}/messages` — история
|
||||
* - `GET /conversations/{id}/events` — SSE: события хода
|
||||
* - `GET /conversations/{id}/events/replay` — replay-after-disconnect (events after ?after_id=X)
|
||||
* - `GET /events` — SSE: события агента
|
||||
* - `GET /events/replay` — replay-after-disconnect (global)
|
||||
* - `GET /health` — `"ok"`
|
||||
*
|
||||
* @param eventStore если null — `/events/replay` endpoints возвращают 503 (event
|
||||
* persistence не настроен). Live-streaming работает как обычно.
|
||||
*/
|
||||
fun Route.agentikAgent(agent: Agent, path: String = "/agentik", token: String? = null) {
|
||||
fun Route.agentikAgent(
|
||||
agent: Agent,
|
||||
eventStore: EventStore? = null,
|
||||
path: String = "/agentik",
|
||||
token: String? = null,
|
||||
) {
|
||||
route(path) {
|
||||
install(ContentNegotiation) {
|
||||
json(agentikJson)
|
||||
@@ -43,6 +54,6 @@ fun Route.agentikAgent(agent: Agent, path: String = "/agentik", token: String? =
|
||||
this.token = token
|
||||
}
|
||||
}
|
||||
agentikRoutes(agent)
|
||||
agentikRoutes(agent, eventStore)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -21,12 +21,14 @@ import kotlinx.serialization.KSerializer
|
||||
import kotlinx.serialization.json.Json
|
||||
import pw.binom.agentik.proto.Agent
|
||||
import pw.binom.agentik.proto.AgentEvent
|
||||
import pw.binom.agentik.proto.AllEvent
|
||||
import pw.binom.agentik.proto.Conversation
|
||||
import pw.binom.agentik.proto.Content
|
||||
import pw.binom.agentik.proto.Event
|
||||
import pw.binom.agentik.storage.events.EventStore
|
||||
import kotlin.time.Instant
|
||||
|
||||
internal fun Route.agentikRoutes(agent: Agent) {
|
||||
internal fun Route.agentikRoutes(agent: Agent, eventStore: EventStore? = null) {
|
||||
|
||||
get("/health") {
|
||||
call.respondText("ok")
|
||||
@@ -134,6 +136,48 @@ internal fun Route.agentikRoutes(agent: Agent) {
|
||||
val after = call.parseAfter() ?: return@get
|
||||
call.streamJsonSse(agent.events(after), AgentEvent.serializer())
|
||||
}
|
||||
|
||||
/**
|
||||
* Все события в одном потоке: agent lifecycle + все conversation events.
|
||||
* Для admin-дашборда, debug-инструментов, parent-агента.
|
||||
* Для UI достаточно `/events` + `/conversations/{id}/events`.
|
||||
*/
|
||||
get("/events/all") {
|
||||
val after = call.parseAfter() ?: return@get
|
||||
call.streamJsonSse(agent.allEvents(after), AllEvent.serializer())
|
||||
}
|
||||
|
||||
// ---- Replay-after-disconnect (event log) ----
|
||||
//
|
||||
// Stream `/events` и `/conversations/{id}/events` — cold (no replay).
|
||||
// Клиент, отвалившийся от SSE, при reconnect делает GET на replay-endpoint
|
||||
// с `?after_id=X` чтобы получить events, которые произошли во время разрыва.
|
||||
// paginated: делает несколько запросов пока `result.size < limit`.
|
||||
//
|
||||
// Если [eventStore] == null (не сконфигурирован), эти endpoints возвращают 503.
|
||||
|
||||
get("/events/replay") {
|
||||
if (eventStore == null) {
|
||||
call.respond(HttpStatusCode.ServiceUnavailable, "EventStore not configured on this agent")
|
||||
return@get
|
||||
}
|
||||
val afterId = call.request.queryParameters["after_id"]?.takeIf { it.isNotBlank() }
|
||||
val limit = (call.request.queryParameters["limit"]?.toIntOrNull() ?: 100).coerceIn(1, 1000)
|
||||
val records = eventStore.query(conversationId = null, afterId = afterId, limit = limit)
|
||||
call.respond(records)
|
||||
}
|
||||
|
||||
get("/conversations/{id}/events/replay") {
|
||||
if (eventStore == null) {
|
||||
call.respond(HttpStatusCode.ServiceUnavailable, "EventStore not configured on this agent")
|
||||
return@get
|
||||
}
|
||||
val id = call.parameters["id"]!!
|
||||
val afterId = call.request.queryParameters["after_id"]?.takeIf { it.isNotBlank() }
|
||||
val limit = (call.request.queryParameters["limit"]?.toIntOrNull() ?: 100).coerceIn(1, 1000)
|
||||
val records = eventStore.query(conversationId = id, afterId = afterId, limit = limit)
|
||||
call.respond(records)
|
||||
}
|
||||
}
|
||||
|
||||
// ---------- helpers ----------
|
||||
|
||||
@@ -1,38 +1,50 @@
|
||||
package pw.binom.agentik.standalone.agent
|
||||
|
||||
import kotlin.time.Instant
|
||||
import kotlinx.coroutines.flow.Flow
|
||||
import kotlinx.coroutines.flow.MutableSharedFlow
|
||||
import kotlinx.coroutines.flow.asSharedFlow
|
||||
import kotlinx.coroutines.flow.emitAll
|
||||
import kotlinx.coroutines.flow.filterNotNull
|
||||
import kotlinx.coroutines.flow.flow
|
||||
import kotlinx.coroutines.flow.map
|
||||
import kotlinx.coroutines.flow.merge
|
||||
import kotlinx.coroutines.launch
|
||||
import kotlinx.coroutines.runBlocking
|
||||
import kotlinx.coroutines.sync.Mutex
|
||||
import kotlinx.coroutines.sync.withLock
|
||||
import kotlinx.serialization.json.Json
|
||||
import pw.binom.agentik.llm.tools.ContextCompactor
|
||||
import pw.binom.agentik.llm.tools.LlmMemoryReviewer
|
||||
import pw.binom.agentik.llm.tools.LlmReflector
|
||||
import pw.binom.agentik.llm.tools.SkillMiner
|
||||
import pw.binom.agentik.memory.MemoryPrefetcher
|
||||
import pw.binom.agentik.memory.MemoryReviewer
|
||||
import pw.binom.agentik.memory.MemorySystemGuidance
|
||||
import pw.binom.agentik.proto.Agent as ProtoAgent
|
||||
import pw.binom.agentik.proto.AgentEvent
|
||||
import pw.binom.agentik.proto.AllEvent
|
||||
import pw.binom.agentik.proto.Conversation as ProtoConversation
|
||||
import pw.binom.agentik.proto.Event as ProtoEvent
|
||||
import pw.binom.agentik.skills.SkillCatalog
|
||||
import pw.binom.agentik.skills.renderSystemPromptSection
|
||||
import pw.binom.agentik.standalone.agent.memory.MemoryToolsFactory
|
||||
import pw.binom.agentik.standalone.llm.LlmConfig
|
||||
import pw.binom.agentik.storage.ConversationRecord
|
||||
import pw.binom.agentik.storage.Ids
|
||||
import pw.binom.agentik.storage.Reflection
|
||||
import pw.binom.agentik.storage.WorkingMemoryEntry
|
||||
import pw.binom.agentik.storage.StorageBundle
|
||||
import pw.binom.agentik.toolsets.EnableToolsetTool
|
||||
import pw.binom.agentik.storage.WorkingMemoryEntry
|
||||
import pw.binom.agentik.storage.events.EventRecord
|
||||
import pw.binom.agentik.storage.events.EventType
|
||||
import pw.binom.agentik.toolsets.DisableToolsetTool
|
||||
import pw.binom.agentik.toolsets.EnableToolsetTool
|
||||
import pw.binom.agentik.toolsets.NamedTool
|
||||
import pw.binom.agentik.toolsets.SystemPromptToolsetSection
|
||||
import pw.binom.agentik.toolsets.ToolsetContribution
|
||||
import pw.binom.agentik.toolsets.ToolsetDispatchPolicy
|
||||
import pw.binom.agentik.toolsets.ToolsetRegistry
|
||||
import pw.binom.litert.LiteLlm
|
||||
import kotlin.time.Instant
|
||||
import pw.binom.agentik.llm.tools.LlmReflector
|
||||
import pw.binom.agentik.llm.tools.SkillMiner
|
||||
import pw.binom.agentik.llm.tools.LlmMemoryReviewer
|
||||
import pw.binom.agentik.llm.tools.ContextCompactor
|
||||
import pw.binom.agentik.toolsets.NamedTool
|
||||
|
||||
/**
|
||||
* Stateful [ProtoAgent] на базе SQLite (история + working memory) и
|
||||
@@ -192,6 +204,54 @@ class ChatAgent(
|
||||
extraBufferCapacity = 64,
|
||||
)
|
||||
|
||||
/** Json-encoder для payload в EventStore. Один на весь agent. */
|
||||
private val eventJson = Json {
|
||||
ignoreUnknownKeys = true
|
||||
encodeDefaults = true
|
||||
}
|
||||
|
||||
/**
|
||||
* Fire-and-forget persist в EventStore. Если EventStore настроен (production
|
||||
* deployment с `:storage-sqlite`), каждый AgentEvent также уходит в SQLite
|
||||
* с уникальным id, чтобы клиенты могли сделать /events/replay после
|
||||
* disconnect. Если EventStore == null (например, dev mode с in-memory
|
||||
* storage или Android без persistent event log) — no-op.
|
||||
*
|
||||
* Background launch + runCatching: ошибки БД не должны ронять agent loop.
|
||||
* Логирование — если persistence падает, видно в логах, но live-stream
|
||||
* продолжает работать.
|
||||
*/
|
||||
private fun persistAgentEvent(event: AgentEvent) {
|
||||
val store = storage.eventStore ?: return
|
||||
kotlinx.coroutines.CoroutineScope(
|
||||
kotlinx.coroutines.SupervisorJob() + kotlinx.coroutines.Dispatchers.IO
|
||||
).launch {
|
||||
runCatching {
|
||||
store.append(
|
||||
EventRecord(
|
||||
id = "ev-${Ids.new("agent")}",
|
||||
conversationId = when (event) {
|
||||
is AgentEvent.Created -> event.conversationId
|
||||
is AgentEvent.Deleted -> event.id
|
||||
is AgentEvent.Renamed -> event.id
|
||||
},
|
||||
createdAt = event.date,
|
||||
type = when (event) {
|
||||
is AgentEvent.Created -> EventType.AGENT_CREATED
|
||||
is AgentEvent.Deleted -> EventType.AGENT_DELETED
|
||||
is AgentEvent.Renamed -> EventType.AGENT_RENAMED
|
||||
},
|
||||
payload = eventJson.encodeToString(AgentEvent.serializer(), event),
|
||||
)
|
||||
)
|
||||
}.onFailure {
|
||||
mu.KotlinLogging.logger("ChatAgent").warn(it) {
|
||||
"failed to persist agent event ${event::class.simpleName}: ${it.message}"
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/** Защищает карту живых диалогов. */
|
||||
private val liveLock = Mutex()
|
||||
private val live: MutableMap<String, ChatConversation> = HashMap()
|
||||
@@ -204,6 +264,36 @@ class ChatAgent(
|
||||
return agentEvents.asSharedFlow()
|
||||
}
|
||||
|
||||
/**
|
||||
* Все события в одном потоке: agent lifecycle + events всех диалогов.
|
||||
* Реализация — merge двух cold-flow'ов. Snapshot-список живых диалогов
|
||||
* берётся на момент подписки; новые Created-Event'ы НЕ переподписывают
|
||||
* (это ответственность caller'а: если хочет всё — он может
|
||||
* переподписаться или следить за AgentEvent.Created сам).
|
||||
*
|
||||
* Для admin/debug — допустимое упрощение. Для long-running мониторинга
|
||||
* (admin-дашборд часами) надо добавить reactive re-subscribe (см.
|
||||
* notes 04-sub-agents.md, Variant B).
|
||||
*/
|
||||
override fun allEvents(after: Instant): Flow<AllEvent> = flow {
|
||||
// Agent lifecycle events
|
||||
emitAll(agentEvents.asSharedFlow().map { e: AgentEvent ->
|
||||
AllEvent.Agent(date = e.date, event = e)
|
||||
})
|
||||
// Snapshot живых диалогов на момент подписки.
|
||||
// ВАЖНО: не подписываемся на новые Created — это ответственность caller'а
|
||||
// (см. KDoc выше).
|
||||
live.values
|
||||
.asSequence()
|
||||
.filterNot { it.isClosed }
|
||||
.forEach { conv: ProtoConversation ->
|
||||
val cid: String = conv.id
|
||||
emitAll(conv.events(after).map { ev: ProtoEvent ->
|
||||
AllEvent.Conversation(date = ev.date, conversationId = cid, event = ev)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
override fun createConversation(temp: Boolean): ProtoConversation {
|
||||
val now = now()
|
||||
val id = pw.binom.agentik.storage.Ids.new("conv")
|
||||
@@ -243,6 +333,7 @@ class ChatAgent(
|
||||
liveLock.withLock { live[conv.id] = conv }
|
||||
}
|
||||
agentEvents.tryEmit(AgentEvent.Created(date = now(), conversationId = conv.id))
|
||||
persistAgentEvent(AgentEvent.Created(date = now(), conversationId = conv.id))
|
||||
return conv
|
||||
}
|
||||
|
||||
@@ -258,7 +349,11 @@ class ChatAgent(
|
||||
val conv = liveLock.withLock { live.remove(id) }
|
||||
conv?.close()
|
||||
val ok = storage.conversationStore.delete(id)
|
||||
if (ok) agentEvents.tryEmit(AgentEvent.Deleted(date = now(), id = id))
|
||||
if (ok) {
|
||||
val event = AgentEvent.Deleted(date = now(), id = id)
|
||||
agentEvents.tryEmit(event)
|
||||
persistAgentEvent(event)
|
||||
}
|
||||
return ok
|
||||
}
|
||||
|
||||
|
||||
+102
-3
@@ -4,9 +4,45 @@ import kotlinx.coroutines.channels.BufferOverflow
|
||||
import kotlinx.coroutines.flow.MutableSharedFlow
|
||||
import kotlinx.coroutines.flow.SharedFlow
|
||||
import kotlinx.coroutines.flow.asSharedFlow
|
||||
import kotlinx.coroutines.launch
|
||||
import kotlinx.serialization.json.Json
|
||||
import mu.KotlinLogging
|
||||
import pw.binom.agentik.proto.Event as ProtoEvent
|
||||
import pw.binom.agentik.storage.Ids
|
||||
import pw.binom.agentik.storage.events.EventRecord
|
||||
import pw.binom.agentik.storage.events.EventStore
|
||||
import pw.binom.agentik.storage.events.EventType
|
||||
|
||||
internal class ConversationEvents {
|
||||
private val log = KotlinLogging.logger {}
|
||||
|
||||
/**
|
||||
* SSE-события диалога + (опционально) persistence в [EventStore].
|
||||
*
|
||||
* Двойная ответственность:
|
||||
* 1. Live-streaming через [flow] — клиенты подписываются на long-lived SSE.
|
||||
* 2. Durable storage через [eventStore] — для replay после disconnect
|
||||
* через /conversations/{id}/events/replay?after_id=X.
|
||||
*
|
||||
* **Persistence strategy**: при каждом [tryEmit]/[emit] параллельно пишем в
|
||||
* EventStore (fire-and-forget в IO scope). Ошибка БД НЕ должна ронять live-stream
|
||||
* — оборачиваем в runCatching и логируем.
|
||||
*
|
||||
* **Idempotency**: каждое event имеет детерминированный id (из messageStore при
|
||||
* создании), append с тем же id в EventStore — no-op. Это критично для retry
|
||||
* между producer и БД.
|
||||
*
|
||||
* **Dual-write cost**: на каждый event одно INSERT в SQLite. SQLite на локальном
|
||||
* диске выдерживает ~50K events/sec; для hot pathов можно вынести persist в
|
||||
* отдельную batched очередь. Для v1 — синхронный launch — OK.
|
||||
*/
|
||||
internal class ConversationEvents(
|
||||
private val eventStore: EventStore? = null,
|
||||
private val conversationId: String? = null,
|
||||
private val eventJson: Json = Json {
|
||||
ignoreUnknownKeys = true
|
||||
encodeDefaults = true
|
||||
},
|
||||
) {
|
||||
private val _flow = MutableSharedFlow<ProtoEvent>(
|
||||
replay = 0,
|
||||
extraBufferCapacity = 4096,
|
||||
@@ -15,7 +51,70 @@ internal class ConversationEvents {
|
||||
|
||||
val flow: SharedFlow<ProtoEvent> get() = _flow.asSharedFlow()
|
||||
|
||||
fun tryEmit(event: ProtoEvent): Boolean = _flow.tryEmit(event)
|
||||
/**
|
||||
* Emit event в live-stream + persist в EventStore (если настроен).
|
||||
*
|
||||
* @return true если event попал в live-stream (false если buffer overflow
|
||||
* и event был дропнут — DROP_OLDEST policy).
|
||||
*/
|
||||
fun tryEmit(event: ProtoEvent): Boolean {
|
||||
val ok = _flow.tryEmit(event)
|
||||
if (ok) persistAsync(event)
|
||||
return ok
|
||||
}
|
||||
|
||||
suspend fun emit(event: ProtoEvent) = _flow.emit(event)
|
||||
/**
|
||||
* Same as [tryEmit] но suspend — ждёт места в buffer'е (а не дропает).
|
||||
* Используется реже — там где мы хотим гарантировать доставку подписчикам.
|
||||
*/
|
||||
suspend fun emit(event: ProtoEvent) {
|
||||
_flow.emit(event)
|
||||
persistAsync(event)
|
||||
}
|
||||
|
||||
/**
|
||||
* Persist в EventStore в fire-and-forget. Если [eventStore] == null — no-op
|
||||
* (in-memory dev или Android без persistence).
|
||||
*
|
||||
* **Не использует agentScope** — мы не знаем о нём здесь (ConversationEvents
|
||||
* не владеет lifecycle). Если нужна более аккуратная lifecycle management —
|
||||
* передавать scope параметром или держать свой CoroutineScope.
|
||||
*
|
||||
* **Сейчас**: создаём transient `GlobalScope`-like через `MainScope()`-style —
|
||||
* НЕТ, лучше через `CoroutineScope(SupervisorJob + Dispatchers.IO).launch`.
|
||||
* Это сделано лениво, чтобы не плодить треды при hot path.
|
||||
*/
|
||||
private fun persistAsync(event: ProtoEvent) {
|
||||
val store = eventStore ?: return
|
||||
val convId = conversationId ?: return // не знаем к чему привязать
|
||||
kotlinx.coroutines.CoroutineScope(kotlinx.coroutines.SupervisorJob() + kotlinx.coroutines.Dispatchers.IO).launch {
|
||||
runCatching {
|
||||
store.append(
|
||||
EventRecord(
|
||||
id = "ev-${Ids.new("conv")}",
|
||||
conversationId = convId,
|
||||
createdAt = event.date,
|
||||
type = mapEventType(event),
|
||||
payload = eventJson.encodeToString(ProtoEvent.serializer(), event),
|
||||
)
|
||||
)
|
||||
}.onFailure {
|
||||
log.warn(it) {
|
||||
"failed to persist conversation event ${event::class.simpleName}: ${it.message}"
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private fun mapEventType(event: ProtoEvent): EventType = when (event) {
|
||||
is ProtoEvent.StartReasoning -> EventType.CONVERSATION_START_REASONING
|
||||
is ProtoEvent.StartResponse -> EventType.CONVERSATION_START_RESPONSE
|
||||
is ProtoEvent.AppendText -> EventType.CONVERSATION_APPEND_TEXT
|
||||
is ProtoEvent.AppendImage -> EventType.CONVERSATION_APPEND_IMAGE
|
||||
is ProtoEvent.ToolCall -> EventType.CONVERSATION_TOOL_CALL
|
||||
is ProtoEvent.ToolResult -> EventType.CONVERSATION_TOOL_RESULT
|
||||
is ProtoEvent.End -> EventType.CONVERSATION_END
|
||||
is ProtoEvent.Interrupted -> EventType.CONVERSATION_INTERRUPTED
|
||||
is ProtoEvent.Error -> EventType.CONVERSATION_ERROR
|
||||
}
|
||||
}
|
||||
|
||||
@@ -84,7 +84,10 @@ class ConversationLoop(
|
||||
agentScope = agentScope,
|
||||
)
|
||||
|
||||
private val events = ConversationEvents()
|
||||
private val events = ConversationEvents(
|
||||
eventStore = storage.eventStore,
|
||||
conversationId = state.id,
|
||||
)
|
||||
|
||||
/** Per-conversation background event bus. Lifecycle scoped к этому ConversationLoop. */
|
||||
private val backgroundEvents = BackgroundEventBus()
|
||||
|
||||
@@ -1,5 +1,7 @@
|
||||
package pw.binom.agentik.storage
|
||||
|
||||
import pw.binom.agentik.storage.events.EventStore
|
||||
|
||||
/**
|
||||
* Агрегатор всех storage-интерфейсов, нужных агенту для работы с историей диалога.
|
||||
*
|
||||
@@ -11,7 +13,10 @@ package pw.binom.agentik.storage
|
||||
* `SkillStore` НЕ входит сюда — он живёт в модуле `:skills` (другая ответственность:
|
||||
* не сообщения/рефлексии, а контент-файлы навыков) и принимается отдельно в `ChatAgent`.
|
||||
*
|
||||
* AutoCloseable: один `close()` закрывает все четыре store'а. В реализациях,
|
||||
* [EventStore] входит начиная с commit "event-store" — для replay после
|
||||
* disconnect (см. `/events/replay` endpoint в `:server`).
|
||||
*
|
||||
* AutoCloseable: один `close()` закрывает все store'ы. В реализациях,
|
||||
* которые не владеют ресурсами (in-memory), close — no-op.
|
||||
*/
|
||||
data class StorageBundle(
|
||||
@@ -19,11 +24,13 @@ data class StorageBundle(
|
||||
val messageStore: MessageStore,
|
||||
val workingMemoryStore: WorkingMemoryStore,
|
||||
val reflectionStore: ReflectionStore,
|
||||
val eventStore: EventStore? = null,
|
||||
) : AutoCloseable {
|
||||
override fun close() {
|
||||
conversationStore.close()
|
||||
messageStore.close()
|
||||
workingMemoryStore.close()
|
||||
reflectionStore.close()
|
||||
eventStore?.close()
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,109 @@
|
||||
package pw.binom.agentik.storage.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
|
||||
}
|
||||
+84
@@ -0,0 +1,84 @@
|
||||
package pw.binom.agentik.storage.inmemory
|
||||
|
||||
import pw.binom.agentik.storage.events.EventRecord
|
||||
import pw.binom.agentik.storage.events.EventStore
|
||||
import pw.binom.agentik.storage.events.EventType
|
||||
import kotlin.time.Instant
|
||||
import kotlinx.coroutines.sync.Mutex
|
||||
import kotlinx.coroutines.sync.withLock
|
||||
|
||||
/**
|
||||
* Thread-safe in-memory [EventStore]. Используется в тестах и в dev-режиме
|
||||
* `:standalone` (когда persistent events не нужны — например, для agentik-cli).
|
||||
*
|
||||
* Хранит все events в одном sorted-list (по createdAt). Поиск по afterId —
|
||||
* бинарный (list отсортирован). На 10K events работает за доли ms — для тестов
|
||||
* хватает. На production нужен Sqlite-impl.
|
||||
*
|
||||
* Thread-safety: один [Mutex] на все операции. Не concurrent-write-optimised —
|
||||
* для высоких нагрузок заменить на concurrent skip-list.
|
||||
*/
|
||||
class InMemoryEventStore : EventStore {
|
||||
|
||||
private val all: MutableList<EventRecord> = mutableListOf()
|
||||
private val byId: MutableMap<String, EventRecord> = mutableMapOf()
|
||||
private val mutex = Mutex()
|
||||
|
||||
override suspend fun append(record: EventRecord) {
|
||||
mutex.withLock {
|
||||
// Идемпотентность по id
|
||||
if (byId.containsKey(record.id)) return
|
||||
byId[record.id] = record
|
||||
// Insert maintaining ASC order by createdAt, ties broken by id.
|
||||
// binarySearch returns negative (-insertionPoint - 1) if not found,
|
||||
// or non-negative index if equal element found.
|
||||
val idx = all.binarySearch {
|
||||
val cmp = it.createdAt.compareTo(record.createdAt)
|
||||
if (cmp != 0) cmp else it.id.compareTo(record.id)
|
||||
}
|
||||
if (idx < 0) {
|
||||
all.add(-idx - 1, record)
|
||||
} else {
|
||||
// Found equal element — insert AFTER it to keep insertion order.
|
||||
all.add(idx + 1, record)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
override suspend fun query(
|
||||
conversationId: String?,
|
||||
afterId: String?,
|
||||
limit: Int,
|
||||
): List<EventRecord> {
|
||||
mutex.withLock {
|
||||
val startIdx = if (afterId == null) 0 else {
|
||||
val afterIdx = all.indexOfFirst { it.id == afterId }
|
||||
if (afterIdx < 0) return emptyList()
|
||||
afterIdx + 1
|
||||
}
|
||||
val filtered = if (conversationId == null) {
|
||||
all.subList(startIdx.coerceAtMost(all.size), all.size)
|
||||
} else {
|
||||
all.subList(startIdx.coerceAtMost(all.size), all.size)
|
||||
.filter { it.conversationId == conversationId }
|
||||
}
|
||||
return filtered.take(limit)
|
||||
}
|
||||
}
|
||||
|
||||
override suspend fun pruneOlderThan(olderThan: Instant): Int {
|
||||
mutex.withLock {
|
||||
val toRemove = all.filter { it.createdAt < olderThan }.map { it.id }
|
||||
if (toRemove.isEmpty()) return 0
|
||||
all.removeAll { it.id in toRemove }
|
||||
toRemove.forEach { byId.remove(it) }
|
||||
return toRemove.size
|
||||
}
|
||||
}
|
||||
|
||||
override suspend fun count(): Int = mutex.withLock { all.size }
|
||||
|
||||
override fun close() {
|
||||
// no-op: nothing to release
|
||||
}
|
||||
}
|
||||
+1
@@ -19,5 +19,6 @@ object InMemoryStorage {
|
||||
messageStore = InMemoryMessageStore(),
|
||||
workingMemoryStore = InMemoryWorkingMemoryStore(),
|
||||
reflectionStore = InMemoryReflectionStore(),
|
||||
eventStore = InMemoryEventStore(),
|
||||
)
|
||||
}
|
||||
|
||||
+165
@@ -0,0 +1,165 @@
|
||||
package pw.binom.agentik.storage.inmemory
|
||||
|
||||
import kotlinx.coroutines.coroutineScope
|
||||
import kotlinx.coroutines.test.runTest
|
||||
import pw.binom.agentik.storage.events.EventRecord
|
||||
import pw.binom.agentik.storage.events.EventType
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertNull
|
||||
import kotlin.test.assertTrue
|
||||
import kotlin.time.Clock
|
||||
import kotlin.time.Instant
|
||||
|
||||
class InMemoryEventStoreTest {
|
||||
|
||||
private fun rec(
|
||||
id: String,
|
||||
ts: Long,
|
||||
conv: String? = null,
|
||||
type: EventType = EventType.CONVERSATION_APPEND_TEXT,
|
||||
): EventRecord = EventRecord(
|
||||
id = id,
|
||||
conversationId = conv,
|
||||
createdAt = Instant.fromEpochMilliseconds(ts),
|
||||
type = type,
|
||||
payload = "{\"i\":$id}",
|
||||
)
|
||||
|
||||
@Test
|
||||
fun `append then query returns the record`() = runTest {
|
||||
val store = InMemoryEventStore()
|
||||
val r = rec("ev-1", ts = 1000)
|
||||
store.append(r)
|
||||
val result = store.query()
|
||||
assertEquals(listOf(r), result)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `query with afterId returns events strictly after the cursor`() = runTest {
|
||||
val store = InMemoryEventStore()
|
||||
store.append(rec("ev-1", ts = 1000))
|
||||
store.append(rec("ev-2", ts = 2000))
|
||||
store.append(rec("ev-3", ts = 3000))
|
||||
|
||||
// afterId = "ev-1" → only ev-2, ev-3 (exclusive)
|
||||
assertEquals(listOf("ev-2", "ev-3"), store.query(afterId = "ev-1").map { it.id })
|
||||
// afterId = "ev-2" → only ev-3
|
||||
assertEquals(listOf("ev-3"), store.query(afterId = "ev-2").map { it.id })
|
||||
// afterId = null → all
|
||||
assertEquals(listOf("ev-1", "ev-2", "ev-3"), store.query(afterId = null).map { it.id })
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `query with afterId pointing at unknown id returns empty`() = runTest {
|
||||
val store = InMemoryEventStore()
|
||||
store.append(rec("ev-1", ts = 1000))
|
||||
assertEquals(emptyList(), store.query(afterId = "ev-unknown"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `query with conversationId filters to that conversation only`() = runTest {
|
||||
val store = InMemoryEventStore()
|
||||
store.append(rec("ev-1", ts = 1000, conv = "c-1"))
|
||||
store.append(rec("ev-2", ts = 2000, conv = "c-2"))
|
||||
store.append(rec("ev-3", ts = 3000, conv = "c-1"))
|
||||
|
||||
val c1 = store.query(conversationId = "c-1")
|
||||
assertEquals(listOf("ev-1", "ev-3"), c1.map { it.id })
|
||||
|
||||
val c2 = store.query(conversationId = "c-2")
|
||||
assertEquals(listOf("ev-2"), c2.map { it.id })
|
||||
|
||||
val all = store.query(conversationId = null)
|
||||
assertEquals(listOf("ev-1", "ev-2", "ev-3"), all.map { it.id })
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `query with limit caps the result size`() = runTest {
|
||||
val store = InMemoryEventStore()
|
||||
repeat(10) { i -> store.append(rec("ev-$i", ts = (i * 1000).toLong())) }
|
||||
val first5 = store.query(limit = 5)
|
||||
assertEquals(5, first5.size)
|
||||
assertEquals(listOf("ev-0", "ev-1", "ev-2", "ev-3", "ev-4"), first5.map { it.id })
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `append is idempotent on id`() = runTest {
|
||||
val store = InMemoryEventStore()
|
||||
val r1 = rec("ev-1", ts = 1000, type = EventType.CONVERSATION_APPEND_TEXT)
|
||||
val r2 = rec("ev-1", ts = 1000, type = EventType.CONVERSATION_TOOL_CALL) // тот же id, разный тип
|
||||
store.append(r1)
|
||||
store.append(r2)
|
||||
// Second append no-op (idempotent by id) — первая запись побеждает
|
||||
val result = store.query()
|
||||
assertEquals(1, result.size)
|
||||
assertEquals(EventType.CONVERSATION_APPEND_TEXT, result[0].type)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `records returned in createdAt ascending order`() = runTest {
|
||||
val store = InMemoryEventStore()
|
||||
// Insert out of order
|
||||
store.append(rec("ev-b", ts = 2000))
|
||||
store.append(rec("ev-a", ts = 1000))
|
||||
store.append(rec("ev-c", ts = 3000))
|
||||
|
||||
val result = store.query()
|
||||
assertEquals(listOf("ev-a", "ev-b", "ev-c"), result.map { it.id })
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `records with same timestamp ordered by id ascending (stable sort)`() = runTest {
|
||||
val store = InMemoryEventStore()
|
||||
store.append(rec("ev-c", ts = 1000))
|
||||
store.append(rec("ev-a", ts = 1000))
|
||||
store.append(rec("ev-b", ts = 1000))
|
||||
|
||||
val result = store.query()
|
||||
// id lexicographic order: a < b < c
|
||||
assertEquals(listOf("ev-a", "ev-b", "ev-c"), result.map { it.id })
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `pruneOlderThan removes records before cutoff`() = runTest {
|
||||
val store = InMemoryEventStore()
|
||||
store.append(rec("ev-1", ts = 1000))
|
||||
store.append(rec("ev-2", ts = 2000))
|
||||
store.append(rec("ev-3", ts = 3000))
|
||||
|
||||
val removed = store.pruneOlderThan(Instant.fromEpochMilliseconds(2500))
|
||||
assertEquals(2, removed)
|
||||
assertEquals(listOf("ev-3"), store.query().map { it.id })
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `pruneOlderThan returns 0 when nothing to remove`() = runTest {
|
||||
val store = InMemoryEventStore()
|
||||
store.append(rec("ev-1", ts = 5000))
|
||||
val removed = store.pruneOlderThan(Instant.fromEpochMilliseconds(1000))
|
||||
assertEquals(0, removed)
|
||||
assertEquals(1, store.count())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `count returns number of stored records`() = runTest {
|
||||
val store = InMemoryEventStore()
|
||||
assertEquals(0, store.count())
|
||||
store.append(rec("ev-1", ts = 1000))
|
||||
store.append(rec("ev-2", ts = 2000))
|
||||
assertEquals(2, store.count())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `concurrent append from multiple coroutines all succeed`() = runTest {
|
||||
// Не strict — InMemoryEventStore использует Mutex, поэтому concurrent calls
|
||||
// сериализуются. Тест проверяет, что при последовательных append все ids
|
||||
// попадают в store (для concurrent test нужна отдельная TestScope — это
|
||||
// покрыто integration-тестами в :standalone).
|
||||
val store = InMemoryEventStore()
|
||||
repeat(51) { i ->
|
||||
store.append(rec("ev-$i", ts = i.toLong()))
|
||||
}
|
||||
assertEquals(51, store.count())
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,86 @@
|
||||
package pw.binom.agentik.storage.sqlite
|
||||
|
||||
import kotlin.time.Instant
|
||||
import mu.KotlinLogging
|
||||
import pw.binom.agentik.storage.events.EventRecord
|
||||
import pw.binom.agentik.storage.events.EventStore
|
||||
import pw.binom.agentik.storage.events.EventType
|
||||
import pw.binom.agentik.storage.sqlite.Agent_event as DbAgentEvent
|
||||
|
||||
private val log = KotlinLogging.logger {}
|
||||
|
||||
/**
|
||||
* SQLDelight-реализация [EventStore] поверх таблицы `agent_event`.
|
||||
*
|
||||
* Использует [EventStoreQueries] (генерируется SQLDelight из EventStore.sq).
|
||||
* Все запросы готовы — мы только маппим `Agent_event` (DB) ↔ `EventRecord` (domain).
|
||||
*
|
||||
* **Idempotency**: `append()` использует `INSERT OR IGNORE` — повторный append
|
||||
* с тем же id (network retry) — no-op. Это критично для producer'а, который
|
||||
* может retry при transient failure.
|
||||
*
|
||||
* **Pruning**: вызывай [pruneOlderThan] раз в час из background scope. Типичный
|
||||
* TTL = 24h. Если eventStore разрастётся (миллионы записей), индексы
|
||||
* (created_at, conversation_id+created_at) обеспечат O(log N) для query.
|
||||
*/
|
||||
class SqliteEventStore(
|
||||
private val db: AgentikDatabase,
|
||||
) : EventStore {
|
||||
|
||||
private val queries: EventStoreQueries get() = db.eventStoreQueries
|
||||
|
||||
override suspend fun append(record: EventRecord) {
|
||||
// payload — opaque JSON, хранится как UTF-8 bytes. Используем ByteArray,
|
||||
// потому что BLOB-колонка эффективнее TEXT для >100KB строк, и
|
||||
// API sqldelight нативно работает с ByteArray.
|
||||
val payloadBytes = record.payload.encodeToByteArray()
|
||||
queries.insert(
|
||||
id = record.id,
|
||||
conversation_id = record.conversationId,
|
||||
created_at = record.createdAt.toEpochMilliseconds(),
|
||||
type = record.type.name,
|
||||
payload = payloadBytes,
|
||||
)
|
||||
log.debug { "appended event id=${record.id} type=${record.type} conv=${record.conversationId}" }
|
||||
}
|
||||
|
||||
override suspend fun query(
|
||||
conversationId: String?,
|
||||
afterId: String?,
|
||||
limit: Int,
|
||||
): List<EventRecord> {
|
||||
// Если conversationId == null — используем queryAfter без фильтра
|
||||
// (он сам обрабатывает :convId IS NULL внутри SQL).
|
||||
// Если задан — queryAfterByConv (тогда SQL имеет WHERE conversation_id = :convId).
|
||||
val rows: List<DbAgentEvent> = if (conversationId == null) {
|
||||
queries.queryAfter(convId = null, afterId = afterId, limit = limit.toLong()).executeAsList()
|
||||
} else {
|
||||
queries.queryAfterByConv(convId = conversationId, afterId = afterId, limit = limit.toLong())
|
||||
.executeAsList()
|
||||
}
|
||||
return rows.map { it.toDomain() }
|
||||
}
|
||||
|
||||
override suspend fun pruneOlderThan(olderThan: Instant): Int {
|
||||
val deleted = queries.pruneOlderThan(olderThan.toEpochMilliseconds()).value
|
||||
if (deleted > 0) log.info { "pruned $deleted events older than $olderThan" }
|
||||
return deleted.toInt()
|
||||
}
|
||||
|
||||
override suspend fun count(): Int = queries.countAll().executeAsOne().toInt()
|
||||
|
||||
override fun close() {
|
||||
// no-op: lifecycle owned by AgentikDatabase / SqliteStores
|
||||
}
|
||||
|
||||
private fun DbAgentEvent.toDomain(): EventRecord = EventRecord(
|
||||
id = id,
|
||||
conversationId = conversation_id,
|
||||
createdAt = Instant.fromEpochMilliseconds(created_at),
|
||||
// type name → enum. Если в БД оказался неизвестный тип (новая версия,
|
||||
// unknown старому коду) — fallback на AGENT_CREATED (нейтральное значение).
|
||||
// Это безопаснее чем throw: клиент просто получит event с минимальным payload.
|
||||
type = runCatching { EventType.valueOf(type) }.getOrDefault(EventType.AGENT_CREATED),
|
||||
payload = payload.decodeToString(),
|
||||
)
|
||||
}
|
||||
@@ -7,12 +7,13 @@ import pw.binom.agentik.storage.ConversationStore
|
||||
import pw.binom.agentik.storage.MessageStore
|
||||
import pw.binom.agentik.storage.ReflectionStore
|
||||
import pw.binom.agentik.storage.StorageBundle
|
||||
import pw.binom.agentik.storage.events.EventStore
|
||||
import pw.binom.agentik.storage.sqlite.SqliteReflectionStore
|
||||
import pw.binom.agentik.storage.WorkingMemoryStore
|
||||
|
||||
/**
|
||||
* Корневой объект SQLite-слоя: держит [SqlDriver] и три [WorkingMemoryStore]/[MessageStore]/[ConversationStore].
|
||||
* Закрывается вместе с приложением.
|
||||
* Корневой объект SQLite-слоя: держит [SqlDriver] и пять store'ов (включая
|
||||
* [EventStore] — replay-after-disconnect).
|
||||
*/
|
||||
class SqliteStores private constructor(
|
||||
val driver: SqlDriver,
|
||||
@@ -20,6 +21,7 @@ class SqliteStores private constructor(
|
||||
val messages: MessageStore,
|
||||
val workingMemory: WorkingMemoryStore,
|
||||
val reflections: ReflectionStore,
|
||||
val events: EventStore,
|
||||
) : AutoCloseable {
|
||||
|
||||
/**
|
||||
@@ -32,9 +34,11 @@ class SqliteStores private constructor(
|
||||
messageStore = messages,
|
||||
workingMemoryStore = workingMemory,
|
||||
reflectionStore = reflections,
|
||||
eventStore = events,
|
||||
)
|
||||
|
||||
override fun close() {
|
||||
events.close()
|
||||
conversations.close()
|
||||
messages.close()
|
||||
workingMemory.close()
|
||||
@@ -55,6 +59,7 @@ class SqliteStores private constructor(
|
||||
messages = SqliteMessageStore(db),
|
||||
workingMemory = SqliteWorkingMemoryStore(db),
|
||||
reflections = SqliteReflectionStore(db),
|
||||
events = SqliteEventStore(db),
|
||||
)
|
||||
}
|
||||
|
||||
@@ -69,6 +74,7 @@ class SqliteStores private constructor(
|
||||
messages = SqliteMessageStore(db),
|
||||
workingMemory = SqliteWorkingMemoryStore(db),
|
||||
reflections = SqliteReflectionStore(db),
|
||||
events = SqliteEventStore(db),
|
||||
)
|
||||
}
|
||||
|
||||
@@ -118,6 +124,18 @@ class SqliteStores private constructor(
|
||||
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(),
|
||||
// v3: event log для replay-after-disconnect (см. EventStore.sq)
|
||||
"""
|
||||
CREATE TABLE IF NOT EXISTS agent_event (
|
||||
id TEXT NOT NULL PRIMARY KEY,
|
||||
conversation_id TEXT,
|
||||
created_at INTEGER NOT NULL,
|
||||
type TEXT NOT NULL,
|
||||
payload BLOB NOT NULL
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS idx_agent_event_conv_time ON agent_event(conversation_id, created_at);
|
||||
CREATE INDEX IF NOT EXISTS idx_agent_event_time ON agent_event(created_at);
|
||||
""".trimIndent(),
|
||||
)
|
||||
for (sql in migrations) {
|
||||
driver.execute(null, sql, 0)
|
||||
|
||||
@@ -0,0 +1,49 @@
|
||||
-- Event log для replay после disconnect (см. EventStore.kt в :storage-core).
|
||||
-- Каждая запись — один event из ChatAgent (AgentEvent) или ConversationLoop
|
||||
-- (Conversation.Event). payload — opaque JSON, сериализуется в :standalone перед append.
|
||||
--
|
||||
-- Cursor для pagination: query() принимает afterId, возвращает events строго
|
||||
-- после него (exclusive). Используется /events/replay endpoint в :server.
|
||||
|
||||
CREATE TABLE agent_event (
|
||||
id TEXT NOT NULL PRIMARY KEY,
|
||||
conversation_id TEXT, -- NULL для agent-level events (Created/Deleted/Renamed)
|
||||
created_at INTEGER NOT NULL, -- epoch millis (UTC)
|
||||
type TEXT NOT NULL, -- EventType.name (см. :storage-core/events/EventStore.kt)
|
||||
payload BLOB NOT NULL -- serialized JSON
|
||||
);
|
||||
|
||||
CREATE INDEX idx_agent_event_conv_time ON agent_event(conversation_id, created_at);
|
||||
CREATE INDEX idx_agent_event_time ON agent_event(created_at);
|
||||
|
||||
insert:
|
||||
INSERT OR IGNORE INTO agent_event (id, conversation_id, created_at, type, payload)
|
||||
VALUES (?, ?, ?, ?, ?);
|
||||
|
||||
queryById:
|
||||
SELECT * FROM agent_event WHERE id = ?;
|
||||
|
||||
queryAfter:
|
||||
-- Catchup по conversationId (или все если null). afterId exclusive.
|
||||
-- Сортировка: created_at ASC, id ASC (стабильный tie-break для events с одинаковым timestamp).
|
||||
SELECT * FROM agent_event
|
||||
WHERE (:convId IS NULL OR conversation_id = :convId)
|
||||
AND id > COALESCE(:afterId, '')
|
||||
ORDER BY created_at ASC, id ASC
|
||||
LIMIT :limit;
|
||||
|
||||
queryAfterByConv:
|
||||
SELECT * FROM agent_event
|
||||
WHERE conversation_id = :convId
|
||||
AND id > COALESCE(:afterId, '')
|
||||
ORDER BY created_at ASC, id ASC
|
||||
LIMIT :limit;
|
||||
|
||||
countAll:
|
||||
SELECT COUNT(*) FROM agent_event;
|
||||
|
||||
countByConv:
|
||||
SELECT COUNT(*) FROM agent_event WHERE conversation_id = :convId;
|
||||
|
||||
pruneOlderThan:
|
||||
DELETE FROM agent_event WHERE created_at < :cutoffEpochMillis;
|
||||
+152
@@ -0,0 +1,152 @@
|
||||
package pw.binom.agentik.storage.sqlite
|
||||
|
||||
import kotlinx.coroutines.test.runTest
|
||||
import pw.binom.agentik.storage.events.EventRecord
|
||||
import pw.binom.agentik.storage.events.EventType
|
||||
import kotlin.test.AfterTest
|
||||
import kotlin.test.BeforeTest
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertTrue
|
||||
import kotlin.time.Instant
|
||||
|
||||
class SqliteEventStoreTest {
|
||||
|
||||
private lateinit var driver: app.cash.sqldelight.driver.jdbc.sqlite.JdbcSqliteDriver
|
||||
private lateinit var db: AgentikDatabase
|
||||
private lateinit var store: SqliteEventStore
|
||||
|
||||
@BeforeTest
|
||||
fun setup() {
|
||||
driver = app.cash.sqldelight.driver.jdbc.sqlite.JdbcSqliteDriver(
|
||||
app.cash.sqldelight.driver.jdbc.sqlite.JdbcSqliteDriver.IN_MEMORY,
|
||||
)
|
||||
AgentikDatabase.Schema.create(driver)
|
||||
db = AgentikDatabase(driver)
|
||||
store = SqliteEventStore(db)
|
||||
}
|
||||
|
||||
@AfterTest
|
||||
fun tearDown() {
|
||||
driver.close()
|
||||
}
|
||||
|
||||
private fun rec(
|
||||
id: String,
|
||||
ts: Long,
|
||||
conv: String? = null,
|
||||
type: EventType = EventType.CONVERSATION_APPEND_TEXT,
|
||||
): EventRecord = EventRecord(
|
||||
id = id,
|
||||
conversationId = conv,
|
||||
createdAt = Instant.fromEpochMilliseconds(ts),
|
||||
type = type,
|
||||
payload = """{"i":"$id"}""",
|
||||
)
|
||||
|
||||
@Test
|
||||
fun `append then query returns the record`() = runTest {
|
||||
val r = rec("ev-1", ts = 1000)
|
||||
store.append(r)
|
||||
assertEquals(listOf(r), store.query())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `query with conversationId filters to that conversation only`() = runTest {
|
||||
store.append(rec("ev-1", ts = 1000, conv = "c-1"))
|
||||
store.append(rec("ev-2", ts = 2000, conv = "c-2"))
|
||||
store.append(rec("ev-3", ts = 3000, conv = "c-1"))
|
||||
|
||||
val c1 = store.query(conversationId = "c-1")
|
||||
assertEquals(listOf("ev-1", "ev-3"), c1.map { it.id })
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `query with afterId returns events strictly after cursor`() = runTest {
|
||||
store.append(rec("ev-1", ts = 1000))
|
||||
store.append(rec("ev-2", ts = 2000))
|
||||
store.append(rec("ev-3", ts = 3000))
|
||||
|
||||
assertEquals(listOf("ev-2", "ev-3"), store.query(afterId = "ev-1").map { it.id })
|
||||
assertEquals(listOf("ev-3"), store.query(afterId = "ev-2").map { it.id })
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `append is idempotent (INSERT OR IGNORE)`() = runTest {
|
||||
val r1 = rec("ev-1", ts = 1000, type = EventType.CONVERSATION_APPEND_TEXT)
|
||||
val r2 = rec("ev-1", ts = 1000, type = EventType.CONVERSATION_TOOL_CALL)
|
||||
store.append(r1)
|
||||
store.append(r2)
|
||||
// INSERT OR IGNORE — вторая попытка no-op
|
||||
val result = store.query()
|
||||
assertEquals(1, result.size)
|
||||
assertEquals(EventType.CONVERSATION_APPEND_TEXT, result[0].type)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `records ordered by createdAt then id ascending`() = runTest {
|
||||
store.append(rec("ev-c", ts = 2000))
|
||||
store.append(rec("ev-a", ts = 1000))
|
||||
store.append(rec("ev-d", ts = 1000)) // same ts as a, different id
|
||||
store.append(rec("ev-b", ts = 1500))
|
||||
|
||||
val result = store.query()
|
||||
assertEquals(listOf("ev-a", "ev-d", "ev-b", "ev-c"), result.map { it.id })
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `query with limit caps the result`() = runTest {
|
||||
repeat(10) { i -> store.append(rec("ev-$i", ts = (i * 100).toLong())) }
|
||||
val first3 = store.query(limit = 3)
|
||||
assertEquals(3, first3.size)
|
||||
assertEquals(listOf("ev-0", "ev-1", "ev-2"), first3.map { it.id })
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `pruneOlderThan removes records before cutoff`() = runTest {
|
||||
store.append(rec("ev-1", ts = 1000))
|
||||
store.append(rec("ev-2", ts = 2000))
|
||||
store.append(rec("ev-3", ts = 3000))
|
||||
|
||||
val removed = store.pruneOlderThan(Instant.fromEpochMilliseconds(2500))
|
||||
assertEquals(2, removed)
|
||||
assertEquals(listOf("ev-3"), store.query().map { it.id })
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `count returns total record count`() = runTest {
|
||||
assertEquals(0, store.count())
|
||||
store.append(rec("ev-1", ts = 1000))
|
||||
store.append(rec("ev-2", ts = 2000))
|
||||
assertEquals(2, store.count())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `unknown EventType in DB is loaded as AGENT_CREATED fallback`() = runTest {
|
||||
// Insert record with type that doesn't exist in current enum (simulating
|
||||
// a future enum value that older code doesn't know about).
|
||||
store.append(rec("ev-future", ts = 1000, type = EventType.CONVERSATION_END))
|
||||
// Manually mutate DB to use a fake type name (simulating forward-compat)
|
||||
db.eventStoreQueries.insert(
|
||||
id = "ev-fake",
|
||||
conversation_id = null,
|
||||
created_at = 2000,
|
||||
type = "FUTURE_TYPE_NOT_IN_ENUM",
|
||||
payload = "{}".encodeToByteArray(),
|
||||
)
|
||||
// Оба должны загрузиться (первый с правильным type, второй — с fallback)
|
||||
val result = store.query()
|
||||
assertEquals(2, result.size)
|
||||
assertEquals(EventType.CONVERSATION_END, result[0].type)
|
||||
assertEquals(EventType.AGENT_CREATED, result[1].type) // fallback для неизвестного type
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `close is no-op (does not close shared driver)`() = runTest {
|
||||
// store.close() НЕ должен закрывать driver — driver shared с другими store'ами.
|
||||
store.close()
|
||||
// Если бы close закрыл driver — следующий запрос упал бы. Проверяем что работает.
|
||||
store.append(rec("ev-1", ts = 1000))
|
||||
assertEquals(1, store.count())
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user