diff options
| author | Adam Malczewski <[email protected]> | 2026-06-10 08:29:59 +0900 |
|---|---|---|
| committer | Adam Malczewski <[email protected]> | 2026-06-10 08:29:59 +0900 |
| commit | 6db12ff70acb3333d05a5020ab66da4172a5225a (patch) | |
| tree | de5cc6314a3a6dd966d7c4fdb9b20adb04ae8307 /packages/transport-http/src | |
| parent | 4248cd1d546a4c1fb4e68940c11b5e309c2c2736 (diff) | |
| download | dispatch-6db12ff70acb3333d05a5020ab66da4172a5225a.tar.gz dispatch-6db12ff70acb3333d05a5020ab66da4172a5225a.zip | |
feat(metrics): durable per-turn/step token+timing metrics (observability spans + persisted replay)
Two-part token-data improvement:
#2 Observability spans (kernel run-turn): turn & step span-close now stamp
ALL four Usage fields — added usage.cacheReadTokens/cacheWriteTokens (were
silently dropped) and normalized usage_* -> usage.* to match the
provider.request span (consistent D9 GROUP BY). No contract change.
#3 Persisted replay metrics (conversation-store + read endpoint): new
StepMetrics/TurnMetrics wire types; conversation-store persists per-turn
metrics in a separate key space (appendMetrics/loadMetrics, turn-append
order); session-orchestrator accumulates per-step+turn metrics from the
event stream (pure metrics.ts) and persists after seal; transport-http
serves GET /conversations/:id/metrics -> ConversationMetricsResponse.
Contracts: @dispatch/wire + @dispatch/transport-contract bumped 0.3.0->0.4.0
(additive). GLOSSARY: turn metrics / step metrics.
typecheck EXIT 0, biome clean, 546 vitest + 89 bun = 635 tests.
Diffstat (limited to 'packages/transport-http/src')
| -rw-r--r-- | packages/transport-http/src/app.test.ts | 128 | ||||
| -rw-r--r-- | packages/transport-http/src/app.ts | 23 | ||||
| -rw-r--r-- | packages/transport-http/src/server.bun.test.ts | 4 |
3 files changed, 153 insertions, 2 deletions
diff --git a/packages/transport-http/src/app.test.ts b/packages/transport-http/src/app.test.ts index aa47dce..0a6c5b0 100644 --- a/packages/transport-http/src/app.test.ts +++ b/packages/transport-http/src/app.test.ts @@ -1,4 +1,4 @@ -import type { AgentEvent, Logger, StoredChunk } from "@dispatch/kernel"; +import type { AgentEvent, Logger, StepId, StoredChunk, TurnMetrics } from "@dispatch/kernel"; import { describe, expect, it } from "vitest"; import { createApp } from "./app.js"; import type { ConversationStore, CredentialStore, SessionOrchestrator } from "./seam.js"; @@ -47,6 +47,7 @@ function createFakeLogger(): Logger & { readonly records: readonly CapturedLog[] function createFakeConversationStore( store: Map<string, StoredChunk[]> = new Map(), + metricsStore: Map<string, TurnMetrics[]> = new Map(), ): ConversationStore { return { async append() {}, @@ -58,6 +59,10 @@ function createFakeConversationStore( const minSeq = sinceSeq ?? 0; return chunks.filter((c) => c.seq > minSeq); }, + async appendMetrics() {}, + async loadMetrics(conversationId) { + return metricsStore.get(conversationId) ?? []; + }, }; } @@ -465,6 +470,127 @@ describe("GET /conversations/:id", () => { }); }); +describe("GET /conversations/:id/metrics", () => { + const sampleMetrics: TurnMetrics[] = [ + { + turnId: "turn1", + usage: { inputTokens: 100, outputTokens: 50, cacheReadTokens: 0, cacheWriteTokens: 0 }, + durationMs: 1000, + steps: [ + { + stepId: "step1" as StepId, + usage: { inputTokens: 100, outputTokens: 50, cacheReadTokens: 0, cacheWriteTokens: 0 }, + ttftMs: 200, + decodeMs: 300, + genTotalMs: 500, + }, + ], + }, + { + turnId: "turn2", + usage: { inputTokens: 200, outputTokens: 80, cacheReadTokens: 10, cacheWriteTokens: 5 }, + durationMs: 1500, + steps: [ + { + stepId: "step2" as StepId, + usage: { inputTokens: 200, outputTokens: 80, cacheReadTokens: 10, cacheWriteTokens: 5 }, + ttftMs: 300, + decodeMs: 500, + genTotalMs: 800, + }, + ], + }, + ]; + + it("returns persisted turn metrics as { turns }", async () => { + const metricsStore = new Map<string, TurnMetrics[]>([["conv1", sampleMetrics]]); + const app = createApp({ + conversationStore: createFakeConversationStore(new Map(), metricsStore), + orchestrator: createFakeOrchestrator([]), + credentialStore: createFakeCredentialStore([]), + }); + + const res = await app.request("/conversations/conv1/metrics"); + expect(res.status).toBe(200); + const body = (await res.json()) as { turns: readonly TurnMetrics[] }; + expect(body.turns).toHaveLength(2); + expect(body.turns[0]?.turnId).toBe("turn1"); + expect(body.turns[1]?.turnId).toBe("turn2"); + }); + + it("returns { turns: [] } for an unknown conversation", async () => { + const app = createApp({ + conversationStore: createFakeConversationStore(), + orchestrator: createFakeOrchestrator([]), + credentialStore: createFakeCredentialStore([]), + }); + + const res = await app.request("/conversations/unknown/metrics"); + expect(res.status).toBe(200); + const body = (await res.json()) as { turns: readonly TurnMetrics[] }; + expect(body.turns).toHaveLength(0); + }); + + it("the metrics route does not collide with GET /conversations/:id history route", async () => { + const sampleChunks: StoredChunk[] = [ + { seq: 1, role: "user", chunk: { type: "text", text: "hello" } }, + ]; + const store = new Map<string, StoredChunk[]>([["conv1", sampleChunks]]); + const metricsStore = new Map<string, TurnMetrics[]>([["conv1", sampleMetrics]]); + const app = createApp({ + conversationStore: createFakeConversationStore(store, metricsStore), + orchestrator: createFakeOrchestrator([]), + credentialStore: createFakeCredentialStore([]), + }); + + const metricsRes = await app.request("/conversations/conv1/metrics"); + expect(metricsRes.status).toBe(200); + const metricsBody = (await metricsRes.json()) as { turns: readonly TurnMetrics[] }; + expect(metricsBody.turns).toHaveLength(2); + + const historyRes = await app.request("/conversations/conv1"); + expect(historyRes.status).toBe(200); + const historyBody = (await historyRes.json()) as { + chunks: readonly StoredChunk[]; + latestSeq: number; + }; + expect(historyBody.chunks).toHaveLength(1); + }); + + it("a store failure on the metrics read returns an error status + logs an error", async () => { + const logger = createFakeLogger(); + const brokenStore: ConversationStore = { + async append() {}, + async load() { + return []; + }, + async loadSince() { + return []; + }, + async appendMetrics() {}, + async loadMetrics() { + throw new Error("storage exploded"); + }, + }; + const app = createApp({ + conversationStore: brokenStore, + orchestrator: createFakeOrchestrator([]), + credentialStore: createFakeCredentialStore([]), + logger, + }); + + const res = await app.request("/conversations/conv1/metrics"); + expect(res.status).toBe(500); + const body = (await res.json()) as { error: string }; + expect(body.error).toContain("Failed to load conversation metrics"); + + const errorLogs = logger.records.filter((r) => r.level === "error"); + expect(errorLogs).toHaveLength(1); + expect(errorLogs[0]?.msg).toBe("conversations: metrics store failure"); + expect(errorLogs[0]?.attrs?.err).toBeInstanceOf(Error); + }); +}); + describe("POST /chat logging", () => { it("POST /chat logs an info line when a request is accepted", async () => { const logger = createFakeLogger(); diff --git a/packages/transport-http/src/app.ts b/packages/transport-http/src/app.ts index a8edc71..4002e23 100644 --- a/packages/transport-http/src/app.ts +++ b/packages/transport-http/src/app.ts @@ -1,5 +1,9 @@ import type { AgentEvent, Logger } from "@dispatch/kernel"; -import type { ConversationHistoryResponse, ModelsResponse } from "@dispatch/transport-contract"; +import type { + ConversationHistoryResponse, + ConversationMetricsResponse, + ModelsResponse, +} from "@dispatch/transport-contract"; import { Hono } from "hono"; import { cors } from "hono/cors"; import { @@ -57,6 +61,23 @@ export function createApp(opts: CreateServerOptions): Hono { app.get("/health", (c) => c.json({ ok: true })); + app.get("/conversations/:id/metrics", async (c) => { + const conversationId = c.req.param("id"); + + try { + const turns = await opts.conversationStore.loadMetrics(conversationId); + log.info("conversations: metrics read", { + conversationId, + count: turns.length, + }); + const body: ConversationMetricsResponse = { turns }; + return c.json(body, 200); + } catch (err) { + log.error("conversations: metrics store failure", { err }); + return c.json({ error: "Failed to load conversation metrics" }, 500); + } + }); + app.get("/conversations/:id", async (c) => { const conversationId = c.req.param("id"); const sinceSeqResult = parseSinceSeq(c.req.query("sinceSeq")); diff --git a/packages/transport-http/src/server.bun.test.ts b/packages/transport-http/src/server.bun.test.ts index e824d18..b2a978a 100644 --- a/packages/transport-http/src/server.bun.test.ts +++ b/packages/transport-http/src/server.bun.test.ts @@ -37,6 +37,10 @@ function fakeConversationStore(): ConversationStore { async loadSince() { return []; }, + async appendMetrics() {}, + async loadMetrics() { + return []; + }, }; } |
