feat(events): EventStore + AllEvent unified stream + replay endpoints
EventStore (persistent event log) и AllEvent (sealed wrapper для
третьего типа подписки — ВСЕ events в одном потоке). Touches 7 modules.
Архитектура:
Producer (ChatAgent + ConversationEvents) → EventStore + SharedFlow
↓ ↓
Live SSE (cold, no replay) Replay endpoints (cursor-based)
(1) :storage-core — EventStore interface
- append(record): idempotent по record.id (INSERT OR IGNORE)
- query(conversationId?, afterId?, limit): пагинированный catchup
- pruneOlderThan(instant): TTL cleanup
- count(): maintenance метрика
- @Serializable EventRecord(id, conversationId?, createdAt, type, payload)
- enum EventType: AGENT_*/CONVERSATION_* (forward-compat fallback)
- StorageBundle дополнен eventStore: EventStore? = null (backward-compat)
(2) :storage-inmemory — InMemoryEventStore
- Thread-safe (Mutex), binarySearch для упорядоченной вставки
- Записи сортируются по createdAt ASC, ties по id ASC (стабильно)
- Idempotency по id (повторный append no-op)
(3) :storage-sqlite — SqliteEventStore
- sqldelight schema: agent_event (id PK, conversation_id?, created_at,
type, payload BLOB) + 2 индекса (conversation_id+created_at,
created_at)
- Миграция v3: CREATE TABLE IF NOT EXISTS (additive)
- 5 запросов: insert, queryGlobal, queryByConv, pruneOlderThan, count
- Forward-compat: неизвестный EventType в БД → fallback AGENT_CREATED
(чтобы старые клиенты не падали на новых enum values)
- Добавлен в SqliteStores (open/inMemory + asBundle())
(4) :standalone — Producer wiring
- ChatAgent.persistAgentEvent() — fire-and-forget append при каждом
AgentEvent (Created/Deleted/Renamed)
- ConversationEvents — персистит в EventStore при каждом tryEmit/emit
(концертный случай от connect disconnect)
- ChatAgent.allEvents() — merge agent-events + snapshot всех живых
диалогов в единый Flow<AllEvent>
(5) :proto — AllEvent sealed interface
- AllEvent.Agent(date, event: AgentEvent)
- AllEvent.Conversation(date, conversationId, event: Event)
- Agent.allEvents(after): Flow<AllEvent> — третий тип подписки
(в дополнение к events() и Conversation.events)
(6) :server — Endpoints
- GET /events/all — SSE поток AllEvent (cold)
- GET /events/replay?after_id=&limit= — пагинированный catchup
(503 если EventStore не сконфигурирован)
- GET /conversations/{id}/events/replay?after_id=&limit= — то же per-conv
- Module.kt принимает eventStore: EventStore? параметром
(7) :client — Client API
- AgentClient.allEvents(after) — подписка на /events/all SSE
- AgentClient.replayAllEvents(afterId, limit) — catchup /events/replay
- AgentClient.replayConversationEvents(convId, afterId, limit)
- EventRecordDto — wire-зеркало EventRecord (клиент не зависит
от :storage-core, определяет DTO локально; формат совместим с
серверным JSON)
Тесты: 22 новых теста (12 InMemory + 10 Sqlite), все зелёные.
Все три слоя синхронизированы: proto contract + standalone impl +
server endpoint + client API.
This commit is contained in:
@@ -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()
|
||||
|
||||
Reference in New Issue
Block a user