standalone: token accounting per assistant turn (input/output → SQLite)

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.
This commit is contained in:
2026-09-15 05:45:40 +03:00
parent 9afa877e39
commit f4ce82957b
8 changed files with 321 additions and 18 deletions
+21
View File
@@ -200,6 +200,27 @@ OpenAI). Это происходит в фоне (`Dispatchers.IO`), основ
повторения. Используется как cheap "auto-improving prompt feedback" повторения. Используется как cheap "auto-improving prompt feedback"
без ручного переписывания system prompt. без ручного переписывания 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) ## Куратор памяти (Curator)
Фоновая корутина (запускается автоматически, если `AGENTIK_MEMORY_DIR != off`): Фоновая корутина (запускается автоматически, если `AGENTIK_MEMORY_DIR != off`):
@@ -4,10 +4,36 @@ import kotlinx.serialization.SerialName
import kotlinx.serialization.Serializable import kotlinx.serialization.Serializable
import kotlin.time.Instant 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). * Запись в таблице `message` (append-only audit) и `working_memory` (mutable view).
* *
* Используется sealed-иерархия: подтипы `User`/`Assistant`/`ToolCall`/`ToolResult` * Используется sealed-иерархию: подтипы `User`/`Assistant`/`ToolCall`/`ToolResult`
* живут и там, и там. `Summary`/`System` — только в `working_memory` * живут и там, и там. `Summary`/`System` — только в `working_memory`
* (синтетические строки, созданные при суммаризации или как system-prompt). * (синтетические строки, созданные при суммаризации или как system-prompt).
* *
@@ -53,6 +79,12 @@ sealed interface MessageRecord {
override val conversationId: String, override val conversationId: String,
override val content: List<Content>, override val content: List<Content>,
override val createdAt: Instant, override val createdAt: Instant,
/**
* Token usage этого turn'а: сколько input+output токенов обработала
* модель. Заполняется в [ChatConversation.runTurn] через
* `LiteConversation.tokenCount()` (до/после send).
*/
val tokens: TurnTokens? = null,
) : Body ) : Body
@Serializable @Serializable
@@ -2,6 +2,23 @@ package pw.binom.agentik.standalone.persistence
import kotlin.time.Instant 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). * Append-only audit log сообщений (`message` table).
* *
@@ -21,4 +38,10 @@ interface MessageStore : AutoCloseable {
/** Все сообщения диалога, отсортированные по `createdAt ASC` (для rebuild working memory). */ /** Все сообщения диалога, отсортированные по `createdAt ASC` (для rebuild working memory). */
suspend fun listAll(conversationId: String): List<MessageRecord> suspend fun listAll(conversationId: String): List<MessageRecord>
/**
* Суммарная token-статистика по диалогу: input/output/turns. Один проход
* по всем assistant-сообщениям. Дёшево (на практике < 1мс на SQLite).
*/
suspend fun tokenStats(conversationId: String): TokenStats
} }
@@ -7,16 +7,18 @@ import kotlinx.serialization.json.Json
/** /**
* JSON-формат для тел user/assistant сообщений: список [Content], опционально * JSON-формат для тел user/assistant сообщений: список [Content], опционально
* с [MessageContext] (для user — кто инициировал ход). * с [MessageContext] (для user — кто инициировал ход) и [TurnTokens]
* (для assistant — сколько токенов стоил этот turn).
* *
* Encoded-формат: * Encoded-формат:
* ``` * ```
* {"content": [ ...Content ], "context": {...MessageContext?}} * {"content": [ ...Content ], "context": {...MessageContext?}, "tokens": {...TurnTokens?}}
* ``` * ```
* *
* Backward compat: при чтении старых строк, где payload был просто * Backward compat: при чтении старых строк, где payload был просто
* `[ ... ]` (без обёртки), парсер падает на wrapper-формат и fallback'ит * `[ ... ]` (без обёртки), парсер падает на wrapper-формат и fallback'ит
* к `ListSerializer<Content>` — такие строки возвращаются с `context = null`. * к `ListSerializer<Content>` — такие строки возвращаются с `context = null`,
* `tokens = null`.
*/ */
private val bodyJson = Json { private val bodyJson = Json {
ignoreUnknownKeys = true ignoreUnknownKeys = true
@@ -29,35 +31,49 @@ internal data class MessageBodyPayload(
val content: List<Content>, val content: List<Content>,
@SerialName("context") @SerialName("context")
val context: MessageContext? = null, val context: MessageContext? = null,
/**
* Token usage для assistant (input + output). `null` для user-сообщений,
* для исторических assistant-сообщений без метрики и для бэкендов без
* tokenCount() (off-line модели).
*/
val tokens: TurnTokens? = null,
) )
/** /**
* Сериализует тело user (или assistant) сообщения в JSON-строку для * Сериализует тело user (или assistant) сообщения в JSON-строку для
* `payload_json` SQLite. Для user может нести [context] — кто инициировал ход. * `payload_json` SQLite. Для user может нести [context] — кто инициировал ход;
* для assistant может нести [tokens] — token usage этого turn'а.
*/ */
fun encodeBodyPayload(content: List<Content>, context: MessageContext? = null): String = fun encodeBodyPayload(
bodyJson.encodeToString( content: List<Content>,
MessageBodyPayload.serializer(), context: MessageContext? = null,
MessageBodyPayload(content = content, context = context), tokens: TurnTokens? = null,
) ): String = bodyJson.encodeToString(
MessageBodyPayload.serializer(),
MessageBodyPayload(content = content, context = context, tokens = tokens),
)
/** /**
* Десериализует тело сообщения: возвращает пару `(content, context)`. * Десериализует тело сообщения: возвращает тройку `(content, context, tokens)`.
* Контекст null если: * Контекст и токены — null если:
* - в строке нет поля `context` (новый формат, обычный user); * - поля отсутствуют в новом формате;
* - payload в старом plain-array формате (миграция не нужна — fallback). * - payload в старом plain-array формате (миграция не нужна — fallback).
*/ */
fun decodeBodyPayload(json: String): BodyDecoded = readPayload(json) fun decodeBodyPayload(json: String): BodyDecoded = readPayload(json)
data class BodyDecoded(val content: List<Content>, val context: MessageContext?) data class BodyDecoded(
val content: List<Content>,
val context: MessageContext?,
val tokens: TurnTokens? = null,
)
private fun readPayload(json: String): BodyDecoded { private fun readPayload(json: String): BodyDecoded {
return try { return try {
val p = bodyJson.decodeFromString(MessageBodyPayload.serializer(), json) val p = bodyJson.decodeFromString(MessageBodyPayload.serializer(), json)
BodyDecoded(p.content, p.context) BodyDecoded(p.content, p.context, p.tokens)
} catch (e: kotlinx.serialization.SerializationException) { } catch (e: kotlinx.serialization.SerializationException) {
// Старый формат: голый JSON-массив Content, без обёртки. // Старый формат: голый JSON-массив Content, без обёртки.
val arr = bodyJson.decodeFromString(ListSerializer(Content.serializer()), json) val arr = bodyJson.decodeFromString(ListSerializer(Content.serializer()), json)
BodyDecoded(arr, null) BodyDecoded(arr, null, null)
} }
} }
@@ -215,6 +215,25 @@ fun main() {
Runtime.getRuntime().addShutdownHook(Thread { curator.stop() }) 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})") 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 { Runtime.getRuntime().addShutdownHook(Thread {
agent.close() agent.close()
mcpRegistry.close() mcpRegistry.close()
@@ -32,6 +32,7 @@ import pw.binom.agentik.standalone.persistence.ConversationStore
import pw.binom.agentik.standalone.persistence.MessageContext import pw.binom.agentik.standalone.persistence.MessageContext
import pw.binom.agentik.standalone.persistence.MessageOrigin import pw.binom.agentik.standalone.persistence.MessageOrigin
import pw.binom.agentik.standalone.persistence.MessageRecord 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.MessageStore
import pw.binom.agentik.standalone.persistence.WorkingMemoryEntry import pw.binom.agentik.standalone.persistence.WorkingMemoryEntry
import pw.binom.agentik.standalone.persistence.WorkingMemoryRow import pw.binom.agentik.standalone.persistence.WorkingMemoryRow
@@ -269,6 +270,15 @@ class ChatConversation(
var currentParts: List<LiteContentPart> = initialParts var currentParts: List<LiteContentPart> = initialParts
var loopGuard = 0 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) { while (loopGuard++ < MAX_TOOL_LOOPS) {
val collectedCalls = mutableListOf<LiteToolCall>() val collectedCalls = mutableListOf<LiteToolCall>()
try { try {
@@ -306,6 +316,15 @@ class ChatConversation(
log.warn { "tool loop hit MAX_TOOL_LOOPS=$MAX_TOOL_LOOPS for $id — bailing" } 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 assistantId = newId("msg")
val assistantAt = now() val assistantAt = now()
val assistantContent = listOf(Content.Text(reply.toString())) val assistantContent = listOf(Content.Text(reply.toString()))
@@ -314,6 +333,7 @@ class ChatConversation(
conversationId = id, conversationId = id,
content = assistantContent, content = assistantContent,
createdAt = assistantAt, createdAt = assistantAt,
tokens = turnTokens,
) )
if (!record.isTemporal) { if (!record.isTemporal) {
@@ -915,3 +935,16 @@ internal fun MessageRecord.toProto(): ProtoMessage = when (this) {
content = listOf(ProtoContent.Text(body = text)), 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
}
@@ -3,6 +3,7 @@ package pw.binom.agentik.standalone.persistence.sqlite
import kotlinx.serialization.json.Json import kotlinx.serialization.json.Json
import pw.binom.agentik.standalone.persistence.MessageRecord import pw.binom.agentik.standalone.persistence.MessageRecord
import pw.binom.agentik.standalone.persistence.MessageStore 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.decodeBodyPayload
import pw.binom.agentik.standalone.persistence.encodeBodyPayload import pw.binom.agentik.standalone.persistence.encodeBodyPayload
import kotlin.time.Instant import kotlin.time.Instant
@@ -48,12 +49,42 @@ class SqliteMessageStore(private val db: AgentikDatabase) : MessageStore {
override suspend fun listAll(conversationId: String): List<MessageRecord> = override suspend fun listAll(conversationId: String): List<MessageRecord> =
q.listByConversationAll(conversation_id = conversationId).executeAsList().map { it.toRecord() } 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() {} override fun close() {}
} }
private fun encodeRecord(record: MessageRecord): Pair<String, String> = when (record) { private fun encodeRecord(record: MessageRecord): Pair<String, String> = when (record) {
is MessageRecord.UserMessage -> "user" to encodeBodyPayload(record.content, record.context) is MessageRecord.UserMessage -> "user" to encodeBodyPayload(
is MessageRecord.AssistantMessage -> "assistant" to encodeBodyPayload(record.content) 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( is MessageRecord.ToolCall -> "tool_call" to Json.encodeToString(
CallPayload.serializer(), CallPayload.serializer(),
CallPayload(name = record.toolName, title = record.toolTitle, argsJson = record.toolArgsJson), CallPayload(name = record.toolName, title = record.toolTitle, argsJson = record.toolArgsJson),
@@ -101,6 +132,7 @@ private fun Message.toRecord(): MessageRecord {
conversationId = convId, conversationId = convId,
content = decoded.content, content = decoded.content,
createdAt = createdAt, createdAt = createdAt,
tokens = decoded.tokens,
) )
} }
"tool_call" -> { "tool_call" -> {
@@ -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),
)
}