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 3e1cbc0..8a54970 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 @@ -1,6 +1,8 @@ package pw.binom.agentik.eventStore import kotlinx.coroutines.flow.Flow +import kotlinx.coroutines.flow.filter +import kotlinx.coroutines.flow.filterIsInstance import kotlin.time.Instant import pw.binom.agentik.proto.CommonEvent @@ -72,13 +74,23 @@ interface EventStore : AutoCloseable { * - [conversationId] != null → события **только этого** диалога. * * Семантика `after` идентична [events] (catchup + live). - * Возвращаемый тип — конкретный subtype [CommonEvent.Conversation], - * так что caller'у не нужно `.filterIsInstance<...>()` на стороне клиента. - * - * Для реализации: самый дешёвый путь — `events(after).filterIsInstance<...>()`, - * но persistent impl'ы могут делать `WHERE conversation_id = ?` на уровне БД. + * Возвращаемый тип — конкретный subtype [CommonEvent.Conversation]. */ - fun conversationEvents(after: Instant?, conversationId: String? = null): Flow + /** + * **Default implementation** (читает все events + фильтрует). + * + * Простая реализация через [events] + filterIsInstance. Реализации + * могут override'нуть для эффективности (например, добавить SQL + * `WHERE conversation_id = ?` чтобы не тянуть всё в память), но + * контракт корректен и без override. + */ + fun conversationEvents(after: Instant?, conversationId: String? = null): Flow = + events(after) + .filterIsInstance() + .let { filtered -> + if (conversationId == null) filtered + else filtered.filter { it.conversationId == conversationId } + } /** * Subscribe на **только agent events** ([CommonEvent.Agent] — @@ -90,7 +102,15 @@ interface EventStore : AutoCloseable { * Полезно для admin-дашборда, который хочет видеть только lifecycle * диалогов без деталей ходов. */ - fun agentEvents(after: Instant?): Flow + /** + * **Default implementation** (читает все events + фильтрует по типу). + * + * Простая реализация через [events] + filterIsInstance. Реализации + * могут override'нуть для эффективности (например, читать только agent + * row'ы из БД), но контракт корректен и без override. + */ + fun agentEvents(after: Instant?): Flow = + events(after).filterIsInstance() /** * Date **стартовой точки** буфера.