feat(event-store): add conversationEvents and agentEvents filters
ci / JVM build + tests (push) Has been cancelled
ci / JVM build + tests (push) Has been cancelled
Расширяет EventStore тремя вариантами подписки (вместо одного events()):
- events(after) → Flow<CommonEvent>
весь поток (микс Agent + Conversation)
- conversationEvents(after, conversationId?)
→ Flow<CommonEvent.Conversation>
опциональный фильтр по conversationId (null = все диалоги)
- agentEvents(after)
→ Flow<CommonEvent.Agent>
только lifecycle (Created/Deleted/Renamed)
Типизированные subtype'ы вместо Flow<CommonEvent> + .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 не сломаны (добавление, не изменение).
This commit is contained in:
+20
-7
@@ -65,11 +65,24 @@ eventStore.events(after = client.lastSeen).collect { apply(it) }
|
|||||||
```kotlin
|
```kotlin
|
||||||
interface EventStore : AutoCloseable {
|
interface EventStore : AutoCloseable {
|
||||||
fun events(after: Instant?): Flow<CommonEvent>
|
fun events(after: Instant?): Flow<CommonEvent>
|
||||||
|
fun conversationEvents(
|
||||||
|
after: Instant?,
|
||||||
|
conversationId: String? = null, // null = все диалоги
|
||||||
|
): Flow<CommonEvent.Conversation>
|
||||||
|
fun agentEvents(after: Instant?): Flow<CommonEvent.Agent>
|
||||||
suspend fun earliestEventDate(): Instant // non-null: now() для пустого буфера
|
suspend fun earliestEventDate(): Instant // non-null: now() для пустого буфера
|
||||||
override fun close()
|
override fun close()
|
||||||
}
|
}
|
||||||
```
|
```
|
||||||
|
|
||||||
|
**Семантика фильтров**:
|
||||||
|
- `conversationEvents(null)` — все диалоги.
|
||||||
|
- `conversationEvents("c-123")` — один конкретный диалог.
|
||||||
|
- `agentEvents(...)` — только lifecycle (Created/Deleted/Renamed).
|
||||||
|
|
||||||
|
Все три возвращают **типизированные** subtype'ы [CommonEvent], так что
|
||||||
|
caller'у не нужно `.filterIsInstance` на клиентской стороне.
|
||||||
|
|
||||||
### `MutableEventStore : EventStore` (для producer'ов)
|
### `MutableEventStore : EventStore` (для producer'ов)
|
||||||
|
|
||||||
```kotlin
|
```kotlin
|
||||||
@@ -88,14 +101,14 @@ append по TTL/cap. Producer не должен полагаться на то,
|
|||||||
|
|
||||||
## Когда использовать какой интерфейс
|
## Когда использовать какой интерфейс
|
||||||
|
|
||||||
| Caller | Interface |
|
| Caller | Method |
|
||||||
|---|---|
|
|---|---|
|
||||||
| ChatAgent (producer) | `MutableEventStore` |
|
| Server `/events/all` SSE (mixed) | `events(after)` |
|
||||||
| Sub-agents (producer) | `MutableEventStore` |
|
| Server `/conversations/{id}/events` SSE | `conversationEvents(after, conversationId)` |
|
||||||
| Server SSE endpoint | `EventStore` |
|
| Server `/agent/events` SSE (lifecycle only) | `agentEvents(after)` |
|
||||||
| Admin dashboard | `EventStore` |
|
| Admin dashboard (lifecycle) | `agentEvents(after)` |
|
||||||
| Parent orchestrator | `EventStore` |
|
| Parent orchestrator (mixed) | `events(after)` |
|
||||||
| Тесты | `EventStore` (read-only) |
|
| Тесты | `events(after)` + projection через фильтр |
|
||||||
|
|
||||||
## Как добавить новый implementation
|
## Как добавить новый implementation
|
||||||
|
|
||||||
|
|||||||
@@ -65,6 +65,33 @@ interface EventStore : AutoCloseable {
|
|||||||
*/
|
*/
|
||||||
fun events(after: Instant?): Flow<CommonEvent>
|
fun events(after: Instant?): Flow<CommonEvent>
|
||||||
|
|
||||||
|
/**
|
||||||
|
* 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<CommonEvent.Conversation>
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Subscribe на **только agent events** ([CommonEvent.Agent] —
|
||||||
|
* создание/удаление/переименование диалога).
|
||||||
|
*
|
||||||
|
* Семантика `after` идентична [events] (catchup + live).
|
||||||
|
* Возвращаемый тип — конкретный subtype [CommonEvent.Agent].
|
||||||
|
*
|
||||||
|
* Полезно для admin-дашборда, который хочет видеть только lifecycle
|
||||||
|
* диалогов без деталей ходов.
|
||||||
|
*/
|
||||||
|
fun agentEvents(after: Instant?): Flow<CommonEvent.Agent>
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Date **стартовой точки** буфера.
|
* Date **стартовой точки** буфера.
|
||||||
*
|
*
|
||||||
|
|||||||
Reference in New Issue
Block a user