refactor(event-store): default impls for conversationEvents/agentEvents
ci / JVM build + tests (push) Failing after 1m24s
ci / JVM build + tests (push) Failing after 1m24s
Методы conversationEvents() и agentEvents() теперь default в интерфейсе:
реализуют фильтрацию через [events] + filterIsInstance.
Зачем:
- Минимальный контракт для impl — достаточно реализовать только events().
- InMemoryEventStore и любой новый backend получают работающие
specialized views автоматически, без копипасты filterIsInstance.
- Persistent impl'ы (SQL/ksqlite) могут override'нуть для эффективности
(WHERE conversation_id = ? — не тянуть все events в память), но
контракт корректен и без override.
Семантика идентична:
- conversationEvents(after, null) = events().filterIsInstance<Conversation>()
- conversationEvents(after, "c-1") = events()...filter { it.conversationId == "c-1" }
- agentEvents(after) = events().filterIsInstance<Agent>()
Imports добавлены: kotlinx.coroutines.flow.filter, filterIsInstance.
This commit is contained in:
@@ -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<CommonEvent.Conversation>
|
||||
/**
|
||||
* **Default implementation** (читает все events + фильтрует).
|
||||
*
|
||||
* Простая реализация через [events] + filterIsInstance. Реализации
|
||||
* могут override'нуть для эффективности (например, добавить SQL
|
||||
* `WHERE conversation_id = ?` чтобы не тянуть всё в память), но
|
||||
* контракт корректен и без override.
|
||||
*/
|
||||
fun conversationEvents(after: Instant?, conversationId: String? = null): Flow<CommonEvent.Conversation> =
|
||||
events(after)
|
||||
.filterIsInstance<CommonEvent.Conversation>()
|
||||
.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<CommonEvent.Agent>
|
||||
/**
|
||||
* **Default implementation** (читает все events + фильтрует по типу).
|
||||
*
|
||||
* Простая реализация через [events] + filterIsInstance. Реализации
|
||||
* могут override'нуть для эффективности (например, читать только agent
|
||||
* row'ы из БД), но контракт корректен и без override.
|
||||
*/
|
||||
fun agentEvents(after: Instant?): Flow<CommonEvent.Agent> =
|
||||
events(after).filterIsInstance<CommonEvent.Agent>()
|
||||
|
||||
/**
|
||||
* Date **стартовой точки** буфера.
|
||||
|
||||
Reference in New Issue
Block a user