From f4ce82957b8f26bb5881080ed412302e67e2d1ec Mon Sep 17 00:00:00 2001 From: subochev Date: Tue, 15 Sep 2026 05:45:40 +0300 Subject: [PATCH] =?UTF-8?q?standalone:=20token=20accounting=20per=20assist?= =?UTF-8?q?ant=20turn=20(input/output=20=E2=86=92=20SQLite)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit LiteLlm API не отдаёт split prompt/completion наружу через send() (внутренний OpenAI Usage сидит в pw.binom.litert.openai и недоступен), поэтому измеряем через LiteConversation.tokenCount(): input = tokenCount() до первого send в turn'е (= system + вся история + tools + только что добавленное user-сообщение) output = tokenCount() после завершения turn'а - input (= assistant text + tool calls + tool results за tool loop) Пишем в assistant-запись как TurnTokens(input, output) в payload_json. Никаких schema-миграций: payload-формат уже обёрнут в MessageBodyPayload, просто добавлено опциональное поле tokens. MessageStore.tokenStats(conversationId) → TokenStats(turns, inputTokens, outputTokens). На старте агент печатает сводку по всем диалогам: tokens: 17 convs, 134 turns, in=523844, out=58290, total=582134 Бэкенды без tokenCount() (off-line LiteRT-LM модели) → tokens=null, старые assistant-записи без метрики → пропускаются в tokenStats без ошибок. Tests: TokenStatsTest (5 green) + PersistenceTest (unchanged) → 165 total. Backward compat: legacy plain-array payload всё ещё читается, tokens=null. --- standalone/README.md | 21 +++ .../standalone/persistence/MessageRecord.kt | 34 ++++- .../standalone/persistence/MessageStore.kt | 23 ++++ .../agentik/standalone/persistence/Payload.kt | 46 ++++--- .../pw/binom/agentik/standalone/Main.kt | 19 +++ .../standalone/agent/ChatConversation.kt | 33 +++++ .../persistence/sqlite/SqliteMessageStore.kt | 36 ++++- .../standalone/persistence/TokenStatsTest.kt | 127 ++++++++++++++++++ 8 files changed, 321 insertions(+), 18 deletions(-) create mode 100644 standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/persistence/TokenStatsTest.kt diff --git a/standalone/README.md b/standalone/README.md index c7d0770..f076f01 100644 --- a/standalone/README.md +++ b/standalone/README.md @@ -200,6 +200,27 @@ OpenAI). Это происходит в фоне (`Dispatchers.IO`), основ повторения. Используется как cheap "auto-improving prompt feedback" без ручного переписывания system prompt. +## Учёт токенов (token accounting) + +Каждый assistant-ход после LiteLlm.send помечает assistant-запись `TurnTokens(input, output)`: + +- **`input`** — снимок `LiteConversation.tokenCount()` **до** первого send в turn'е + (system + вся история + tools + только что добавленное user-сообщение). +- **`output`** — дельта после завершения turn'а (assistant text + tool calls + + tool results, всё что LiteConversation добавила за весь tool loop). +- Хранится в `payload_json` assistant-сообщения (без schema-миграций). Бэкенды + без `tokenCount()` (off-line модели LiteRT-LM счётчик не отдают) дают `tokens=null`. + +На старте агент печатает сводку по всем существующим диалогам: + +``` +tokens: 17 convs, 134 turns, in=523844, out=58290, total=582134 +``` + +`MessageStore.tokenStats(conversationId)` отдаёт `TokenStats(turns, inputTokens, outputTokens)` +для одного диалога — можно использовать из HTTP фасада или клиентских дашбордов +для оценки cost. + ## Куратор памяти (Curator) Фоновая корутина (запускается автоматически, если `AGENTIK_MEMORY_DIR != off`): 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 74db00d..82373f1 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 @@ -4,10 +4,36 @@ import kotlinx.serialization.SerialName import kotlinx.serialization.Serializable import kotlin.time.Instant +/** + * Token usage одного assistant turn'а: сколько токенов модель обработала + * на входе (system + history + tools + user message) и сколько сгенерировала + * (assistant text + tool calls + tool results, всё что LiteConversation + * добавила к истории за этот turn). + * + * `input` — снимок [LiteConversation.tokenCount] перед первым send() в turn'е + * (после подготовки user-сообщения). `output` — дельта после завершения turn'а + * (включая все tool loop итерации). + * + * Persisted в `message.payload_json` — никаких schema-миграций при добавлении + * полей. Optional: `null` для исторических сообщений или для бэкендов, не + * отдающих tokenCount (например off-line embedded LLM без контекст-счётчика). + */ +@Serializable +data class TurnTokens( + val input: Int, + val output: Int, +) { + val total: Int get() = input + output + init { + require(input >= 0) { "input tokens must be non-negative, got $input" } + require(output >= 0) { "output tokens must be non-negative, got $output" } + } +} + /** * Запись в таблице `message` (append-only audit) и `working_memory` (mutable view). * - * Используется sealed-иерархия: подтипы `User`/`Assistant`/`ToolCall`/`ToolResult` + * Используется sealed-иерархию: подтипы `User`/`Assistant`/`ToolCall`/`ToolResult` * живут и там, и там. `Summary`/`System` — только в `working_memory` * (синтетические строки, созданные при суммаризации или как system-prompt). * @@ -53,6 +79,12 @@ sealed interface MessageRecord { override val conversationId: String, override val content: List, override val createdAt: Instant, + /** + * Token usage этого turn'а: сколько input+output токенов обработала + * модель. Заполняется в [ChatConversation.runTurn] через + * `LiteConversation.tokenCount()` (до/после send). + */ + val tokens: TurnTokens? = null, ) : Body @Serializable diff --git a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/persistence/MessageStore.kt b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/persistence/MessageStore.kt index a5daa72..44f9c10 100644 --- a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/persistence/MessageStore.kt +++ b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/persistence/MessageStore.kt @@ -2,6 +2,23 @@ package pw.binom.agentik.standalone.persistence import kotlin.time.Instant +/** + * Суммарная статистика токенов диалога — aggregate по всем assistant-сообщениям + * в audit log. Делит input/output и считает число assistant-ходов. + * + * Используется: + * - В startup banner'е агента (см. `Main.kt` → "conversation stats"). + * - На HTTP фасаде `/agentik/conversations/{id}/stats` (если будет endpoint). + * - В клиентских дашбордах для оценки cost. + */ +data class TokenStats( + val turns: Int, + val inputTokens: Long, + val outputTokens: Long, +) { + val totalTokens: Long get() = inputTokens + outputTokens +} + /** * Append-only audit log сообщений (`message` table). * @@ -21,4 +38,10 @@ interface MessageStore : AutoCloseable { /** Все сообщения диалога, отсортированные по `createdAt ASC` (для rebuild working memory). */ suspend fun listAll(conversationId: String): List + + /** + * Суммарная token-статистика по диалогу: input/output/turns. Один проход + * по всем assistant-сообщениям. Дёшево (на практике < 1мс на SQLite). + */ + suspend fun tokenStats(conversationId: String): TokenStats } diff --git a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/persistence/Payload.kt b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/persistence/Payload.kt index ff5a7bc..a1a66a8 100644 --- a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/persistence/Payload.kt +++ b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/persistence/Payload.kt @@ -7,16 +7,18 @@ import kotlinx.serialization.json.Json /** * JSON-формат для тел user/assistant сообщений: список [Content], опционально - * с [MessageContext] (для user — кто инициировал ход). + * с [MessageContext] (для user — кто инициировал ход) и [TurnTokens] + * (для assistant — сколько токенов стоил этот turn). * * Encoded-формат: * ``` - * {"content": [ ...Content ], "context": {...MessageContext?}} + * {"content": [ ...Content ], "context": {...MessageContext?}, "tokens": {...TurnTokens?}} * ``` * * Backward compat: при чтении старых строк, где payload был просто * `[ ... ]` (без обёртки), парсер падает на wrapper-формат и fallback'ит - * к `ListSerializer` — такие строки возвращаются с `context = null`. + * к `ListSerializer` — такие строки возвращаются с `context = null`, + * `tokens = null`. */ private val bodyJson = Json { ignoreUnknownKeys = true @@ -29,35 +31,49 @@ internal data class MessageBodyPayload( val content: List, @SerialName("context") val context: MessageContext? = null, + /** + * Token usage для assistant (input + output). `null` для user-сообщений, + * для исторических assistant-сообщений без метрики и для бэкендов без + * tokenCount() (off-line модели). + */ + val tokens: TurnTokens? = null, ) /** * Сериализует тело user (или assistant) сообщения в JSON-строку для - * `payload_json` SQLite. Для user может нести [context] — кто инициировал ход. + * `payload_json` SQLite. Для user может нести [context] — кто инициировал ход; + * для assistant может нести [tokens] — token usage этого turn'а. */ -fun encodeBodyPayload(content: List, context: MessageContext? = null): String = - bodyJson.encodeToString( - MessageBodyPayload.serializer(), - MessageBodyPayload(content = content, context = context), - ) +fun encodeBodyPayload( + content: List, + context: MessageContext? = null, + tokens: TurnTokens? = null, +): String = bodyJson.encodeToString( + MessageBodyPayload.serializer(), + MessageBodyPayload(content = content, context = context, tokens = tokens), +) /** - * Десериализует тело сообщения: возвращает пару `(content, context)`. - * Контекст null если: - * - в строке нет поля `context` (новый формат, обычный user); + * Десериализует тело сообщения: возвращает тройку `(content, context, tokens)`. + * Контекст и токены — null если: + * - поля отсутствуют в новом формате; * - payload в старом plain-array формате (миграция не нужна — fallback). */ fun decodeBodyPayload(json: String): BodyDecoded = readPayload(json) -data class BodyDecoded(val content: List, val context: MessageContext?) +data class BodyDecoded( + val content: List, + val context: MessageContext?, + val tokens: TurnTokens? = null, +) private fun readPayload(json: String): BodyDecoded { return try { val p = bodyJson.decodeFromString(MessageBodyPayload.serializer(), json) - BodyDecoded(p.content, p.context) + BodyDecoded(p.content, p.context, p.tokens) } catch (e: kotlinx.serialization.SerializationException) { // Старый формат: голый JSON-массив Content, без обёртки. val arr = bodyJson.decodeFromString(ListSerializer(Content.serializer()), json) - BodyDecoded(arr, null) + BodyDecoded(arr, null, null) } } diff --git a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/Main.kt b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/Main.kt index 19f7b0a..cf33840 100644 --- a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/Main.kt +++ b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/Main.kt @@ -215,6 +215,25 @@ fun main() { Runtime.getRuntime().addShutdownHook(Thread { curator.stop() }) println(" curator: enabled (interval=${pw.binom.agentik.standalone.agent.memory.Curator.DEFAULT_INTERVAL}, maxAge=${pw.binom.agentik.standalone.agent.memory.Curator.DEFAULT_MAX_AGE})") } + // Token stats по существующим диалогам (агрегат на старте — каждая запись + // парсится из payload_json, ну >100 turns и БД приличная — но в рамках + // стартапа это терпимо). + val existingConvs = kotlinx.coroutines.runBlocking { stores.conversations.list(offset = 0, limit = 1000) } + if (existingConvs.isNotEmpty()) { + var totalTurns = 0 + var totalIn = 0L + var totalOut = 0L + for (c in existingConvs) { + if (c.isTemporal) continue + val s = kotlinx.coroutines.runBlocking { stores.messages.tokenStats(c.id) } + totalTurns += s.turns + totalIn += s.inputTokens + totalOut += s.outputTokens + } + if (totalTurns > 0) { + println(" tokens: ${existingConvs.size} convs, $totalTurns turns, in=${totalIn}, out=${totalOut}, total=${totalIn + totalOut}") + } + } Runtime.getRuntime().addShutdownHook(Thread { agent.close() mcpRegistry.close() 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 1e0b3dc..7f9a502 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 @@ -32,6 +32,7 @@ import pw.binom.agentik.standalone.persistence.ConversationStore import pw.binom.agentik.standalone.persistence.MessageContext import pw.binom.agentik.standalone.persistence.MessageOrigin import pw.binom.agentik.standalone.persistence.MessageRecord +import pw.binom.agentik.standalone.persistence.TurnTokens import pw.binom.agentik.standalone.persistence.MessageStore import pw.binom.agentik.standalone.persistence.WorkingMemoryEntry import pw.binom.agentik.standalone.persistence.WorkingMemoryRow @@ -269,6 +270,15 @@ class ChatConversation( var currentParts: List = initialParts var loopGuard = 0 + // Token accounting: снимаем tokenCount() ДО первого send (это будет + // наш `input` для этого turn'а — вся история conversation + только что + // добавленное user-сообщение). ПОСЛЕ цикла снимаем ещё раз — дельта + // даёт нам `output` (то что добавила модель: assistant text + tool + // call args + tool results, естественно накопленные за tool loop). + // Если tokenCount() не поддерживается бэкендом или кидает — tokens останется null. + val tokensAtTurnStart: Int? = readTokenCount(liteConv) + var turnTokens: TurnTokens? = null + while (loopGuard++ < MAX_TOOL_LOOPS) { val collectedCalls = mutableListOf() try { @@ -306,6 +316,15 @@ class ChatConversation( log.warn { "tool loop hit MAX_TOOL_LOOPS=$MAX_TOOL_LOOPS for $id — bailing" } } + // Считаем дельту после цикла (defensive: turnTokens может остаться null). + if (tokensAtTurnStart != null) { + val tokensAtTurnEnd = readTokenCount(liteConv) + if (tokensAtTurnEnd != null) { + val output = (tokensAtTurnEnd - tokensAtTurnStart).coerceAtLeast(0) + turnTokens = TurnTokens(input = tokensAtTurnStart, output = output) + } + } + val assistantId = newId("msg") val assistantAt = now() val assistantContent = listOf(Content.Text(reply.toString())) @@ -314,6 +333,7 @@ class ChatConversation( conversationId = id, content = assistantContent, createdAt = assistantAt, + tokens = turnTokens, ) if (!record.isTemporal) { @@ -915,3 +935,16 @@ internal fun MessageRecord.toProto(): ProtoMessage = when (this) { content = listOf(ProtoContent.Text(body = text)), ) } + +/** + * Defensive-обёртка вокруг [LiteConversation.tokenCount]: некоторые бэкенды + * (например LiteRT-LM с off-line моделью без контекст-счётчика) могут + * кидать или возвращать невалидное значение. Возвращаем `null` в таких + * случаях — лучше не иметь token-accounting'а за один turn, чем сломать turn. + */ +private fun readTokenCount(liteConv: LiteConversation): Int? = try { + val n = liteConv.tokenCount() + if (n < 0) null else n +} catch (_: Throwable) { + null +} 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 1fc35a0..d7a1540 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 @@ -3,6 +3,7 @@ package pw.binom.agentik.standalone.persistence.sqlite import kotlinx.serialization.json.Json import pw.binom.agentik.standalone.persistence.MessageRecord import pw.binom.agentik.standalone.persistence.MessageStore +import pw.binom.agentik.standalone.persistence.TokenStats import pw.binom.agentik.standalone.persistence.decodeBodyPayload import pw.binom.agentik.standalone.persistence.encodeBodyPayload import kotlin.time.Instant @@ -48,12 +49,42 @@ class SqliteMessageStore(private val db: AgentikDatabase) : MessageStore { override suspend fun listAll(conversationId: String): List = q.listByConversationAll(conversation_id = conversationId).executeAsList().map { it.toRecord() } + override suspend fun tokenStats(conversationId: String): TokenStats { + // Загружаем все assistant-сообщения и считаем локально. Для >10K ходов + // это можно оптимизировать SQL aggregation с JSON_EXTRACT, но пока + // узких мест нет — реальные диалоги редко длиннее нескольких сотен ходов. + val messages = q.listByConversationAll(conversation_id = conversationId).executeAsList() + var turns = 0 + var inputTotal = 0L + var outputTotal = 0L + for (m in messages) { + if (m.kind != "assistant") continue + val decoded = try { + decodeBodyPayload(m.payload_json) + } catch (_: Throwable) { + // skip malformed legacy payloads + continue + } + val tokens = decoded.tokens ?: continue + turns++ + inputTotal += tokens.input + outputTotal += tokens.output + } + return TokenStats(turns = turns, inputTokens = inputTotal, outputTokens = outputTotal) + } + override fun close() {} } private fun encodeRecord(record: MessageRecord): Pair = when (record) { - is MessageRecord.UserMessage -> "user" to encodeBodyPayload(record.content, record.context) - is MessageRecord.AssistantMessage -> "assistant" to encodeBodyPayload(record.content) + is MessageRecord.UserMessage -> "user" to encodeBodyPayload( + content = record.content, + context = record.context, + ) + is MessageRecord.AssistantMessage -> "assistant" to encodeBodyPayload( + content = record.content, + tokens = record.tokens, + ) is MessageRecord.ToolCall -> "tool_call" to Json.encodeToString( CallPayload.serializer(), CallPayload(name = record.toolName, title = record.toolTitle, argsJson = record.toolArgsJson), @@ -101,6 +132,7 @@ private fun Message.toRecord(): MessageRecord { conversationId = convId, content = decoded.content, createdAt = createdAt, + tokens = decoded.tokens, ) } "tool_call" -> { diff --git a/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/persistence/TokenStatsTest.kt b/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/persistence/TokenStatsTest.kt new file mode 100644 index 0000000..b57a802 --- /dev/null +++ b/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/persistence/TokenStatsTest.kt @@ -0,0 +1,127 @@ +package pw.binom.agentik.standalone.persistence + +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertNull +import kotlin.time.Instant +import pw.binom.agentik.standalone.persistence.sqlite.SqliteStores + +class TokenStatsTest { + + @Test + fun `tokenStats sums across assistant messages of same conversation`() { + val stores = SqliteStores.inMemory() + try { + kotlinx.coroutines.runBlocking { + stores.conversations.upsert( + ConversationRecord( + id = "c1", + title = null, + isTemporal = false, + createdAt = Instant.parse("2026-09-15T12:00:00Z"), + updatedAt = Instant.parse("2026-09-15T12:00:00Z"), + ) + )} + kotlinx.coroutines.runBlocking { + stores.messages.append( + assistant("a1", "c1", 100, 50, Instant.parse("2026-09-15T12:01:00Z")) + ) + stores.messages.append( + assistant("a2", "c1", 200, 80, Instant.parse("2026-09-15T12:02:00Z")) + ) + // User без tokens — не должны считаться. + stores.messages.append( + MessageRecord.UserMessage( + id = "u1", + conversationId = "c1", + content = listOf(Content.Text("hi")), + createdAt = Instant.parse("2026-09-15T12:00:30Z"), + ) + ) + } + val stats = kotlinx.coroutines.runBlocking { stores.messages.tokenStats("c1") } + assertEquals(2, stats.turns) + assertEquals(300L, stats.inputTokens) + assertEquals(130L, stats.outputTokens) + assertEquals(430L, stats.totalTokens) + } finally { stores.close() } + } + + @Test + fun `tokenStats skips legacy assistant messages without tokens field`() { + val stores = SqliteStores.inMemory() + try { + // Вставляем запись со СТАРЫМ payload-форматом (plain array через + // прямой SQL апдейт, минуя типизированный encoder). + val id = "legacy" + val msgId = "m1" + stores.driver.execute(null, "INSERT INTO conversation(id, title, is_temporal, created_at, updated_at) VALUES('c1', NULL, 0, 0, 0)", 0) + stores.driver.execute( + null, + "INSERT INTO message(id, conversation_id, kind, payload_json, created_at) VALUES(?, ?, ?, ?, ?)", + 5, + ) { + bindString(0, msgId) + bindString(1, "c1") + bindString(2, "assistant") + bindString(3, """[{"kind":"text","body":"legacy"}]""") + bindLong(4, Instant.parse("2026-09-15T12:00:00Z").toEpochMilliseconds()) + } + val stats = kotlinx.coroutines.runBlocking { stores.messages.tokenStats("c1") } + assertEquals(0, stats.turns, "legacy без tokens не должен считаться") + assertEquals(0L, stats.inputTokens) + assertEquals(0L, stats.outputTokens) + } finally { stores.close() } + } + + @Test + fun `TurnTokens rejects negative values`() { + val tokens = TurnTokens(input = 100, output = 50) + assertEquals(100, tokens.input) + assertEquals(150, tokens.total) + // Sanity для init{} + try { + TurnTokens(input = -1, output = 50) + error("should have thrown") + } catch (_: IllegalArgumentException) {} + } + + @Test + fun `backwards compatible encode-decode round trip preserves both context and tokens`() { + val original = MessageBodyPayload( + content = listOf(Content.Text("hi")), + context = MessageContext( + origin = MessageOrigin.USER, + description = "x", + ), + tokens = TurnTokens(input = 100, output = 50), + ) + val json = kotlinx.serialization.json.Json.encodeToString(MessageBodyPayload.serializer(), original) + val decoded = kotlinx.serialization.json.Json.decodeFromString(MessageBodyPayload.serializer(), json) + assertEquals(100, decoded.tokens?.input) + assertEquals(50, decoded.tokens?.output) + assertEquals(1, decoded.content.size) + } + + @Test + fun `decodeBodyPayload round-trip preserves tokens`() { + val tokens = TurnTokens(input = 250, output = 80) + val encoded = encodeBodyPayload( + content = listOf(Content.Text("ok")), + tokens = tokens, + ) + val decoded = decodeBodyPayload(encoded) + assertEquals(250, decoded.tokens?.input) + assertEquals(80, decoded.tokens?.output) + } + + private fun assistant( + msgId: String, convId: String, input: Int, output: Int, at: Instant, + ) = MessageRecord.AssistantMessage( + id = msgId, + conversationId = convId, + content = listOf(Content.Text("reply")), + createdAt = at, + tokens = TurnTokens(input = input, output = output), + ) +}