refactor(protocol): add toolName to ToolResult, remove proto typealiases, rename id→toolCallId, drop Conversation.events()
Three protocol-level changes from Android-client review (items 1-3, 5-6):
1) toolName denormalization in ToolResult (3 layers):
- :outbox-api/Event.ToolResult: +toolName: String? = null
- :journal-api/MessageRecord.ToolResult: +toolName: String? = null
- :proto/Message.ToolResult: +toolName: String? = null
- :storage-ksqlite, :journal-ksqlite ResultPayload codec: +toolName
- :standalone/ToolDispatcher, ConversationLoop: thread toolName = call.name
Nullable + default = backward-compat for already-persisted histories
and existing clients.
2) Drop proto/Event.kt, AgentEvent.kt, CommonEvent.kt typealiases.
is proto.Event.End failed with 'Unresolved reference End' (alias
loses nested-class access). Use pw.binom.agentik.outbox.{Event,
AgentEvent, CommonEvent} directly everywhere — :proto already has
api(:outbox-api), the package is visible to consumers, no shim
needed. 21 files rewired, 3 files deleted.
3) Rename Event.ToolResult.id → toolCallId (option B per user).
In :outbox-api Event.ToolResult.id == Event.ToolCall.id (one value,
one name); the persistent journal keeps MessageRecord.ToolResult.id
as its own PK + toolCallId as FK to the call — different semantics,
left untouched. Fixed ToolDispatcher bug: emitted id = resultId
while KDoc claimed id == ToolCall.id; now emits toolCallId = callId.
4) Remove Conversation.events() from :proto; OutboxStore is sole event source.
Conversation is a pure per-conversation abstraction (send/getMessages/
rename/close). Live events only via agent.outbox.conversationEvents/
agentEvents/events. HTTP route /conversations/{id}/events stays for
wire-compat but routes through outbox internally (map { it.event }).
jvmTest green (95 tasks).
This commit is contained in:
@@ -27,9 +27,10 @@ class SendSubcommand : AgentikSubcommand("send", "Отправить user-ход
|
||||
// Подписываемся на поток событий ДО send: события, отправленные
|
||||
// до подписки, не реплеятся (shared-flow без replay).
|
||||
val eventsJob = launch {
|
||||
conv.events(Instant.DISTANT_PAST)
|
||||
agent.outbox.conversationEvents(Instant.DISTANT_PAST, conv.id)
|
||||
// onEach печатает и терминальный event, takeWhile лишь
|
||||
// завершает сбор после него.
|
||||
.map { it.event }
|
||||
.onEach { ev -> emit(ev) }
|
||||
.takeWhile { ev -> !isTerminal(ev) }
|
||||
.collect { }
|
||||
@@ -53,7 +54,7 @@ class SendSubcommand : AgentikSubcommand("send", "Отправить user-ход
|
||||
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.id} ${escape(ev.result ?: "")}")
|
||||
is Event.ToolResult -> println("event ToolResult ${ev.toolCallId} ${escape(ev.result ?: "")}")
|
||||
is Event.End -> println("event End")
|
||||
is Event.Interrupted -> println("event Interrupted")
|
||||
is Event.Error -> println("event Error ${ev.code ?: ""} ${escape(ev.message)}")
|
||||
|
||||
@@ -7,7 +7,7 @@ import kotlinx.coroutines.launch
|
||||
import pw.binom.agentik.proto.Agent
|
||||
import pw.binom.agentik.proto.Content
|
||||
import pw.binom.agentik.proto.Conversation
|
||||
import pw.binom.agentik.proto.Event
|
||||
import pw.binom.agentik.outbox.Event
|
||||
import kotlin.coroutines.CoroutineContext
|
||||
import kotlin.time.Instant
|
||||
|
||||
@@ -92,12 +92,12 @@ internal class TuiBackend(
|
||||
}
|
||||
|
||||
/**
|
||||
* Подписывается на [Conversation.events] и перенаправляет их в [state].
|
||||
* Подписывается на `outbox.conversationEvents(after, conv.id)` и перенаправляет их в [state].
|
||||
*/
|
||||
private fun subscribeEvents(conv: Conversation, from: Instant) {
|
||||
eventsJob?.cancel()
|
||||
eventsJob = scope.launch {
|
||||
conv.events(from).collect { ev -> dispatch(ev) }
|
||||
agent.outbox.conversationEvents(from, conv.id).collect { ce -> dispatch(ce.event) }
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -8,7 +8,7 @@ import pw.binom.agentik.outbox.OutboxStore
|
||||
import pw.binom.agentik.proto.Agent
|
||||
import pw.binom.agentik.proto.Content
|
||||
import pw.binom.agentik.proto.Conversation
|
||||
import pw.binom.agentik.proto.Event
|
||||
import pw.binom.agentik.outbox.Event
|
||||
import pw.binom.agentik.proto.Message
|
||||
import pw.binom.agentik.proto.MessageContext
|
||||
import kotlin.time.Instant
|
||||
@@ -31,7 +31,7 @@ internal class FakeAgent(
|
||||
// emptyFlow, journal — error-on-access (никто не должен его трогать).
|
||||
override val journal: JournalStore = error("journal not used in TuiBackend tests")
|
||||
override val outbox: OutboxStore = object : OutboxStore {
|
||||
override fun events(after: Instant?) = emptyFlow<pw.binom.agentik.proto.CommonEvent>()
|
||||
override fun events(after: Instant?) = emptyFlow<pw.binom.agentik.outbox.CommonEvent>()
|
||||
override suspend fun earliestEventDate(): Instant = Instant.DISTANT_PAST
|
||||
override fun close() {}
|
||||
}
|
||||
|
||||
@@ -4,7 +4,7 @@ import kotlinx.coroutines.ExperimentalCoroutinesApi
|
||||
import kotlinx.coroutines.test.runCurrent
|
||||
import kotlinx.coroutines.test.runTest
|
||||
import pw.binom.agentik.proto.Content
|
||||
import pw.binom.agentik.proto.Event
|
||||
import pw.binom.agentik.outbox.Event
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertFalse
|
||||
@@ -170,7 +170,7 @@ class TuiBackendTest {
|
||||
runCurrent()
|
||||
val now = kotlin.time.Clock.System.now()
|
||||
conv.emit(Event.ToolCall(date = now, id = "1", title = null, toolName = "echo", toolArgs = """{"x":1}"""))
|
||||
conv.emit(Event.ToolResult(date = now, id = "1", result = "ok"))
|
||||
conv.emit(Event.ToolResult(date = now, toolCallId = "1", result = "ok"))
|
||||
runCurrent()
|
||||
|
||||
val toolMsgs = state.messages.value.filterIsInstance<TuiMessage.ToolCall>()
|
||||
|
||||
@@ -6,18 +6,12 @@ import io.ktor.client.request.get
|
||||
import io.ktor.client.request.parameter
|
||||
import io.ktor.client.request.patch
|
||||
import io.ktor.client.request.post
|
||||
import io.ktor.client.request.prepareGet
|
||||
import io.ktor.client.request.setBody
|
||||
import io.ktor.client.statement.bodyAsChannel
|
||||
import io.ktor.http.ContentType
|
||||
import io.ktor.http.HttpStatusCode
|
||||
import io.ktor.http.contentType
|
||||
import kotlinx.coroutines.flow.Flow
|
||||
import kotlinx.coroutines.flow.flow
|
||||
import kotlinx.serialization.Serializable
|
||||
import pw.binom.agentik.proto.Content
|
||||
import pw.binom.agentik.proto.Conversation
|
||||
import pw.binom.agentik.proto.Event
|
||||
import pw.binom.agentik.proto.Message
|
||||
import pw.binom.agentik.proto.MessageContext
|
||||
import kotlin.time.Instant
|
||||
@@ -72,22 +66,6 @@ internal class ConversationClient(
|
||||
httpClient.post("$convUrl/interrupt")
|
||||
}
|
||||
|
||||
override fun events(after: Instant): Flow<Event> = flow {
|
||||
// prepareGet + execute (а не get) обязателен: `get` дожидается полного
|
||||
// тела ответа, а SSE-поток не заканчивается никогда — вызов висел бы
|
||||
// вечно. `execute` отдаёт HttpResponse со стриминговым bodyAsChannel.
|
||||
httpClient.prepareGet("$convUrl/events?after=$after") { noSseReadTimeout() }
|
||||
.execute { response ->
|
||||
check(response.status == HttpStatusCode.OK) {
|
||||
"events: server returned ${response.status}"
|
||||
}
|
||||
readSse(response.bodyAsChannel())
|
||||
.collect { payload ->
|
||||
emit(agentikJson.decodeFromString(Event.serializer(), payload))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
override suspend fun getMessages(after: Instant, offset: Int, limit: Int): List<Message> =
|
||||
httpClient.get("$convUrl/messages") {
|
||||
parameter("after", after.toString())
|
||||
|
||||
@@ -55,6 +55,14 @@ sealed interface MessageRecord {
|
||||
override val id: String,
|
||||
override val conversationId: String,
|
||||
val toolCallId: String,
|
||||
/**
|
||||
* Имя тула, денормализованное из соответствующего `MessageRecord.ToolCall.toolName`.
|
||||
* Денормализация экономна (одна строка в SQLite) и снимает с UI
|
||||
* необходимость сопоставления `toolCallId → toolName`. `null` —
|
||||
* безопасный backfill для записей до миграции или для сиротливых
|
||||
* результатов без предшествующего `ToolCall`.
|
||||
*/
|
||||
val toolName: String? = null,
|
||||
val result: String?,
|
||||
override val createdAt: Instant,
|
||||
) : MessageRecord
|
||||
|
||||
+13
-3
@@ -30,7 +30,7 @@ internal fun encodeRecord(record: MessageRecord): Pair<String, String> = when (r
|
||||
)
|
||||
is MessageRecord.ToolResult -> "tool_result" to Json.encodeToString(
|
||||
ResultPayload.serializer(),
|
||||
ResultPayload(toolCallId = record.toolCallId, result = record.result),
|
||||
ResultPayload(toolCallId = record.toolCallId, toolName = record.toolName, result = record.result),
|
||||
)
|
||||
is MessageRecord.Error -> "error" to Json.encodeToString(
|
||||
ErrorPayload.serializer(),
|
||||
@@ -59,7 +59,7 @@ internal fun SQLiteResultSet.toMessageRecord(json: Json): MessageRecord {
|
||||
}
|
||||
"tool_result" -> {
|
||||
val p = Json.decodeFromString(ResultPayload.serializer(), payload)
|
||||
MessageRecord.ToolResult(id = id, conversationId = convId, toolCallId = p.toolCallId, result = p.result, createdAt = createdAt)
|
||||
MessageRecord.ToolResult(id = id, conversationId = convId, toolCallId = p.toolCallId, toolName = p.toolName, result = p.result, createdAt = createdAt)
|
||||
}
|
||||
"error" -> {
|
||||
val p = Json.decodeFromString(ErrorPayload.serializer(), payload)
|
||||
@@ -72,8 +72,18 @@ internal fun SQLiteResultSet.toMessageRecord(json: Json): MessageRecord {
|
||||
@kotlinx.serialization.Serializable
|
||||
internal data class CallPayload(val name: String, val title: String?, val argsJson: String)
|
||||
|
||||
/**
|
||||
* Тулрезалт-сериализация для SQLite. [toolName] денормализован из
|
||||
* соответствующего `ToolCall.name` для упрощения UI (нет нужды в
|
||||
* локальной `Map<id, name>`). Nullable с дефолтом — старые записи
|
||||
* без поля десериализуются как `null`.
|
||||
*/
|
||||
@kotlinx.serialization.Serializable
|
||||
internal data class ResultPayload(val toolCallId: String, val result: String?)
|
||||
internal data class ResultPayload(
|
||||
val toolCallId: String,
|
||||
val toolName: String? = null,
|
||||
val result: String?,
|
||||
)
|
||||
|
||||
@kotlinx.serialization.Serializable
|
||||
internal data class ErrorPayload(val message: String, val code: String?)
|
||||
|
||||
@@ -14,9 +14,9 @@ import kotlin.time.Instant
|
||||
* идентична `OutboxStore.events`: поток **не реплеит** прошлое, для бэкфилла
|
||||
* используются `Agent.getConversations` / `getConversation`.
|
||||
*
|
||||
* **История**: раньше жил в `:proto` (как `pw.binom.agentik.proto.AgentEvent`).
|
||||
* После миграции в `:outbox-api` — `:proto.AgentEvent` стал typealias'ом,
|
||||
* backward-compat для существующих импортов сохранён.
|
||||
* **История**: до 2026-09-21 жил в `:proto` как `pw.binom.agentik.proto.AgentEvent`;
|
||||
* typealias удалён 2026-09-21 (стирал nested-типы в `is`/`when`) — потребители
|
||||
* импортируют напрямую из `pw.binom.agentik.outbox.AgentEvent`.
|
||||
*/
|
||||
@Serializable
|
||||
sealed interface AgentEvent {
|
||||
|
||||
@@ -9,17 +9,17 @@ import kotlinx.serialization.Serializable
|
||||
*
|
||||
* Useful for admin dashboards, debug tools, parent agents: one subscription
|
||||
* instead of N+1. For regular UI use two separate SSE feeds
|
||||
* ([AgentEvent] via `/events` и `Event` via `/conversations/{id}/events`);
|
||||
* ([AgentEvent] via `/events` и [Event] via `/conversations/{id}/events`);
|
||||
* [CommonEvent] — for those who need everything in one place.
|
||||
*
|
||||
* Server endpoint: `GET /events/all` (SSE), or replay via `OutboxStore.events(after)`.
|
||||
*
|
||||
* **История**: раньше жил в `:proto` (как `pw.binom.agentik.proto.CommonEvent`).
|
||||
* После миграции в `:outbox-api` — `:proto.CommonEvent` стал typealias'ом,
|
||||
* backward-compat для существующих импортов сохранён. `CommonEvent.Conversation`
|
||||
* ссылается на [Event] (тоже в `:outbox-api` теперь) — раньше был
|
||||
* `pw.binom.agentik.proto.Event`, теперь это `pw.binom.agentik.outbox.Event`
|
||||
* (он тоже typealias-нут в `:proto.Event`).
|
||||
* **История**: до 2026-09-21 жил в `:proto` как `pw.binom.agentik.proto.CommonEvent`;
|
||||
* при миграции в `:outbox-api` был оставлен typealias в `:proto` для backward-compat,
|
||||
* но он стирал nested-типы (`CommonEvent.Agent`, `CommonEvent.Conversation`),
|
||||
* что ломало `is CommonEvent.Agent` на стороне клиента. Typealias'ы
|
||||
* `Event`/`AgentEvent`/`CommonEvent` из `:proto` удалены — потребители
|
||||
* импортируют напрямую из `pw.binom.agentik.outbox.*`.
|
||||
*/
|
||||
@Serializable
|
||||
sealed interface CommonEvent {
|
||||
|
||||
@@ -15,9 +15,12 @@ import kotlin.time.Instant
|
||||
* `StartReasoning?` → `StartResponse(TEXT|IMAGE)` → ...контент... → `End` | `Interrupted` | `Error`.
|
||||
* `StartReasoning` может отсутствовать, если агент не показывал рассуждения.
|
||||
*
|
||||
* **История**: раньше жил в `:proto` (как `pw.binom.agentik.proto.Event`).
|
||||
* После миграции в `:outbox-api` — `:proto.Event` стал typealias'ом,
|
||||
* backward-compat для существующих импортов сохранён.
|
||||
* **История**: до 2026-09-21 жил в `:proto` как `pw.binom.agentik.proto.Event`;
|
||||
* при миграции в `:outbox-api` был оставлен typealias в `:proto` для
|
||||
* backward-compat, но он стирал nested-типы (`Event.End`, `Event.ToolCall`,
|
||||
* `Event.ToolResult`), что ломало `is Event.End` на стороне клиента.
|
||||
* Typealias удалён 2026-09-21 — потребители импортируют напрямую из
|
||||
* `pw.binom.agentik.outbox.Event`.
|
||||
*/
|
||||
@Serializable
|
||||
sealed interface Event {
|
||||
@@ -75,12 +78,26 @@ sealed interface Event {
|
||||
|
||||
/**
|
||||
* Результат вызова тула. Приходит целиком после завершения исполнения.
|
||||
* [id] совпадает с [ToolCall.id], к которому относится результат, и
|
||||
* с id `Message.ToolResult` в истории.
|
||||
*
|
||||
* [toolCallId] = id [ToolCall], к которому относится результат, и
|
||||
* `MessageRecord.ToolResult.toolCallId` в истории. Один Call → один Result,
|
||||
* пара `(date, toolCallId)` уникальна — отдельный `id` в live-событии
|
||||
* не нужен (PK живёт в персистентном журнале).
|
||||
*
|
||||
* [toolName] денормализован из соответствующего [ToolCall.toolName] —
|
||||
* UI рендерит имя тула без локальной `Map<toolCallId, name>` и без риска
|
||||
* «Result пришёл до Call». `null` допустим для backfill'а старых
|
||||
* записей, у которых поле отсутствует, или теоретического случая
|
||||
* Result без предшествующего Call (orphan).
|
||||
*/
|
||||
@Serializable
|
||||
@SerialName("tool_result")
|
||||
data class ToolResult(override val date: Instant, val id: String, val result: String?) : Event
|
||||
data class ToolResult(
|
||||
override val date: Instant,
|
||||
val toolCallId: String,
|
||||
val toolName: String? = null,
|
||||
val result: String?,
|
||||
) : Event
|
||||
|
||||
/**
|
||||
* Ошибка хода. После неё поток завершается; дальнейшие события могут
|
||||
|
||||
@@ -1,10 +0,0 @@
|
||||
package pw.binom.agentik.proto
|
||||
|
||||
/**
|
||||
* Backward-compat typealias: `AgentEvent` теперь живёт в `:outbox-api`
|
||||
* (логически принадлежит сущности outbox, не wire-протоколу `:proto`).
|
||||
*
|
||||
* Существующие импорты `pw.binom.agentik.proto.AgentEvent` продолжают
|
||||
* работать транспарентно. Использовать typealias в новом коде.
|
||||
*/
|
||||
typealias AgentEvent = pw.binom.agentik.outbox.AgentEvent
|
||||
@@ -1,9 +0,0 @@
|
||||
package pw.binom.agentik.proto
|
||||
|
||||
/**
|
||||
* Backward-compat typealias: `CommonEvent` теперь живёт в `:outbox-api`.
|
||||
*
|
||||
* Существующие импорты `pw.binom.agentik.proto.CommonEvent` продолжают
|
||||
* работать транспарентно.
|
||||
*/
|
||||
typealias CommonEvent = pw.binom.agentik.outbox.CommonEvent
|
||||
@@ -8,6 +8,19 @@ import kotlin.time.Instant
|
||||
* Stateful-диалог клиента и [Agent]. Хранит собственную историю: на каждый
|
||||
* [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 -> ... }
|
||||
* ```
|
||||
* Для cross-conversation view (admin / parent-agent / debug):
|
||||
* `agent.outbox.events(after)`. Для lifecycle агента (created/deleted/renamed):
|
||||
* `agent.outbox.agentEvents(after)`.
|
||||
*/
|
||||
interface Conversation : AutoCloseable {
|
||||
val id: String
|
||||
@@ -33,7 +46,7 @@ interface Conversation : AutoCloseable {
|
||||
|
||||
/**
|
||||
* Ставит новый user-ход в очередь. Возвращает управление сразу — поток
|
||||
* событий ответа приходит через [events].
|
||||
* событий ответа приходит через `agent.outbox.conversationEvents(...)`.
|
||||
*
|
||||
* Если в момент вызова выполняется другой ход, новый встаёт в очередь
|
||||
* за ним. Чтобы отменить текущий — вызови [interrupt] перед [send].
|
||||
@@ -48,25 +61,13 @@ interface Conversation : AutoCloseable {
|
||||
|
||||
/**
|
||||
* Прерывает текущий исполняемый ход (best-effort: LLM-stream прибивается,
|
||||
* in-flight tool может доехать или отвалиться). В [events] эмитится
|
||||
* [Event.Interrupted], затем может начаться следующий ход из очереди.
|
||||
* in-flight tool может доехать или отвалиться). В `agent.outbox.conversationEvents`
|
||||
* эмитится `Event.Interrupted`, затем может начаться следующий ход из очереди.
|
||||
*
|
||||
* Если хода нет — no-op.
|
||||
*/
|
||||
suspend fun interrupt()
|
||||
|
||||
/**
|
||||
* Live-подписка на всё, что происходит в диалоге, начиная с [after].
|
||||
*
|
||||
* **Не реплеит** события, произошедшие до [after] — для бэкфилла
|
||||
* используй [getMessages]. Если [after] — момент последнего виденного
|
||||
* клиентом события, поток продолжается «с того места».
|
||||
*
|
||||
* Подписки независимы: каждый вызов возвращает свой [Flow], отмена одного
|
||||
* не влияет на других подписчиков и на сам диалог.
|
||||
*/
|
||||
fun events(after: Instant): Flow<Event>
|
||||
|
||||
/** Страница истории: не более [limit] сообщений после [after], начиная с [offset]-го. */
|
||||
suspend fun getMessages(after: Instant, offset: Int, limit: Int): List<Message>
|
||||
|
||||
@@ -83,7 +84,7 @@ interface Conversation : AutoCloseable {
|
||||
|
||||
/**
|
||||
* Освобождает ресурсы диалога (подписки, сетевые хэндлы). Идемпотентно.
|
||||
* После [close] дальнейшие вызовы [send]/[interrupt]/[events]/[getMessages]/[rename] не определены.
|
||||
* После [close] дальнейшие вызовы [send]/[interrupt]/[getMessages]/[rename] не определены.
|
||||
*/
|
||||
override fun close()
|
||||
|
||||
|
||||
@@ -1,11 +0,0 @@
|
||||
package pw.binom.agentik.proto
|
||||
|
||||
/**
|
||||
* Backward-compat typealias: `Event` теперь живёт в `:outbox-api`.
|
||||
*
|
||||
* Существующие импорты `pw.binom.agentik.proto.Event` продолжают
|
||||
* работать транспарентно. `Conversation.events(after): Flow<Event>` в
|
||||
* `:proto.Conversation` теперь фактически возвращает
|
||||
* `pw.binom.agentik.outbox.Event` — тот же тип, другое имя.
|
||||
*/
|
||||
typealias Event = pw.binom.agentik.outbox.Event
|
||||
@@ -45,9 +45,22 @@ sealed interface Message {
|
||||
override val date: Instant
|
||||
) : Message
|
||||
|
||||
/**
|
||||
* Результат вызова тула. Приходит в историю `getMessages` после завершения хода.
|
||||
* [id] совпадает с [ToolCall.id], к которому относится результат, и
|
||||
* с id соответствующего `MessageRecord.ToolResult` в journal.
|
||||
*
|
||||
* [toolName] денормализован из [ToolCall.toolName] — UI рендерит
|
||||
* имя тула в строке результата без отдельной `Map<id, name>`.
|
||||
*/
|
||||
@Serializable
|
||||
@SerialName("tool_result")
|
||||
class ToolResult(override val id: String, val result: String?, override val date: Instant) : Message
|
||||
class ToolResult(
|
||||
override val id: String,
|
||||
val toolName: String? = null,
|
||||
val result: String?,
|
||||
override val date: Instant,
|
||||
) : Message
|
||||
|
||||
/**
|
||||
* Ход завершился ошибкой (LLM, инициализация движка или иная отказоустойчивая
|
||||
|
||||
@@ -3,8 +3,8 @@ package pw.binom.agentik.server
|
||||
import io.ktor.server.routing.Route
|
||||
import io.ktor.server.routing.get
|
||||
import io.ktor.server.routing.route
|
||||
import pw.binom.agentik.outbox.CommonEvent
|
||||
import pw.binom.agentik.outbox.OutboxStore
|
||||
import pw.binom.agentik.proto.CommonEvent
|
||||
|
||||
/**
|
||||
* HTTP-фасад для [OutboxStore] (bounded-tail live event stream агента).
|
||||
|
||||
@@ -20,11 +20,11 @@ import kotlinx.coroutines.flow.map
|
||||
import kotlinx.serialization.KSerializer
|
||||
import kotlinx.serialization.json.Json
|
||||
import pw.binom.agentik.proto.Agent
|
||||
import pw.binom.agentik.proto.AgentEvent
|
||||
import pw.binom.agentik.proto.CommonEvent
|
||||
import pw.binom.agentik.proto.Conversation
|
||||
import pw.binom.agentik.proto.Content
|
||||
import pw.binom.agentik.proto.Event
|
||||
import pw.binom.agentik.outbox.AgentEvent
|
||||
import pw.binom.agentik.outbox.CommonEvent
|
||||
import pw.binom.agentik.outbox.Event
|
||||
import kotlin.time.Instant
|
||||
|
||||
internal fun Route.agentikRoutes(agent: Agent) {
|
||||
@@ -128,7 +128,10 @@ internal fun Route.agentikRoutes(agent: Agent) {
|
||||
return@get
|
||||
}
|
||||
val after = call.parseAfter() ?: return@get
|
||||
call.streamJsonSse(c.events(after), Event.serializer())
|
||||
// Live-источник событий — `OutboxStore` (единая точка истины);
|
||||
// разворачиваем `CommonEvent.Conversation` → `Event` для совместимости
|
||||
// wire-формата (клиент десериализует как `Event`, не как `CommonEvent.Conversation`).
|
||||
call.streamJsonSse(agent.outbox.conversationEvents(after, id).map { it.event }, Event.serializer())
|
||||
}
|
||||
|
||||
get("/events") {
|
||||
|
||||
@@ -49,7 +49,8 @@ class A2aBridge(private val agent: Agent) : AgentHandler {
|
||||
val reply = StringBuilder()
|
||||
val turnDone = CompletableDeferred<Unit>()
|
||||
val subscription = async {
|
||||
conv.events(since).collect { e ->
|
||||
agent.outbox.conversationEvents(since, conv.id).collect { ce ->
|
||||
val e = ce.event
|
||||
when (e) {
|
||||
is Event.AppendText -> reply.append(e.body)
|
||||
is Event.End, is Event.Interrupted -> turnDone.complete(Unit)
|
||||
|
||||
@@ -18,11 +18,11 @@ import pw.binom.agentik.memory.MemorySystemGuidance
|
||||
import pw.binom.agentik.proto.Agent as ProtoAgent
|
||||
import pw.binom.agentik.outbox.AgentEvent
|
||||
import pw.binom.agentik.outbox.CommonEvent
|
||||
import pw.binom.agentik.outbox.Event as ProtoEvent
|
||||
import pw.binom.agentik.outbox.MutableOutboxStore
|
||||
import pw.binom.agentik.journal.JournalStore
|
||||
import pw.binom.agentik.outbox.OutboxStore
|
||||
import pw.binom.agentik.proto.Conversation as ProtoConversation
|
||||
import pw.binom.agentik.proto.Event as ProtoEvent
|
||||
import pw.binom.agentik.skills.SkillCatalog
|
||||
import pw.binom.agentik.skills.renderSystemPromptSection
|
||||
import pw.binom.agentik.standalone.agent.memory.MemoryToolsFactory
|
||||
|
||||
+4
-4
@@ -3,15 +3,15 @@ package pw.binom.agentik.standalone.agent
|
||||
import kotlinx.coroutines.runBlocking
|
||||
import kotlinx.coroutines.flow.Flow
|
||||
import kotlinx.coroutines.flow.map
|
||||
import pw.binom.agentik.outbox.MutableOutboxStore
|
||||
import pw.binom.agentik.outbox.CommonEvent
|
||||
import pw.binom.agentik.proto.Event as ProtoEvent
|
||||
import pw.binom.agentik.outbox.Event
|
||||
import pw.binom.agentik.outbox.MutableOutboxStore
|
||||
|
||||
internal class ConversationEvents(
|
||||
private val globalEventStore: MutableOutboxStore,
|
||||
private val conversationId: String,
|
||||
) {
|
||||
fun tryEmit(event: ProtoEvent): Boolean {
|
||||
fun tryEmit(event: Event): Boolean {
|
||||
runBlocking {
|
||||
globalEventStore.append(
|
||||
CommonEvent.Conversation(
|
||||
@@ -24,7 +24,7 @@ internal class ConversationEvents(
|
||||
return true
|
||||
}
|
||||
|
||||
fun events(after: kotlin.time.Instant?): Flow<ProtoEvent> =
|
||||
fun events(after: kotlin.time.Instant?): Flow<Event> =
|
||||
globalEventStore.conversationEvents(after = after, conversationId = conversationId)
|
||||
.map { it.event }
|
||||
}
|
||||
|
||||
+1
-3
@@ -220,9 +220,6 @@ class ConversationLoop(
|
||||
toolDispatcher.currentToolJob?.cancel()
|
||||
}
|
||||
|
||||
override fun events(after: Instant): Flow<ProtoEvent> =
|
||||
events.events(after)
|
||||
|
||||
override suspend fun getMessages(after: Instant, offset: Int, limit: Int): List<ProtoMessage> =
|
||||
messageStore.list(conversationId = id, after = after, offset = offset, limit = limit)
|
||||
.map { it.toProto() }
|
||||
@@ -568,6 +565,7 @@ internal fun MessageRecord.toProto(): ProtoMessage = when (this) {
|
||||
is MessageRecord.ToolResult -> ProtoMessage.ToolResult(
|
||||
id = id,
|
||||
date = createdAt,
|
||||
toolName = toolName,
|
||||
result = result,
|
||||
)
|
||||
is MessageRecord.Error -> ProtoMessage.Error(
|
||||
|
||||
+2
-1
@@ -93,7 +93,7 @@ internal class ToolDispatcher(
|
||||
}
|
||||
|
||||
val resultAt = now()
|
||||
events.tryEmit(ProtoEvent.ToolResult(date = resultAt, id = resultId, result = resultText))
|
||||
events.tryEmit(ProtoEvent.ToolResult(date = resultAt, toolCallId = callId, toolName = call.name, result = resultText))
|
||||
|
||||
// Эмитим background event — другие компоненты (BackgroundScheduler)
|
||||
// решают, делать ли что-то. Cancellation = not a failure (не эмитим Failed).
|
||||
@@ -110,6 +110,7 @@ internal class ToolDispatcher(
|
||||
id = resultId,
|
||||
conversationId = state.id,
|
||||
toolCallId = callId,
|
||||
toolName = call.name,
|
||||
result = resultText,
|
||||
createdAt = resultAt,
|
||||
),
|
||||
|
||||
@@ -319,7 +319,7 @@ class ChatAgentTest {
|
||||
|
||||
val events = mutableListOf<ProtoEvent>()
|
||||
val job = launch(start = kotlinx.coroutines.CoroutineStart.UNDISPATCHED) {
|
||||
conv.events(Instant.DISTANT_PAST).collect { events.add(it) }
|
||||
agent.outbox.conversationEvents(Instant.DISTANT_PAST, conv.id).collect { events.add(it.event) }
|
||||
}
|
||||
conv.send(listOf(Content.Text("hi")))
|
||||
delay(50)
|
||||
@@ -356,7 +356,7 @@ class ChatAgentTest {
|
||||
// отправки событий подписка ничего не увидит.
|
||||
val events = mutableListOf<ProtoEvent>()
|
||||
val eventsJob = launch(start = kotlinx.coroutines.CoroutineStart.UNDISPATCHED) {
|
||||
conv.events(Instant.DISTANT_PAST).collect { events.add(it) }
|
||||
agent.outbox.conversationEvents(Instant.DISTANT_PAST, conv.id).collect { events.add(it.event) }
|
||||
}
|
||||
|
||||
val sendJob = launch {
|
||||
@@ -414,7 +414,7 @@ class ChatAgentTest {
|
||||
// Подписываемся ДО send — SharedFlow без replay
|
||||
val events = mutableListOf<ProtoEvent>()
|
||||
val eventsJob = launch(start = kotlinx.coroutines.CoroutineStart.UNDISPATCHED) {
|
||||
conv.events(Instant.DISTANT_PAST).collect { events.add(it) }
|
||||
agent.outbox.conversationEvents(Instant.DISTANT_PAST, conv.id).collect { events.add(it.event) }
|
||||
}
|
||||
|
||||
val sendJob = launch {
|
||||
|
||||
+13
-3
@@ -30,7 +30,7 @@ internal fun encodeRecord(record: MessageRecord): Pair<String, String> = when (r
|
||||
)
|
||||
is MessageRecord.ToolResult -> "tool_result" to Json.encodeToString(
|
||||
ResultPayload.serializer(),
|
||||
ResultPayload(toolCallId = record.toolCallId, result = record.result),
|
||||
ResultPayload(toolCallId = record.toolCallId, toolName = record.toolName, result = record.result),
|
||||
)
|
||||
is MessageRecord.Error -> "error" to Json.encodeToString(
|
||||
ErrorPayload.serializer(),
|
||||
@@ -59,7 +59,7 @@ internal fun SQLiteResultSet.toMessageRecord(json: Json): MessageRecord {
|
||||
}
|
||||
"tool_result" -> {
|
||||
val p = Json.decodeFromString(ResultPayload.serializer(), payload)
|
||||
MessageRecord.ToolResult(id = id, conversationId = convId, toolCallId = p.toolCallId, result = p.result, createdAt = createdAt)
|
||||
MessageRecord.ToolResult(id = id, conversationId = convId, toolCallId = p.toolCallId, toolName = p.toolName, result = p.result, createdAt = createdAt)
|
||||
}
|
||||
"error" -> {
|
||||
val p = Json.decodeFromString(ErrorPayload.serializer(), payload)
|
||||
@@ -72,8 +72,18 @@ internal fun SQLiteResultSet.toMessageRecord(json: Json): MessageRecord {
|
||||
@kotlinx.serialization.Serializable
|
||||
internal data class CallPayload(val name: String, val title: String?, val argsJson: String)
|
||||
|
||||
/**
|
||||
* Тулрезалт-сериализация для SQLite. [toolName] денормализован из
|
||||
* соответствующего `ToolCall.name` для упрощения UI (нет нужды в
|
||||
* локальной `Map<id, name>`). Nullable с дефолтом — старые записи
|
||||
* без поля десериализуются как `null`.
|
||||
*/
|
||||
@kotlinx.serialization.Serializable
|
||||
internal data class ResultPayload(val toolCallId: String, val result: String?)
|
||||
internal data class ResultPayload(
|
||||
val toolCallId: String,
|
||||
val toolName: String? = null,
|
||||
val result: String?,
|
||||
)
|
||||
|
||||
@kotlinx.serialization.Serializable
|
||||
internal data class ErrorPayload(val message: String, val code: String?)
|
||||
|
||||
Reference in New Issue
Block a user