refactor(web): publish transient model request capacity (round 1)

This commit is contained in:
Hypatia May
2026-07-28 18:35:39 +08:00
parent 3f0ba77bfa
commit 35b9c454e5
54 changed files with 765 additions and 1352 deletions
+19 -64
View File
@@ -6,7 +6,6 @@
import { randomUUID } from 'node:crypto'
import { mkdir, stat } from 'node:fs/promises'
import { join } from 'node:path'
import { FiberState } from 'cordis'
import type { Context } from 'cordis'
import { installAgentLlmTarget } from '@deepseek-ai/dsh-agent'
import type {
@@ -409,26 +408,6 @@ export function createApiProxy(ctx: Context, defaults: ApiProxyDefaults): ApiPro
return target
}
/**
* Read the best capacity route without taking ownership of foreign routing.
* Web agents expose their live selection; other agents expose only a route
* that already crossed the durable request-header boundary.
*/
function metricsRouteFor(agent: Agent): Pick<AgentLlmTarget, 'provider' | 'model'> | undefined {
const installed = targets.get(agent)
if (installed !== undefined) return installed.current
const logged = agent.session.requestHeader()?.config
return logged === undefined
? undefined
: { provider: logged.provider, model: logged.model }
}
/** Pair a registry agent only with the exact Session lifecycle it owns. */
function metricsAgentFor(session: Session): Agent | undefined {
const agent = ctx.get('agents')?.get(session.id)
return agent?.session === session ? agent : undefined
}
/** Pre-publication setup used by both fresh and resumed Web agents. */
function installTarget(agentCtx: Context): void {
const agent = agentCtx.agent
@@ -444,39 +423,23 @@ export function createApiProxy(ctx: Context, defaults: ApiProxyDefaults): ApiPro
const pendingMetricSessions = new Set<Session>()
let metricFlushScheduled = false
let metricsDisposed = false
const metricsProjector = new SessionMetricsProjector(
ctx,
metricsRouteFor,
(agent) => {
if (metricsDisposed) return
const agents = ctx.get('agents')
if (agents?.get(agent.id) !== agent) return
const sessions = ctx.get('sessions')
if (sessions?.get(agent.id) !== agent.session) return
scheduleMetrics(agent.session)
},
)
const metricsProjector = new SessionMetricsProjector(ctx)
/** Queue one full-log metrics publication after synchronous session listeners drain. */
function scheduleMetrics(session: Session): void {
if (metricsDisposed || muxQueues.size === 0) return
if (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, metricsAgentFor(current)),
metrics: metricsProjector.snapshot(current),
})
}
})
@@ -488,25 +451,22 @@ export function createApiProxy(ctx: Context, defaults: ApiProxyDefaults): ApiPro
if (affectsSessionMetrics(event)) scheduleMetrics(session)
}),
ctx.on('agent/created', (agent: Agent) => { scheduleMetrics(agent.session) }),
ctx.on('agent/model-request', (agent, turn, step, request) => {
broadcast({
type: 'session/model-request',
sessionId: agent.session.id,
turn,
step,
provider: request.provider,
model: request.model,
...request.contextWindow === undefined
? {}
: { contextWindow: request.contextWindow },
})
}),
ctx.on('session/disposed', (session: Session) => { pendingMetricSessions.delete(session) }),
ctx.on('internal/status', (fiber) => {
if (metricsDisposed) return
if (fiber.state === FiberState.UNLOADING) {
metricsProjector.invalidateCapacities()
return
}
if (fiber.state !== FiberState.ACTIVE
&& fiber.state !== FiberState.FAILED
&& fiber.state !== FiberState.DISPOSED) return
metricsProjector.invalidateCapacities()
const sessions = ctx.get('sessions')
if (sessions === undefined) return
for (const session of sessions.list()) scheduleMetrics(session)
}, { global: true }),
]
return () => {
metricsDisposed = true
metricsProjector.dispose()
pendingMetricSessions.clear()
for (const dispose of disposers) dispose()
}
@@ -819,7 +779,7 @@ export function createApiProxy(ctx: Context, defaults: ApiProxyDefaults): ApiPro
// client cannot reconstruct session-level state from it).
const todos = beforeSeq === undefined ? backscanTodos(found.agent.session.events) : undefined
const metrics = beforeSeq === undefined
? metricsProjector.snapshot(found.agent.session, found.agent)
? metricsProjector.snapshot(found.agent.session)
: undefined
return ok(request, {
events: entries,
@@ -920,11 +880,6 @@ 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, {
@@ -1235,7 +1190,7 @@ export function createApiProxy(ctx: Context, defaults: ApiProxyDefaults): ApiPro
queue.push(frame({
type: 'session/metrics',
sessionId: session.id,
metrics: metricsProjector.snapshot(session, metricsAgentFor(session)),
metrics: metricsProjector.snapshot(session),
}))
}
for (const pending of pendingQuestions.values()) {
@@ -1292,7 +1247,7 @@ export function createApiProxy(ctx: Context, defaults: ApiProxyDefaults): ApiPro
queue.push(frame({
type: 'session/metrics',
sessionId: session.id,
metrics: metricsProjector.snapshot(session, metricsAgentFor(session)),
metrics: metricsProjector.snapshot(session),
}))
}),
ctx.on('session/disposed', (session: Session) => {
@@ -30,6 +30,15 @@ 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/model-request'),
sessionId: sessionIdSchema,
turn: z.number().int().positive(),
step: z.number().int().positive(),
provider: z.string().min(1),
model: z.string().min(1),
contextWindow: z.number().int().positive().optional(),
}),
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')]) }),
+16
View File
@@ -59,6 +59,22 @@ 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 }
/**
* One model request observed by this already-open mux connection after its
* final route and stream handle were resolved. This frame is transient: mux
* baselines, reconnects, and session history never replay it. An absent
* `contextWindow` explicitly clears a capacity observed from an earlier
* request on the same connection.
*/
| {
type: 'session/model-request'
sessionId: SessionId
turn: number
step: number
provider: string
model: string
contextWindow?: number
}
| { 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 }
@@ -145,7 +145,7 @@ 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. */
/** Host-owned durable usage and current-pressure projection. */
export const sessionMetricsSchema = z.object({
logRevision: z.number().int().nonnegative(),
projectionRevision: z.number().int().nonnegative(),
@@ -154,7 +154,6 @@ export const sessionMetricsSchema = z.object({
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. */
+7 -7
View File
@@ -34,9 +34,9 @@ export interface HistoryEntry {
/**
* 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.
* buckets are cumulative across the full log; current context pressure
* describes the replayed request surface at this revision and is absent when
* the Host cannot measure it.
*/
export interface SessionMetrics {
/** Number of durable events included in this projection. */
@@ -53,8 +53,6 @@ export interface SessionMetrics {
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. */
@@ -178,8 +176,10 @@ 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.
* The same tail-only rule carries `metrics`, whose cumulative usage and
* current pressure are Host projections over the full log rather than
* products of the returned page. Live model capacity is connection-local
* telemetry and is never reconstructed here.
*/
history(request: RpcRequest<{ sessionId: SessionId; beforeSeq?: number; maxMessages?: number }>):
Promise<RpcResponse<{ events: HistoryEntry[]; hasMore: boolean; todos?: TodoItem[]; metrics?: SessionMetrics }>>
+5 -140
View File
@@ -5,7 +5,6 @@
*/
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'
@@ -20,27 +19,10 @@ interface UsageState {
byStep: Map<string, TokenUsage>
}
interface CapacityState {
routeKey: string | undefined
generation: number
epoch: number
status: 'pending' | 'ready' | 'retryable'
contextWindow?: number
controller?: AbortController
}
type CapacityTarget = Pick<AgentLlmTarget, 'provider' | 'model'>
interface TokenMeterLike {
measure(session: Session): { totalTokens: number }
}
interface LlmLike {
resolveModelInfo(provider: string, model: string, signal?: AbortSignal): 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 }
@@ -79,51 +61,19 @@ function recordUsage(state: UsageState, turn: number, step: number, usage: Token
state.cacheWriteTokens += usage.cacheWriteTokens ?? 0
}
function routeKeyFor(target: CapacityTarget | undefined): string | undefined {
return target === undefined ? undefined : `${target.provider}\u0000${target.model}`
}
/**
* Projects durable cumulative usage and route-aware current context without
* awaiting model metadata on the session append path.
*/
/** Projects durable cumulative usage and synchronous current context pressure. */
export class SessionMetricsProjector {
private readonly usage = new WeakMap<Session, UsageState>()
private readonly capacities = new WeakMap<Agent, CapacityState>()
private readonly pendingCapacities = new Set<CapacityState>()
private capacityEpoch = 0
private disposed = false
/**
* @param ctx - Host context providing optional token-meter and LLM services.
* @param targetFor - side-effect-free selected or logged route lookup for one attached agent.
* @param onCapacityResolved - schedules a fresh live projection after exact-route metadata resolves.
*/
constructor(
private readonly ctx: Context,
private readonly targetFor: (agent: Agent) => CapacityTarget | undefined,
private readonly onCapacityResolved: (agent: Agent) => void,
) {}
/** Retire adapter-owned metadata and fence every resolution already in flight. */
invalidateCapacities(): void {
this.capacityEpoch++
for (const pending of this.pendingCapacities) this.abortCapacityResolution(pending)
}
/** Permanently retire capacity projection and cancel every adapter-owned lookup. */
dispose(): void {
this.disposed = true
this.invalidateCapacities()
}
/** @param ctx - Host context providing an optional token-meter service. */
constructor(private readonly ctx: Context) {}
/**
* 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.
* @returns cumulative usage and any currently measurable pressure.
*/
snapshot(session: Session, agent?: Agent): SessionMetrics {
snapshot(session: Session): SessionMetrics {
const state = this.syncUsage(session)
const tokenMeter = this.ctx.get('tokenMeter') as TokenMeterLike | undefined
let contextTokens: number | undefined
@@ -134,7 +84,6 @@ export class SessionMetricsProjector {
// 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++,
@@ -143,7 +92,6 @@ export class SessionMetricsProjector {
cacheReadTokens: state.cacheReadTokens,
cacheWriteTokens: state.cacheWriteTokens,
...contextTokens === undefined ? {} : { contextTokens },
...contextWindow === undefined ? {} : { contextWindow },
}
}
@@ -171,87 +119,4 @@ export class SessionMetricsProjector {
}
return state
}
private capacityFor(agent: Agent): number | undefined {
if (this.disposed) return undefined
const target = this.targetFor(agent)
const routeKey = routeKeyFor(target)
let state = this.capacities.get(agent)
if (state === undefined
|| state.routeKey !== routeKey
|| state.epoch !== this.capacityEpoch
|| state.status === 'retryable') {
this.abortCapacityResolution(state)
state = {
routeKey,
generation: (state?.generation ?? 0) + 1,
epoch: this.capacityEpoch,
status: target === undefined ? 'ready' : 'pending',
}
this.capacities.set(agent, state)
if (target !== undefined) this.resolveCapacity(agent, target, state)
}
return state.status === 'ready' ? state.contextWindow : undefined
}
private resolveCapacity(
agent: Agent,
target: CapacityTarget,
pending: CapacityState,
): void {
const llm = this.ctx.get('llm') as LlmLike | undefined
if (llm === undefined) {
pending.status = 'retryable'
return
}
const controller = new AbortController()
pending.controller = controller
this.pendingCapacities.add(pending)
void Promise.resolve()
.then(() => {
controller.signal.throwIfAborted()
return llm.resolveModelInfo(target.provider, target.model, controller.signal)
})
.then(
(resolved) => {
this.finishCapacityResolution(pending, controller)
if (this.capacityResolutionIsStale(agent, pending)) return
pending.status = 'ready'
if (resolved.context !== undefined) pending.contextWindow = resolved.context.contextWindow
this.onCapacityResolved(agent)
},
() => {
this.finishCapacityResolution(pending, controller)
if (!this.capacityResolutionIsStale(agent, pending)) pending.status = 'retryable'
},
)
}
private abortCapacityResolution(pending: CapacityState | undefined): void {
if (pending === undefined || pending.controller === undefined) return
const controller = pending.controller
delete pending.controller
this.pendingCapacities.delete(pending)
controller.abort()
}
private finishCapacityResolution(pending: CapacityState, controller: AbortController): void {
this.pendingCapacities.delete(pending)
if (pending.controller === controller) delete pending.controller
}
private capacityResolutionIsStale(agent: Agent, pending: CapacityState): boolean {
if (pending.epoch !== this.capacityEpoch) return true
if (this.capacities.get(agent)?.generation !== pending.generation) return true
if (routeKeyFor(this.targetFor(agent)) === pending.routeKey) return false
// Unknown is the neutral generation; the next observed concrete route
// starts a fresh resolution even when it equals the route that disappeared.
this.capacities.set(agent, {
routeKey: undefined,
generation: pending.generation + 1,
epoch: this.capacityEpoch,
status: 'ready',
})
return true
}
}