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

This commit is contained in:
2026-09-30 12:41:10 +03:00
parent bd7e079780
commit 059c23e3ab
60 changed files with 613 additions and 400 deletions
@@ -13,7 +13,7 @@ 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.content.Content
import pw.binom.agentik.proto.Conversation
import java.util.concurrent.ConcurrentHashMap
@@ -31,7 +31,7 @@ private val log = KotlinLogging.logger {}
* Ответ A2A = склеенные [OnlineEvent.AppendText] нашего хода. Подписку на онлайн-поток
* ([pw.binom.agentik.outbox.OnlineOutbox]) открываем ДО [Conversation.send] (live-only,
* без catchup — события начала хода иначе можно упустить), завершение хода ждём
* по durable-событиям [Event.End] / [Event.Interrupted] / [Event.Error].
* по онлайн [OnlineEvent.End] и durable [Event.AssistantMessage] / [Event.Interrupted] / [Event.Error].
*
* Ограничение v1: tool-события и картинки в A2A-ответ не транслируются;
* при нескольких ходов в очереди за контекстом текст предыдущего хода
@@ -53,14 +53,20 @@ class A2aBridge(private val agent: Agent) : AgentHandler {
// Онлайн-поток — дельты ответа (live-only, без catchup).
val onlineJob = async {
agent.onlineOutbox.onlineEvents(conv.id).collect { e ->
if (e is OnlineEvent.AppendText) reply.append(e.body)
when (e) {
is OnlineEvent.AppendText -> reply.append(e.body)
is OnlineEvent.End -> turnDone.complete(Unit)
else -> {}
}
}
}
// Durable-поток — терминатор хода (catchup + live).
// Durable-поток — терминатор хода (catchup + live): целый ответ
// приходит [Event.AssistantMessage], обрыв — [Event.Interrupted],
// провал — [Event.Error].
val turnJob = async {
agent.outbox.conversationEvents(since, conv.id).collect { ce ->
when (val e = ce.event) {
is Event.End, is Event.Interrupted -> turnDone.complete(Unit)
is Event.AssistantMessage, is Event.Interrupted -> turnDone.complete(Unit)
is Event.Error ->
turnDone.completeExceptionally(
IllegalStateException("agent turn failed: ${e.message}")
@@ -137,5 +137,5 @@ internal suspend fun recentTurns(workingMemoryStore: ContextStore, conversationI
}
/** Текстовое содержимое записей working memory (Text-контент, без картинок). */
internal fun List<pw.binom.agentik.journal.Content>.text(): String =
filterIsInstance<pw.binom.agentik.journal.Content.Text>().joinToString("\n") { it.body }
internal fun List<pw.binom.agentik.content.Content>.text(): String =
filterIsInstance<pw.binom.agentik.content.Content.Text>().joinToString("\n") { it.body }
@@ -380,10 +380,10 @@ private val onlineEventStore: MutableOnlineOutbox = pw.binom.agentik.outbox.inme
for (entry in filtered.takeLast(limit * 2)) {
val text = when (entry) {
is pw.binom.agentik.context.WorkingMemoryEntry.User ->
entry.content.filterIsInstance<pw.binom.agentik.journal.Content.Text>()
entry.content.filterIsInstance<pw.binom.agentik.content.Content.Text>()
.joinToString("\n") { it.body }
is pw.binom.agentik.context.WorkingMemoryEntry.Assistant ->
entry.content.filterIsInstance<pw.binom.agentik.journal.Content.Text>()
entry.content.filterIsInstance<pw.binom.agentik.content.Content.Text>()
.joinToString("\n") { it.body }
else -> continue
}
@@ -6,7 +6,7 @@ import pw.binom.agentik.memory.ConversationTurn
import pw.binom.agentik.memory.MemoryReviewer
import pw.binom.agentik.memory.MemoryStore
import pw.binom.agentik.standalone.agent.memory.materializeReviewNote
import pw.binom.agentik.journal.Content
import pw.binom.agentik.content.Content
import pw.binom.agentik.context.WorkingMemoryEntry
import pw.binom.agentik.context.WorkingMemoryRow
import pw.binom.agentik.context.ContextStore
@@ -3,8 +3,8 @@ package pw.binom.agentik.standalone.agent
import mu.KotlinLogging
import pw.binom.agentik.memory.MemoryPrefetcher
import pw.binom.litert.LiteContentPart
import pw.binom.agentik.journal.MessageContext
import pw.binom.agentik.journal.MessageOrigin
import pw.binom.agentik.content.MessageContext
import pw.binom.agentik.content.MessageOrigin
internal class ContextBuilder(
private val memoryPrefetcher: MemoryPrefetcher?,
@@ -19,22 +19,22 @@ import mu.KotlinLogging
import pw.binom.agentik.memory.MemoryPrefetcher
import pw.binom.agentik.memory.MemoryReviewer
import pw.binom.agentik.memory.MemoryStore
import pw.binom.agentik.proto.Content as ProtoContent
import pw.binom.agentik.content.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
import pw.binom.agentik.journal.Content
import pw.binom.agentik.content.MessageContext as ProtoMessageContext
import pw.binom.agentik.content.Content
import pw.binom.agentik.journal.ConversationRecord
import pw.binom.agentik.journal.MutableConversationStore
import pw.binom.agentik.journal.MessageContext
import pw.binom.agentik.journal.MessageOrigin
import pw.binom.agentik.content.MessageContext
import pw.binom.agentik.content.MessageOrigin
import pw.binom.agentik.journal.MessageRecord
import pw.binom.agentik.journal.MutableJournalStore as MutableJournalStore
import pw.binom.agentik.journal.TurnTokens
import pw.binom.agentik.content.TurnTokens
import pw.binom.agentik.context.WorkingMemoryEntry
import pw.binom.agentik.context.ContextStore
import pw.binom.agentik.toolsets.ToolsetDispatchPolicy
@@ -179,11 +179,11 @@ class ConversationLoop(
check(!state.isClosed) { "Conversation closed: $id" }
val turnStarted = now()
// Working-маркер — самый первый event хода. Эмитим синхронно через
// tryEmit (events.tryEmit → outbox.append, без сетевого I/O), чтобы
// клиент увидел «агент работает» ещё до turnLock.withLock { launch }
// и до первого токена от LLM. Терминатор — End/Interrupted/Error
// (см. KDoc Event.Working).
emitEvent(pw.binom.agentik.outbox.Event.Working(date = turnStarted))
// tryEmit (events.tryEmitOnline → без сетевого I/O), чтобы клиент
// увидел «агент работает» ещё до turnLock.withLock { launch } и до
// первого токена от LLM. Working/End — онлайн-маркеры (live-only),
// терминатор durable-части — AssistantMessage/Interrupted/Error.
emitOnline(OnlineEvent.Working(date = turnStarted))
val userMessageId = newId("msg")
val storageContext = context?.toStorage()
@@ -206,6 +206,16 @@ class ConversationLoop(
),
now = turnStarted,
)
// Durable-событие user-сообщения: позволяет восстановить историю
// по курсору outbox без отдельного запроса в journal.
emitEvent(
ProtoEvent.UserMessage(
date = turnStarted,
id = userMessageId,
content = userRecord.content,
context = storageContext,
)
)
}
turnLock.withLock {
@@ -448,6 +458,17 @@ class ConversationLoop(
)
messageStore.append(assistantRecord)
// Durable-событие готового ответа агента.
emitEvent(
ProtoEvent.AssistantMessage(
date = assistantAt,
id = assistantId,
content = assistantContent,
reasoning = null,
tokens = turnTokens,
)
)
workingMemoryStore.append(
conversationId = id,
entry = WorkingMemoryEntry.Assistant(
@@ -489,7 +510,7 @@ class ConversationLoop(
if (wasInterrupted || interrupted.get()) {
emitEvent(ProtoEvent.Interrupted(date = now()))
}
emitEvent(ProtoEvent.End(date = now()))
emitOnline(OnlineEvent.End(date = now()))
interrupted.set(false)
}
@@ -541,9 +562,9 @@ internal fun ProtoContent.toStorage(): Content = when (this) {
internal fun ProtoMessageContext.toStorage(): MessageContext = MessageContext(
origin = when (origin) {
pw.binom.agentik.proto.MessageOrigin.USER -> MessageOrigin.USER
pw.binom.agentik.proto.MessageOrigin.SYSTEM -> MessageOrigin.SYSTEM
pw.binom.agentik.proto.MessageOrigin.EVENT -> MessageOrigin.EVENT
pw.binom.agentik.content.MessageOrigin.USER -> MessageOrigin.USER
pw.binom.agentik.content.MessageOrigin.SYSTEM -> MessageOrigin.SYSTEM
pw.binom.agentik.content.MessageOrigin.EVENT -> MessageOrigin.EVENT
},
description = description,
sourceId = sourceId,
@@ -552,9 +573,9 @@ internal fun ProtoMessageContext.toStorage(): MessageContext = MessageContext(
internal fun MessageContext.toProto(): ProtoMessageContext {
val protoOrigin = when (origin) {
MessageOrigin.USER -> pw.binom.agentik.proto.MessageOrigin.USER
MessageOrigin.SYSTEM -> pw.binom.agentik.proto.MessageOrigin.SYSTEM
MessageOrigin.EVENT -> pw.binom.agentik.proto.MessageOrigin.EVENT
MessageOrigin.USER -> pw.binom.agentik.content.MessageOrigin.USER
MessageOrigin.SYSTEM -> pw.binom.agentik.content.MessageOrigin.SYSTEM
MessageOrigin.EVENT -> pw.binom.agentik.content.MessageOrigin.EVENT
}
return ProtoMessageContext(
origin = protoOrigin,
@@ -9,7 +9,7 @@ import kotlinx.coroutines.flow.onEach
import kotlinx.coroutines.launch
import mu.KotlinLogging
import pw.binom.agentik.memory.ConversationTurn
import pw.binom.agentik.journal.Content
import pw.binom.agentik.content.Content
import pw.binom.agentik.reflection.ReflectionStore
import pw.binom.agentik.context.WorkingMemoryEntry
import pw.binom.agentik.context.ContextStore
@@ -9,7 +9,7 @@ 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.content.Content
import pw.binom.agentik.outbox.Event as ProtoEvent
import pw.binom.agentik.skill.mining.SkillReadTool
import pw.binom.agentik.skills.SkillCatalog
@@ -186,7 +186,7 @@ class ChatAgentTest {
pw.binom.agentik.journal.MessageRecord.UserMessage(
id = "m1",
conversationId = id,
content = listOf(pw.binom.agentik.journal.Content.Text("hi")),
content = listOf(pw.binom.agentik.content.Content.Text("hi")),
createdAt = Instant.fromEpochMilliseconds(1_700_000_000_000),
),
)
@@ -214,10 +214,10 @@ class ChatAgentTest {
val msgs = sqliteStores.messages.listFlow(conv.id, Instant.DISTANT_PAST).toList()
assertEquals(2, msgs.size)
assertEquals("hi", (msgs[0] as pw.binom.agentik.journal.MessageRecord.UserMessage).content.let {
(it[0] as pw.binom.agentik.journal.Content.Text).body
(it[0] as pw.binom.agentik.content.Content.Text).body
})
assertEquals("hello back", (msgs[1] as pw.binom.agentik.journal.MessageRecord.AssistantMessage).content.let {
(it[0] as pw.binom.agentik.journal.Content.Text).body
(it[0] as pw.binom.agentik.content.Content.Text).body
})
}
@@ -341,12 +341,18 @@ class ChatAgentTest {
val job = launch(start = kotlinx.coroutines.CoroutineStart.UNDISPATCHED) {
agent.outbox.conversationEvents(Instant.DISTANT_PAST, conv.id).collect { events.add(it.event) }
}
val online = mutableListOf<OnlineEvent>()
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()
onlineJob.cancel()
assertTrue(events.any { it is ProtoEvent.Error && it.message == "boom from llm" }, "events=$events")
assertTrue(events.any { it is ProtoEvent.End }, "events=$events")
assertTrue(events.any { it is ProtoEvent.UserMessage }, "events=$events")
assertTrue(online.any { it is OnlineEvent.End }, "online=$online")
val msgs = sqliteStores.messages.listFlow(conv.id, Instant.DISTANT_PAST).toList()
assertEquals(2, msgs.size)
@@ -363,14 +369,11 @@ class ChatAgentTest {
}
@Test
fun `Working is first durable event, streaming goes to online channel`() = runTest {
// Working-маркер обязан прийти самым первым durable-событием хода, до
// End. Это позволяет UI показать спиннер сразу же при отправке, не
// дожидаясь первого токена от LLM.
//
// Стриминг (StartReasoning / StartResponse / AppendText) — теперь
// онлайн-события (live-only, не сохраняются) и приходят из
// [OnlineOutbox], а НЕ из durable-потока.
fun `message events are durable, Working-and-End and streaming go online`() = runTest {
// Durable-поток несёт целые сообщения (UserMessage/AssistantMessage).
// Working/End — маркеры фаз (онлайн, live-only), как и стриминг
// (StartReasoning / StartResponse / AppendText), который приходит
// из [OnlineOutbox], а НЕ из durable-потока.
val agent = newAgent()
fakeLlm.reply = "ok"
val conv = agent.createConversation(temp = false)
@@ -388,17 +391,20 @@ class ChatAgentTest {
durableJob.cancel()
onlineJob.cancel()
// Working — первый durable event хода (индекс 0).
// durable: первое событие хода — user-сообщение; терминатор успешного
// хода — готовое assistant-сообщение.
assertTrue(durable.isNotEmpty(), "no events captured: $durable")
val first = durable.first()
assertIs<ProtoEvent.Working>(first)
assertTrue(durable.any { it is ProtoEvent.End }, "no End: $durable")
// Ранее стриминговые маркеры — теперь в онлайн-канале.
assertIs<ProtoEvent.UserMessage>(first)
assertTrue(durable.any { it is ProtoEvent.AssistantMessage }, "no AssistantMessage: $durable")
// Working/End и стриминговые маркеры — онлайн-канал (live-only).
assertTrue(online.any { it is OnlineEvent.Working }, "no Working: $online")
assertTrue(online.any { it is OnlineEvent.StartReasoning }, "no StartReasoning: $online")
assertTrue(online.any { it is OnlineEvent.StartResponse }, "no StartResponse: $online")
assertTrue(online.any { it is OnlineEvent.End }, "no End: $online")
// Working не дублируется (emit'ится один раз на send()).
assertEquals(1, durable.count { it is ProtoEvent.Working }, "Working emitted >1 times: $durable")
assertEquals(1, online.count { it is OnlineEvent.Working }, "Working emitted >1 times: $online")
}
@Test
@@ -417,6 +423,10 @@ class ChatAgentTest {
val eventsJob = launch(start = kotlinx.coroutines.CoroutineStart.UNDISPATCHED) {
agent.outbox.conversationEvents(Instant.DISTANT_PAST, conv.id).collect { events.add(it.event) }
}
val online = mutableListOf<OnlineEvent>()
val onlineJob = launch(start = kotlinx.coroutines.CoroutineStart.UNDISPATCHED) {
agent.onlineOutbox.onlineEvents(conv.id).collect { online.add(it) }
}
val sendJob = launch {
try {
@@ -430,6 +440,7 @@ class ChatAgentTest {
conv.interrupt()
sendJob.join()
eventsJob.cancel()
onlineJob.cancel()
// audit: только user (assistant не успел сгенериться)
val msgs = sqliteStores.messages.listFlow(conv.id, Instant.DISTANT_PAST).toList()
@@ -441,9 +452,9 @@ class ChatAgentTest {
assertEquals(1, wm.size)
assertTrue(wm[0].entry is WorkingMemoryEntry.User)
// events: должны включать Interrupted + End
// durable: включая Interrupted; End — онлайн-маркер.
assertTrue(events.any { it is ProtoEvent.Interrupted }, "events=$events")
assertTrue(events.any { it is ProtoEvent.End }, "events=$events")
assertTrue(online.any { it is OnlineEvent.End }, "online=$online")
}
@Test
@@ -504,10 +515,10 @@ class ChatAgentTest {
assertFalse(exchanges[0].wasCancelled, "tool реально выполнился, не был отменён")
assertTrue(exchanges[0].resultText.contains("echo"))
// events должны включать ToolCall + ToolResult. End — обязательно (turn завершился).
// events должны включать ToolCall + ToolResult. AssistantMessage — обязательно (turn завершился).
assertTrue(events.any { it is ProtoEvent.ToolCall }, "events=$events")
assertTrue(events.any { it is ProtoEvent.ToolResult }, "events=$events")
assertTrue(events.any { it is ProtoEvent.End }, "events=$events")
assertTrue(events.any { it is ProtoEvent.AssistantMessage }, "events=$events")
}
@Test
@@ -615,7 +626,7 @@ class ChatAgentTest {
val agent = newAgent(llm = toolLlm, tools = listOf(echoTool))
val conv = agent.createConversation(temp = false)
conv.send(listOf(pw.binom.agentik.proto.Content.Text("call the tool")))
conv.send(listOf(pw.binom.agentik.content.Content.Text("call the tool")))
// sendStreamContents вызывается дважды: первый раз с user-сообщением
// (LLM отвечает tool_call), второй раз — после addToolResult — для
@@ -12,7 +12,7 @@ import pw.binom.agentik.standalone.llm.LlmBackend
import pw.binom.agentik.standalone.llm.LlmConfig
import pw.binom.agentik.standalone.llm.OpenAiConfig
import pw.binom.agentik.standalone.persistence.SqliteStores
import pw.binom.agentik.proto.Content as ProtoContent
import pw.binom.agentik.content.Content as ProtoContent
import pw.binom.litert.LiteConversation
import pw.binom.litert.LiteConversationConfig
import pw.binom.litert.LiteLlm
@@ -1,9 +1,9 @@
package pw.binom.agentik.standalone.agent
import pw.binom.agentik.journal.MessageContext
import pw.binom.agentik.journal.MessageOrigin.EVENT
import pw.binom.agentik.journal.MessageOrigin.SYSTEM
import pw.binom.agentik.journal.MessageOrigin.USER
import pw.binom.agentik.content.MessageContext
import pw.binom.agentik.content.MessageOrigin.EVENT
import pw.binom.agentik.content.MessageOrigin.SYSTEM
import pw.binom.agentik.content.MessageOrigin.USER
import pw.binom.litert.LiteContentPart
import kotlin.test.Test
import kotlin.test.assertEquals
@@ -25,7 +25,7 @@ import pw.binom.agentik.memory.MemorySystemGuidance
import pw.binom.agentik.memory.NewMemoryNote
import pw.binom.agentik.memory.ReviewedTurn
import pw.binom.agentik.memory.md.openMdMemorySystem
import pw.binom.agentik.proto.Content
import pw.binom.agentik.content.Content
import pw.binom.agentik.standalone.agent.memory.MemoryToolsFactory
import pw.binom.agentik.standalone.llm.LlmBackend
import pw.binom.agentik.standalone.llm.LlmConfig
@@ -1,9 +1,9 @@
package pw.binom.agentik.standalone.persistence
import pw.binom.agentik.journal.MessageContext
import pw.binom.agentik.journal.MessageOrigin
import pw.binom.agentik.content.MessageContext
import pw.binom.agentik.content.MessageOrigin
import pw.binom.agentik.journal.ConversationRecord
import pw.binom.agentik.journal.MessageRecord
import pw.binom.agentik.journal.Content
import pw.binom.agentik.content.Content
import pw.binom.agentik.context.WorkingMemoryEntry
import kotlinx.coroutines.flow.toList