refactor(event-store): split into EventStore (read-only) + MutableEventStore
ci / JVM build + tests (push) Failing after 1m49s
ci / JVM build + tests (push) Failing after 1m49s
Разделяет интерфейс на read-only (EventStore) и write (MutableEventStore).
EventStore (read-only, для consumer'ов):
- events(after: Instant?): Flow<CommonEvent>
- earliestEventDate(): Instant
- close()
MutableEventStore : EventStore (для producer'ов):
- + append(event: CommonEvent)
- suspend, не идемпотентный, может быть silently evicted
Зачем:
- Consumer'ы (server SSE, admin dashboard, parent agents) принимают
EventStore — compile-time гарантия что они не могут писать в store.
- Producer'ы (ChatAgent, sub-agents, A2A-bridge) принимают MutableEventStore.
- Тесты могут использовать EventStore без опасности случайной модификации.
Миграция:
- :event-store пока без implementations, поэтому ничего не сломалось.
- Когда добавим InMemoryEventStore — он будет реализовывать оба
(MutableEventStore = EventStore + append). Подписки получают только
read-only projection через приведение типа.
Также: импорт обновлён AllEvent → CommonEvent (по rename в :proto).
This commit is contained in:
@@ -17,7 +17,7 @@ import kotlinx.coroutines.flow.flow
|
|||||||
import kotlinx.coroutines.runBlocking
|
import kotlinx.coroutines.runBlocking
|
||||||
import pw.binom.agentik.proto.Agent
|
import pw.binom.agentik.proto.Agent
|
||||||
import pw.binom.agentik.proto.AgentEvent
|
import pw.binom.agentik.proto.AgentEvent
|
||||||
import pw.binom.agentik.proto.AllEvent
|
import pw.binom.agentik.proto.CommonEvent
|
||||||
import pw.binom.agentik.proto.Conversation
|
import pw.binom.agentik.proto.Conversation
|
||||||
import kotlin.time.Instant
|
import kotlin.time.Instant
|
||||||
|
|
||||||
@@ -85,7 +85,7 @@ internal class AgentClient(
|
|||||||
* Подписка на ВСЕ события: agent lifecycle + все conversation events.
|
* Подписка на ВСЕ события: agent lifecycle + все conversation events.
|
||||||
* Использует SSE endpoint /events/all.
|
* Использует SSE endpoint /events/all.
|
||||||
*/
|
*/
|
||||||
override fun allEvents(after: Instant): Flow<AllEvent> = flow {
|
override fun allEvents(after: Instant): Flow<CommonEvent> = flow {
|
||||||
httpClient.prepareGet("$agentUrl/events/all?after=$after") { noSseReadTimeout() }
|
httpClient.prepareGet("$agentUrl/events/all?after=$after") { noSseReadTimeout() }
|
||||||
.execute { response ->
|
.execute { response ->
|
||||||
check(response.status == HttpStatusCode.OK) {
|
check(response.status == HttpStatusCode.OK) {
|
||||||
@@ -93,7 +93,7 @@ internal class AgentClient(
|
|||||||
}
|
}
|
||||||
readSse(response.bodyAsChannel())
|
readSse(response.bodyAsChannel())
|
||||||
.collect { payload ->
|
.collect { payload ->
|
||||||
emit(agentikJson.decodeFromString(AllEvent.serializer(), payload))
|
emit(agentikJson.decodeFromString(CommonEvent.serializer(), payload))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -2,7 +2,7 @@ package pw.binom.agentik.eventStore
|
|||||||
|
|
||||||
import kotlinx.coroutines.flow.Flow
|
import kotlinx.coroutines.flow.Flow
|
||||||
import kotlin.time.Instant
|
import kotlin.time.Instant
|
||||||
import pw.binom.agentik.proto.AllEvent
|
import pw.binom.agentik.proto.CommonEvent
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Bounded-tail event log с автоматическим управлением TTL.
|
* Bounded-tail event log с автоматическим управлением TTL.
|
||||||
@@ -30,11 +30,8 @@ import pw.binom.agentik.proto.AllEvent
|
|||||||
* - Позволяет impl выбирать retention strategy (TTL, size cap, sliding window).
|
* - Позволяет impl выбирать retention strategy (TTL, size cap, sliding window).
|
||||||
* - Сохраняет контракт clean: интерфейс только о put/get.
|
* - Сохраняет контракт clean: интерфейс только о put/get.
|
||||||
*
|
*
|
||||||
* **Append is NOT idempotent**: [AllEvent] не имеет уникального id,
|
* **Read-only**: этот интерфейс предоставляет только read-операции.
|
||||||
* поэтому retry с тем же logical event приведёт к дубликату в tail'е.
|
* Для записи см. [MutableEventStore].
|
||||||
* Для at-least-once → exactly-once нужна дедупликация на стороне
|
|
||||||
* consumer'а (catchup через `:message-store-api` audit log имеет
|
|
||||||
* монотонный `id` и работает как dedup anchor).
|
|
||||||
*
|
*
|
||||||
* **Подписки нереентрантные**: каждый вызов [events] создаёт **новую
|
* **Подписки нереентрантные**: каждый вызов [events] создаёт **новую
|
||||||
* подписку** (cold Flow). Один [events] НЕ видит события, добавленные до
|
* подписку** (cold Flow). Один [events] НЕ видит события, добавленные до
|
||||||
@@ -46,18 +43,6 @@ import pw.binom.agentik.proto.AllEvent
|
|||||||
*/
|
*/
|
||||||
interface EventStore : AutoCloseable {
|
interface EventStore : AutoCloseable {
|
||||||
|
|
||||||
/**
|
|
||||||
* Положить event в log.
|
|
||||||
*
|
|
||||||
* - **Идемпотентно по [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: AllEvent)
|
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Subscribe на events.
|
* Subscribe на events.
|
||||||
*
|
*
|
||||||
@@ -78,7 +63,7 @@ interface EventStore : AutoCloseable {
|
|||||||
* полноты клиент обязан cross-check с [earliestEventDate] и fallback
|
* полноты клиент обязан cross-check с [earliestEventDate] и fallback
|
||||||
* в message store при gap'е (см. KDoc интерфейса).
|
* в message store при gap'е (см. KDoc интерфейса).
|
||||||
*/
|
*/
|
||||||
fun events(after: Instant?): Flow<AllEvent>
|
fun events(after: Instant?): Flow<CommonEvent>
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Date **стартовой точки** буфера.
|
* Date **стартовой точки** буфера.
|
||||||
|
|||||||
@@ -0,0 +1,52 @@
|
|||||||
|
package pw.binom.agentik.eventStore
|
||||||
|
|
||||||
|
import pw.binom.agentik.proto.CommonEvent
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Mutable вариант [EventStore] — добавляет producer-операцию [append].
|
||||||
|
*
|
||||||
|
* Этот интерфейс предназначен **только для producer'ов** (ChatAgent,
|
||||||
|
* sub-agents, A2A-bridge). Consumer'ы (server SSE endpoints, admin
|
||||||
|
* dashboards, parent agents) должны принимать **read-only** [EventStore]
|
||||||
|
* — тогда невозможно случайно писать в store из observer'а.
|
||||||
|
*
|
||||||
|
* Типичное использование:
|
||||||
|
* ```
|
||||||
|
* // Producer
|
||||||
|
* class ChatAgent(private val events: MutableEventStore) {
|
||||||
|
* suspend fun doSomething() {
|
||||||
|
* events.append(CommonEvent.Agent(date = now, event = AgentEvent.Created(...)))
|
||||||
|
* }
|
||||||
|
* }
|
||||||
|
*
|
||||||
|
* // Consumer
|
||||||
|
* class EventStreamEndpoint(private val events: EventStore) {
|
||||||
|
* fun stream() = events.events(after = null)
|
||||||
|
* // Ошибка компиляции если раскомментировать:
|
||||||
|
* // events.append(...) // ← нельзя, MutableEventStore нет в типе
|
||||||
|
* }
|
||||||
|
* ```
|
||||||
|
*
|
||||||
|
* **Append НЕ идемпотентен**: [CommonEvent] не имеет уникального id,
|
||||||
|
* поэтому retry с тем же logical event (например, после network failure
|
||||||
|
* между producer и store) приведёт к дубликату в tail'е. Это OK для
|
||||||
|
* use case'a bounded-tail — клиент, делающий catchup через [events](after),
|
||||||
|
* получит свой диапазон ровно один раз при подключении, а последующие
|
||||||
|
* retry producer'а просто насытят tail повторами, не задевая уже
|
||||||
|
* обработанные. Для гарантированной exactly-once — dedup через
|
||||||
|
* [message-store] (там есть монотонный `id`).
|
||||||
|
*
|
||||||
|
* **Silently evicted**: implementation может выкинуть этот event сразу
|
||||||
|
* после append (TTL/cap) без уведомления producer'а. Producer **не
|
||||||
|
* должен** полагаться на то, что event дойдёт до клиента, если он
|
||||||
|
* вне retention window.
|
||||||
|
*/
|
||||||
|
interface MutableEventStore : EventStore {
|
||||||
|
/**
|
||||||
|
* Положить event в log.
|
||||||
|
*
|
||||||
|
* - **Не идемпотентно** — см. KDoc интерфейса.
|
||||||
|
* - **Suspend** для KMP I/O impl'ов (SQLite через JNI).
|
||||||
|
*/
|
||||||
|
suspend fun append(event: CommonEvent)
|
||||||
|
}
|
||||||
@@ -59,7 +59,7 @@ public interface Agent {
|
|||||||
*
|
*
|
||||||
* Cold (no replay). For catchup use EventStore.
|
* Cold (no replay). For catchup use EventStore.
|
||||||
*/
|
*/
|
||||||
fun allEvents(after: Instant): Flow<AllEvent>
|
fun allEvents(after: Instant): Flow<CommonEvent>
|
||||||
|
|
||||||
companion object {
|
companion object {
|
||||||
|
|
||||||
|
|||||||
+4
-4
@@ -10,7 +10,7 @@ import kotlinx.serialization.Serializable
|
|||||||
* Useful for admin dashboards, debug tools, parent agents: one subscription
|
* Useful for admin dashboards, debug tools, parent agents: one subscription
|
||||||
* instead of N+1. For regular UI use two separate SSE feeds
|
* instead of N+1. For regular UI use two separate SSE feeds
|
||||||
* ([AgentEvent] via /events and [Event] via /conversations/{id}/events);
|
* ([AgentEvent] via /events and [Event] via /conversations/{id}/events);
|
||||||
* [AllEvent] is for those who need everything in one place.
|
* [CommonEvent] is for those who need everything in one place.
|
||||||
*
|
*
|
||||||
* Server endpoint: GET /events/all (SSE), or replay via EventStore.
|
* Server endpoint: GET /events/all (SSE), or replay via EventStore.
|
||||||
*
|
*
|
||||||
@@ -19,7 +19,7 @@ import kotlinx.serialization.Serializable
|
|||||||
* only.
|
* only.
|
||||||
*/
|
*/
|
||||||
@Serializable
|
@Serializable
|
||||||
sealed interface AllEvent {
|
sealed interface CommonEvent {
|
||||||
val date: Instant
|
val date: Instant
|
||||||
|
|
||||||
@Serializable
|
@Serializable
|
||||||
@@ -27,7 +27,7 @@ sealed interface AllEvent {
|
|||||||
data class Agent(
|
data class Agent(
|
||||||
override val date: Instant,
|
override val date: Instant,
|
||||||
val event: AgentEvent,
|
val event: AgentEvent,
|
||||||
) : AllEvent
|
) : CommonEvent
|
||||||
|
|
||||||
@Serializable
|
@Serializable
|
||||||
@SerialName("conversation")
|
@SerialName("conversation")
|
||||||
@@ -35,5 +35,5 @@ sealed interface AllEvent {
|
|||||||
override val date: Instant,
|
override val date: Instant,
|
||||||
val conversationId: String,
|
val conversationId: String,
|
||||||
val event: Event,
|
val event: Event,
|
||||||
) : AllEvent
|
) : CommonEvent
|
||||||
}
|
}
|
||||||
@@ -3,7 +3,6 @@ package pw.binom.agentik.server
|
|||||||
import io.ktor.http.ContentType
|
import io.ktor.http.ContentType
|
||||||
import io.ktor.http.HttpStatusCode
|
import io.ktor.http.HttpStatusCode
|
||||||
import io.ktor.server.application.ApplicationCall
|
import io.ktor.server.application.ApplicationCall
|
||||||
import io.ktor.server.application.call
|
|
||||||
import io.ktor.server.request.receive
|
import io.ktor.server.request.receive
|
||||||
import io.ktor.server.request.receiveText
|
import io.ktor.server.request.receiveText
|
||||||
import io.ktor.server.response.respond
|
import io.ktor.server.response.respond
|
||||||
@@ -21,7 +20,7 @@ import kotlinx.serialization.KSerializer
|
|||||||
import kotlinx.serialization.json.Json
|
import kotlinx.serialization.json.Json
|
||||||
import pw.binom.agentik.proto.Agent
|
import pw.binom.agentik.proto.Agent
|
||||||
import pw.binom.agentik.proto.AgentEvent
|
import pw.binom.agentik.proto.AgentEvent
|
||||||
import pw.binom.agentik.proto.AllEvent
|
import pw.binom.agentik.proto.CommonEvent
|
||||||
import pw.binom.agentik.proto.Conversation
|
import pw.binom.agentik.proto.Conversation
|
||||||
import pw.binom.agentik.proto.Content
|
import pw.binom.agentik.proto.Content
|
||||||
import pw.binom.agentik.proto.Event
|
import pw.binom.agentik.proto.Event
|
||||||
@@ -144,7 +143,7 @@ internal fun Route.agentikRoutes(agent: Agent, eventStore: EventStore? = null) {
|
|||||||
*/
|
*/
|
||||||
get("/events/all") {
|
get("/events/all") {
|
||||||
val after = call.parseAfter() ?: return@get
|
val after = call.parseAfter() ?: return@get
|
||||||
call.streamJsonSse(agent.allEvents(after), AllEvent.serializer())
|
call.streamJsonSse(agent.allEvents(after), CommonEvent.serializer())
|
||||||
}
|
}
|
||||||
|
|
||||||
// ---- Replay-after-disconnect (event log) ----
|
// ---- Replay-after-disconnect (event log) ----
|
||||||
|
|||||||
@@ -5,17 +5,14 @@ import kotlinx.coroutines.flow.Flow
|
|||||||
import kotlinx.coroutines.flow.MutableSharedFlow
|
import kotlinx.coroutines.flow.MutableSharedFlow
|
||||||
import kotlinx.coroutines.flow.asSharedFlow
|
import kotlinx.coroutines.flow.asSharedFlow
|
||||||
import kotlinx.coroutines.flow.emitAll
|
import kotlinx.coroutines.flow.emitAll
|
||||||
import kotlinx.coroutines.flow.filterNotNull
|
|
||||||
import kotlinx.coroutines.flow.flow
|
import kotlinx.coroutines.flow.flow
|
||||||
import kotlinx.coroutines.flow.map
|
import kotlinx.coroutines.flow.map
|
||||||
import kotlinx.coroutines.flow.merge
|
|
||||||
import kotlinx.coroutines.launch
|
import kotlinx.coroutines.launch
|
||||||
import kotlinx.coroutines.runBlocking
|
import kotlinx.coroutines.runBlocking
|
||||||
import kotlinx.coroutines.sync.Mutex
|
import kotlinx.coroutines.sync.Mutex
|
||||||
import kotlinx.coroutines.sync.withLock
|
import kotlinx.coroutines.sync.withLock
|
||||||
import kotlinx.serialization.json.Json
|
import kotlinx.serialization.json.Json
|
||||||
import pw.binom.agentik.llm.tools.ContextCompactor
|
import pw.binom.agentik.llm.tools.ContextCompactor
|
||||||
import pw.binom.agentik.llm.tools.LlmMemoryReviewer
|
|
||||||
import pw.binom.agentik.llm.tools.LlmReflector
|
import pw.binom.agentik.llm.tools.LlmReflector
|
||||||
import pw.binom.agentik.llm.tools.SkillMiner
|
import pw.binom.agentik.llm.tools.SkillMiner
|
||||||
import pw.binom.agentik.memory.MemoryPrefetcher
|
import pw.binom.agentik.memory.MemoryPrefetcher
|
||||||
@@ -23,7 +20,7 @@ import pw.binom.agentik.memory.MemoryReviewer
|
|||||||
import pw.binom.agentik.memory.MemorySystemGuidance
|
import pw.binom.agentik.memory.MemorySystemGuidance
|
||||||
import pw.binom.agentik.proto.Agent as ProtoAgent
|
import pw.binom.agentik.proto.Agent as ProtoAgent
|
||||||
import pw.binom.agentik.proto.AgentEvent
|
import pw.binom.agentik.proto.AgentEvent
|
||||||
import pw.binom.agentik.proto.AllEvent
|
import pw.binom.agentik.proto.CommonEvent
|
||||||
import pw.binom.agentik.proto.Conversation as ProtoConversation
|
import pw.binom.agentik.proto.Conversation as ProtoConversation
|
||||||
import pw.binom.agentik.proto.Event as ProtoEvent
|
import pw.binom.agentik.proto.Event as ProtoEvent
|
||||||
import pw.binom.agentik.skills.SkillCatalog
|
import pw.binom.agentik.skills.SkillCatalog
|
||||||
@@ -33,8 +30,6 @@ import pw.binom.agentik.standalone.llm.LlmConfig
|
|||||||
import pw.binom.agentik.messageStore.ConversationRecord
|
import pw.binom.agentik.messageStore.ConversationRecord
|
||||||
import pw.binom.agentik.messageStore.Ids
|
import pw.binom.agentik.messageStore.Ids
|
||||||
import pw.binom.agentik.messageStore.Reflection
|
import pw.binom.agentik.messageStore.Reflection
|
||||||
import pw.binom.agentik.storageBundle.StorageBundle
|
|
||||||
import pw.binom.agentik.workingMemory.WorkingMemoryEntry
|
|
||||||
import pw.binom.agentik.messageStore.events.EventRecord
|
import pw.binom.agentik.messageStore.events.EventRecord
|
||||||
import pw.binom.agentik.messageStore.events.EventType
|
import pw.binom.agentik.messageStore.events.EventType
|
||||||
import pw.binom.agentik.toolsets.DisableToolsetTool
|
import pw.binom.agentik.toolsets.DisableToolsetTool
|
||||||
@@ -275,10 +270,10 @@ class ChatAgent(
|
|||||||
* (admin-дашборд часами) надо добавить reactive re-subscribe (см.
|
* (admin-дашборд часами) надо добавить reactive re-subscribe (см.
|
||||||
* notes 04-sub-agents.md, Variant B).
|
* notes 04-sub-agents.md, Variant B).
|
||||||
*/
|
*/
|
||||||
override fun allEvents(after: Instant): Flow<AllEvent> = flow {
|
override fun allEvents(after: Instant): Flow<CommonEvent> = flow {
|
||||||
// Agent lifecycle events
|
// Agent lifecycle events
|
||||||
emitAll(agentEvents.asSharedFlow().map { e: AgentEvent ->
|
emitAll(agentEvents.asSharedFlow().map { e: AgentEvent ->
|
||||||
AllEvent.Agent(date = e.date, event = e)
|
CommonEvent.Agent(date = e.date, event = e)
|
||||||
})
|
})
|
||||||
// Snapshot живых диалогов на момент подписки.
|
// Snapshot живых диалогов на момент подписки.
|
||||||
// ВАЖНО: не подписываемся на новые Created — это ответственность caller'а
|
// ВАЖНО: не подписываемся на новые Created — это ответственность caller'а
|
||||||
@@ -289,7 +284,7 @@ class ChatAgent(
|
|||||||
.forEach { conv: ProtoConversation ->
|
.forEach { conv: ProtoConversation ->
|
||||||
val cid: String = conv.id
|
val cid: String = conv.id
|
||||||
emitAll(conv.events(after).map { ev: ProtoEvent ->
|
emitAll(conv.events(after).map { ev: ProtoEvent ->
|
||||||
AllEvent.Conversation(date = ev.date, conversationId = cid, event = ev)
|
CommonEvent.Conversation(date = ev.date, conversationId = cid, event = ev)
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user