docs(event-store): README explaining module purpose and contract
ci / JVM build + tests (push) Failing after 1m16s
ci / JVM build + tests (push) Failing after 1m16s
Документирует:
- Три принципа дизайна (TTL внутри, catchup+live в одном Flow,
read-only контракт для observer'ов)
- Архитектуру двухуровневого хранилища событий со схемой
- Reconnect pattern с gap detection
- API EventStore + MutableEventStore (когда какой использовать)
- Таблица: какой caller принимает какой интерфейс
- Текущее состояние: interfaces готовы, implementations в работе
- Зависимости (минимальные: :proto + kotlinx-coroutines)
В том же стиле что и :agent-toolsets/README.md.
This commit is contained in:
@@ -0,0 +1,124 @@
|
|||||||
|
# `: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 через `MutableEventStore` (заменяет
|
||||||
|
текущий `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 делает):
|
||||||
|
|
||||||
|
```kotlin
|
||||||
|
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
|
||||||
|
|
||||||
|
### `EventStore` (read-only, для consumer'ов)
|
||||||
|
|
||||||
|
```kotlin
|
||||||
|
interface EventStore : AutoCloseable {
|
||||||
|
fun events(after: Instant?): Flow<CommonEvent>
|
||||||
|
suspend fun earliestEventDate(): Instant // non-null: now() для пустого буфера
|
||||||
|
override fun close()
|
||||||
|
}
|
||||||
|
```
|
||||||
|
|
||||||
|
### `MutableEventStore : EventStore` (для producer'ов)
|
||||||
|
|
||||||
|
```kotlin
|
||||||
|
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 | Interface |
|
||||||
|
|---|---|
|
||||||
|
| ChatAgent (producer) | `MutableEventStore` |
|
||||||
|
| Sub-agents (producer) | `MutableEventStore` |
|
||||||
|
| Server SSE endpoint | `EventStore` |
|
||||||
|
| Admin dashboard | `EventStore` |
|
||||||
|
| Parent orchestrator | `EventStore` |
|
||||||
|
| Тесты | `EventStore` (read-only) |
|
||||||
|
|
||||||
|
## Как добавить новый 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 дизайн (`EventStore` + `MutableEventStore`)
|
||||||
|
- ✅ 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.
|
||||||
Reference in New Issue
Block a user