summaryrefslogtreecommitdiffhomepage
path: root/packages/transport-http/src
diff options
context:
space:
mode:
authorAdam Malczewski <[email protected]>2026-06-10 08:29:59 +0900
committerAdam Malczewski <[email protected]>2026-06-10 08:29:59 +0900
commit6db12ff70acb3333d05a5020ab66da4172a5225a (patch)
treede5cc6314a3a6dd966d7c4fdb9b20adb04ae8307 /packages/transport-http/src
parent4248cd1d546a4c1fb4e68940c11b5e309c2c2736 (diff)
downloaddispatch-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.ts128
-rw-r--r--packages/transport-http/src/app.ts23
-rw-r--r--packages/transport-http/src/server.bun.test.ts4
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 [];
+ },
};
}