feat(event-store): new :event-store module with interfaces only
ci / JVM build + tests (push) Failing after 1m54s
ci / JVM build + tests (push) Failing after 1m54s
Выделяет bounded-tail event log в отдельный KMP-модуль. Это **новый
контракт** (не замена :message-store-api/events/EventStore — тот пока жив).
**Двухуровневое хранилище событий**:
1. :event-store (этот PR) — короткий bounded tail, авто-TTL.
Для live SSE и recent replay (catchup после короткого disconnect).
2. :message-store-api (MessageStore) — полный audit log, никогда не
эвиктится. Source of truth для длинного disconnect / audit query.
**API**:
- append(Event) — put, идемпотентный по id
- events(after: Instant?): Flow<Event> — catchup + live в одном Flow
- earliestEventDate(): Instant? — для gap detection у клиента
**Чего НЕТ в API** (by design):
- delete/prune методов — TTL/cap eviction полностью на стороне impl.
Caller'ы не должны забыть вызвать cleanup (single source of truth).
- conversationId / type в Event — opaque payload, тип envelope'а
решает producer.
Пока interfaces only — implementations (InMemoryEventStore, persistent)
появятся в следующих коммитах. Никаких изменений в существующем
EventStore в :message-store-api, чтобы не ломать зависимости.
KMP targets: jvm + linuxX64 + mingwX64 (Apple auto-disabled на Linux).
Modules:
+ :event-store — новый, commonMain only, ~150 строк
This commit is contained in:
@@ -0,0 +1,28 @@
|
|||||||
|
plugins {
|
||||||
|
alias(libs.plugins.kotlin.multiplatform)
|
||||||
|
}
|
||||||
|
|
||||||
|
// KMP-интерфейс bounded-tail event log'а. Implementation-specific TTL/cap
|
||||||
|
// eviction — caller's responsibility НЕ вызывать cleanup() (метод не существует).
|
||||||
|
//
|
||||||
|
// Зависимости минимальные — только kotlinx-coroutines для Flow. Никаких
|
||||||
|
// :message-store-api, :proto, kotlinx-serialization — тип [Event] намеренно
|
||||||
|
// generic (payload: String), JSON-сериализацию делает producer.
|
||||||
|
|
||||||
|
kotlin {
|
||||||
|
jvmToolchain(21)
|
||||||
|
|
||||||
|
jvm()
|
||||||
|
linuxX64()
|
||||||
|
mingwX64()
|
||||||
|
|
||||||
|
sourceSets {
|
||||||
|
commonMain.dependencies {
|
||||||
|
api(libs.kotlinx.coroutines.core)
|
||||||
|
}
|
||||||
|
commonTest.dependencies {
|
||||||
|
implementation(kotlin("test"))
|
||||||
|
implementation(libs.kotlinx.coroutines.test)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,27 @@
|
|||||||
|
package pw.binom.agentik.eventStore
|
||||||
|
|
||||||
|
import kotlin.time.Instant
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Элемент bounded-tail event log'а.
|
||||||
|
*
|
||||||
|
* **Дизайн — generic envelope**: `payload` это opaque JSON-строка.
|
||||||
|
* Типизация envelope'а (что внутри) — на стороне producer'а.
|
||||||
|
* Это позволяет `:event-store` оставаться generic и не зависеть от
|
||||||
|
* `:proto` или каких-либо доменных типов.
|
||||||
|
*
|
||||||
|
* @property id монотонный по времени prefix (`ev-<uuid>`), используется для:
|
||||||
|
* - **Idempotency**: повторный [EventStore.append] с тем же `id` — no-op.
|
||||||
|
* - **Cursor-based pagination**: `id` последнего seen event'а.
|
||||||
|
*
|
||||||
|
* @property date момент эмиссии в UTC. Используется для:
|
||||||
|
* - Ordering в [EventStore.events] catchup replay.
|
||||||
|
* - [EventStore.earliestEventDate] для gap detection у клиента.
|
||||||
|
*
|
||||||
|
* @property payload opaque JSON-строка, описывающая само событие.
|
||||||
|
*/
|
||||||
|
data class Event(
|
||||||
|
val id: String,
|
||||||
|
val date: Instant,
|
||||||
|
val payload: String,
|
||||||
|
)
|
||||||
@@ -0,0 +1,98 @@
|
|||||||
|
package pw.binom.agentik.eventStore
|
||||||
|
|
||||||
|
import kotlinx.coroutines.flow.Flow
|
||||||
|
import kotlin.time.Instant
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Bounded-tail event log с автоматическим управлением TTL.
|
||||||
|
*
|
||||||
|
* **Архитектура двухуровневого хранилища событий**:
|
||||||
|
* 1. **Этот store** = короткий bounded tail (live SSE + недавний replay).
|
||||||
|
* События автоматически эвиктятся по TTL/cap (implementation-defined).
|
||||||
|
* 2. **Message store (`:message-store-api`)** = полный audit log, никогда не
|
||||||
|
* эвиктится. Source of truth для всего прошлого.
|
||||||
|
*
|
||||||
|
* **Паттерн reconnect** (caller'ы):
|
||||||
|
* ```
|
||||||
|
* val earliest = store.earliestEventDate()
|
||||||
|
* if (earliest != null && client.lastSeen < earliest) {
|
||||||
|
* // gap обнаружен — идём в message store за прошлым
|
||||||
|
* val gap = messageStore.query(after = client.lastSeen, before = earliest)
|
||||||
|
* applyAll(gap)
|
||||||
|
* }
|
||||||
|
* store.events(after = client.lastSeen).collect { apply(it) }
|
||||||
|
* ```
|
||||||
|
*
|
||||||
|
* **Нет delete/cleanup методов** — TTL/cap eviction полностью на стороне
|
||||||
|
* implementation. Это:
|
||||||
|
* - Убирает single source of truth дублирование (caller не может забыть cleanup).
|
||||||
|
* - Позволяет impl выбирать retention strategy (TTL, size cap, sliding window).
|
||||||
|
* - Сохраняет контракт clean: интерфейс только о put/get.
|
||||||
|
*
|
||||||
|
* **Idempotent append**: повторный [append] с тем же [Event.id] — no-op.
|
||||||
|
* Критично для retry при network failure между producer и store.
|
||||||
|
*
|
||||||
|
* **Подписки нереентрантные**: каждый вызов [events] создаёт **новую
|
||||||
|
* подписку** (cold Flow). Один [events] НЕ видит события, добавленные до
|
||||||
|
* его вызова, если [after] == null. Если нужен catchup — передавайте
|
||||||
|
* `after = lastSeenDate` явно.
|
||||||
|
*
|
||||||
|
* **Multi-consumer**: разные [events] подписки видят одно и то же live
|
||||||
|
* tail. Каждая подписка — независимая projection.
|
||||||
|
*/
|
||||||
|
interface EventStore : AutoCloseable {
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Положить event в log.
|
||||||
|
*
|
||||||
|
* - **Идемпотентно по [Event.id]**: повторный append с тем же id — no-op.
|
||||||
|
* - **Silently evicted**: implementation может выкинуть этот event сразу
|
||||||
|
* после append (TTL/cap) без уведомления producer'а. Producer **не
|
||||||
|
* должен** полагаться на то, что event дойдёт до клиента, если он
|
||||||
|
* вне retention window.
|
||||||
|
* - **Suspend**: для KMP I/O impl'ов (SQLite через JNI).
|
||||||
|
*/
|
||||||
|
suspend fun append(event: Event)
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Subscribe на events.
|
||||||
|
*
|
||||||
|
* **`after == null`** → только **live** (события с момента вызова
|
||||||
|
* `events()`). Каждое новое событие от любого producer'а немедленно
|
||||||
|
* появится в Flow. Буфер replay не отдаётся.
|
||||||
|
*
|
||||||
|
* **`after != null`** → сначала **catchup**: эмитт все буферизованные
|
||||||
|
* события с `date > after`, порядок `date ASC` (ties по `id ASC`).
|
||||||
|
* Затем **live** (как null-case).
|
||||||
|
*
|
||||||
|
* Cold Flow: каждый вызов — новая подписка. Вызов **после** append'а
|
||||||
|
* не увидит этот конкретный event (если `after == null`); для catchup
|
||||||
|
* передавайте явный `after`.
|
||||||
|
*
|
||||||
|
* ВАЖНО: `Flow` НЕ бросает ошибку при потере сети между producer и
|
||||||
|
* store — такие события просто не дойдут до этого Flow. Для гарантии
|
||||||
|
* полноты клиент обязан cross-check с [earliestEventDate] и fallback
|
||||||
|
* в message store при gap'е (см. KDoc интерфейса).
|
||||||
|
*/
|
||||||
|
fun events(after: Instant?): Flow<Event>
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Date самого старого event'а, всё ещё хранящегося в буфере.
|
||||||
|
*
|
||||||
|
* `null` если буфер пуст (или store только что стартовал — ни одного
|
||||||
|
* события ещё не было).
|
||||||
|
*
|
||||||
|
* **Используется клиентом для gap detection**:
|
||||||
|
* - `lastSeen < earliest` → есть дыра в покрытии, нужен fallback
|
||||||
|
* в message store за диапазоном `[lastSeen, earliest)`.
|
||||||
|
* - `lastSeen >= earliest` → всё доступно через [events](after),
|
||||||
|
* fallback не нужен.
|
||||||
|
* - `earliest == null` → store пуст, первый live event сам станет
|
||||||
|
* `earliest` для следующего клиента.
|
||||||
|
*
|
||||||
|
* Suspend потому что в persistent impl'ах требует SQL query (`MIN(date)`).
|
||||||
|
*/
|
||||||
|
suspend fun earliestEventDate(): Instant?
|
||||||
|
|
||||||
|
override fun close()
|
||||||
|
}
|
||||||
@@ -67,6 +67,11 @@ include(":memory-vector")
|
|||||||
// IO-зависимостей. Используется в тестах (быстрый setup, без JDBC) и будет
|
// IO-зависимостей. Используется в тестах (быстрый setup, без JDBC) и будет
|
||||||
// использоваться в Android-сборке (JVector/SQLite не подходят для ART out-of-box).
|
// использоваться в Android-сборке (JVector/SQLite не подходят для ART out-of-box).
|
||||||
include(":message-store-api")
|
include(":message-store-api")
|
||||||
|
// Bounded-tail event log с auto-TTL. Двухуровневое хранилище: этот модуль —
|
||||||
|
// короткий live tail + recent replay; полный audit log живёт в :message-store-api
|
||||||
|
// (там — MessageStore + ConversationStore). EventStore сам управляет eviction,
|
||||||
|
// никаких prune-методов наружу. KMP, без implementations пока.
|
||||||
|
include(":event-store")
|
||||||
include(":working-memory-api")
|
include(":working-memory-api")
|
||||||
include(":storage-bundle")
|
include(":storage-bundle")
|
||||||
include(":storage-inmemory")
|
include(":storage-inmemory")
|
||||||
|
|||||||
Reference in New Issue
Block a user