diff --git a/agentik-cli/src/commonMain/kotlin/pw/binom/agentik/cli/commands/SendSubcommand.kt b/agentik-cli/src/commonMain/kotlin/pw/binom/agentik/cli/commands/SendSubcommand.kt index 5ad10fb..909f094 100644 --- a/agentik-cli/src/commonMain/kotlin/pw/binom/agentik/cli/commands/SendSubcommand.kt +++ b/agentik-cli/src/commonMain/kotlin/pw/binom/agentik/cli/commands/SendSubcommand.kt @@ -3,6 +3,7 @@ package pw.binom.agentik.cli.commands import kotlinx.cli.ArgType import kotlinx.cli.vararg import kotlinx.coroutines.delay +import kotlinx.coroutines.flow.map import kotlinx.coroutines.flow.onEach import kotlinx.coroutines.flow.takeWhile import kotlinx.coroutines.launch @@ -10,6 +11,7 @@ import pw.binom.agentik.cli.AgentikSubcommand import pw.binom.agentik.cli.defaultCliHttpClient import pw.binom.agentik.client.AgentikAgent import pw.binom.agentik.outbox.Event +import pw.binom.agentik.outbox.OnlineEvent import pw.binom.agentik.proto.Content import kotlin.time.Instant @@ -24,8 +26,8 @@ class SendSubcommand : AgentikSubcommand("send", "Отправить user-ход return@runBlocking } try { - // Подписываемся на поток событий ДО send: события, отправленные - // до подписки, не реплеятся (shared-flow без replay). + // Durable-поток (End/Interrupted/Error + Tool*) — ловит терминатор хода. + // Подписываемся ДО send: события, отправленные до подписки, не реплеятся. val eventsJob = launch { agent.outbox.conversationEvents(Instant.DISTANT_PAST, conv.id) // onEach печатает и терминальный event, takeWhile лишь @@ -35,10 +37,19 @@ class SendSubcommand : AgentikSubcommand("send", "Отправить user-ход .takeWhile { ev -> !isTerminal(ev) } .collect { } } - // Даём SSE-подписке установиться, затем шлём ход. + // Онлайн-поток (дельты стриминга ответа) — live-only, без терминатора. + val onlineJob = launch { + agent.onlineOutbox.onlineEvents(conv.id) + .onEach { ev -> emitOnline(ev) } + .collect { } + } + // Даём SSE-подпискам установиться, затем шлём ход. delay(200) conv.send(listOf(Content.Text(text.joinToString(" ")))) eventsJob.join() + // Даём онлайн-потоку дослать хвостовые дельты, эмитнутые до End. + delay(100) + onlineJob.cancel() } finally { conv.close() } @@ -49,10 +60,6 @@ class SendSubcommand : AgentikSubcommand("send", "Отправить user-ход private fun emit(ev: Event) { when (ev) { - is Event.StartReasoning -> println("event StartReasoning") - is Event.StartResponse -> println("event StartResponse ${ev.responseType}") - is Event.AppendText -> println("event AppendText ${escape(ev.body)}") - is Event.AppendImage -> println("event AppendImage <${ev.body.size}B ${ev.mime}>") is Event.ToolCall -> println("event ToolCall ${ev.id} ${ev.toolName} ${escape(ev.toolArgs)}") is Event.ToolResult -> println("event ToolResult ${ev.toolCallId} ${escape(ev.result ?: "")}") is Event.End -> println("event End") @@ -61,5 +68,14 @@ class SendSubcommand : AgentikSubcommand("send", "Отправить user-ход } } + private fun emitOnline(ev: OnlineEvent) { + when (ev) { + is OnlineEvent.StartReasoning -> println("event StartReasoning") + is OnlineEvent.StartResponse -> println("event StartResponse ${ev.responseType}") + is OnlineEvent.AppendText -> println("event AppendText ${escape(ev.body)}") + is OnlineEvent.AppendImage -> println("event AppendImage <${ev.body.size}B ${ev.mime}>") + } + } + private fun escape(s: String): String = s.replace("\n", "\\n").replace("\r", "\\r") } diff --git a/agentik-tui/src/commonMain/kotlin/pw/binom/agentik/tui/TuiBackend.kt b/agentik-tui/src/commonMain/kotlin/pw/binom/agentik/tui/TuiBackend.kt index a887fe4..317c163 100644 --- a/agentik-tui/src/commonMain/kotlin/pw/binom/agentik/tui/TuiBackend.kt +++ b/agentik-tui/src/commonMain/kotlin/pw/binom/agentik/tui/TuiBackend.kt @@ -8,6 +8,7 @@ import pw.binom.agentik.proto.Agent import pw.binom.agentik.proto.Content import pw.binom.agentik.proto.Conversation import pw.binom.agentik.outbox.Event +import pw.binom.agentik.outbox.OnlineEvent import kotlin.coroutines.CoroutineContext import kotlin.time.Instant @@ -34,6 +35,9 @@ internal class TuiBackend( /** Активная джоба подписки на [Conversation.events]. */ private var eventsJob: Job? = null + /** Активная джоба подписки на онлайн-поток (стриминг ответа). */ + private var onlineJob: Job? = null + /** Последний виденный момент событий — для переподписки при reconnect. */ private var lastSeenAt: Instant = Instant.DISTANT_PAST @@ -92,41 +96,42 @@ internal class TuiBackend( } /** - * Подписывается на `outbox.conversationEvents(after, conv.id)` и перенаправляет их в [state]. + * Подписывается на durable-поток `outbox.conversationEvents(after, conv.id)` + * и live-поток `onlineOutbox.onlineEvents(conv.id)`; оба перенаправляет в [state]. + * + * Онлайн-поток live-only (без catchup), поэтому подписку открываем ДО [Conversation.send] + * (см. [ensureConversation] → [onUserMessage]), чтобы не упустить начало хода. */ private fun subscribeEvents(conv: Conversation, from: Instant) { eventsJob?.cancel() eventsJob = scope.launch { agent.outbox.conversationEvents(from, conv.id).collect { ce -> dispatch(ce.event) } } + onlineJob?.cancel() + onlineJob = scope.launch { + agent.onlineOutbox.onlineEvents(conv.id).collect { ev -> dispatchOnline(ev) } + } } /** - * Маппинг [Event] → [AppState] (что показать в TUI). + * Маппинг [Event] (durable) → [AppState] (что показать в TUI). * - * - AppendText → дописывает в последний ассистентский чанк - * - StartReasoning / StartResponse → новый streaming-чанк * - End → закрывает streaming * - Interrupted → закрывает streaming + системное сообщение * - ToolCall / ToolResult → сообщения в историю * - Error → системное сообщение + * + * Стриминг ответа (дельты текста/картинок) приходит отдельным потоком — + * см. [dispatchOnline]. */ private fun dispatch(ev: Event) { lastSeenAt = ev.date when (ev) { - is Event.AppendText -> state.appendAssistant(ev.body) - is Event.StartReasoning -> { - state.postSystem("… думаю") - } - is Event.StartResponse -> state.setStreaming(true) is Event.End -> state.finishAssistant() is Event.Interrupted -> { state.finishAssistant() state.postSystem("прервано") } - is Event.AppendImage -> { - state.postSystem("[картинка: ${ev.mime}, ${ev.body.size} байт]") - } is Event.ToolCall -> { state.postToolCall(toolName = ev.toolName, title = null, args = ev.toolArgs) } @@ -139,4 +144,20 @@ internal class TuiBackend( } } } + + /** + * Маппинг [OnlineEvent] (стриминг ответа, live-only) → [AppState]. + * + * - AppendText → дописывает в последний ассистентский чанк + * - StartReasoning / StartResponse → новый streaming-чанк + * - AppendImage → системное сообщение-заглушка + */ + private fun dispatchOnline(ev: OnlineEvent) { + when (ev) { + is OnlineEvent.AppendText -> state.appendAssistant(ev.body) + is OnlineEvent.StartReasoning -> state.postSystem("… думаю") + is OnlineEvent.StartResponse -> state.setStreaming(true) + is OnlineEvent.AppendImage -> state.postSystem("[картинка: ${ev.mime}, ${ev.body.size} байт]") + } + } } 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 ae7f932..fc9c358 100644 --- a/client/src/commonMain/kotlin/pw/binom/agentik/client/AgentClient.kt +++ b/client/src/commonMain/kotlin/pw/binom/agentik/client/AgentClient.kt @@ -14,6 +14,7 @@ import io.ktor.http.contentType import kotlinx.coroutines.runBlocking import pw.binom.agentik.journal.ConversationStore import pw.binom.agentik.journal.JournalStore +import pw.binom.agentik.outbox.OnlineOutbox import pw.binom.agentik.outbox.OutboxStore import pw.binom.agentik.proto.Agent import pw.binom.agentik.proto.AgentInfo @@ -45,6 +46,7 @@ internal class AgentClient private constructor( private val agentUrl: String = baseUrl.trimEnd('/') override val outbox: OutboxStore = HttpEventStore(httpClient = httpClient, baseUrl = agentUrl) + override val onlineOutbox: OnlineOutbox = HttpOnlineOutbox(httpClient = httpClient, baseUrl = agentUrl) override val journal: JournalStore = HttpJournalStore(httpClient = httpClient, baseUrl = agentUrl) override val conversationStore: ConversationStore = HttpConversationStore(httpClient = httpClient, baseUrl = agentUrl) diff --git a/client/src/commonMain/kotlin/pw/binom/agentik/client/HttpOnlineOutbox.kt b/client/src/commonMain/kotlin/pw/binom/agentik/client/HttpOnlineOutbox.kt new file mode 100644 index 0000000..e260f25 --- /dev/null +++ b/client/src/commonMain/kotlin/pw/binom/agentik/client/HttpOnlineOutbox.kt @@ -0,0 +1,46 @@ +package pw.binom.agentik.client + +import io.ktor.client.HttpClient +import io.ktor.client.request.prepareGet +import io.ktor.client.statement.bodyAsChannel +import io.ktor.http.HttpStatusCode +import kotlinx.coroutines.flow.Flow +import kotlinx.coroutines.flow.flow +import pw.binom.agentik.outbox.OnlineEvent +import pw.binom.agentik.outbox.OnlineOutbox + +/** + * HTTP-реализация [OnlineOutbox] (= [pw.binom.agentik.outbox.OnlineOutbox]), + * ходящая в `:server`-фасад. + * + * **Endpoint**: [onlineEvents] → `GET {baseUrl}/conversations/{id}/online` + * (live-only SSE, без `after`). Сервер не реплеит — подписка получает + * только то, что эмитится после подключения. Reconnect-логика здесь не + * нужна: при обрыве поток просто закрывается, а потерянные дельты + * восстанавливаются из durable-истории ([HttpEventStore] + журнал). + */ +internal class HttpOnlineOutbox( + private val httpClient: HttpClient, + private val baseUrl: String, +) : OnlineOutbox { + + private val agentUrl: String = baseUrl.trimEnd('/') + + override fun onlineEvents(conversationId: String): Flow = flow { + val url = "$agentUrl/conversations/$conversationId/online" + httpClient.prepareGet(url) { noReadTimeout() } + .execute { response -> + check(response.status == HttpStatusCode.OK) { + "onlineEvents: server returned ${response.status}" + } + readSse(response.bodyAsChannel()) + .collect { payload -> + emit(agentikJson.decodeFromString(OnlineEvent.serializer(), payload)) + } + } + } + + override fun close() { + // HttpClient закрывает владелец (AgentClient / AgentikAgent). + } +} diff --git a/client/src/commonTest/kotlin/pw/binom/agentik/client/SseTimeoutTest.kt b/client/src/commonTest/kotlin/pw/binom/agentik/client/SseTimeoutTest.kt index 67b9b6c..15edfd7 100644 --- a/client/src/commonTest/kotlin/pw/binom/agentik/client/SseTimeoutTest.kt +++ b/client/src/commonTest/kotlin/pw/binom/agentik/client/SseTimeoutTest.kt @@ -26,7 +26,7 @@ import kotlin.test.fail /** * Репродукция бага Ktor CIO: дефолтный [io.ktor.client.engine.cio.CIOEngineConfig.requestTimeout] - * = 15 с убивает SSE read. Наш fix — [noSseReadTimeout] ставит capability + * = 15 с убивает SSE read. Наш fix — [noReadTimeout] ставит capability * [io.ktor.client.plugins.HttpTimeoutCapability] со всеми таймаутами = INFINITE * перед каждым read-стримом. * @@ -43,7 +43,7 @@ class SseTimeoutTest { private fun freePort(): Int = ServerSocket(0).use { it.localPort } @Test - fun `sse read survives past default cio timeout with noSseReadTimeout`(): Unit = runBlocking { + fun `sse read survives past default cio timeout with noReadTimeout`(): Unit = runBlocking { val port = freePort() val server = embeddedServer(io.ktor.server.cio.CIO, port = port) { routing { @@ -65,7 +65,7 @@ class SseTimeoutTest { val received = mutableListOf() client.prepareGet("http://127.0.0.1:$port/sse") { header("Accept", "text/event-stream") - noSseReadTimeout() + noReadTimeout() }.execute { resp -> val ch = resp.bodyAsChannel() // 19 с запас: ждём, пока сервер пошлёт "done" после 17 с. @@ -95,14 +95,14 @@ class SseTimeoutTest { } /** - * Контр-тест: убеждаемся что БЕЗ [noSseReadTimeout] дефолтный + * Контр-тест: убеждаемся что БЕЗ [noReadTimeout] дефолтный * CIO requestTimeout = 15 с действительно убивает SSE-стрим. * Сервер держит stream 17 с; если клиент не выставил capability — * мы должны получить [HttpRequestTimeoutException] на ~15 с, не * дожидаясь "done". */ @Test - fun `without noSseReadTimeout default cio requestTimeout kills the stream`(): Unit = runBlocking { + fun `without noReadTimeout default cio requestTimeout kills the stream`(): Unit = runBlocking { val port = freePort() val server = embeddedServer(io.ktor.server.cio.CIO, port = port) { routing { @@ -123,7 +123,7 @@ class SseTimeoutTest { try { client.prepareGet("http://127.0.0.1:$port/sse") { header("Accept", "text/event-stream") - // НАМЕРЕННО без noSseReadTimeout. + // НАМЕРЕННО без noReadTimeout. }.execute { resp -> val ch = resp.bodyAsChannel() // Читаем строки, пока не придёт "data: done" — без capability diff --git a/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/Event.kt b/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/Event.kt index c9d756f..e265260 100644 --- a/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/Event.kt +++ b/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/Event.kt @@ -12,8 +12,14 @@ import kotlin.time.Instant * порядка при равных timestamps. * * Базовая структура хода: - * `Working` → `StartReasoning?` → `StartResponse(TEXT|IMAGE)` → ...контент... → `End` | `Interrupted` | `Error`. - * `StartReasoning` может отсутствовать, если агент не показывал рассуждения. + * `Working` → `End` | `Interrupted` | `Error`, между ними — целые события + * ([ToolCall]/[ToolResult]/[ToolFailed]). + * + * **Стриминг ответа (дельты текста/картинок и маркеры фаз) — НЕ здесь.** + * Дельты токенов живут в [OnlineEvent] (live-only, не сохраняются). + * [Event] — только «целые» (durable) события, пригодные к перезапросу + * по курсору. + * * `Working` — маркер «агент принял запрос и пошёл обрабатывать», эмитится * синхронно в `Conversation.send()` ДО старта LLM-цикла (и возможной долгой * очереди turnLock'а). Парный терминатор не нужен: `End`/`Interrupted`/`Error` @@ -32,15 +38,9 @@ sealed interface Event { /** Момент эмиссии события в UTC. */ val date: Instant - @Serializable - enum class ResponseType { - @SerialName("text") TEXT, - @SerialName("image") IMAGE - } - /** * Маркер «агент принял запрос и пошёл обрабатывать». Эмитится **до** - * [StartReasoning]/[StartResponse], синхронно из `Conversation.send()`, + * [End]/[Interrupted]/[Error], синхронно из `Conversation.send()`, * чтобы UI мог показать спиннер ещё до первого токена ответа. * Терминатор хода ([End]/[Interrupted]/[Error]) — парный. */ @@ -48,16 +48,6 @@ sealed interface Event { @SerialName("working") data class Working(override val date: Instant) : Event - /** Ассистент начал рассуждение (опциональный маркер; контент рассуждения приходит через [AppendText]). */ - @Serializable - @SerialName("start_reasoning") - data class StartReasoning(override val date: Instant) : Event - - /** Начало ответа ассистента заданного типа. После него идут соответствующие `Append*`/`Tool*`-события, потом [End]/[Interrupted]/[Error]. */ - @Serializable - @SerialName("start_response") - data class StartResponse(override val date: Instant, val responseType: ResponseType) : Event - /** Ход завершён нормально. Соответствующий `Message.AssistantMessage` появится в `getMessages`. */ @Serializable @SerialName("end") @@ -68,14 +58,6 @@ sealed interface Event { @SerialName("interrupted") data class Interrupted(override val date: Instant) : Event - @Serializable - @SerialName("append_text") - data class AppendText(override val date: Instant, val body: String) : Event - - @Serializable - @SerialName("append_image") - data class AppendImage(override val date: Instant, val body: ByteArray, val mime: String) : Event - /** * Агент начал вызов тула. Аргументы приходят целиком — стриминга нет. * [id] совпадает с id соответствующего `Message.ToolCall` в истории diff --git a/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/MutableOnlineOutbox.kt b/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/MutableOnlineOutbox.kt new file mode 100644 index 0000000..f0426ab --- /dev/null +++ b/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/MutableOnlineOutbox.kt @@ -0,0 +1,26 @@ +package pw.binom.agentik.outbox + +/** + * Write-сторона [OnlineOutbox]: используется продюсерами + * ([pw.binom.agentik.standalone.agent.ConversationLoop] и т.п.) для эмиссии + * live-дельт. Наружу (в [pw.binom.agentik.proto.Agent]) отдаётся + * read-only [OnlineOutbox]. + * + * **Non-suspend и best-effort**: онлайн-события по определению нигде не + * персистятся, I/O нет — блокировать продюсера незачем. [tryAppendOnline] + * не буферизует и не ждёт (см. [OnlineOutbox]): медленный подписчик может + * потерять дельту, это допустимо. + */ +interface MutableOnlineOutbox : OnlineOutbox { + + suspend fun appendOnline(conversationId: String, event: OnlineEvent) + + /** + * Эмитит [event] в live-канал диалога [conversationId]. Не сохраняется. + * + * Возвращает `true`, если событие принято live-каналом. Возврат `false` + * (нет активных подписчиков / буфер переполнен с DROP-политикой) — + * не ошибка: у онлайн-событий нет гарантии доставки. + */ + fun tryAppendOnline(conversationId: String, event: OnlineEvent): Boolean +} diff --git a/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/OnlineEvent.kt b/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/OnlineEvent.kt new file mode 100644 index 0000000..0ef7f8f --- /dev/null +++ b/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/OnlineEvent.kt @@ -0,0 +1,56 @@ +package pw.binom.agentik.outbox + +import kotlinx.serialization.SerialName +import kotlinx.serialization.Serializable +import kotlin.time.Instant + +/** + * **Онлайн-события** диалога: стриминг ответа агента «в моменте» — + * дельты текста/картинок и маркеры фаз хода. + * + * Принципиальное отличие от [Event] (durable): + * - **Никогда и нигде не сохраняются** — ни в буфер [OnlineOutbox], + * ни в journal. Это чистый live-канал. + * - **Только онлайн-подписка**: события, эмитнутые до подписки + * (или в момент обрыва соединения), не реплеятся и не восстанавливаются. + * Потерянный фрагмент не страшен — целый результат хода приходит + * durable-событием ([Event.End]) и/или лежит в журнале. + * - **Нет курсора**: у потока нет `after`/`lastSeen` — курсор там, где + * есть что реплеить. + * + * Зачем разделять: дельты токенов — высокочастотный мусор, который, + * попав в durable store, копится в RAM (standalone-outbox растёт unbounded) + * и засоряет историю. В [Event] остаются только «целые» события, + * пригодные к перезапросу по курсору. + */ +@Serializable +sealed interface OnlineEvent { + /** Момент эмиссии события в UTC (для упорядочивания в рамках стрима). */ + val date: Instant + + @Serializable + enum class ResponseType { + @SerialName("text") TEXT, + @SerialName("image") IMAGE + } + + /** Ассистент начал рассуждение (опциональный маркер; контент идёт через [AppendText]). */ + @Serializable + @SerialName("start_reasoning") + data class StartReasoning(override val date: Instant) : OnlineEvent + + /** Начало ответа ассистента заданного типа. Далее идут соответствующие `Append*`. */ + @Serializable + @SerialName("start_response") + data class StartResponse(override val date: Instant, val responseType: ResponseType) : OnlineEvent + + /** Очередная дельта текста ответа. */ + @Serializable + @SerialName("append_text") + data class AppendText(override val date: Instant, val body: String) : OnlineEvent + + /** Очередная дельта картинки ответа. */ + @Serializable + @SerialName("append_image") + data class AppendImage(override val date: Instant, val body: ByteArray, val mime: String) : OnlineEvent +} diff --git a/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/OnlineOutbox.kt b/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/OnlineOutbox.kt new file mode 100644 index 0000000..3b813b3 --- /dev/null +++ b/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/OnlineOutbox.kt @@ -0,0 +1,34 @@ +package pw.binom.agentik.outbox + +import kotlinx.coroutines.flow.Flow + +/** + * Live-канал **онлайн-событий** ([OnlineEvent]) диалога — стриминга ответа + * агента «в моменте». + * + * Контракт принципиально проще [OutboxStore]: + * - **без catchup**: `onlineEvents` отдаёт только то, что эмитится *после* + * подписки. Прошлое не реплеится — его и нет (нечего хранить). + * - **без курсора**: нет `after`/`lastSeen`. + * - **без TTL/buffer**: [OnlineEvent] не буферизуются. + * + * Это осознанный компромисс: дельты токенов — высокочастотный мусор, + * который в durable-сторе копился бы в RAM и засорял историю. Потеря + * фрагмента при обрыве не критична — целый ответ приходит [Event.End] + * и/или лежит в [pw.binom.agentik.journal.JournalStore]. + * + * Read-only view: запись — через [MutableOnlineOutbox]. + */ +interface OnlineOutbox : AutoCloseable { + + fun onlineEvents(): Flow + + /** + * Подписка на live-поток онлайн-событий диалога [conversationId]. + * События, эмитнутые до подписки, не приходят. + */ + fun onlineEvents(conversationId: String): Flow + + /** Освобождает ресурсы. Idempotent. */ + override fun close() +} diff --git a/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/OutboxStore.kt b/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/OutboxStore.kt index b68174e..df4a2ed 100644 --- a/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/OutboxStore.kt +++ b/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/OutboxStore.kt @@ -8,6 +8,11 @@ import kotlin.time.Instant /** * Bounded-tail event log с автоматическим управлением TTL. * + * Хранит **только durable-события [Event]** — «целые» факты хода + * (Working/End/Interrupted/Error, ToolCall/ToolResult/ToolFailed). + * Высокочастотный **стриминг ответа** (дельты текста/картинок) сюда + * НЕ попадает — он живёт в [OnlineOutbox] (live-only, не сохраняется). + * * **Архитектура двухуровневого хранилища событий**: * 1. **Этот store** = короткий bounded tail (live SSE + недавний replay). * События автоматически эвиктятся по TTL/cap (implementation-defined). diff --git a/outbox-inmemory/src/commonMain/kotlin/pw/binom/agentik/outbox/inmemory/InMemoryOnlineOutbox.kt b/outbox-inmemory/src/commonMain/kotlin/pw/binom/agentik/outbox/inmemory/InMemoryOnlineOutbox.kt new file mode 100644 index 0000000..0d5953d --- /dev/null +++ b/outbox-inmemory/src/commonMain/kotlin/pw/binom/agentik/outbox/inmemory/InMemoryOnlineOutbox.kt @@ -0,0 +1,65 @@ +package pw.binom.agentik.outbox.inmemory + +import kotlinx.coroutines.channels.BufferOverflow +import kotlinx.coroutines.flow.Flow +import kotlinx.coroutines.flow.MutableSharedFlow +import kotlinx.coroutines.flow.filter +import kotlinx.coroutines.flow.map +import pw.binom.agentik.outbox.MutableOnlineOutbox +import pw.binom.agentik.outbox.OnlineEvent + +/** + * In-memory реализация [MutableOnlineOutbox] — единственный live-канал + * без хранения. + * + * **Никакого буфера событий**: [OnlineEvent] нигде не накапливаются + * (в этом весь смысл — дельты токенов не должны течь в durable-стор). + * Единственное, что живёт в памяти, — [MutableSharedFlow] с BOUNDED + * internal buffer'ом для развязки продюсера/подписчиков; при переполнении + * **старые дропаются** ([BufferOverflow.DROP_OLDEST]), [appendOnline] не + * блокируется. Потеря дельты допустима (см. [OnlineOutbox]). + * + * **Маршрутизация**: один общий [MutableSharedFlow] c `conversationId` + * в envelope; [onlineEvents] фильтрует по диалогу. Отдельный flow-на-диалог + * не держим, чтобы не плодить per-conversation подписки, которые надо + * чистить вручную. + */ +class InMemoryOnlineOutbox( + liveBufferCapacity: Int = DEFAULT_LIVE_BUFFER_CAPACITY, +) : MutableOnlineOutbox { + + private data class Envelope(val conversationId: String, val event: OnlineEvent) + + private val liveFlow = MutableSharedFlow( + replay = 0, + extraBufferCapacity = liveBufferCapacity, + onBufferOverflow = BufferOverflow.DROP_OLDEST, + ) + + init { + require(liveBufferCapacity > 0) { + "liveBufferCapacity must be > 0, got $liveBufferCapacity" + } + } + + override fun onlineEvents(): Flow = + liveFlow.map { it.event } + + override fun onlineEvents(conversationId: String): Flow = + liveFlow.filter { it.conversationId == conversationId }.map { it.event } + + override suspend fun appendOnline(conversationId: String, event: OnlineEvent)= + liveFlow.emit(Envelope(conversationId, event)) + + override fun tryAppendOnline(conversationId: String, event: OnlineEvent): Boolean = + liveFlow.tryEmit(Envelope(conversationId, event)) + + override fun close() { + // replay = 0 — чистить нечего; сам flow соберётся GC'ом при выходе ссылки. + // Идемпотентно: повторный close() безопасен. + } + + private companion object { + private const val DEFAULT_LIVE_BUFFER_CAPACITY = 1 + } +} diff --git a/outbox-inmemory/src/commonTest/kotlin/pw/binom/agentik/outbox/inmemory/InMemoryOutboxStoreTest.kt b/outbox-inmemory/src/commonTest/kotlin/pw/binom/agentik/outbox/inmemory/InMemoryOutboxStoreTest.kt index 6611bf4..9eb5f34 100644 --- a/outbox-inmemory/src/commonTest/kotlin/pw/binom/agentik/outbox/inmemory/InMemoryOutboxStoreTest.kt +++ b/outbox-inmemory/src/commonTest/kotlin/pw/binom/agentik/outbox/inmemory/InMemoryOutboxStoreTest.kt @@ -164,7 +164,7 @@ class InMemoryOutboxStoreTest { store.append(CommonEvent.Conversation( date = now, conversationId = "c-1", - event = Event.AppendText(date = now, body = "hi"), + event = Event.Working(date = now), )) // Snapshot-based test of the default impl (uses events() + filterIsInstance). @@ -180,9 +180,9 @@ class InMemoryOutboxStoreTest { fun `conversationEvents with conversationId filters to that conversation`() = runBlocking { val store = InMemoryOutboxStore(maxMessages = null, ttl = null) val now = Instant.fromEpochSeconds(0) - store.append(CommonEvent.Conversation(now, "c-1", Event.AppendText(now, "a"))) - store.append(CommonEvent.Conversation(now, "c-2", Event.AppendText(now, "b"))) - store.append(CommonEvent.Conversation(now, "c-1", Event.AppendText(now, "c"))) + store.append(CommonEvent.Conversation(now, "c-1", Event.Working(now))) + store.append(CommonEvent.Conversation(now, "c-2", Event.Working(now))) + store.append(CommonEvent.Conversation(now, "c-1", Event.Working(now))) // Test the filter logic by manually filtering snapshot. val c1 = store.snapshot() @@ -200,7 +200,7 @@ class InMemoryOutboxStoreTest { date = now, event = AgentEvent.Created(date = now, conversationId = "created"), )) - store.append(CommonEvent.Conversation(now, "c-1", Event.AppendText(now, "hi"))) + store.append(CommonEvent.Conversation(now, "c-1", Event.Working(now))) val all = store.snapshot() val agents = all.filterIsInstance() 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 f4d2061..ec8839b 100644 --- a/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Agent.kt +++ b/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Agent.kt @@ -2,6 +2,8 @@ package pw.binom.agentik.proto import pw.binom.agentik.journal.ConversationStore import pw.binom.agentik.journal.JournalStore +import pw.binom.agentik.outbox.OnlineEvent +import pw.binom.agentik.outbox.OnlineOutbox import pw.binom.agentik.outbox.OutboxStore import kotlin.time.Instant @@ -12,9 +14,10 @@ import kotlin.time.Instant * [createConversation] возвращает [Conversation], который сам хранит историю * и которому отправляют ходы через [Conversation.send]. * - * **Хранилища вынесены в [Agent.journal], [Agent.outbox] и - * [Agent.conversationStore]**: все три read-only views. События живут - * в [outbox] как `OutboxStore.events(after)` / `outbox.agentEvents(after)`. + * **Хранилища вынесены в [Agent.journal], [Agent.outbox], + * [Agent.onlineOutbox] и [Agent.conversationStore]**: read-only views. + * Durable-события живут в [outbox] как `OutboxStore.events(after)` / + * `outbox.agentEvents(after)`; live-стриминг ответа — в [onlineOutbox]. * Это даёт единый путь для всех read-операций по хранилищу и убирает * дублирование между протоколом и хранилищем. * @@ -83,6 +86,24 @@ interface Agent : AutoCloseable { */ val outbox: OutboxStore + /** + * Live-канал **онлайн-событий** ([OnlineEvent] — стриминг ответа агента). + * + * Отделён от [outbox] принципиально: + * - [outbox] — durable: «целые» события, перезапрашиваемые по курсору + * (`after`), с catchup + live; + * - [onlineOutbox] — live-only: дельты токенов, **никогда не сохраняются**, + * без catchup/курсора, подписка работает только «онлайн». + * + * Используется HTTP-фасадом `:server` для endpoint'а + * `GET /{path}/conversations/{id}/online` (SSE). Клиент рендерит ответ + * из этого потока в реальном времени, а durable-историю берёт из [journal]. + * + * **Read-only**: write-доступ только через `MutableOnlineOutbox` внутри + * ChatAgent / ConversationLoop, не через [Agent] interface. + */ + val onlineOutbox: OnlineOutbox + /** * Read-only view на `conversation` table (id + title + timestamps). * diff --git a/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Conversation.kt b/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Conversation.kt index 5bad0c0..edcf82e 100644 --- a/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Conversation.kt +++ b/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Conversation.kt @@ -9,15 +9,21 @@ import kotlin.time.Instant * [send] агенту не нужно пересылать транскрипт — он уже живёт внутри * [Conversation]. * - * **Live-события** диалога (turn stream: StartReasoning / AppendText / End / - * ToolCall / ToolResult / ...) НЕ часть этого интерфейса — единственный - * источник live-событий это [pw.binom.agentik.outbox.OutboxStore]. - * Подписаться на события конкретного диалога: - * ``` - * agent.outbox.conversationEvents(after = lastSeen, conversationId = id) - * .map { it.event } - * .collect { e -> ... } - * ``` + * **Live-события** диалога НЕ часть этого интерфейса. Их два независимых + * потока: + * - **durable** ([pw.binom.agentik.outbox.Event]: Working / End / Interrupted / + * Error / ToolCall / ToolResult / ToolFailed) — из + * [pw.binom.agentik.outbox.OutboxStore], перезапрашивается по курсору: + * ``` + * agent.outbox.conversationEvents(after = lastSeen, conversationId = id) + * .map { it.event } + * .collect { e -> ... } + * ``` + * - **online** ([pw.binom.agentik.outbox.OnlineEvent]: StartReasoning / + * StartResponse / AppendText / AppendImage) — live-only стриминг ответа + * из [pw.binom.agentik.outbox.OnlineOutbox], без catchup/курсора: + * `agent.onlineOutbox.onlineEvents(id)`. + * * Для cross-conversation view (admin / parent-agent / debug): * `agent.outbox.events(after)`. Для lifecycle агента (created/deleted/renamed): * `agent.outbox.agentEvents(after)`. 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 047c057..17c628d 100644 --- a/server/src/commonMain/kotlin/pw/binom/agentik/server/Routes.kt +++ b/server/src/commonMain/kotlin/pw/binom/agentik/server/Routes.kt @@ -25,6 +25,7 @@ import pw.binom.agentik.proto.Content import pw.binom.agentik.outbox.AgentEvent import pw.binom.agentik.outbox.CommonEvent import pw.binom.agentik.outbox.Event +import pw.binom.agentik.outbox.OnlineEvent import kotlin.time.Instant internal fun Route.agentikRoutes(agent: Agent) { @@ -167,6 +168,27 @@ internal fun Route.agentikRoutes(agent: Agent) { call.streamJsonSse(agent.outbox.conversationEvents(after, id).map { it.event }, Event.serializer()) } + /** + * `GET /conversations/{id}/online` — **live-only** SSE поток + * онлайн-событий ([pw.binom.agentik.outbox.OnlineEvent]): дельты + * текста/картинок стриминга ответа. + * + * Отличия от `/conversations/{id}/events`: + * - нет `after` — catchup невозможен, события не хранятся; + * - подписка получает только то, что эмитится после подключения; + * - потерянные дельты восстанавливаются из durable-истории + * (`/events` + журнал), а не реплеятся здесь. + */ + get("/conversations/{id}/online") { + val id = call.parameters["id"]!! + val c = agent.getConversation(id) + if (c == null) { + call.respond(HttpStatusCode.NotFound) + return@get + } + call.streamJsonSse(agent.onlineOutbox.onlineEvents(id), OnlineEvent.serializer()) + } + get("/events") { val after = call.parseAfter() ?: return@get // agent.outbox.agentEvents(after) возвращает Flow; diff --git a/server/src/commonTest/kotlin/pw/binom/agentik/server/AgentInfoRouteTest.kt b/server/src/commonTest/kotlin/pw/binom/agentik/server/AgentInfoRouteTest.kt index ae86809..9f66c1a 100644 --- a/server/src/commonTest/kotlin/pw/binom/agentik/server/AgentInfoRouteTest.kt +++ b/server/src/commonTest/kotlin/pw/binom/agentik/server/AgentInfoRouteTest.kt @@ -21,6 +21,7 @@ import pw.binom.agentik.journal.ConversationStore import pw.binom.agentik.journal.JournalStore import pw.binom.agentik.journal.MessageRecord import pw.binom.agentik.outbox.CommonEvent +import pw.binom.agentik.outbox.OnlineOutbox import pw.binom.agentik.outbox.OutboxStore import pw.binom.agentik.proto.Agent import pw.binom.agentik.proto.AgentInfo @@ -74,6 +75,7 @@ class AgentInfoRouteTest { override suspend fun earliestEventDate(): Instant = Instant.DISTANT_PAST override fun close() {} } + override val onlineOutbox: OnlineOutbox = emptyOnlineOutbox() override val conversationStore: ConversationStore = object : ConversationStore { override suspend fun get(id: String) = null override suspend fun list(offset: Int, limit: Int) = emptyList() @@ -143,6 +145,7 @@ class AgentInfoRouteTest { override suspend fun earliestEventDate(): Instant = Instant.DISTANT_PAST override fun close() {} } + override val onlineOutbox: OnlineOutbox = emptyOnlineOutbox() override val conversationStore: ConversationStore = object : ConversationStore { override suspend fun get(id: String) = null override suspend fun list(offset: Int, limit: Int) = emptyList() diff --git a/server/src/commonTest/kotlin/pw/binom/agentik/server/BearerTokenTest.kt b/server/src/commonTest/kotlin/pw/binom/agentik/server/BearerTokenTest.kt index 423c59f..5697862 100644 --- a/server/src/commonTest/kotlin/pw/binom/agentik/server/BearerTokenTest.kt +++ b/server/src/commonTest/kotlin/pw/binom/agentik/server/BearerTokenTest.kt @@ -17,6 +17,7 @@ import pw.binom.agentik.journal.ConversationStore import pw.binom.agentik.journal.JournalStore import pw.binom.agentik.journal.MessageRecord import pw.binom.agentik.outbox.CommonEvent +import pw.binom.agentik.outbox.OnlineOutbox import pw.binom.agentik.outbox.OutboxStore import pw.binom.agentik.proto.Agent import pw.binom.agentik.proto.AgentInfo @@ -58,6 +59,7 @@ class BearerTokenTest { override suspend fun earliestEventDate(): Instant = Instant.DISTANT_PAST override fun close() {} } + override val onlineOutbox: OnlineOutbox = emptyOnlineOutbox() override val conversationStore: ConversationStore = object : ConversationStore { override suspend fun get(id: String) = null override suspend fun list(offset: Int, limit: Int) = emptyList() diff --git a/server/src/commonTest/kotlin/pw/binom/agentik/server/ConversationRoutesTest.kt b/server/src/commonTest/kotlin/pw/binom/agentik/server/ConversationRoutesTest.kt index e745881..1604bc0 100644 --- a/server/src/commonTest/kotlin/pw/binom/agentik/server/ConversationRoutesTest.kt +++ b/server/src/commonTest/kotlin/pw/binom/agentik/server/ConversationRoutesTest.kt @@ -20,6 +20,7 @@ import pw.binom.agentik.journal.JournalStore import pw.binom.agentik.journal.inmemory.InMemoryJournalStore import pw.binom.agentik.journal.inmemory.InMemoryMutableConversationStore import pw.binom.agentik.outbox.CommonEvent +import pw.binom.agentik.outbox.OnlineOutbox import pw.binom.agentik.outbox.OutboxStore import pw.binom.agentik.proto.Agent import pw.binom.agentik.proto.AgentInfo @@ -168,6 +169,7 @@ class ConversationRoutesTest { override suspend fun earliestEventDate(): Instant = Instant.DISTANT_PAST override fun close() {} } + override val onlineOutbox: OnlineOutbox = emptyOnlineOutbox() override val conversationStore: ConversationStore = cs override fun createConversation(temp: Boolean): pw.binom.agentik.proto.Conversation = diff --git a/server/src/commonTest/kotlin/pw/binom/agentik/server/JournalRoutesCountTest.kt b/server/src/commonTest/kotlin/pw/binom/agentik/server/JournalRoutesCountTest.kt index 5a2e343..42e4c8e 100644 --- a/server/src/commonTest/kotlin/pw/binom/agentik/server/JournalRoutesCountTest.kt +++ b/server/src/commonTest/kotlin/pw/binom/agentik/server/JournalRoutesCountTest.kt @@ -24,6 +24,7 @@ import pw.binom.agentik.journal.JournalStore import pw.binom.agentik.journal.MessageRecord import pw.binom.agentik.journal.inmemory.InMemoryJournalStore import pw.binom.agentik.outbox.CommonEvent +import pw.binom.agentik.outbox.OnlineOutbox import pw.binom.agentik.outbox.OutboxStore import pw.binom.agentik.proto.Agent import pw.binom.agentik.proto.AgentInfo @@ -187,6 +188,7 @@ class JournalRoutesCountTest { override suspend fun earliestEventDate(): Instant = Instant.DISTANT_PAST override fun close() {} } + override val onlineOutbox: OnlineOutbox = emptyOnlineOutbox() override val conversationStore: ConversationStore = object : ConversationStore { override suspend fun get(id: String) = null override suspend fun list(offset: Int, limit: Int) = emptyList() diff --git a/server/src/commonTest/kotlin/pw/binom/agentik/server/TestOnlineOutbox.kt b/server/src/commonTest/kotlin/pw/binom/agentik/server/TestOnlineOutbox.kt new file mode 100644 index 0000000..48be547 --- /dev/null +++ b/server/src/commonTest/kotlin/pw/binom/agentik/server/TestOnlineOutbox.kt @@ -0,0 +1,14 @@ +package pw.binom.agentik.server + +import kotlinx.coroutines.flow.emptyFlow +import pw.binom.agentik.outbox.OnlineEvent +import pw.binom.agentik.outbox.OnlineOutbox + +/** + * Пустой [OnlineOutbox] для тестовых [pw.binom.agentik.proto.Agent]-заглушек: + * онлайн-канал тестам не нужен, но интерфейс обязывает его отдать. + */ +internal fun emptyOnlineOutbox(): OnlineOutbox = object : OnlineOutbox { + override fun onlineEvents(conversationId: String) = emptyFlow() + override fun close() {} +} diff --git a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/A2aBridge.kt b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/A2aBridge.kt index 9c906c5..c4d9802 100644 --- a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/A2aBridge.kt +++ b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/A2aBridge.kt @@ -11,6 +11,7 @@ import pw.binom.a2a.model.Role import pw.binom.a2a.model.TextPart import pw.binom.a2a.server.AgentHandler import pw.binom.agentik.outbox.Event +import pw.binom.agentik.outbox.OnlineEvent import pw.binom.agentik.proto.Agent import pw.binom.agentik.proto.Content import pw.binom.agentik.proto.Conversation @@ -27,9 +28,10 @@ private val log = KotlinLogging.logger {} * -> новый диалог, контекст пересоздаётся. * - id внутреннего диалога отдаётся клиенту в `metadata.agentikConversationId` ответа. * - * Ответ A2A = склеенные [Event.AppendText] нашего хода. Подписку на [Conversation.events] - * открываем ДО [Conversation.send] (иначе события начала хода могут быть упущены), - * завершение хода ждём по [Event.End] / [Event.Interrupted] / [Event.Error]. + * Ответ A2A = склеенные [OnlineEvent.AppendText] нашего хода. Подписку на онлайн-поток + * ([pw.binom.agentik.outbox.OnlineOutbox]) открываем ДО [Conversation.send] (live-only, + * без catchup — события начала хода иначе можно упустить), завершение хода ждём + * по durable-событиям [Event.End] / [Event.Interrupted] / [Event.Error]. * * Ограничение v1: tool-события и картинки в A2A-ответ не транслируются; * при нескольких ходов в очереди за контекстом текст предыдущего хода @@ -48,17 +50,21 @@ class A2aBridge(private val agent: Agent) : AgentHandler { val since = conv.updatedAt val reply = StringBuilder() val turnDone = CompletableDeferred() - val subscription = async { + // Онлайн-поток — дельты ответа (live-only, без catchup). + val onlineJob = async { + agent.onlineOutbox.onlineEvents(conv.id).collect { e -> + if (e is OnlineEvent.AppendText) reply.append(e.body) + } + } + // Durable-поток — терминатор хода (catchup + live). + val turnJob = async { agent.outbox.conversationEvents(since, conv.id).collect { ce -> - val e = ce.event - when (e) { - is Event.AppendText -> reply.append(e.body) + when (val e = ce.event) { is Event.End, is Event.Interrupted -> turnDone.complete(Unit) is Event.Error -> - if (!turnDone.completeExceptionally( - IllegalStateException("agent turn failed: ${e.message}") - ) - ) {} + turnDone.completeExceptionally( + IllegalStateException("agent turn failed: ${e.message}") + ) else -> {} } } @@ -67,7 +73,8 @@ class A2aBridge(private val agent: Agent) : AgentHandler { try { turnDone.await() } finally { - subscription.cancel() + onlineJob.cancel() + turnJob.cancel() } log.info { "a2a context=$contextId conv=${conv.id} reply=${reply.length} chars" } Message( diff --git a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ChatAgent.kt b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ChatAgent.kt index 07dac27..1916a23 100644 --- a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ChatAgent.kt +++ b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ChatAgent.kt @@ -29,6 +29,8 @@ import pw.binom.agentik.outbox.AgentEvent import pw.binom.agentik.outbox.Event as OutboxEvent import pw.binom.agentik.outbox.CommonEvent import pw.binom.agentik.outbox.MutableOutboxStore +import pw.binom.agentik.outbox.MutableOnlineOutbox +import pw.binom.agentik.outbox.OnlineOutbox import pw.binom.agentik.journal.JournalStore import pw.binom.agentik.outbox.OutboxStore import pw.binom.agentik.proto.Conversation as ProtoConversation @@ -213,6 +215,13 @@ private val eventStore: MutableOutboxStore = pw.binom.agentik.outbox.inmemory.In ttl = null, ) +/** + * Live-канал стриминга ответа (дельты текста/картинок и фазовые маркеры). + * Онлайн-события никогда не сохраняются и не реплеятся — см. [OnlineOutbox]. + * Durable-события по-прежнему идут в [eventStore]. + */ +private val onlineEventStore: MutableOnlineOutbox = pw.binom.agentik.outbox.inmemory.InMemoryOnlineOutbox() + /** * Собирает **актуальный** список тулов для диспетчеризации: * внешние из [toolProviders] + testTools + встроенные (read_skill, @@ -413,6 +422,15 @@ private val eventStore: MutableOutboxStore = pw.binom.agentik.outbox.inmemory.In override val outbox: OutboxStore get() = eventStore + /** + * Read-only view of [onlineEventStore] для HTTP-фасада в `:server` + * (`Route.agentikAgent` → `/conversations/{id}/online` SSE). + * + * Live-only: без catchup/курсора, события не сохраняются. + */ + override val onlineOutbox: OnlineOutbox + get() = onlineEventStore + /** * Read-only view of [mutableConversationStore] для HTTP-фасада в `:server` * (`GET /conversations` → список [ConversationRecord] для UI). @@ -456,6 +474,7 @@ private val eventStore: MutableOutboxStore = pw.binom.agentik.outbox.inmemory.In workingMemoryStore = workingMemoryStore, reflectionStore = reflectionStore, eventStore = eventStore, + onlineEventStore = onlineEventStore, llm = llm, systemPrompt = systemPrompt, systemPromptResolver = { buildRuntimeSystemPrompt() }, @@ -528,6 +547,7 @@ private val eventStore: MutableOutboxStore = pw.binom.agentik.outbox.inmemory.In workingMemoryStore = workingMemoryStore, reflectionStore = reflectionStore, eventStore = eventStore, + onlineEventStore = onlineEventStore, llm = llm, systemPrompt = systemPrompt, systemPromptResolver = { buildRuntimeSystemPrompt() }, diff --git a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ConversationEvents.kt b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ConversationEvents.kt index d0b9adc..6a1f3c9 100644 --- a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ConversationEvents.kt +++ b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ConversationEvents.kt @@ -5,10 +5,18 @@ import kotlinx.coroutines.flow.Flow import kotlinx.coroutines.flow.map import pw.binom.agentik.outbox.CommonEvent import pw.binom.agentik.outbox.Event +import pw.binom.agentik.outbox.MutableOnlineOutbox import pw.binom.agentik.outbox.MutableOutboxStore +import pw.binom.agentik.outbox.OnlineEvent +/** + * Фасад эмиссии и чтения событий одного диалога. Разводит два канала: + * - durable ([Event]) → [globalEventStore] (`:outbox`), с catchup по `after`; + * - online ([OnlineEvent]) → [onlineStore], live-only (без catchup). + */ internal class ConversationEvents( private val globalEventStore: MutableOutboxStore, + private val onlineStore: MutableOnlineOutbox, private val conversationId: String, ) { fun tryEmit(event: Event): Boolean { @@ -27,4 +35,11 @@ internal class ConversationEvents( fun events(after: kotlin.time.Instant?): Flow = globalEventStore.conversationEvents(after = after, conversationId = conversationId) .map { it.event } + + /** Best-effort эмиссия онлайн-события — без блокировки продюсера и без хранения. */ + fun tryEmitOnline(event: OnlineEvent): Boolean = + onlineStore.tryAppendOnline(conversationId, event) + + /** Live-поток онлайн-событий диалога (без catchup — см. KDoc [pw.binom.agentik.outbox.OnlineOutbox]). */ + fun onlineEvents(): Flow = onlineStore.onlineEvents(conversationId) } diff --git a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ConversationLoop.kt b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ConversationLoop.kt index 5c84f4e..1603156 100644 --- a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ConversationLoop.kt +++ b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ConversationLoop.kt @@ -22,6 +22,8 @@ import pw.binom.agentik.memory.MemoryStore import pw.binom.agentik.proto.Content as ProtoContent import pw.binom.agentik.proto.Conversation as ProtoConversation import pw.binom.agentik.outbox.Event as ProtoEvent +import pw.binom.agentik.outbox.OnlineEvent +import pw.binom.agentik.outbox.MutableOnlineOutbox import pw.binom.agentik.reflection.ReflectionStore import pw.binom.agentik.proto.Message as ProtoMessage import pw.binom.agentik.proto.MessageContext as ProtoMessageContext @@ -53,6 +55,11 @@ class ConversationLoop( private val workingMemoryStore: ContextStore, private val reflectionStore: ReflectionStore?, private val eventStore: pw.binom.agentik.outbox.MutableOutboxStore, + /** + * Live-канал онлайн-событий (дельты ответа). Не сохраняется; подписка + * возможна только «онлайн». Durable-события по-прежнему в [eventStore]. + */ + private val onlineEventStore: MutableOnlineOutbox, private val llm: LiteLlm, /** * Базовый системный промпт, который задаётся беседе при создании. @@ -102,6 +109,7 @@ class ConversationLoop( private val events = ConversationEvents( globalEventStore = eventStore, + onlineStore = onlineEventStore, conversationId = state.id, ) @@ -259,8 +267,8 @@ class ConversationLoop( compactor.compactPreTurnIfNeeded() } - emitEvent(ProtoEvent.StartReasoning(date = turnStarted)) - emitEvent(ProtoEvent.StartResponse(date = now(), responseType = ProtoEvent.ResponseType.TEXT)) + emitOnline(OnlineEvent.StartReasoning(date = turnStarted)) + emitOnline(OnlineEvent.StartResponse(date = now(), responseType = OnlineEvent.ResponseType.TEXT)) val parts = userRecord.content.mapNotNull { c -> when (c) { @@ -327,7 +335,7 @@ class ConversationLoop( lc.sendStreamContents(pendingParts).collect { delta -> if (delta.text.isNotEmpty()) { reply.append(delta.text) - emitEvent(ProtoEvent.AppendText(date = now(), body = delta.text)) + emitOnline(OnlineEvent.AppendText(date = now(), body = delta.text)) } if (delta.toolCalls.isNotEmpty()) { collectedCalls.addAll(delta.toolCalls) @@ -365,7 +373,7 @@ class ConversationLoop( } if (delta.text.isNotEmpty()) { reply.append(delta.text) - emitEvent(ProtoEvent.AppendText(date = now(), body = delta.text)) + emitOnline(OnlineEvent.AppendText(date = now(), body = delta.text)) } if (delta.toolCalls.isNotEmpty()) { nextCalls.addAll(delta.toolCalls) @@ -382,7 +390,7 @@ class ConversationLoop( lc.sendStreamContents(listOf(LiteContentPart.Text(" "))).collect { followUp -> if (followUp.text.isNotEmpty()) { reply.append(followUp.text) - emitEvent(ProtoEvent.AppendText(date = now(), body = followUp.text)) + emitOnline(OnlineEvent.AppendText(date = now(), body = followUp.text)) } if (followUp.toolCalls.isNotEmpty()) { collectedPostTool.addAll(followUp.toolCalls) @@ -491,6 +499,10 @@ class ConversationLoop( events.tryEmit(event) } + private fun emitOnline(event: OnlineEvent) { + events.tryEmitOnline(event) + } + private suspend fun failTurn(message: String, code: String? = null) { val ts = now() if (!state.isTemporal) { diff --git a/standalone/src/commonTest/kotlin/pw/binom/agentik/standalone/agent/ChatAgentTest.kt b/standalone/src/commonTest/kotlin/pw/binom/agentik/standalone/agent/ChatAgentTest.kt index e802ae1..54516a3 100644 --- a/standalone/src/commonTest/kotlin/pw/binom/agentik/standalone/agent/ChatAgentTest.kt +++ b/standalone/src/commonTest/kotlin/pw/binom/agentik/standalone/agent/ChatAgentTest.kt @@ -8,6 +8,7 @@ import kotlinx.coroutines.flow.toList import kotlinx.coroutines.launch import kotlinx.coroutines.test.runTest import pw.binom.agentik.outbox.AgentEvent +import pw.binom.agentik.outbox.OnlineEvent import pw.binom.agentik.proto.Content import pw.binom.agentik.outbox.Event as ProtoEvent import pw.binom.agentik.skill.mining.SkillReadTool @@ -362,33 +363,42 @@ class ChatAgentTest { } @Test - fun `Working event is the first event of a turn`() = runTest { - // Working-маркер обязан прийти самым первым событием хода, до - // StartReasoning/StartResponse/AppendText/End. Это позволяет UI - // показать спиннер сразу же при отправке, не дожидаясь первого - // токена от LLM. + fun `Working is first durable event, streaming goes to online channel`() = runTest { + // Working-маркер обязан прийти самым первым durable-событием хода, до + // End. Это позволяет UI показать спиннер сразу же при отправке, не + // дожидаясь первого токена от LLM. + // + // Стриминг (StartReasoning / StartResponse / AppendText) — теперь + // онлайн-события (live-only, не сохраняются) и приходят из + // [OnlineOutbox], а НЕ из durable-потока. val agent = newAgent() fakeLlm.reply = "ok" val conv = agent.createConversation(temp = false) - val events = mutableListOf() - val job = launch(start = kotlinx.coroutines.CoroutineStart.UNDISPATCHED) { - agent.outbox.conversationEvents(Instant.DISTANT_PAST, conv.id).collect { events.add(it.event) } + val durable = mutableListOf() + val online = mutableListOf() + val durableJob = launch(start = kotlinx.coroutines.CoroutineStart.UNDISPATCHED) { + agent.outbox.conversationEvents(Instant.DISTANT_PAST, conv.id).collect { durable.add(it.event) } + } + val onlineJob = launch(start = kotlinx.coroutines.CoroutineStart.UNDISPATCHED) { + agent.onlineOutbox.onlineEvents(conv.id).collect { online.add(it) } } conv.send(listOf(Content.Text("hi"))) delay(50) - job.cancel() + durableJob.cancel() + onlineJob.cancel() - // Working — первый event хода (индекс 0). - assertTrue(events.isNotEmpty(), "no events captured: $events") - val first = events.first() + // Working — первый durable event хода (индекс 0). + assertTrue(durable.isNotEmpty(), "no events captured: $durable") + val first = durable.first() assertIs(first) - assertTrue(events.any { it is ProtoEvent.StartReasoning }, "no StartReasoning: $events") - assertTrue(events.any { it is ProtoEvent.StartResponse }, "no StartResponse: $events") - assertTrue(events.any { it is ProtoEvent.End }, "no End: $events") + assertTrue(durable.any { it is ProtoEvent.End }, "no End: $durable") + // Ранее стриминговые маркеры — теперь в онлайн-канале. + assertTrue(online.any { it is OnlineEvent.StartReasoning }, "no StartReasoning: $online") + assertTrue(online.any { it is OnlineEvent.StartResponse }, "no StartResponse: $online") // Working не дублируется (emit'ится один раз на send()). - assertEquals(1, events.count { it is ProtoEvent.Working }, "Working emitted >1 times: $events") + assertEquals(1, durable.count { it is ProtoEvent.Working }, "Working emitted >1 times: $durable") } @Test