diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index 5d87f4c..3ffb6e8 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -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. Слои персистентности diff --git a/docs/STANDALONE.md b/docs/STANDALONE.md index 2551787..bb282c3 100644 --- a/docs/STANDALONE.md +++ b/docs/STANDALONE.md @@ -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)) diff --git a/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Message.kt b/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Message.kt index 4a908c8..bfbb8fc 100644 --- a/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Message.kt +++ b/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Message.kt @@ -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 } diff --git a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/persistence/MessageRecord.kt b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/persistence/MessageRecord.kt index 273c8d3..55d9eae 100644 --- a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/persistence/MessageRecord.kt +++ b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/persistence/MessageRecord.kt @@ -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") diff --git a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ChatConversation.kt b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ChatConversation.kt index c7c100f..b299ba7 100644 --- a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ChatConversation.kt +++ b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ChatConversation.kt @@ -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, diff --git a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/persistence/sqlite/SqliteMessageStore.kt b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/persistence/sqlite/SqliteMessageStore.kt index c63f909..ec53367 100644 --- a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/persistence/sqlite/SqliteMessageStore.kt +++ b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/persistence/sqlite/SqliteMessageStore.kt @@ -62,6 +62,10 @@ private fun encodeRecord(record: MessageRecord): Pair = 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") } } diff --git a/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/ChatAgentTest.kt b/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/ChatAgentTest.kt index f63ac23..932a097 100644 --- a/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/ChatAgentTest.kt +++ b/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/ChatAgentTest.kt @@ -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() + 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(msgs[0]) + val err = assertIs(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? = null val conversations = mutableListOf() @@ -486,6 +518,9 @@ private class FakeLiteConversation( override fun sendStreamContents(contents: List): Flow { 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() {} diff --git a/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/persistence/PersistenceTest.kt b/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/persistence/PersistenceTest.kt index 47b3842..ca92cfc 100644 --- a/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/persistence/PersistenceTest.kt +++ b/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/persistence/PersistenceTest.kt @@ -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(all[0]) + assertEquals("e1", first.id) + assertEquals("boom", first.message) + assertEquals("E42", first.code) + assertEquals(t0, first.createdAt) + assertNull(assertIs(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(