:event-store — bounded-tail event log (KMP)
Что это
Двухуровневое хранилище событий агента. Этот модуль — короткий
bounded tail для live-SSE и недавнего replay. Полный audit log
живёт в :message-store-api (никогда не эвиктится, source of truth).
Три принципа:
- Tail управляет TTL сам. Никаких
prune/cleanupметодов наружу — implementation решает, когда выкинуть старый event. Caller'ы не могут забыть cleanup. - Catchup + live в одном Flow.
events(after)сначала отдаёт буферизованный диапазон, потом переключается на live tail — клиент не должен знать, где у него "разрыв". - Read-only контракт для consumer'ов. Запись через [MutableEventStore], чтение через [EventStore]. Compile-time гарантия что observer не сможет писать в store.
Где используется
:standaloneChatAgent — append черезMutableOutboxStore(заменяет текущийagentEvents: MutableSharedFlow+persistAgentEvent).:serverRoutes.kt —/events/allSSE endpoint читает черезEventStore.events(after).- Будущий
:android-agentcore — same интерфейс для локального bounded tail без dedicated server connection.
Архитектура
┌─ :event-store (этот модуль) ────────────────────────┐
│ Bounded tail с auto-TTL: │
│ • append(event) ← producer │
│ • events(after): Flow ← consumer │
│ • earliestEventDate() для gap detection │
│ TTL/cap eviction — внутри impl │
└───────────────────────────────────────────────────┘
▲ gap detected
│
┌─ :message-store-api (полный audit log) ───────────┐
│ MessageStore: query(after, before, limit) │
│ Никогда не эвиктится. Source of truth. │
└───────────────────────────────────────────────────┘
Reconnect pattern (caller делает):
val earliest = eventStore.earliestEventDate()
if (client.lastSeen < earliest) {
// gap: догоняем через :message-store-api
val gap = messageStore.query(after = client.lastSeen, before = earliest)
applyAll(gap)
client.lastSeen = gap.last().createdAt
}
eventStore.events(after = client.lastSeen).collect { apply(it) }
API
OutboxStore (read-only, для consumer'ов)
interface EventStore : AutoCloseable {
fun events(after: Instant?): Flow<CommonEvent>
fun conversationEvents(
after: Instant?,
conversationId: String? = null, // null = все диалоги
): Flow<CommonEvent.Conversation>
fun agentEvents(after: Instant?): Flow<CommonEvent.Agent>
suspend fun earliestEventDate(): Instant // non-null: now() для пустого буфера
override fun close()
}
Семантика фильтров:
conversationEvents(null)— все диалоги.conversationEvents("c-123")— один конкретный диалог.agentEvents(...)— только lifecycle (Created/Deleted/Renamed).
Все три возвращают типизированные subtype'ы [CommonEvent], так что
caller'у не нужно .filterIsInstance на клиентской стороне.
MutableEventStore : EventStore (для producer'ов)
interface MutableEventStore : EventStore {
suspend fun append(event: CommonEvent)
}
Append НЕ идемпотентен: [CommonEvent] не имеет уникального id,
retry даст дубликат. Для exactly-once — dedup через
:message-store-api (там есть монотонный id).
Silently evicted: implementation может выкинуть event сразу после append по TTL/cap. Producer не должен полагаться на то, что event дойдёт до клиента, если он вне retention window.
Когда использовать какой интерфейс
| Caller | Method |
|---|---|
Server /events/all SSE (mixed) |
events(after) |
Server /conversations/{id}/events SSE |
conversationEvents(after, conversationId) |
Server /agent/events SSE (lifecycle only) |
agentEvents(after) |
| Admin dashboard (lifecycle) | agentEvents(after) |
| Parent orchestrator (mixed) | events(after) |
| Тесты | events(after) + projection через фильтр |
Как добавить новый implementation
- Создать класс с конструктором и lifecycle (
close()обязан освободить ресурсы). - Реализовать минимум: append (с TTL eviction), events (Flow с catchup + live), earliestEventDate (non-null Instant, now() если буфер пуст).
- Для persistent impl: SQL/ksqlite таблица с индексом по date,
вставка =
INSERT OR IGNOREдля дедупликации на уровне БД (если в схеме будет id).
Текущее состояние
- ✅ Interface дизайн (
OutboxStore+MutableOutboxStore) - ✅ KMP build (jvm + linuxX64 + mingwX64)
- ⏳ Нет implementations (next:
InMemoryEventStoreдля тестов) - ⏳ Не интегрирован в
:standalone/:server
Зависимости
:proto(api) — типCommonEvent(3 AgentEvent + 9 Conversation.Event вариантов).kotlinx-coroutines-core(api) —Flow.
Никаких kotlinx-serialization, kotlin-logging, platform-specific
зависимостей — этот модуль намеренно minimal.