feat(web): add durable session metrics (round 1)

This commit is contained in:
Hypatia May
2026-07-28 12:10:48 +08:00
parent 2a46685414
commit 9ca0241d5d
43 changed files with 1139 additions and 97 deletions

View File

@@ -40,6 +40,7 @@ import type {
} from '@deepseek-ai/dsh-user-interaction'
import { UserInteractionError } from '@deepseek-ai/dsh-user-interaction'
import { pickNativeDirectory } from './native-directory-picker.ts'
import { affectsSessionMetrics, SessionMetricsProjector } from './session-metrics.ts'
/** Page size when history is called without maxMessages. */
const DEFAULT_MAX_MESSAGES = 50
@@ -417,6 +418,54 @@ export function createApiProxy(ctx: Context, defaults: ApiProxyDefaults): ApiPro
for (const queue of muxQueues) queue.push(envelope)
}
const pendingMetricSessions = new Set<Session>()
let metricFlushScheduled = false
let metricsDisposed = false
const metricsProjector = new SessionMetricsProjector(
ctx,
agent => targetFor(agent).current,
(agent) => { scheduleMetrics(agent.session) },
)
/** Queue one full-log metrics publication after synchronous session listeners drain. */
function scheduleMetrics(session: Session): void {
if (metricsDisposed || muxQueues.size === 0) return
pendingMetricSessions.add(session)
if (metricFlushScheduled) return
metricFlushScheduled = true
queueMicrotask(() => {
metricFlushScheduled = false
if (metricsDisposed) {
pendingMetricSessions.clear()
return
}
const sessions = [...pendingMetricSessions]
pendingMetricSessions.clear()
for (const current of sessions) {
broadcast({
type: 'session/metrics',
sessionId: current.id,
metrics: metricsProjector.snapshot(current, ctx.agents.get(current.id)),
})
}
})
}
ctx.effect(() => {
const disposers = [
ctx.on('session/event', (session: Session, event: SessionEvent) => {
if (affectsSessionMetrics(event)) scheduleMetrics(session)
}),
ctx.on('agent/created', (agent: Agent) => { scheduleMetrics(agent.session) }),
ctx.on('session/disposed', (session: Session) => { pendingMetricSessions.delete(session) }),
]
return () => {
metricsDisposed = true
pendingMetricSessions.clear()
for (const dispose of disposers) dispose()
}
}, 'api-proxy: session metrics')
/**
* Per-session inbox mirror serving the mux-open queue snapshot (the same
* refresh-recovery baseline as pending questions). Keyed by the stable
@@ -723,7 +772,15 @@ export function createApiProxy(ctx: Context, defaults: ApiProxyDefaults): ApiPro
// log (the page window may not contain the last todo/write; a paged
// client cannot reconstruct session-level state from it).
const todos = beforeSeq === undefined ? backscanTodos(found.agent.session.events) : undefined
return ok(request, { events: entries, hasMore: page.hasMore, ...todos === undefined ? {} : { todos } })
const metrics = beforeSeq === undefined
? metricsProjector.snapshot(found.agent.session, found.agent)
: undefined
return ok(request, {
events: entries,
hasMore: page.hasMore,
...todos === undefined ? {} : { todos },
...metrics === undefined ? {} : { metrics },
})
},
async models(request) {
@@ -817,6 +874,11 @@ export function createApiProxy(ctx: Context, defaults: ApiProxyDefaults): ApiPro
: { reasoningEffort: resolved.reasoningEffort },
}
targetFor(found.agent).current = selected
broadcast({
type: 'session/metrics',
sessionId: found.agent.session.id,
metrics: metricsProjector.snapshot(found.agent.session, found.agent),
})
return ok(request, { selected: { ...selected } })
} catch (error: unknown) {
return err(request, {
@@ -1102,6 +1164,11 @@ export function createApiProxy(ctx: Context, defaults: ApiProxyDefaults): ApiPro
muxQueues.add(queue)
for (const session of ctx.sessions.list()) {
subscribeSession(queue, session)
queue.push(frame({
type: 'session/metrics',
sessionId: session.id,
metrics: metricsProjector.snapshot(session, ctx.agents.get(session.id)),
}))
}
for (const pending of pendingQuestions.values()) {
queue.push({
@@ -1154,6 +1221,11 @@ export function createApiProxy(ctx: Context, defaults: ApiProxyDefaults): ApiPro
}),
ctx.on('session/created', (session: Session) => {
subscribeSession(queue, session)
queue.push(frame({
type: 'session/metrics',
sessionId: session.id,
metrics: metricsProjector.snapshot(session, ctx.agents.get(session.id)),
}))
}),
ctx.on('session/disposed', (session: Session) => {
openCalls.delete(session.id)

View File

@@ -10,7 +10,9 @@ import type { HostFrame, MuxFrame } from './events.ts'
import type { Wire } from './rpc.schema.ts'
import { rpcErrorSchema, rpcIdSchema } from './rpc.schema.ts'
import { approvalRequestIdSchema } from './approvals.schema.ts'
import { contentBlockSchema, sessionEventSchema, sessionIdSchema, toolEventViewSchema } from './sessions.schema.ts'
import {
contentBlockSchema, sessionEventSchema, sessionIdSchema, sessionMetricsSchema, toolEventViewSchema,
} from './sessions.schema.ts'
import { workspaceIdSchema, workspaceViewSchema } from './workspace.schema.ts'
/** Question shape validated strictly against core dsh-user-interaction. */
@@ -27,6 +29,7 @@ export const askUserQuestionItemSchema = z.object({
export const muxFrameSchema = z.discriminatedUnion('type', [
z.object({ type: z.literal('session/event'), sessionId: sessionIdSchema, event: sessionEventSchema, view: toolEventViewSchema.optional() }),
z.object({ type: z.literal('session/subscribed'), sessionId: sessionIdSchema, lastSeq: z.number().int() }),
z.object({ type: z.literal('session/metrics'), sessionId: sessionIdSchema, metrics: sessionMetricsSchema }),
z.object({ type: z.literal('session/title'), sessionId: sessionIdSchema, title: z.string().min(1), eventSeq: z.number().int().nonnegative(), updatedAt: z.number() }),
z.object({ type: z.literal('approval/requested'), sessionId: sessionIdSchema, approvalId: approvalRequestIdSchema, toolName: z.string(), callId: z.string().optional(), reason: z.string().optional() }),
z.object({ type: z.literal('approval/resolved'), sessionId: sessionIdSchema, approvalId: approvalRequestIdSchema, outcome: z.union([z.literal('allowed-once'), z.literal('rejected'), z.literal('cancelled'), z.literal('unavailable')]) }),

View File

@@ -13,6 +13,7 @@ import type { CallId } from '@deepseek-ai/dsh-llm/brand'
import type { SessionEvent, SessionId } from '@deepseek-ai/dsh-session/types'
import type { ToolCallView, ToolResultView } from '@deepseek-ai/dsh-tools/presentation'
import type { RpcError, RpcId, RpcRequest } from './rpc.ts'
import type { SessionMetrics } from './sessions.ts'
import type { WorkspaceView } from './workspace.ts'
// Client-side consumers take the render-intent vocabulary from the contract;
@@ -57,6 +58,7 @@ export interface EventsApi {
export type MuxFrame =
| { type: 'session/event'; sessionId: SessionId; event: SessionEvent; view?: ToolEventView }
| { type: 'session/subscribed'; sessionId: SessionId; lastSeq: number }
| { type: 'session/metrics'; sessionId: SessionId; metrics: SessionMetrics }
| { type: 'session/title'; sessionId: SessionId; title: string; eventSeq: number; updatedAt: number }
| { type: 'approval/requested'; sessionId: SessionId; approvalId: ApprovalRequestId; toolName: string; callId?: CallId; reason?: string }
| { type: 'approval/resolved'; sessionId: SessionId; approvalId: ApprovalRequestId; outcome: ApprovalOutcome }

View File

@@ -27,7 +27,7 @@ export interface ApiProxy {
// ---- Domain interfaces and payload entities ----
export type {
HistoryEntry, ModelCatalogFailure, ModelCatalogModel, ModelProviderGroup, ModelReasoning,
ModelReasoningEffort, ModelTarget, SessionModels, SessionsApi, SessionSummary,
ModelReasoningEffort, ModelTarget, SessionMetrics, SessionModels, SessionsApi, SessionSummary,
} from './sessions.ts'
export type { HostApi } from './host.ts'
export type { WorkspaceApi, WorkspaceId, WorkspaceView } from './workspace.ts'

View File

@@ -11,7 +11,7 @@ import type { RequestPayload, ResponseValue } from './rpc-map.ts'
import type { Wire } from './rpc.schema.ts'
import type {
HistoryEntry, ModelCatalogFailure, ModelCatalogModel, ModelProviderGroup, ModelReasoning,
ModelReasoningEffort, ModelTarget, SessionSummary,
ModelReasoningEffort, ModelTarget, SessionMetrics, SessionSummary,
} from './sessions.ts'
import type { ToolEventView } from './events.ts'
import type { WorkspaceId } from './workspace.ts'
@@ -145,11 +145,24 @@ export const todoItemSchema = z.object({
status: z.union([z.literal('pending'), z.literal('in_progress'), z.literal('completed')]),
})
/** Host-owned durable usage and current-context projection. */
export const sessionMetricsSchema = z.object({
logRevision: z.number().int().nonnegative(),
projectionRevision: z.number().int().nonnegative(),
uncachedInputTokens: z.number().nonnegative(),
outputTokens: z.number().nonnegative(),
cacheReadTokens: z.number().nonnegative(),
cacheWriteTokens: z.number().nonnegative(),
contextTokens: z.number().nonnegative().optional(),
contextWindow: z.number().int().positive().optional(),
}) satisfies z.ZodType<Wire<SessionMetrics>>
/** session.history response value. */
export const sessionHistoryValueSchema = z.object({
events: z.array(historyEntrySchema),
hasMore: z.boolean(),
todos: z.array(todoItemSchema).optional(),
metrics: sessionMetricsSchema.optional(),
}) satisfies z.ZodType<Wire<ResponseValue<'session.history'>>>
/** session.models request payload. */

View File

@@ -32,6 +32,31 @@ export interface HistoryEntry {
view?: ToolEventView
}
/**
* Host-owned token metrics for one durable session revision. Provider usage
* buckets are cumulative across the full log; current context fields describe
* the replayed request surface at this revision and are absent when the Host
* cannot measure pressure or resolve exact-route capacity.
*/
export interface SessionMetrics {
/** Number of durable events included in this projection. */
logRevision: number
/** Monotone ordering within one Host process and mux subscription generation. */
projectionRevision: number
/** Cumulative uncached provider input. */
uncachedInputTokens: number
/** Cumulative provider output. */
outputTokens: number
/** Cumulative provider cache reads. */
cacheReadTokens: number
/** Cumulative provider cache writes; excluded from the Web cache-hit formula. */
cacheWriteTokens: number
/** Current request pressure from `ctx.tokenMeter.measure(session).totalTokens`. */
contextTokens?: number
/** Exact selected-route capacity from `ctx.llm.resolveModelInfo()`. */
contextWindow?: number
}
/** Complete model target selected for one session. */
export interface ModelTarget {
/** Registered provider route. */
@@ -153,9 +178,11 @@ export interface SessionsApi {
* projection (latest `todo/write` over the FULL log, independent of the page window) —
* so a paged client restores the plan without walking history; absent when the session
* never wrote one. Older pages omit it (the projection is session-level, not per-page).
* The same tail-only rule carries `metrics`, whose cumulative usage and current context
* are Host projections over the full log rather than products of the returned page.
*/
history(request: RpcRequest<{ sessionId: SessionId; beforeSeq?: number; maxMessages?: number }>):
Promise<RpcResponse<{ events: HistoryEntry[]; hasMore: boolean; todos?: TodoItem[] }>>
Promise<RpcResponse<{ events: HistoryEntry[]; hasMore: boolean; todos?: TodoItem[]; metrics?: SessionMetrics }>>
/** Reads a fresh advisory model directory for this session. Provider lookups run independently. */
models(request: RpcRequest<{ sessionId: SessionId }>): Promise<RpcResponse<SessionModels>>

View File

@@ -0,0 +1,194 @@
/**
* Full-log usage and current-context projection for Web clients.
*
* @module @deepseek-ai/dsh-host-apiproxy/session-metrics
*/
import type { Context } from 'cordis'
import type { Agent, AgentLlmTarget } from '@deepseek-ai/dsh-agent'
import type { TokenUsage } from '@deepseek-ai/dsh-llm'
import type { Session, SessionEvent } from '@deepseek-ai/dsh-session'
import type { SessionMetrics } from './api/sessions.ts'
interface UsageState {
logRevision: number
projectionRevision: number
uncachedInputTokens: number
outputTokens: number
cacheReadTokens: number
cacheWriteTokens: number
byStep: Map<string, TokenUsage>
}
interface CapacityState {
routeKey: string
generation: number
status: 'pending' | 'ready'
contextWindow?: number
}
interface TokenMeterLike {
measure(session: Session): { totalTokens: number }
}
interface LlmLike {
resolveModelInfo(provider: string, model: string): Promise<{
context?: { contextWindow: number }
}>
}
function usageFrom(event: SessionEvent): { turn: number; step: number; usage: TokenUsage } | undefined {
if (event.type === 'assistant/chunk' && event.data.chunk.type === 'usage') {
return { turn: event.data.turn, step: event.data.step, usage: event.data.chunk.usage }
}
if (event.type === 'assistant/message' && event.data.usage !== undefined) {
return { turn: event.data.turn, step: event.data.step, usage: event.data.usage }
}
return undefined
}
/**
* Whether an appended event can change cumulative usage or token-meter
* pressure. Text/reasoning stream deltas remain outside both projections.
* @param event - appended durable event.
* @returns true when the Host must publish a fresh metrics snapshot.
*/
export function affectsSessionMetrics(event: SessionEvent): boolean {
if (event.type === 'assistant/chunk') return event.data.chunk.type === 'usage'
if (event.type === 'request/header') return true
return 'surfaceOp' in event
}
function recordUsage(state: UsageState, turn: number, step: number, usage: TokenUsage): void {
const key = `${turn}:${step}`
const previous = state.byStep.get(key)
if (previous !== undefined) {
state.uncachedInputTokens -= previous.inputTokens
state.outputTokens -= previous.outputTokens
state.cacheReadTokens -= previous.cacheReadTokens ?? 0
state.cacheWriteTokens -= previous.cacheWriteTokens ?? 0
}
state.byStep.set(key, usage)
state.uncachedInputTokens += usage.inputTokens
state.outputTokens += usage.outputTokens
state.cacheReadTokens += usage.cacheReadTokens ?? 0
state.cacheWriteTokens += usage.cacheWriteTokens ?? 0
}
/**
* Projects durable cumulative usage and route-aware current context without
* awaiting model metadata on the session append path.
*/
export class SessionMetricsProjector {
private readonly usage = new WeakMap<Session, UsageState>()
private readonly capacities = new WeakMap<Agent, CapacityState>()
/**
* @param ctx - Host context providing optional token-meter and LLM services.
* @param targetFor - selected route owner for one attached Web agent.
* @param onCapacityResolved - schedules a fresh live projection after exact-route metadata resolves.
*/
constructor(
private readonly ctx: Context,
private readonly targetFor: (agent: Agent) => Pick<AgentLlmTarget, 'provider' | 'model'>,
private readonly onCapacityResolved: (agent: Agent) => void,
) {}
/**
* Read a fresh detached projection through the session's durable tail.
* @param session - authoritative durable log owner.
* @param agent - attached route owner, when available.
* @returns cumulative usage and any currently available pressure/capacity.
*/
snapshot(session: Session, agent?: Agent): SessionMetrics {
const state = this.syncUsage(session)
const tokenMeter = this.ctx.get('tokenMeter') as TokenMeterLike | undefined
let contextTokens: number | undefined
if (tokenMeter !== undefined) {
try {
contextTokens = tokenMeter.measure(session).totalTokens
} catch {
// A malformed or temporarily unmeasurable replay has no honest pressure value.
}
}
const contextWindow = agent === undefined ? undefined : this.capacityFor(agent)
return {
logRevision: state.logRevision,
projectionRevision: state.projectionRevision++,
uncachedInputTokens: state.uncachedInputTokens,
outputTokens: state.outputTokens,
cacheReadTokens: state.cacheReadTokens,
cacheWriteTokens: state.cacheWriteTokens,
...contextTokens === undefined ? {} : { contextTokens },
...contextWindow === undefined ? {} : { contextWindow },
}
}
private syncUsage(session: Session): UsageState {
let state = this.usage.get(session)
if (state === undefined) {
state = {
logRevision: 0,
projectionRevision: 0,
uncachedInputTokens: 0,
outputTokens: 0,
cacheReadTokens: 0,
cacheWriteTokens: 0,
byStep: new Map(),
}
this.usage.set(session, state)
}
while (state.logRevision < session.events.length) {
const event = session.events[state.logRevision]
/* v8 ignore next -- Session events are append-only and dense; logRevision is bounded by length. */
if (event === undefined) break
const usage = usageFrom(event)
if (usage !== undefined) recordUsage(state, usage.turn, usage.step, usage.usage)
state.logRevision++
}
return state
}
private capacityFor(agent: Agent): number | undefined {
const target = this.targetFor(agent)
const routeKey = `${target.provider}\u0000${target.model}`
let state = this.capacities.get(agent)
if (state === undefined || state.routeKey !== routeKey) {
state = {
routeKey,
generation: (state?.generation ?? 0) + 1,
status: 'pending',
}
this.capacities.set(agent, state)
this.resolveCapacity(agent, target, state)
}
return state.status === 'ready' ? state.contextWindow : undefined
}
private resolveCapacity(
agent: Agent,
target: Pick<AgentLlmTarget, 'provider' | 'model'>,
pending: CapacityState,
): void {
const llm = this.ctx.get('llm') as LlmLike | undefined
if (llm === undefined) {
pending.status = 'ready'
return
}
void Promise.resolve()
.then(() => llm.resolveModelInfo(target.provider, target.model))
.then(
(resolved) => {
if (this.capacities.get(agent)?.generation !== pending.generation) return
const current = this.targetFor(agent)
if (`${current.provider}\u0000${current.model}` !== pending.routeKey) return
pending.status = 'ready'
if (resolved.context !== undefined) pending.contextWindow = resolved.context.contextWindow
this.onCapacityResolved(agent)
},
() => {
if (this.capacities.get(agent)?.generation === pending.generation) pending.status = 'ready'
},
)
}
}