Изменения вводят live-канал OnlineOutbox и фоновую синхронизацию диалогов.
This commit is contained in:
@@ -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")
|
||||
}
|
||||
|
||||
@@ -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} байт]")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
|
||||
|
||||
@@ -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<OnlineEvent> = 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).
|
||||
}
|
||||
}
|
||||
@@ -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<String>()
|
||||
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
|
||||
|
||||
@@ -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` в истории
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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<OnlineEvent>
|
||||
|
||||
/**
|
||||
* Подписка на live-поток онлайн-событий диалога [conversationId].
|
||||
* События, эмитнутые до подписки, не приходят.
|
||||
*/
|
||||
fun onlineEvents(conversationId: String): Flow<OnlineEvent>
|
||||
|
||||
/** Освобождает ресурсы. Idempotent. */
|
||||
override fun close()
|
||||
}
|
||||
@@ -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).
|
||||
|
||||
+65
@@ -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<Envelope>(
|
||||
replay = 0,
|
||||
extraBufferCapacity = liveBufferCapacity,
|
||||
onBufferOverflow = BufferOverflow.DROP_OLDEST,
|
||||
)
|
||||
|
||||
init {
|
||||
require(liveBufferCapacity > 0) {
|
||||
"liveBufferCapacity must be > 0, got $liveBufferCapacity"
|
||||
}
|
||||
}
|
||||
|
||||
override fun onlineEvents(): Flow<OnlineEvent> =
|
||||
liveFlow.map { it.event }
|
||||
|
||||
override fun onlineEvents(conversationId: String): Flow<OnlineEvent> =
|
||||
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
|
||||
}
|
||||
}
|
||||
+5
-5
@@ -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<CommonEvent.Agent>()
|
||||
|
||||
@@ -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).
|
||||
*
|
||||
|
||||
@@ -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)`.
|
||||
|
||||
@@ -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<CommonEvent.Agent>;
|
||||
|
||||
@@ -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<pw.binom.agentik.journal.ConversationRecord>()
|
||||
@@ -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<pw.binom.agentik.journal.ConversationRecord>()
|
||||
|
||||
@@ -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<pw.binom.agentik.journal.ConversationRecord>()
|
||||
|
||||
@@ -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 =
|
||||
|
||||
@@ -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<pw.binom.agentik.journal.ConversationRecord>()
|
||||
|
||||
@@ -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<OnlineEvent>()
|
||||
override fun close() {}
|
||||
}
|
||||
@@ -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<Unit>()
|
||||
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(
|
||||
|
||||
@@ -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() },
|
||||
|
||||
+15
@@ -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<Event> =
|
||||
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<OnlineEvent> = onlineStore.onlineEvents(conversationId)
|
||||
}
|
||||
|
||||
+17
-5
@@ -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) {
|
||||
|
||||
+26
-16
@@ -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<ProtoEvent>()
|
||||
val job = launch(start = kotlinx.coroutines.CoroutineStart.UNDISPATCHED) {
|
||||
agent.outbox.conversationEvents(Instant.DISTANT_PAST, conv.id).collect { events.add(it.event) }
|
||||
val durable = mutableListOf<ProtoEvent>()
|
||||
val online = mutableListOf<OnlineEvent>()
|
||||
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<ProtoEvent.Working>(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
|
||||
|
||||
Reference in New Issue
Block a user