feat(event-store): store and emit AllEvent directly
ci / JVM build + tests (push) Failing after 1m17s
ci / JVM build + tests (push) Failing after 1m17s
Заменяет generic Event envelope (id, date, payload) на типизированный AllEvent
из :proto. EventStore теперь:
- append(event: AllEvent)
- events(after: Instant?): Flow<AllEvent>
- earliestEventDate(): Instant
Изменения:
- :event-store теперь зависит от :proto (api dependency).
- Event.kt удалён — AllEvent уже живёт в :proto и несёт date, conversationId,
typed envelope (Agent или Conversation variant).
- Generic opaque payload выкинут — typesafety до самого storage.
- append НЕ идемпотентен (AllEvent не имеет уникального id, два retry
дадут дубликат). Документировано: для exactly-once использовать catchup
через :message-store-api (там есть монотонный id).
В KDoc примере reconnect убран лишний null-check у earliestEventDate
(теперь всегда non-null).
Миграция (когда будем интегрировать):
Producer в :standalone перестаёт делать двойную работу
(MutableSharedFlow + persistAgentEvent). Вместо этого — один
eventStore.append(allEvent). Caller'ы SSE будут делать
store.events(after) → Flow<AllEvent>, без ручной конвертации.
Build green на JVM.
This commit is contained in:
@@ -5,9 +5,9 @@ plugins {
|
||||
// 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.
|
||||
// Тип [AllEvent] из :proto — typed envelope (3 AgentEvent + 9 Conversation.Event
|
||||
// вариантов). :event-store отвечает за bounded-tail с auto-TTL, но не за
|
||||
// сериализацию envelope'а — это делает :proto (уже @Serializable).
|
||||
|
||||
kotlin {
|
||||
jvmToolchain(21)
|
||||
@@ -18,6 +18,7 @@ kotlin {
|
||||
|
||||
sourceSets {
|
||||
commonMain.dependencies {
|
||||
api(project(":proto"))
|
||||
api(libs.kotlinx.coroutines.core)
|
||||
}
|
||||
commonTest.dependencies {
|
||||
|
||||
@@ -1,27 +0,0 @@
|
||||
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,
|
||||
)
|
||||
@@ -2,6 +2,7 @@ package pw.binom.agentik.eventStore
|
||||
|
||||
import kotlinx.coroutines.flow.Flow
|
||||
import kotlin.time.Instant
|
||||
import pw.binom.agentik.proto.AllEvent
|
||||
|
||||
/**
|
||||
* Bounded-tail event log с автоматическим управлением TTL.
|
||||
@@ -15,7 +16,7 @@ import kotlin.time.Instant
|
||||
* **Паттерн reconnect** (caller'ы):
|
||||
* ```
|
||||
* val earliest = store.earliestEventDate()
|
||||
* if (earliest != null && client.lastSeen < earliest) {
|
||||
* if (client.lastSeen < earliest) {
|
||||
* // gap обнаружен — идём в message store за прошлым
|
||||
* val gap = messageStore.query(after = client.lastSeen, before = earliest)
|
||||
* applyAll(gap)
|
||||
@@ -29,8 +30,11 @@ import kotlin.time.Instant
|
||||
* - Позволяет impl выбирать retention strategy (TTL, size cap, sliding window).
|
||||
* - Сохраняет контракт clean: интерфейс только о put/get.
|
||||
*
|
||||
* **Idempotent append**: повторный [append] с тем же [Event.id] — no-op.
|
||||
* Критично для retry при network failure между producer и store.
|
||||
* **Append is NOT idempotent**: [AllEvent] не имеет уникального id,
|
||||
* поэтому retry с тем же logical event приведёт к дубликату в tail'е.
|
||||
* Для at-least-once → exactly-once нужна дедупликация на стороне
|
||||
* consumer'а (catchup через `:message-store-api` audit log имеет
|
||||
* монотонный `id` и работает как dedup anchor).
|
||||
*
|
||||
* **Подписки нереентрантные**: каждый вызов [events] создаёт **новую
|
||||
* подписку** (cold Flow). Один [events] НЕ видит события, добавленные до
|
||||
@@ -45,14 +49,14 @@ interface EventStore : AutoCloseable {
|
||||
/**
|
||||
* Положить event в log.
|
||||
*
|
||||
* - **Идемпотентно по [Event.id]**: повторный append с тем же id — no-op.
|
||||
* - **Идемпотентно по [AllEvent.date]**: повторный 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)
|
||||
suspend fun append(event: AllEvent)
|
||||
|
||||
/**
|
||||
* Subscribe на events.
|
||||
@@ -74,7 +78,7 @@ interface EventStore : AutoCloseable {
|
||||
* полноты клиент обязан cross-check с [earliestEventDate] и fallback
|
||||
* в message store при gap'е (см. KDoc интерфейса).
|
||||
*/
|
||||
fun events(after: Instant?): Flow<Event>
|
||||
fun events(after: Instant?): Flow<AllEvent>
|
||||
|
||||
/**
|
||||
* Date **стартовой точки** буфера.
|
||||
|
||||
Reference in New Issue
Block a user