standalone: persist turn errors so polling clients see them
Event.Error уже эмитился в live-SSE, но если клиент подключился
после провала хода (или опрашивает историю через getMessages
вместо SSE), он видел только user-сообщение без следа, что ход
провалился. Это и нужно было поправить.
* :proto
- Message.Error(id, message, code?, date) — терминальная
персистентная проекция Event.Error. Audit-only; в live-стриме
по-прежнему приходит Event.Error.
* :standalone
- MessageRecord.Error с тем же контрактом.
- SqliteMessageStore: kind "error", JSON-payload {message, code};
encoding-ошибки round-trip покрыты тестом (с code и без).
- ChatConversation.failTurn(message, code?): пишет MessageRecord.Error
в audit (кроме temp-диалогов), затем эмитит Event.Error + Event.End.
Все три exit-точки из runTurn (пустой parts, init-catch,
stream-catch) теперь через failTurn.
- ChatConversation: при ошибке стрима живой LiteConversation
сбрасывается — следующий send пересоберёт его из working_memory.
Раньше оставляли битую инстанцию вопреки KDoc класса.
- toProto: маппит MessageRecord.Error → Message.Error.
* docs
- STANDALONE.md: добавлен Message.Error в таблицу dual-log + абзац
про персист ошибок (audit/working_memory, сброс liteConv).
- ARCHITECTURE.md: Message.Error в сигнатуре Message, failTurn
в ChatConversation, Message.Error в audit log.
- Поправлен пример describe() в доке тулов (реальный формат —
OpenAI-style function wrapper, а не голый JSON Schema).
Тесты:
- ChatAgentTest: LLM failure → Event.Error + Event.End + Error
record в audit + backfill через getMessages.
- PersistenceTest: MessageRecord.Error round-trip (с code и без).
- :standalone jvmTest 70 (было 69).
E2E: неверный model id (HTTP 400) → polling GET /messages теперь
возвращает user_message + error (без assistant_message).
This commit is contained in:
@@ -66,7 +66,7 @@ agentik
|
||||
- `AutoCloseable` — `close()` идемпотентен.
|
||||
|
||||
`Content = Text(body) | Image(data, mime)`.
|
||||
`Message = UserMessage | AssistantMessage | ToolCall | ToolResult | Summary | System`.
|
||||
`Message = UserMessage | AssistantMessage | ToolCall | ToolResult | Error`.
|
||||
`Event = StartReasoning | StartResponse | End | AppendText | AppendImage | ToolCall | ToolResult | Error`.
|
||||
`AgentEvent = Created(conversationId) | Deleted(id) | Renamed(id, title)`.
|
||||
|
||||
@@ -104,6 +104,9 @@ Main.kt
|
||||
- Tool-loop: на `delta.toolCalls` → emit `Event.ToolCall` → execute
|
||||
(`tool.invoke(argsJson)`) → emit `Event.ToolResult` →
|
||||
`liteConv.addToolResult(callId, name, result)` → продолжение стрима.
|
||||
- На ошибку хода (init/стрим LLM): `failTurn` → `messageStore.append(Error)`
|
||||
(audit; в working_memory не пишется) + `Event.Error` + `Event.End`, живой
|
||||
`LiteConversation` сбрасывается и пересобирается на следующем `send`.
|
||||
- `LiteConversation` живёт **один на весь `ChatConversation`** для обоих
|
||||
бэкендов: KV-cache движка сохраняется между ходами. Это контракт
|
||||
litert-api, не только Google-specific.
|
||||
@@ -116,7 +119,7 @@ Main.kt
|
||||
summarization-вставка отложена (нужен дизайн-проработка).
|
||||
|
||||
**`MessageStore`** — append-only аудит. На каждый ход дописываются
|
||||
`UserMessage`, `AssistantMessage`, `ToolCall`, `ToolResult`. Никаких
|
||||
`UserMessage`, `AssistantMessage`, `ToolCall`, `ToolResult`, `Error`. Никаких
|
||||
update/delete кроме каскада из `ConversationStore.delete`.
|
||||
|
||||
## 5. Слои персистентности
|
||||
|
||||
+4
-2
@@ -166,13 +166,15 @@ fun main() {
|
||||
|
||||
| таблица | назначение | мутации |
|
||||
|---|---|---|
|
||||
| `message` | append-only audit log. Все user/assistant/tool-call/tool-result сообщения. Никогда не редактируется (кроме каскадного `DELETE` при удалении диалога). | только `INSERT` |
|
||||
| `message` | append-only audit log. Все user/assistant/tool-call/tool-result/error сообщения. Никогда не редактируется (кроме каскадного `DELETE` при удалении диалога). | только `INSERT` |
|
||||
| `working_memory` | mutable LLM-контекст. System-prompt + текущая история + (в v2) суммаризации. | `INSERT`, `compact(dropFromIdx, summary)` |
|
||||
|
||||
Маппинг `:proto.Message ↔ MessageRecord` живёт в `ChatConversation.kt` (`toProto`/`toStorage`) — сами `MessageRecord` намеренно НЕ зависят от `:proto`, чтобы можно было сменить транспорт без миграции таблиц.
|
||||
|
||||
Подробный контракт — в комментариях к `MessageRecord.kt` и `WorkingMemoryEntry.kt`.
|
||||
|
||||
**Ошибки хода персистятся.** Если ход провалился (LLM/движок недоступны — например, HTTP 400 от endpoint'а), `ChatConversation.failTurn` пишет терминальную запись `MessageRecord.Error` в audit и эмитит `Event.Error` + `Event.End`. Благодаря audit-записи ошибка видна не только подписчику live-SSE, но и клиенту, который делает backfill через `getMessages` (polling/переподключение): в истории будет `Message.Error(id, message, code?)`, а для этого user-сообщения не будет `AssistantMessage`. При ошибке стрима живой `LiteConversation` сбрасывается — следующий `send` пересоберёт его из `working_memory`. В working_memory `Error` не пишется (модель не должна видеть ошибки прошлых ходов).
|
||||
|
||||
### `MessageStore`
|
||||
|
||||
```kotlin
|
||||
@@ -259,7 +261,7 @@ suspend fun touch(id: String, now: Instant)
|
||||
// 1) In-agent tool (нативный LiteTool):
|
||||
val echoTool = object : LiteTool {
|
||||
override fun describe(): String =
|
||||
"""{"type":"object","properties":{"x":{"type":"string"}},"required":["x"]}"""
|
||||
"""{"type":"function","function":{"name":"echo","description":"Echo a string","parameters":{"type":"object","properties":{"x":{"type":"string"}},"required":["x"]}}}"""
|
||||
override fun invoke(args: String): String = "echoed: $args"
|
||||
}
|
||||
val tools = listOf(NamedTool("echo", echoTool))
|
||||
|
||||
@@ -39,4 +39,19 @@ sealed interface Message {
|
||||
@Serializable
|
||||
@SerialName("tool_result")
|
||||
class ToolResult(override val id: String, val result: String?, override val date: Instant) : Message
|
||||
|
||||
/**
|
||||
* Ход завершился ошибкой (LLM, инициализация движка или иная отказоустойчивая
|
||||
* ситуация). Терминальная запись хода: соответствующего [AssistantMessage]
|
||||
* не будет. В live-стриме приходит [Event.Error]; этот тип — его персистентная
|
||||
* проекция, видимая при backfill через `getMessages`.
|
||||
*/
|
||||
@Serializable
|
||||
@SerialName("error")
|
||||
class Error(
|
||||
override val id: String,
|
||||
val message: String,
|
||||
val code: String? = null,
|
||||
override val date: Instant,
|
||||
) : Message
|
||||
}
|
||||
|
||||
+14
@@ -69,6 +69,20 @@ sealed interface MessageRecord {
|
||||
override val createdAt: Instant,
|
||||
) : MessageRecord
|
||||
|
||||
/**
|
||||
* Терминальная запись провалившегося хода. Только audit log
|
||||
* (в working_memory не пишется — модель не должна видеть ошибки прошлых ходов).
|
||||
*/
|
||||
@Serializable
|
||||
@SerialName("error")
|
||||
data class Error(
|
||||
override val id: String,
|
||||
override val conversationId: String,
|
||||
val message: String,
|
||||
val code: String?,
|
||||
override val createdAt: Instant,
|
||||
) : MessageRecord
|
||||
|
||||
/** Синтетическое: суммаризация старого контекста. Только в working_memory. */
|
||||
@Serializable
|
||||
@SerialName("summary")
|
||||
|
||||
+34
-6
@@ -177,8 +177,7 @@ class ChatConversation(
|
||||
}
|
||||
|
||||
if (parts.isEmpty()) {
|
||||
emitEvent(ProtoEvent.Error(date = now(), message = "Empty user input (no text content)"))
|
||||
emitEvent(ProtoEvent.End(date = now()))
|
||||
failTurn("Empty user input (no text content)")
|
||||
return
|
||||
}
|
||||
|
||||
@@ -186,8 +185,7 @@ class ChatConversation(
|
||||
getOrCreateLiteConversation(excludeUserSourceId = if (record.isTemporal) null else userRecord.id)
|
||||
} catch (e: Throwable) {
|
||||
this.liteConv = null
|
||||
emitEvent(ProtoEvent.Error(date = now(), message = e.message ?: "LiteConversation init failed"))
|
||||
emitEvent(ProtoEvent.End(date = now()))
|
||||
failTurn(e.message ?: "LiteConversation init failed")
|
||||
return
|
||||
}
|
||||
|
||||
@@ -210,8 +208,8 @@ class ChatConversation(
|
||||
} catch (e: kotlinx.coroutines.CancellationException) {
|
||||
throw e
|
||||
} catch (e: Throwable) {
|
||||
emitEvent(ProtoEvent.Error(date = now(), message = e.message ?: e.javaClass.simpleName))
|
||||
emitEvent(ProtoEvent.End(date = now()))
|
||||
this.liteConv = null
|
||||
failTurn(e.message ?: e.javaClass.simpleName)
|
||||
return
|
||||
}
|
||||
|
||||
@@ -354,6 +352,30 @@ class ChatConversation(
|
||||
events.tryEmit(event)
|
||||
}
|
||||
|
||||
/**
|
||||
* Терминальное завершение хода ошибкой: пишет [MessageRecord.Error] в audit
|
||||
* (кроме temp-диалогов), затем эмитит [ProtoEvent.Error] + [ProtoEvent.End].
|
||||
*
|
||||
* Благодаря audit-записи ошибка видна не только в live-стриме, но и при
|
||||
* backfill через `getMessages` (переподключение / polling).
|
||||
*/
|
||||
private suspend fun failTurn(message: String, code: String? = null) {
|
||||
val ts = now()
|
||||
if (!record.isTemporal) {
|
||||
messageStore.append(
|
||||
MessageRecord.Error(
|
||||
id = newId("err"),
|
||||
conversationId = id,
|
||||
message = message,
|
||||
code = code,
|
||||
createdAt = ts,
|
||||
),
|
||||
)
|
||||
}
|
||||
emitEvent(ProtoEvent.Error(date = ts, message = message, code = code))
|
||||
emitEvent(ProtoEvent.End(date = ts))
|
||||
}
|
||||
|
||||
private fun now(): Instant =
|
||||
Instant.fromEpochMilliseconds(System.currentTimeMillis())
|
||||
|
||||
@@ -426,6 +448,12 @@ internal fun MessageRecord.toProto(): ProtoMessage = when (this) {
|
||||
date = createdAt,
|
||||
result = result,
|
||||
)
|
||||
is MessageRecord.Error -> ProtoMessage.Error(
|
||||
id = id,
|
||||
date = createdAt,
|
||||
message = message,
|
||||
code = code,
|
||||
)
|
||||
is MessageRecord.Summary -> ProtoMessage.AssistantMessage(
|
||||
id = id,
|
||||
date = createdAt,
|
||||
|
||||
+17
@@ -62,6 +62,10 @@ private fun encodeRecord(record: MessageRecord): Pair<String, String> = when (re
|
||||
ResultPayload.serializer(),
|
||||
ResultPayload(toolCallId = record.toolCallId, result = record.result),
|
||||
)
|
||||
is MessageRecord.Error -> "error" to Json.encodeToString(
|
||||
ErrorPayload.serializer(),
|
||||
ErrorPayload(message = record.message, code = record.code),
|
||||
)
|
||||
is MessageRecord.Summary,
|
||||
is MessageRecord.System -> error("Summary/System — synthetic, cannot append to audit log")
|
||||
}
|
||||
@@ -72,6 +76,9 @@ internal data class CallPayload(val name: String, val title: String?, val argsJs
|
||||
@kotlinx.serialization.Serializable
|
||||
internal data class ResultPayload(val toolCallId: String, val result: String?)
|
||||
|
||||
@kotlinx.serialization.Serializable
|
||||
internal data class ErrorPayload(val message: String, val code: String?)
|
||||
|
||||
private fun Message.toRecord(): MessageRecord {
|
||||
val id = id
|
||||
val convId = conversation_id
|
||||
@@ -110,6 +117,16 @@ private fun Message.toRecord(): MessageRecord {
|
||||
createdAt = createdAt,
|
||||
)
|
||||
}
|
||||
"error" -> {
|
||||
val p = Json.decodeFromString(ErrorPayload.serializer(), payload_json)
|
||||
MessageRecord.Error(
|
||||
id = id,
|
||||
conversationId = convId,
|
||||
message = p.message,
|
||||
code = p.code,
|
||||
createdAt = createdAt,
|
||||
)
|
||||
}
|
||||
else -> error("Unknown message kind in audit log: $kind")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -286,6 +286,37 @@ class ChatAgentTest {
|
||||
assertEquals(LiteRole.MODEL, reopenedLite.initialMessages[1].role)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `LLM failure emits error event and persists error message`() = runTest {
|
||||
val agent = newAgent()
|
||||
fakeLlm.failMessage = "boom from llm"
|
||||
val conv = agent.createConversation(temp = false) as ChatConversation
|
||||
|
||||
val events = mutableListOf<ProtoEvent>()
|
||||
val job = launch(start = kotlinx.coroutines.CoroutineStart.UNDISPATCHED) {
|
||||
conv.events(Instant.DISTANT_PAST).collect { events.add(it) }
|
||||
}
|
||||
conv.send(listOf(Content.Text("hi")))
|
||||
delay(50)
|
||||
job.cancel()
|
||||
|
||||
assertTrue(events.any { it is ProtoEvent.Error && it.message == "boom from llm" }, "events=$events")
|
||||
assertTrue(events.any { it is ProtoEvent.End }, "events=$events")
|
||||
|
||||
val msgs = stores.messages.listAll(conv.id)
|
||||
assertEquals(2, msgs.size)
|
||||
assertIs<pw.binom.agentik.standalone.persistence.MessageRecord.UserMessage>(msgs[0])
|
||||
val err = assertIs<pw.binom.agentik.standalone.persistence.MessageRecord.Error>(msgs[1])
|
||||
assertEquals("boom from llm", err.message)
|
||||
|
||||
// backfill через getMessages (polling/reconnect) тоже видит ошибку
|
||||
val proto = conv.getMessages(Instant.DISTANT_PAST, offset = 0, limit = 10)
|
||||
assertTrue(
|
||||
proto.any { it is pw.binom.agentik.proto.Message.Error && it.message == "boom from llm" },
|
||||
"history=$proto",
|
||||
)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `interrupt cancels active send`() = runTest {
|
||||
val agent = newAgent()
|
||||
@@ -440,6 +471,7 @@ private class FakeLiteLlm : LiteLlm {
|
||||
var reply: String = ""
|
||||
var rememberHistory: Boolean = false
|
||||
var slow: Boolean = false
|
||||
var failMessage: String? = null
|
||||
var lastConfig: LiteConversationConfig? = null
|
||||
var lastContents: List<LiteContentPart>? = null
|
||||
val conversations = mutableListOf<FakeLiteConversation>()
|
||||
@@ -486,6 +518,9 @@ private class FakeLiteConversation(
|
||||
|
||||
override fun sendStreamContents(contents: List<LiteContentPart>): Flow<LiteDelta> {
|
||||
parent.lastContents = contents
|
||||
parent.failMessage?.let { msg ->
|
||||
return kotlinx.coroutines.flow.flow { throw RuntimeException(msg) }
|
||||
}
|
||||
mutableHistory.add(LiteMessage(LiteRole.USER, contents))
|
||||
if (parent.slow) {
|
||||
return kotlinx.coroutines.flow.flow {
|
||||
@@ -554,7 +589,6 @@ private class ToolLoopFakeLiteLlm : LiteLlm {
|
||||
override fun cancel() {}
|
||||
override fun tokenCount(): Int = hist.size
|
||||
override fun addToolResult(callId: String?, name: String, result: String) {
|
||||
System.err.println("[agentik] DEBUG ToolLoopFakeLiteLlm.addToolResult callId=$callId name=$name result=$result")
|
||||
lastToolResult = result
|
||||
}
|
||||
override fun close() {}
|
||||
|
||||
+32
-1
@@ -7,6 +7,7 @@ import kotlin.test.BeforeTest
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertFalse
|
||||
import kotlin.test.assertIs
|
||||
import kotlin.test.assertNotNull
|
||||
import kotlin.test.assertNull
|
||||
import kotlin.test.assertTrue
|
||||
@@ -198,8 +199,38 @@ class PersistenceTest {
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `image content roundtrip through message payload`() = runTest {
|
||||
fun `error record roundtrip through audit log`() = runTest {
|
||||
val t0 = Instant.fromEpochMilliseconds(1_700_000_000_000)
|
||||
stores.messages.append(
|
||||
MessageRecord.Error(
|
||||
id = "e1",
|
||||
conversationId = "c1",
|
||||
message = "boom",
|
||||
code = "E42",
|
||||
createdAt = t0,
|
||||
),
|
||||
)
|
||||
stores.messages.append(
|
||||
MessageRecord.Error(
|
||||
id = "e2",
|
||||
conversationId = "c1",
|
||||
message = "no code",
|
||||
code = null,
|
||||
createdAt = Instant.fromEpochMilliseconds(1_700_000_001_000),
|
||||
),
|
||||
)
|
||||
val all = stores.messages.listAll("c1")
|
||||
assertEquals(2, all.size)
|
||||
val first = assertIs<MessageRecord.Error>(all[0])
|
||||
assertEquals("e1", first.id)
|
||||
assertEquals("boom", first.message)
|
||||
assertEquals("E42", first.code)
|
||||
assertEquals(t0, first.createdAt)
|
||||
assertNull(assertIs<MessageRecord.Error>(all[1]).code)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `image content roundtrip through message payload`() = runTest { val t0 = Instant.fromEpochMilliseconds(1_700_000_000_000)
|
||||
val bytes = byteArrayOf(0x89.toByte(), 0x50, 0x4E, 0x47) // PNG header
|
||||
stores.messages.append(
|
||||
MessageRecord.UserMessage(
|
||||
|
||||
Reference in New Issue
Block a user