Изменение выносит типы контента в :content-api и разделяет потоки событий.

This commit is contained in:
2026-09-30 12:41:10 +03:00
parent 26f94ea83d
commit 6b7b5ee57b
11 changed files with 95 additions and 63 deletions
@@ -3,7 +3,7 @@ package pw.binom.agentik.desktop.model
import pw.binom.agentik.desktop.persistence.MetaRow
import pw.binom.agentik.desktop.session.LiveBlock
import pw.binom.agentik.desktop.session.LiveTurn
import pw.binom.agentik.journal.Content
import pw.binom.agentik.content.Content
import pw.binom.agentik.journal.ConversationRecord
import pw.binom.agentik.journal.MessageRecord
import kotlin.time.Instant
@@ -20,7 +20,7 @@ import kotlin.time.Instant
* `journal.count(convId, after = lastSeen)` и здесь не лежит.
* - **`preview`** — последний фрагмент текста хода, денормализованный
* в строку. Сервер такого поля не отдаёт (STORAGE.md §6.1), клиент
* поддерживает сам на `OnlineEvent.AppendText` / `Event.End`.
* поддерживает сам на `OnlineEvent.AppendText` / `Event.AssistantMessage`.
*
* **Почему НЕ наследуется от `KsqliteJournalStore`:** разные таблицы
* (`conversation_meta` + `chat_group` против `message`) и разные контракты
@@ -67,7 +67,7 @@ internal class AgentSyncEngine(
outbox.conversationEvents(after = null, conversationId = null).collect { ev ->
if (ev.conversationId in openConversationIds) return@collect
when (ev.event) {
is Event.End, is Event.Interrupted -> {
is Event.AssistantMessage, is Event.Interrupted, is Event.Error -> {
runCatching { syncConversation(ev.conversationId) }
onChanged()
}
@@ -18,7 +18,7 @@ import pw.binom.agentik.journal.MutableJournalStore
import pw.binom.agentik.outbox.Event
import pw.binom.agentik.outbox.OnlineOutbox
import pw.binom.agentik.outbox.OutboxStore
import pw.binom.agentik.proto.Content
import pw.binom.agentik.content.Content
import pw.binom.agentik.proto.Conversation
import kotlin.time.Clock
import kotlin.time.Instant
@@ -33,8 +33,8 @@ import kotlin.time.Instant
* Watermark — `conversation_meta.synced_at`, чтобы не вставлять дубли
* (`message.id` — PRIMARY KEY, повторный INSERT бросает).
* - **Live-ход** — два SSE-потока: durable
* `OutboxStore.conversationEvents` (`Working → … → End`) и онлайн
* `OnlineOutbox.onlineEvents` (стриминг ответа). Оба копятся в
* `OutboxStore.conversationEvents` (`UserMessage → … → AssistantMessage`) и
* онлайн `OnlineOutbox.onlineEvents` (Working/End + стриминг ответа). Оба копятся в
* [LiveTurn] и рисуются поверх истории, пока ход не закрыт.
* - **Наши поля** — [ConversationMetaRepository]: `lastSeen` (бейдж),
* `preview` (строка списка).
@@ -42,7 +42,7 @@ import kotlin.time.Instant
* ## Почему подписка на SSE берётся `after = synced_at`
*
* `synced_at` — `createdAt` последнего сообщения, уже лежащего в кэше.
* У завершённого хода все события (`Working`/`AppendText`/`End`) имеют
* У завершённого хода все события (`UserMessage`/`AppendText`/`AssistantMessage`) имеют
* `date`, не больший `createdAt` итогового сообщения, — значит, они
* отфильтруются и повторно не проиграются. А у **идущего** хода события
* свежее — они и попадут в поток. Так «догоняем» уже начатый ход, не
@@ -51,7 +51,7 @@ import kotlin.time.Instant
*
* ## Поток событий vs `sync()` на `End`
*
* На `End`/`Interrupted` ход закрыт: подтягиваем авторитетную историю
* На `AssistantMessage`/`Interrupted`/`Error` ход закрыт: подтягиваем авторитетную историю
* (`sync()` — инкрементальный бэкфилл + перечитка кэша) и гасим live-блоки,
* чтобы не было дублирования «live + history».
*/
@@ -173,8 +173,10 @@ class ChatSession(
.map { it.event }
.collect { event ->
liveTurn.apply(event)
if (event is Event.End || event is Event.Interrupted) {
// Ход закрыт: авторитетная история важнее live-блоков.
// Ход закрыт durable-событием: успех — готовое
// AssistantMessage, прерывание — Interrupted, провал —
// Error. Авторитетная история важнее live-блоков.
if (event is Event.AssistantMessage || event is Event.Interrupted || event is Event.Error) {
runCatching { sync() }
liveTurn.reset()
}
@@ -11,7 +11,7 @@ enum class TurnPhase {
/** Хода нет — стор чист. */
IDLE,
/** `send()` ушёл, ответа ещё нет (`Event.Working`). Спиннер «агент думает». */
/** `send()` ушёл, ответа ещё нет (`OnlineEvent.Working`). Спиннер «агент думает». */
WORKING,
/** Идут рассуждения (`StartReasoning` + `AppendText` в reasoning-блок). Live-канал. */
@@ -20,7 +20,7 @@ enum class TurnPhase {
/** Стримится ответ (`StartResponse` + `AppendText`/`AppendImage`/tool'ы). Live-канал. */
RESPONDING,
/** Ход закрыт нормально (`End`). Live-блоки пора заменить историей из кэша. */
/** Ход закрыт нормально (`OnlineEvent.End` / durable `Event.AssistantMessage`). Live-блоки пора заменить историей из кэша. */
DONE,
/** Ход прерван пользователем (`Interrupted`). Частичный ответ НЕ сохранён. */
@@ -133,7 +133,7 @@ class LiveTurn {
/**
* Начинает новый ход: чистит блоки и переводит в [TurnPhase.WORKING].
* Вызывать перед `send()` (или на `Event.Working`) — иначе блоки
* Вызывать перед `send()` (или на `OnlineEvent.Working`) — иначе блоки
* предыдущего хода смешаются с новыми.
*/
fun start() {
@@ -159,9 +159,17 @@ class LiveTurn {
*/
fun apply(event: Event) {
when (event) {
is Event.Working -> {
// Синхронный маркер старта хода: новый ход — с чистого листа.
start()
is Event.UserMessage -> {
// Durable-маркер старта хода. Если онлайн-`Working` уже
// открыл ход, повторно не чистим блоки (защита от гонки
// двух SSE-потоков).
if (!phase.isStreaming) start()
}
is Event.AssistantMessage -> {
// Durable-терминатор успешного хода: целый ответ уже в журнале.
phase = TurnPhase.DONE
fire()
}
is Event.ToolCall -> {
@@ -203,11 +211,6 @@ class LiveTurn {
fire()
}
is Event.End -> {
phase = TurnPhase.DONE
fire()
}
is Event.Interrupted -> {
phase = TurnPhase.INTERRUPTED
fire()
@@ -228,10 +231,20 @@ class LiveTurn {
* Применяет одно **онлайн** событие ([OnlineEvent]) — live-стриминг
* ответа. Онлайн-события нигде не сохраняются: при обрыве соединения
* потерянный фрагмент невосстановим (целый результат придёт
* durable-[Event.End] и/или ляжет в журнал).
* durable-[Event.AssistantMessage] и/или ляжет в журнал).
*/
fun applyOnline(event: OnlineEvent) {
when (event) {
is OnlineEvent.Working -> {
// Live-маркер старта хода (эмитится синхронно из send()).
start()
}
is OnlineEvent.End -> {
phase = TurnPhase.DONE
fire()
}
is OnlineEvent.StartReasoning -> {
phase = TurnPhase.REASONING
_blocks.add(LiveBlock.Reasoning(nextKey("reasoning")))
@@ -1,6 +1,6 @@
package pw.binom.agentik.desktop.model
import pw.binom.agentik.journal.Content
import pw.binom.agentik.content.Content
import pw.binom.agentik.journal.MessageRecord
import kotlin.test.Test
import kotlin.test.assertEquals
@@ -12,10 +12,10 @@ import pw.binom.agentik.outbox.inmemory.InMemoryOnlineOutbox
import pw.binom.agentik.outbox.inmemory.InMemoryOutboxStore
import pw.binom.agentik.proto.Agent
import pw.binom.agentik.proto.AgentInfo
import pw.binom.agentik.proto.Content
import pw.binom.agentik.content.Content
import pw.binom.agentik.proto.Conversation
import pw.binom.agentik.proto.Message
import pw.binom.agentik.proto.MessageContext
import pw.binom.agentik.content.MessageContext
import pw.binom.agentik.server.agentikAgent
import kotlin.test.AfterTest
import kotlin.test.Test
@@ -26,10 +26,10 @@ import pw.binom.agentik.outbox.inmemory.InMemoryOnlineOutbox
import pw.binom.agentik.outbox.inmemory.InMemoryOutboxStore
import pw.binom.agentik.proto.Agent
import pw.binom.agentik.proto.AgentInfo
import pw.binom.agentik.proto.Content
import pw.binom.agentik.content.Content
import pw.binom.agentik.proto.Conversation
import pw.binom.agentik.proto.Message
import pw.binom.agentik.proto.MessageContext
import pw.binom.agentik.content.MessageContext
import pw.binom.agentik.server.agentikAgent
import kotlin.test.AfterTest
import kotlin.test.Test
@@ -283,10 +283,15 @@ private class EchoConversation(
override suspend fun send(content: List<Content>, context: MessageContext?) {
val now = Clock.System.now()
val userId = nextMsgId(id, "user")
journal.append(
MessageRecord.UserMessage(nextMsgId(id, "user"), id, content.toJournal(), now, context = null),
MessageRecord.UserMessage(userId, id, content.toJournal(), now, context = null),
)
emit(now, Event.Working(now))
emit(
now,
Event.UserMessage(date = now, id = userId, content = content.toJournal()),
)
onlineOutbox.tryAppendOnline(id, OnlineEvent.Working(now))
if (interruptible) {
// Имитация долгого LLM-цикла: крутимся, пока не придёт interrupt().
@@ -300,15 +305,17 @@ private class EchoConversation(
onlineOutbox.tryAppendOnline(id, OnlineEvent.AppendText(now, reply))
val endAt = Clock.System.now()
val assistantId = nextMsgId(id, "assistant")
journal.append(
MessageRecord.AssistantMessage(
nextMsgId(id, "assistant"),
assistantId,
id,
listOf(pw.binom.agentik.journal.Content.Text(reply)),
listOf(pw.binom.agentik.content.Content.Text(reply)),
endAt,
),
)
emit(endAt, Event.End(endAt))
onlineOutbox.tryAppendOnline(id, OnlineEvent.End(endAt))
emit(endAt, Event.AssistantMessage(date = endAt, id = assistantId, content = listOf(Content.Text(reply))))
updated = endAt
store.touch(id, endAt)
@@ -323,9 +330,9 @@ private class EchoConversation(
}
/** `:proto.Content` → `:journal.Content` (в проде это делает storage-impl). */
private fun List<Content>.toJournal(): List<pw.binom.agentik.journal.Content> = map { c ->
private fun List<Content>.toJournal(): List<pw.binom.agentik.content.Content> = map { c ->
when (c) {
is Content.Text -> pw.binom.agentik.journal.Content.Text(c.body)
is Content.Image -> pw.binom.agentik.journal.Content.Image(c.data, c.mime)
is Content.Text -> pw.binom.agentik.content.Content.Text(c.body)
is Content.Image -> pw.binom.agentik.content.Content.Image(c.data, c.mime)
}
}
@@ -13,7 +13,7 @@ import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.withTimeoutOrNull
import org.junit.jupiter.api.Test
import pw.binom.agentik.desktop.persistence.ConversationMetaRepository
import pw.binom.agentik.journal.Content as JournalContent
import pw.binom.agentik.content.Content as JournalContent
import pw.binom.agentik.journal.ConversationRecord
import pw.binom.agentik.journal.JournalStore
import pw.binom.agentik.journal.MessageRecord
@@ -160,7 +160,9 @@ class AgentSyncEngineTest {
f.outbox.awaitSubscriber()
f.remote.records = listOf(user("u1", 1, "привет"), assistant("a1", 2, "пока"))
f.outbox.emit(Event.End(date = at(3)))
f.outbox.emit(
Event.AssistantMessage(date = at(3), id = "a1", content = listOf<JournalContent>(JournalContent.Text("пока")))
)
await { f.cache.count(convId) == 2L }
await { f.meta.meta(convId)?.preview == "пока" }
@@ -178,7 +180,9 @@ class AgentSyncEngineTest {
f.outbox.awaitSubscriber()
f.remote.records = listOf(user("u1", 1, "привет"))
f.outbox.emit(Event.End(date = at(3)))
f.outbox.emit(
Event.AssistantMessage(date = at(3), id = "a1", content = listOf<JournalContent>(JournalContent.Text("привет")))
)
delay(150)
assertEquals(0L, f.cache.count(convId))
@@ -13,7 +13,7 @@ import kotlinx.coroutines.withTimeoutOrNull
import org.junit.jupiter.api.Test
import pw.binom.agentik.desktop.model.UiMessage
import pw.binom.agentik.desktop.persistence.ConversationMetaRepository
import pw.binom.agentik.journal.Content as JournalContent
import pw.binom.agentik.content.Content as JournalContent
import pw.binom.agentik.journal.JournalStore
import pw.binom.agentik.journal.MessageRecord
import pw.binom.agentik.journal.ksqlite.KsqliteJournalStore
@@ -22,10 +22,10 @@ import pw.binom.agentik.outbox.Event
import pw.binom.agentik.outbox.OnlineEvent
import pw.binom.agentik.outbox.OnlineOutbox
import pw.binom.agentik.outbox.OutboxStore
import pw.binom.agentik.proto.Content
import pw.binom.agentik.content.Content
import pw.binom.agentik.proto.Conversation
import pw.binom.agentik.proto.Message
import pw.binom.agentik.proto.MessageContext
import pw.binom.agentik.content.MessageContext
import pw.binom.db.ksqlite.SQLiteConnection
import kotlin.test.assertEquals
import kotlin.test.assertTrue
@@ -163,7 +163,7 @@ class ChatSessionTest {
f.outbox.awaitSubscriber()
f.online.awaitSubscriber()
f.outbox.emit(Event.Working(at(10)))
f.online.emit(OnlineEvent.Working(at(10)))
f.online.emit(OnlineEvent.StartResponse(at(11), OnlineEvent.ResponseType.TEXT))
f.online.emit(OnlineEvent.AppendText(at(12), "Hel"))
f.online.emit(OnlineEvent.AppendText(at(13), "lo"))
@@ -174,9 +174,11 @@ class ChatSessionTest {
assertEquals("Hello", (session.live.value.single() as UiMessage.Assistant)
.let { (it.content.single() as pw.binom.agentik.desktop.model.UiContent.Text).body })
// Сервер зафиксировал ход в журнале и прислал End.
// Сервер зафиксировал ход в журнале и прислал durable AssistantMessage.
f.remote.records = listOf(user("u1", 9, "вопрос"), assistant("a1", 14, "Hello"))
f.outbox.emit(Event.End(at(14)))
f.outbox.emit(
Event.AssistantMessage(at(14), id = "a1", content = listOf<JournalContent>(JournalContent.Text("Hello")))
)
await { session.live.value.isEmpty() && session.history.value.size == 2 }
assertTrue(session.live.value.isEmpty())
@@ -201,7 +203,9 @@ class ChatSessionTest {
// стоял ровно на 5 мс — отсюда отступ назад на 1 мс + отсечение
// известных id вместо опоры на время.
f.remote.records = listOf(user("u1", 5, "вопрос"), assistant("a1", 5, "ответ"))
f.outbox.emit(Event.End(at(5)))
f.outbox.emit(
Event.AssistantMessage(at(5), id = "a1", content = listOf<JournalContent>(JournalContent.Text("ответ")))
)
await { session.history.value.size == 2 }
// Порядок при равных timestamp'ах задаёт id (контракт SQL
@@ -305,6 +309,8 @@ private class FakeOutboxStore : OutboxStore {
private class FakeOnlineOutbox : OnlineOutbox {
private val shared = MutableSharedFlow<OnlineEvent>(replay = 256, extraBufferCapacity = 256)
override fun onlineEvents(): Flow<OnlineEvent> = shared.asSharedFlow()
override fun onlineEvents(conversationId: String): Flow<OnlineEvent> = shared.asSharedFlow()
override fun close() {}
@@ -28,14 +28,14 @@ class LiveTurnTest {
fun `full happy path accumulates reasoning and answer`() {
val turn = LiveTurn()
turn.feed(
Event.Working(at()),
OnlineEvent.Working(at()),
OnlineEvent.StartReasoning(at()),
OnlineEvent.AppendText(at(), "сначала "),
OnlineEvent.AppendText(at(), "подумаем"),
OnlineEvent.StartResponse(at(), OnlineEvent.ResponseType.TEXT),
OnlineEvent.AppendText(at(), "Привет, "),
OnlineEvent.AppendText(at(), "мир!"),
Event.End(at()),
OnlineEvent.End(at()),
)
assertEquals(TurnPhase.DONE, turn.phase)
@@ -52,11 +52,11 @@ class LiveTurnTest {
fun `tool call and result are stitched into one block`() {
val turn = LiveTurn()
turn.feed(
Event.Working(at()),
OnlineEvent.Working(at()),
OnlineEvent.StartResponse(at(), OnlineEvent.ResponseType.TEXT),
Event.ToolCall(at(), id = "t1", title = "grep", toolName = "grep", toolArgs = "{\"q\":\"x\"}"),
Event.ToolResult(at(), toolCallId = "t1", toolName = "grep", result = "found 3"),
Event.End(at()),
OnlineEvent.End(at()),
)
val tools = turn.blocks.filterIsInstance<LiveBlock.Tool>()
@@ -72,7 +72,7 @@ class LiveTurnTest {
fun `orphan tool result creates a block`() {
val turn = LiveTurn()
turn.feed(
Event.Working(at()),
OnlineEvent.Working(at()),
OnlineEvent.StartResponse(at(), OnlineEvent.ResponseType.TEXT),
Event.ToolResult(at(), toolCallId = "missing", toolName = null, result = "ok"),
)
@@ -87,7 +87,7 @@ class LiveTurnTest {
fun `tool failed marks the block`() {
val turn = LiveTurn()
turn.feed(
Event.Working(at()),
OnlineEvent.Working(at()),
Event.ToolCall(at(), id = "t9", title = null, toolName = "bash", toolArgs = "{}"),
Event.ToolFailed(at(), toolCallId = "t9", toolName = "bash", message = "boom", durationMs = 5),
)
@@ -100,7 +100,7 @@ class LiveTurnTest {
fun `error appends failure block and fails the turn`() {
val turn = LiveTurn()
turn.feed(
Event.Working(at()),
OnlineEvent.Working(at()),
OnlineEvent.StartResponse(at(), OnlineEvent.ResponseType.TEXT),
OnlineEvent.AppendText(at(), "partial"),
Event.Error(at(), message = "llm down", code = "E500"),
@@ -118,7 +118,7 @@ class LiveTurnTest {
fun `interrupted turn is terminal and keeps partial text`() {
val turn = LiveTurn()
turn.feed(
Event.Working(at()),
OnlineEvent.Working(at()),
OnlineEvent.StartResponse(at(), OnlineEvent.ResponseType.TEXT),
OnlineEvent.AppendText(at(), "не до конца"),
Event.Interrupted(at()),
@@ -133,7 +133,7 @@ class LiveTurnTest {
fun `text after tool call starts a new answer block`() {
val turn = LiveTurn()
turn.feed(
Event.Working(at()),
OnlineEvent.Working(at()),
OnlineEvent.StartResponse(at(), OnlineEvent.ResponseType.TEXT),
OnlineEvent.AppendText(at(), "до"),
Event.ToolCall(at(), id = "t1", title = null, toolName = "grep", toolArgs = "{}"),
@@ -152,10 +152,10 @@ class LiveTurnTest {
val bytes = byteArrayOf(1, 2, 3)
val turn = LiveTurn()
turn.feed(
Event.Working(at()),
OnlineEvent.Working(at()),
OnlineEvent.StartResponse(at(), OnlineEvent.ResponseType.IMAGE),
OnlineEvent.AppendImage(at(), body = bytes, mime = "image/png"),
Event.End(at()),
OnlineEvent.End(at()),
)
assertTrue(turn.blocks.filterIsInstance<LiveBlock.Answer>().isEmpty())
@@ -168,12 +168,12 @@ class LiveTurnTest {
fun `working resets the previous turn`() {
val turn = LiveTurn()
turn.feed(
Event.Working(at()),
OnlineEvent.Working(at()),
OnlineEvent.StartResponse(at(), OnlineEvent.ResponseType.TEXT),
OnlineEvent.AppendText(at(), "старое"),
Event.End(at()),
OnlineEvent.End(at()),
)
turn.feed(Event.Working(at()), OnlineEvent.AppendText(at(), "новое"))
turn.feed(OnlineEvent.Working(at()), OnlineEvent.AppendText(at(), "новое"))
assertEquals(1, turn.blocks.size)
assertEquals("новое", turn.answerText())
@@ -182,7 +182,7 @@ class LiveTurnTest {
@Test
fun `append without explicit start response still streams as answer`() {
val turn = LiveTurn()
turn.feed(Event.Working(at()), OnlineEvent.AppendText(at(), "без старта"))
turn.feed(OnlineEvent.Working(at()), OnlineEvent.AppendText(at(), "без старта"))
assertEquals(TurnPhase.RESPONDING, turn.phase)
assertEquals("без старта", turn.answerText())
@@ -193,7 +193,7 @@ class LiveTurnTest {
val turn = LiveTurn()
var count = 0
turn.onChange = { count++ }
turn.feed(Event.Working(at()), OnlineEvent.AppendText(at(), "x"))
turn.feed(OnlineEvent.Working(at()), OnlineEvent.AppendText(at(), "x"))
assertEquals(2, count)
}
@@ -202,13 +202,13 @@ class LiveTurnTest {
fun `preview text collapses whitespace and truncates`() {
val turn = LiveTurn()
turn.feed(
Event.Working(at()),
OnlineEvent.Working(at()),
OnlineEvent.AppendText(at(), "Первая строка\n\nвторая строка "),
)
assertEquals("Первая строка вторая строка", turn.previewText())
val long = LiveTurn()
long.feed(Event.Working(at()), OnlineEvent.AppendText(at(), "слово ".repeat(60)))
long.feed(OnlineEvent.Working(at()), OnlineEvent.AppendText(at(), "слово ".repeat(60)))
val preview = long.previewText(maxLength = 30)
assertTrue(preview.length <= 31)
assertTrue(preview.endsWith("…"))