From e2f0e434d1381588d0c9d511a87e6afae0b64e7d Mon Sep 17 00:00:00 2001 From: subochev Date: Sun, 20 Sep 2026 15:46:59 +0300 Subject: [PATCH] refactor(event-store): split into EventStore (read-only) + MutableEventStore MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Разделяет интерфейс на read-only (EventStore) и write (MutableEventStore). EventStore (read-only, для consumer'ов): - events(after: Instant?): Flow - 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). --- .../pw/binom/agentik/client/AgentClient.kt | 6 +-- .../pw/binom/agentik/eventStore/EventStore.kt | 23 ++------ .../agentik/eventStore/MutableEventStore.kt | 52 +++++++++++++++++++ .../kotlin/pw/binom/agentik/proto/Agent.kt | 2 +- .../proto/{AllEvent.kt => CommonEvent.kt} | 8 +-- .../kotlin/pw/binom/agentik/server/Routes.kt | 5 +- .../agentik/standalone/agent/ChatAgent.kt | 13 ++--- 7 files changed, 70 insertions(+), 39 deletions(-) create mode 100644 event-store/src/commonMain/kotlin/pw/binom/agentik/eventStore/MutableEventStore.kt rename proto/src/commonMain/kotlin/pw/binom/agentik/proto/{AllEvent.kt => CommonEvent.kt} (88%) diff --git a/client/src/commonMain/kotlin/pw/binom/agentik/client/AgentClient.kt b/client/src/commonMain/kotlin/pw/binom/agentik/client/AgentClient.kt index 41ffd1c..f62b4db 100644 --- a/client/src/commonMain/kotlin/pw/binom/agentik/client/AgentClient.kt +++ b/client/src/commonMain/kotlin/pw/binom/agentik/client/AgentClient.kt @@ -17,7 +17,7 @@ import kotlinx.coroutines.flow.flow import kotlinx.coroutines.runBlocking import pw.binom.agentik.proto.Agent 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 kotlin.time.Instant @@ -85,7 +85,7 @@ internal class AgentClient( * Подписка на ВСЕ события: agent lifecycle + все conversation events. * Использует SSE endpoint /events/all. */ - override fun allEvents(after: Instant): Flow = flow { + override fun allEvents(after: Instant): Flow = flow { httpClient.prepareGet("$agentUrl/events/all?after=$after") { noSseReadTimeout() } .execute { response -> check(response.status == HttpStatusCode.OK) { @@ -93,7 +93,7 @@ internal class AgentClient( } readSse(response.bodyAsChannel()) .collect { payload -> - emit(agentikJson.decodeFromString(AllEvent.serializer(), payload)) + emit(agentikJson.decodeFromString(CommonEvent.serializer(), payload)) } } } diff --git a/event-store/src/commonMain/kotlin/pw/binom/agentik/eventStore/EventStore.kt b/event-store/src/commonMain/kotlin/pw/binom/agentik/eventStore/EventStore.kt index 7b3f873..a07db41 100644 --- a/event-store/src/commonMain/kotlin/pw/binom/agentik/eventStore/EventStore.kt +++ b/event-store/src/commonMain/kotlin/pw/binom/agentik/eventStore/EventStore.kt @@ -2,7 +2,7 @@ package pw.binom.agentik.eventStore import kotlinx.coroutines.flow.Flow import kotlin.time.Instant -import pw.binom.agentik.proto.AllEvent +import pw.binom.agentik.proto.CommonEvent /** * Bounded-tail event log с автоматическим управлением TTL. @@ -30,11 +30,8 @@ import pw.binom.agentik.proto.AllEvent * - Позволяет impl выбирать retention strategy (TTL, size cap, sliding window). * - Сохраняет контракт clean: интерфейс только о put/get. * - * **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). + * **Read-only**: этот интерфейс предоставляет только read-операции. + * Для записи см. [MutableEventStore]. * * **Подписки нереентрантные**: каждый вызов [events] создаёт **новую * подписку** (cold Flow). Один [events] НЕ видит события, добавленные до @@ -46,18 +43,6 @@ import pw.binom.agentik.proto.AllEvent */ 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. * @@ -78,7 +63,7 @@ interface EventStore : AutoCloseable { * полноты клиент обязан cross-check с [earliestEventDate] и fallback * в message store при gap'е (см. KDoc интерфейса). */ - fun events(after: Instant?): Flow + fun events(after: Instant?): Flow /** * Date **стартовой точки** буфера. diff --git a/event-store/src/commonMain/kotlin/pw/binom/agentik/eventStore/MutableEventStore.kt b/event-store/src/commonMain/kotlin/pw/binom/agentik/eventStore/MutableEventStore.kt new file mode 100644 index 0000000..0856320 --- /dev/null +++ b/event-store/src/commonMain/kotlin/pw/binom/agentik/eventStore/MutableEventStore.kt @@ -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) +} diff --git a/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Agent.kt b/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Agent.kt index bcdf196..c6edc11 100644 --- a/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Agent.kt +++ b/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Agent.kt @@ -59,7 +59,7 @@ public interface Agent { * * Cold (no replay). For catchup use EventStore. */ - fun allEvents(after: Instant): Flow + fun allEvents(after: Instant): Flow companion object { diff --git a/proto/src/commonMain/kotlin/pw/binom/agentik/proto/AllEvent.kt b/proto/src/commonMain/kotlin/pw/binom/agentik/proto/CommonEvent.kt similarity index 88% rename from proto/src/commonMain/kotlin/pw/binom/agentik/proto/AllEvent.kt rename to proto/src/commonMain/kotlin/pw/binom/agentik/proto/CommonEvent.kt index 4ce16a5..fc95e8c 100644 --- a/proto/src/commonMain/kotlin/pw/binom/agentik/proto/AllEvent.kt +++ b/proto/src/commonMain/kotlin/pw/binom/agentik/proto/CommonEvent.kt @@ -10,7 +10,7 @@ import kotlinx.serialization.Serializable * Useful for admin dashboards, debug tools, parent agents: one subscription * instead of N+1. For regular UI use two separate SSE feeds * ([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. * @@ -19,7 +19,7 @@ import kotlinx.serialization.Serializable * only. */ @Serializable -sealed interface AllEvent { +sealed interface CommonEvent { val date: Instant @Serializable @@ -27,7 +27,7 @@ sealed interface AllEvent { data class Agent( override val date: Instant, val event: AgentEvent, - ) : AllEvent + ) : CommonEvent @Serializable @SerialName("conversation") @@ -35,5 +35,5 @@ sealed interface AllEvent { override val date: Instant, val conversationId: String, val event: Event, - ) : AllEvent + ) : CommonEvent } diff --git a/server/src/commonMain/kotlin/pw/binom/agentik/server/Routes.kt b/server/src/commonMain/kotlin/pw/binom/agentik/server/Routes.kt index 6f901b9..3da8183 100644 --- a/server/src/commonMain/kotlin/pw/binom/agentik/server/Routes.kt +++ b/server/src/commonMain/kotlin/pw/binom/agentik/server/Routes.kt @@ -3,7 +3,6 @@ package pw.binom.agentik.server import io.ktor.http.ContentType import io.ktor.http.HttpStatusCode import io.ktor.server.application.ApplicationCall -import io.ktor.server.application.call import io.ktor.server.request.receive import io.ktor.server.request.receiveText import io.ktor.server.response.respond @@ -21,7 +20,7 @@ import kotlinx.serialization.KSerializer import kotlinx.serialization.json.Json import pw.binom.agentik.proto.Agent 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.Content import pw.binom.agentik.proto.Event @@ -144,7 +143,7 @@ internal fun Route.agentikRoutes(agent: Agent, eventStore: EventStore? = null) { */ get("/events/all") { 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) ---- diff --git a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ChatAgent.kt b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ChatAgent.kt index 9b320c4..050898c 100644 --- a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ChatAgent.kt +++ b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ChatAgent.kt @@ -5,17 +5,14 @@ import kotlinx.coroutines.flow.Flow import kotlinx.coroutines.flow.MutableSharedFlow import kotlinx.coroutines.flow.asSharedFlow import kotlinx.coroutines.flow.emitAll -import kotlinx.coroutines.flow.filterNotNull import kotlinx.coroutines.flow.flow import kotlinx.coroutines.flow.map -import kotlinx.coroutines.flow.merge import kotlinx.coroutines.launch import kotlinx.coroutines.runBlocking import kotlinx.coroutines.sync.Mutex import kotlinx.coroutines.sync.withLock import kotlinx.serialization.json.Json 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.SkillMiner 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.proto.Agent as ProtoAgent 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.Event as ProtoEvent 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.Ids 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.EventType import pw.binom.agentik.toolsets.DisableToolsetTool @@ -275,10 +270,10 @@ class ChatAgent( * (admin-дашборд часами) надо добавить reactive re-subscribe (см. * notes 04-sub-agents.md, Variant B). */ - override fun allEvents(after: Instant): Flow = flow { + override fun allEvents(after: Instant): Flow = flow { // Agent lifecycle events emitAll(agentEvents.asSharedFlow().map { e: AgentEvent -> - AllEvent.Agent(date = e.date, event = e) + CommonEvent.Agent(date = e.date, event = e) }) // Snapshot живых диалогов на момент подписки. // ВАЖНО: не подписываемся на новые Created — это ответственность caller'а @@ -289,7 +284,7 @@ class ChatAgent( .forEach { conv: ProtoConversation -> val cid: String = conv.id 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) }) } }