From 7e66baf9e30e197a4ee102d880b283589d4cb557 Mon Sep 17 00:00:00 2001 From: subochev Date: Sun, 20 Sep 2026 15:55:14 +0300 Subject: [PATCH] feat(event-store): add conversationEvents and agentEvents filters MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Расширяет EventStore тремя вариантами подписки (вместо одного events()): - events(after) → Flow весь поток (микс Agent + Conversation) - conversationEvents(after, conversationId?) → Flow опциональный фильтр по conversationId (null = все диалоги) - agentEvents(after) → Flow только lifecycle (Created/Deleted/Renamed) Типизированные subtype'ы вместо Flow + .filterIsInstance: - compile-time safety на клиенте (нет cast'ов в CommonEvent.Conversation) - persistent impl'ы могут делать WHERE conversation_id = ? на уровне БД Соответствует HTTP-маршрутам в :server: - GET /events/all ↔ events(after) - GET /conversations/{id}/events ↔ conversationEvents(after, id) - GET /events (только agent lifecycle) ↔ agentEvents(after) README обновлён — таблица caller→method показывает маппинг. Совместимость: signals не сломаны (добавление, не изменение). --- event-store/README.md | 27 ++++++++++++++----- .../pw/binom/agentik/eventStore/EventStore.kt | 27 +++++++++++++++++++ 2 files changed, 47 insertions(+), 7 deletions(-) diff --git a/event-store/README.md b/event-store/README.md index 18f78ee..7bae567 100644 --- a/event-store/README.md +++ b/event-store/README.md @@ -65,11 +65,24 @@ eventStore.events(after = client.lastSeen).collect { apply(it) } ```kotlin interface EventStore : AutoCloseable { fun events(after: Instant?): Flow + fun conversationEvents( + after: Instant?, + conversationId: String? = null, // null = все диалоги + ): Flow + fun agentEvents(after: Instant?): Flow suspend fun earliestEventDate(): Instant // non-null: now() для пустого буфера override fun close() } ``` +**Семантика фильтров**: +- `conversationEvents(null)` — все диалоги. +- `conversationEvents("c-123")` — один конкретный диалог. +- `agentEvents(...)` — только lifecycle (Created/Deleted/Renamed). + +Все три возвращают **типизированные** subtype'ы [CommonEvent], так что +caller'у не нужно `.filterIsInstance` на клиентской стороне. + ### `MutableEventStore : EventStore` (для producer'ов) ```kotlin @@ -88,14 +101,14 @@ append по TTL/cap. Producer не должен полагаться на то, ## Когда использовать какой интерфейс -| Caller | Interface | +| Caller | Method | |---|---| -| ChatAgent (producer) | `MutableEventStore` | -| Sub-agents (producer) | `MutableEventStore` | -| Server SSE endpoint | `EventStore` | -| Admin dashboard | `EventStore` | -| Parent orchestrator | `EventStore` | -| Тесты | `EventStore` (read-only) | +| Server `/events/all` SSE (mixed) | `events(after)` | +| Server `/conversations/{id}/events` SSE | `conversationEvents(after, conversationId)` | +| Server `/agent/events` SSE (lifecycle only) | `agentEvents(after)` | +| Admin dashboard (lifecycle) | `agentEvents(after)` | +| Parent orchestrator (mixed) | `events(after)` | +| Тесты | `events(after)` + projection через фильтр | ## Как добавить новый implementation 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 a07db41..3e1cbc0 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 @@ -65,6 +65,33 @@ interface EventStore : AutoCloseable { */ fun events(after: Instant?): Flow + /** + * Subscribe на **только conversation events** (т.е. [CommonEvent.Conversation]). + * + * - [conversationId] == null → события **всех** диалогов. + * - [conversationId] != null → события **только этого** диалога. + * + * Семантика `after` идентична [events] (catchup + live). + * Возвращаемый тип — конкретный subtype [CommonEvent.Conversation], + * так что caller'у не нужно `.filterIsInstance<...>()` на стороне клиента. + * + * Для реализации: самый дешёвый путь — `events(after).filterIsInstance<...>()`, + * но persistent impl'ы могут делать `WHERE conversation_id = ?` на уровне БД. + */ + fun conversationEvents(after: Instant?, conversationId: String? = null): Flow + + /** + * Subscribe на **только agent events** ([CommonEvent.Agent] — + * создание/удаление/переименование диалога). + * + * Семантика `after` идентична [events] (catchup + live). + * Возвращаемый тип — конкретный subtype [CommonEvent.Agent]. + * + * Полезно для admin-дашборда, который хочет видеть только lifecycle + * диалогов без деталей ходов. + */ + fun agentEvents(after: Instant?): Flow + /** * Date **стартовой точки** буфера. *