Files
agentik/event-store
subochev aef5083801
ci / JVM build + tests (push) Failing after 1m9s
refactor: migrate EventStore and MessageStore to :outbox-api and :journal-api
- Replaced usages of `:message-store-api` and `:working-memory-api` with `:journal-api`, `:outbox-api`, and `:context-api`.
- Deprecated legacy `EventStore` and `MessageStore` interfaces, added `typealias` for backward compatibility.
- Updated imports across all modules with references to `:journal-api` and `:outbox-api`.
- Introduced `journalRoutes` and `outboxRoutes` in `:server` for audit log and live event stream endpoints.
- Adjusted `Agent` to expose read-only `journal` and `outbox` stores for improved modularity and clarity.
- Removed legacy Event and AgentEvent definitions from `:proto`, migrated to `:outbox-api`.
- Storage-related modules have been updated to support the new APIs consistently.
2026-09-21 02:38:40 +03:00
..

:event-store — bounded-tail event log (KMP)

Что это

Двухуровневое хранилище событий агента. Этот модуль — короткий bounded tail для live-SSE и недавнего replay. Полный audit log живёт в :message-store-api (никогда не эвиктится, source of truth).

Три принципа:

  1. Tail управляет TTL сам. Никаких prune/cleanup методов наружу — implementation решает, когда выкинуть старый event. Caller'ы не могут забыть cleanup.
  2. Catchup + live в одном Flow. events(after) сначала отдаёт буферизованный диапазон, потом переключается на live tail — клиент не должен знать, где у него "разрыв".
  3. Read-only контракт для consumer'ов. Запись через [MutableEventStore], чтение через [EventStore]. Compile-time гарантия что observer не сможет писать в store.

Где используется

  • :standalone ChatAgent — append через MutableOutboxStore (заменяет текущий agentEvents: MutableSharedFlow + persistAgentEvent).
  • :server Routes.kt — /events/all SSE endpoint читает через EventStore.events(after).
  • Будущий :android-agent core — 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

  1. Создать класс с конструктором и lifecycle (close() обязан освободить ресурсы).
  2. Реализовать минимум: append (с TTL eviction), events (Flow с catchup + live), earliestEventDate (non-null Instant, now() если буфер пуст).
  3. Для 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.