From f4d398d1b992081e9f08ad9a0fa04e2c41f51f1c Mon Sep 17 00:00:00 2001 From: doudouOUC Date: Thu, 16 Jul 2026 09:39:19 +0800 Subject: [PATCH 1/4] feat(serve): Complete legacy session workspace telemetry Co-authored-by: Qwen-Coder --- ...emon-legacy-session-workspace-telemetry.md | 114 ++++ docs/developers/daemon/02-serve-runtime.md | 20 +- docs/developers/daemon/19-observability.md | 37 +- .../src/serve/routes/session-runtime.test.ts | 185 ++++++ .../cli/src/serve/routes/session-runtime.ts | 6 +- .../serve/routes/session-telemetry.test.ts | 185 ++++++ packages/cli/src/serve/routes/session.ts | 7 + packages/cli/src/serve/server.ts | 49 +- .../serve/server/telemetry-catalog.test.ts | 104 ++++ .../cli/src/serve/server/telemetry.test.ts | 494 ++++++++++++++-- packages/cli/src/serve/server/telemetry.ts | 550 ++++++++++++++---- 11 files changed, 1539 insertions(+), 212 deletions(-) create mode 100644 docs/design/daemon-legacy-session-workspace-telemetry.md create mode 100644 packages/cli/src/serve/routes/session-runtime.test.ts create mode 100644 packages/cli/src/serve/routes/session-telemetry.test.ts create mode 100644 packages/cli/src/serve/server/telemetry-catalog.test.ts diff --git a/docs/design/daemon-legacy-session-workspace-telemetry.md b/docs/design/daemon-legacy-session-workspace-telemetry.md new file mode 100644 index 00000000000..464e875bb4e --- /dev/null +++ b/docs/design/daemon-legacy-session-workspace-telemetry.md @@ -0,0 +1,114 @@ +# Legacy Session Workspace Telemetry + +## Context + +The daemon telemetry middleware classifies HTTP requests before Express route +handlers run. Legacy singular session routes can resolve to any registered +workspace, but the middleware cannot know the selected runtime from the URL +alone. Resolving the live owner in both middleware and the handler duplicates +work and can disagree if the registry changes between the two lookups. + +This design gives every explicit legacy `/session`, `/sessions`, and +`/permission` route a stable request span while attributing dynamic routes to +the runtime selected by the handler. + +## Route inventory + +The route catalog contains all 48 explicit legacy routes. Each entry declares +its HTTP method, Express path template, canonical route label, and one of two +attribution modes: + +- `handler_resolved` (41 routes): `POST /session`, load/resume, the legacy + transcript route, and every singular session route that resolves a live + owner. The handler publishes the selected runtime workspace to telemetry. +- `pre_resolved` (7 routes): legacy export, A2UI action, legacy organization, + the three global batch mutations, and the global permission vote. These + routes remain bound to the primary workspace. + +The catalog matcher follows the relevant Express 5 defaults: static segments +are case-insensitive, one trailing slash is accepted, and parameter segments +are decoded only after their raw path boundary has been captured. A malformed +session id is retained as its raw value. Permission request ids are decoded +before their existing length and character-set validation. The emitted +`http.route` always uses the canonical catalog template. + +## Deferred attribution + +Handler-resolved requests start without `qwen-code.workspace.hash`. The +middleware stores a private context on the Express response. Route code calls +`setDaemonTelemetryWorkspace(res, runtime.workspaceCwd)` after a unique runtime +has been selected. The setter is best-effort and first-selection-wins: a +repeated identical value is idempotent and a later different value is ignored. + +The four publication seams are: + +1. `requireSessionRuntime`, shared by live-owner routes. +2. Session creation after workspace selection. +3. Session load/resume after target runtime selection. +4. Legacy transcript resolution after a unique live or persisted owner is + found. + +Publication precedes later trust, unsupported-secondary, conflict, and request +validation checks. Consequently those failures retain the uniquely selected +runtime. Requests that fail before unique selection, including not-found, +ambiguous, and workspace-mismatch cases, omit the workspace hash. Attribution +uses `runtime.workspaceCwd`, not a session's requested or temporary cwd. + +On response `finish` or `close`, the middleware hashes the published workspace, +sets the span attribute, records the response, and ends the span. Resolution, +hashing, and span updates are best-effort and cannot affect request handling or +metrics settlement. The context is cleared after one settlement. + +Pre-resolved requests continue to hash the middleware-selected workspace when +the span starts. Removing the middleware's live-owner callback ensures a live +owner is resolved no more than once per request. + +## Streaming and metrics + +All 48 catalog routes create request spans. A successful +`GET /session/:id/events` response ends its span when the SSE connection closes, +but is excluded from the ordinary HTTP request count/duration and the Web Shell +status metrics ring because its duration is the connection lifetime. SSE +handshake failures are recorded as ordinary short HTTP requests. + +`POST /session/:id/generate` is a bounded request-scoped SSE operation. Its +connection ends when generation completes, so its duration remains meaningful +request latency and continues to enter ordinary HTTP metrics. + +Heartbeat requests remain in OpenTelemetry HTTP metrics but stay excluded from +the status metrics ring. `GET /daemon/status` also remains excluded only from +that ring. A shared settlement guard prevents duplicate recording when both +`finish` and `close` fire. + +HTTP metrics and the Web Shell metrics ring remain daemon-global. Adding a +workspace metric dimension requires a separate cardinality and dashboard +compatibility review. + +## Compatibility and boundaries + +This change does not alter routes, request or response schemas, SDKs, +capabilities, persistence, authentication, trust ordering, archive leases, +bridge error mapping, or session execution. It does not add public telemetry +attributes. + +The telemetry middleware is installed after bearer authentication, rate +limiting, and JSON parsing, so requests rejected by those earlier gates remain +outside this request-span coverage. Implicit HEAD/OPTIONS, access-log behavior, +rate-limit path normalization, workspace session-group routes, +workspace-qualified organization, ACP/WebSocket telemetry, and enabling +secondary branch/fork/cd execution are out of scope. + +## Verification + +- A drift guard compares the explicit legacy routes registered with Express to + the catalog and asserts the 48/41/7 inventory. +- Matcher tests cover case, trailing slash, encoded slash, Unicode, malformed + encoding, permission id validation, method/path mismatch, and canonical + labels. +- Middleware tests cover deferred attribution, first-selection-wins, hash + caching, telemetry failures, one-time settlement, SSE metrics, heartbeat, and + status exclusions. +- Route tests cover live-owner, creation, restore, and transcript publication + for primary, secondary, untrusted, missing, ambiguous, and conflict cases. +- A dual-workspace outfile test verifies secondary, primary-bound, and omitted + hashes without exposing raw workspace paths. diff --git a/docs/developers/daemon/02-serve-runtime.md b/docs/developers/daemon/02-serve-runtime.md index 15ec8bad8d1..df5d1690dd0 100644 --- a/docs/developers/daemon/02-serve-runtime.md +++ b/docs/developers/daemon/02-serve-runtime.md @@ -24,16 +24,16 @@ **Middleware** (`packages/cli/src/serve/auth.ts` and `server.ts`): -| Middleware, in registration order | Purpose | Notes | -| ------------------------------------------- | -------------------------------------------------------------------------------------------------------------------------- | ----------------------------------------------------------------------------------------------------------------- | -| `denyBrowserOriginCors` / `allowOriginCors` | Deny all `Origin` headers by default; switch to an allowlist when `--allow-origin ` is configured. | See [`12-auth-security.md`](./12-auth-security.md). | -| `hostAllowlist(bind, getPort)` | On loopback, validate `Host` belongs to `localhost`, `127.0.0.1`, `[::1]`, or `host.docker.internal` plus the actual port. | Defense against DNS rebinding. Comparison is case-insensitive and cached per port. | -| Access-log middleware | Records method, path, status, durationMs, sessionId, and clientId to `DaemonLogger` when a request finishes. | Registered **before** `bearerAuth`, so 401 denials are logged too. Skips `/health` and heartbeat. | -| `bearerAuth(token)` | SHA-256 plus `timingSafeEqual` constant-time bearer comparison. | Open passthrough when no token is configured (loopback dev default). `Bearer` scheme is case-insensitive. | -| Rate-limit middleware | Optional per-tier token bucket for prompt, mutation, and read routes. | Registered after `bearerAuth` and before JSON parsing; returns 429 before parsing when a bucket is exhausted. | -| `express.json({ limit: '10mb' })` | JSON body parsing. | Parse errors return 400. | -| `daemonTelemetryMiddleware` | Wraps each HTTP request in an OpenTelemetry span through `withDaemonRequestSpan`. | Attributes include route, sessionId, clientId, and status code. | -| `createMutationGate` (per-route) | Route-level opt-in gate for mutation routes that require token even on loopback. | Returns `401 { code: 'token_required' }`. Not global `app.use`; routes call `mutate({ strict: true })` as needed. | +| Middleware, in registration order | Purpose | Notes | +| ------------------------------------------- | -------------------------------------------------------------------------------------------------------------------------- | --------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | +| `denyBrowserOriginCors` / `allowOriginCors` | Deny all `Origin` headers by default; switch to an allowlist when `--allow-origin ` is configured. | See [`12-auth-security.md`](./12-auth-security.md). | +| `hostAllowlist(bind, getPort)` | On loopback, validate `Host` belongs to `localhost`, `127.0.0.1`, `[::1]`, or `host.docker.internal` plus the actual port. | Defense against DNS rebinding. Comparison is case-insensitive and cached per port. | +| Access-log middleware | Records method, path, status, durationMs, sessionId, and clientId to `DaemonLogger` when a request finishes. | Registered **before** `bearerAuth`, so 401 denials are logged too. Skips `/health` and heartbeat. | +| `bearerAuth(token)` | SHA-256 plus `timingSafeEqual` constant-time bearer comparison. | Open passthrough when no token is configured (loopback dev default). `Bearer` scheme is case-insensitive. | +| Rate-limit middleware | Optional per-tier token bucket for prompt, mutation, and read routes. | Registered after `bearerAuth` and before JSON parsing; returns 429 before parsing when a bucket is exhausted. | +| `express.json({ limit: '10mb' })` | JSON body parsing. | Parse errors return 400. | +| `daemonTelemetryMiddleware` | Wraps classified daemon API requests that reach this point in an OpenTelemetry span through `withDaemonRequestSpan`. | Attributes include canonical route, resolved workspace hash, sessionId, clientId, and status code. Earlier auth, rate-limit, and body-parser rejections are outside this span boundary. | +| `createMutationGate` (per-route) | Route-level opt-in gate for mutation routes that require token even on loopback. | Returns `401 { code: 'token_required' }`. Not global `app.use`; routes call `mutate({ strict: true })` as needed. | **Subsystems**: diff --git a/docs/developers/daemon/19-observability.md b/docs/developers/daemon/19-observability.md index 6fa38ab4caa..0cae8617d7f 100644 --- a/docs/developers/daemon/19-observability.md +++ b/docs/developers/daemon/19-observability.md @@ -6,23 +6,23 @@ ## What exists today -| Surface | Location | Purpose | -| ------------------------------------------- | ---------------------------------------------- | ----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | -| `QWEN_SERVE_DEBUG` stderr logs | `bridge.ts` and call sites | Env values `1` / `true` / `on` / `yes` (case-insensitive) print `qwen serve debug: ...` lines to stderr. | -| OpenTelemetry span instrumentation | `server.ts` `daemonTelemetryMiddleware` | Each HTTP request is wrapped in `withDaemonRequestSpan`; attributes include route, sessionId, clientId, and status code. Permission routes have dedicated spans. Prompt lifecycle is traced end-to-end. Configuration lives in `settings.json` `telemetry`. | -| OpenTelemetry daemon perf metrics | `telemetry/*event-loop-lag*`, `daemon-metrics` | Event loop lag gauges for daemon and ACP child processes, plus daemon-child pipe message byte histograms. | -| `DaemonLogger` structured file logs | `serve/daemon-logger.ts` | Structured JSON-like log lines are written to a file. Boot prints `daemon log -> `. Supports `info` / `warn` / `error` levels, with structured fields such as `route`, `sessionId`, `clientId`, `childPid`, and `channelId`. | -| Per-request access-log middleware | `server.ts`, registered before `bearerAuth` | Logs `method`, `path`, `status`, `durationMs`, `sessionId`, and `clientId` after each request. Skips `GET /health` and heartbeat. 4xx+ uses `warn`; success uses `info`. | -| `/health` | `server.ts` route | Liveness probe; `?deep=1` returns extended details. | -| `/capabilities` | `server.ts` route | Preflight feature discovery. See [`11-capabilities-versioning.md`](./11-capabilities-versioning.md). | -| `/workspace/preflight` | Route -> `DaemonStatusProvider` | Structured readiness cells: Node version, CLI entry, ripgrep, git, npm, plus ACP-level cells once a child is alive. | -| `/workspace/env` | Route -> `DaemonStatusProvider` | Daemon process env snapshot. Secret env vars report only presence; proxy URL credentials are stripped. | -| `/workspace/mcp` | Route -> bridge extMethod | Pool, budget, and refusal snapshot. | -| `/workspace/skills`, `/workspace/providers` | Routes | ACP-side live snapshots; return empty idle data when no session exists. | -| Per-session SSE | `GET /session/:id/events` | Real-time event stream. | -| `/demo` debug console | `GET /demo` (`packages/cli/src/serve/demo.ts`) | Browser-accessible single-page console: chat, event log, workspace inspector, and permission UX. On loopback, `http://127.0.0.1:4170/demo` is the quickest end-to-end validation path without writing SDK code. Registration rules are in [`02-serve-runtime.md`](./02-serve-runtime.md). | -| `PermissionAuditRing` | `permission-audit.ts` | In-memory FIFO of 512 permission decisions. | -| Mediator `decisionReason` audit | `permissionMediator.ts` | Internal structured record explaining why a permission request resolved the way it did. | +| Surface | Location | Purpose | +| ------------------------------------------- | ---------------------------------------------- | ----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | +| `QWEN_SERVE_DEBUG` stderr logs | `bridge.ts` and call sites | Env values `1` / `true` / `on` / `yes` (case-insensitive) print `qwen serve debug: ...` lines to stderr. | +| OpenTelemetry span instrumentation | `server.ts` `daemonTelemetryMiddleware` | Classified daemon API requests that reach the telemetry middleware are wrapped in `withDaemonRequestSpan`; attributes include canonical route, workspace hash when resolved, sessionId, clientId, and status code. Permission routes have dedicated spans. Prompt lifecycle is traced end-to-end. Configuration lives in `settings.json` `telemetry`. | +| OpenTelemetry daemon perf metrics | `telemetry/*event-loop-lag*`, `daemon-metrics` | Event loop lag gauges for daemon and ACP child processes, plus daemon-child pipe message byte histograms. | +| `DaemonLogger` structured file logs | `serve/daemon-logger.ts` | Structured JSON-like log lines are written to a file. Boot prints `daemon log -> `. Supports `info` / `warn` / `error` levels, with structured fields such as `route`, `sessionId`, `clientId`, `childPid`, and `channelId`. | +| Per-request access-log middleware | `server.ts`, registered before `bearerAuth` | Logs `method`, `path`, `status`, `durationMs`, `sessionId`, and `clientId` after each request. Skips `GET /health` and heartbeat. 4xx+ uses `warn`; success uses `info`. | +| `/health` | `server.ts` route | Liveness probe; `?deep=1` returns extended details. | +| `/capabilities` | `server.ts` route | Preflight feature discovery. See [`11-capabilities-versioning.md`](./11-capabilities-versioning.md). | +| `/workspace/preflight` | Route -> `DaemonStatusProvider` | Structured readiness cells: Node version, CLI entry, ripgrep, git, npm, plus ACP-level cells once a child is alive. | +| `/workspace/env` | Route -> `DaemonStatusProvider` | Daemon process env snapshot. Secret env vars report only presence; proxy URL credentials are stripped. | +| `/workspace/mcp` | Route -> bridge extMethod | Pool, budget, and refusal snapshot. | +| `/workspace/skills`, `/workspace/providers` | Routes | ACP-side live snapshots; return empty idle data when no session exists. | +| Per-session SSE | `GET /session/:id/events` | Real-time event stream. | +| `/demo` debug console | `GET /demo` (`packages/cli/src/serve/demo.ts`) | Browser-accessible single-page console: chat, event log, workspace inspector, and permission UX. On loopback, `http://127.0.0.1:4170/demo` is the quickest end-to-end validation path without writing SDK code. Registration rules are in [`02-serve-runtime.md`](./02-serve-runtime.md). | +| `PermissionAuditRing` | `permission-audit.ts` | In-memory FIFO of 512 permission decisions. | +| Mediator `decisionReason` audit | `permissionMediator.ts` | Internal structured record explaining why a permission request resolved the way it did. | ## What does not exist today @@ -173,7 +173,8 @@ flowchart TD ## Caveats and known limits - **DaemonLogger file logs are structured** and can be filtered by `route`, `sessionId`, and `clientId`. `QWEN_SERVE_DEBUG` stderr logs remain unstructured text. -- **OpenTelemetry spans include per-request correlation.** Each HTTP request span carries route, sessionId, and clientId attributes that can be joined in a tracing backend. +- **OpenTelemetry spans include per-request correlation.** Classified daemon API requests that pass bearer authentication, rate limiting, and body parsing carry canonical route, sessionId, clientId, and (when uniquely resolved) `qwen-code.workspace.hash` attributes. Requests rejected by an earlier middleware gate do not have these request spans. +- **HTTP metrics are daemon-global.** OpenTelemetry HTTP request metrics and the Web Shell status metrics ring do not include a workspace dimension. A successful session SSE connection has a request span but is excluded from ordinary request count/duration metrics because its lifetime is not request latency; failed SSE handshakes are counted normally. - **`runtime.perf` is daemon-only.** Child event loop lag is not reported there by design; use OTel or forwarded stderr stall warnings for ACP child stalls. - **ACP-level `/workspace/preflight` cells require a live session.** On an idle daemon, auth / MCP / skills / providers may show `status: 'not_started'`; this is expected. - **`/workspace/env` only reports secret presence, not values.** Do not expose the response where the mere presence of a secret is sensitive. diff --git a/packages/cli/src/serve/routes/session-runtime.test.ts b/packages/cli/src/serve/routes/session-runtime.test.ts new file mode 100644 index 00000000000..8bce1f7dd68 --- /dev/null +++ b/packages/cli/src/serve/routes/session-runtime.test.ts @@ -0,0 +1,185 @@ +/** + * @license + * Copyright 2026 Qwen Team + * SPDX-License-Identifier: Apache-2.0 + */ + +import type { Response } from 'express'; +import { beforeEach, describe, expect, it, vi } from 'vitest'; +import type { + WorkspaceRegistry, + WorkspaceRuntime, +} from '../workspace-registry.js'; + +const telemetryMocks = vi.hoisted(() => ({ + setDaemonTelemetryWorkspace: vi.fn(), +})); + +vi.mock('../server/telemetry.js', () => telemetryMocks); + +import { requireSessionRuntime } from './session-runtime.js'; + +function runtime( + workspaceCwd: string, + opts: { primary?: boolean; trusted?: boolean } = {}, +): WorkspaceRuntime { + return { + workspaceCwd, + workspaceId: workspaceCwd.split('/').at(-1) ?? workspaceCwd, + primary: opts.primary === true, + trusted: opts.trusted !== false, + } as WorkspaceRuntime; +} + +function response(): Response { + const res = { + statusCode: 200, + status: vi.fn((statusCode: number) => { + res.statusCode = statusCode; + return res; + }), + json: vi.fn(() => res), + }; + return res as unknown as Response; +} + +function registry(opts: { + primary: WorkspaceRuntime; + runtimes: WorkspaceRuntime[]; + resolution: + | { kind: 'found'; runtime: WorkspaceRuntime } + | { kind: 'not_found' } + | { kind: 'ambiguous'; runtimes: WorkspaceRuntime[] }; +}): { + registry: WorkspaceRegistry; + resolveLiveSessionOwner: ReturnType; +} { + const resolveLiveSessionOwner = vi.fn(() => opts.resolution); + return { + registry: { + primary: opts.primary, + list: () => opts.runtimes, + resolveLiveSessionOwner, + } as unknown as WorkspaceRegistry, + resolveLiveSessionOwner, + }; +} + +describe('requireSessionRuntime telemetry attribution', () => { + beforeEach(() => { + vi.clearAllMocks(); + }); + + it('publishes the primary runtime without scanning in single-workspace mode', () => { + const primary = runtime('/workspace/primary', { primary: true }); + const setup = registry({ + primary, + runtimes: [primary], + resolution: { kind: 'not_found' }, + }); + const res = response(); + + expect( + requireSessionRuntime({ + sessionId: 'session-1', + route: 'POST /session/:id/prompt', + res, + workspaceRegistry: setup.registry, + }), + ).toBe(primary); + expect(setup.resolveLiveSessionOwner).not.toHaveBeenCalled(); + expect(telemetryMocks.setDaemonTelemetryWorkspace).toHaveBeenCalledWith( + res, + '/workspace/primary', + ); + }); + + it.each([ + 'GET /session/:id/rewind/snapshots', + 'POST /session/:id/shell', + 'POST /session/:id/prompt', + ])('publishes a uniquely resolved secondary runtime once for %s', (route) => { + const primary = runtime('/workspace/primary', { primary: true }); + const secondary = runtime('/workspace/secondary'); + const setup = registry({ + primary, + runtimes: [primary, secondary], + resolution: { kind: 'found', runtime: secondary }, + }); + const res = response(); + + expect( + requireSessionRuntime({ + sessionId: 'session-2', + route, + res, + workspaceRegistry: setup.registry, + }), + ).toBe(secondary); + expect(setup.resolveLiveSessionOwner).toHaveBeenCalledTimes(1); + expect(telemetryMocks.setDaemonTelemetryWorkspace).toHaveBeenCalledOnce(); + expect(telemetryMocks.setDaemonTelemetryWorkspace).toHaveBeenCalledWith( + res, + '/workspace/secondary', + ); + }); + + it('publishes an untrusted unique runtime before rejecting it', () => { + const primary = runtime('/workspace/primary', { primary: true }); + const secondary = runtime('/workspace/untrusted', { trusted: false }); + const setup = registry({ + primary, + runtimes: [primary, secondary], + resolution: { kind: 'found', runtime: secondary }, + }); + const res = response(); + + expect( + requireSessionRuntime({ + sessionId: 'session-3', + route: 'GET /session/:id/events', + res, + workspaceRegistry: setup.registry, + }), + ).toBeUndefined(); + expect(res.statusCode).toBe(403); + expect(telemetryMocks.setDaemonTelemetryWorkspace).toHaveBeenCalledWith( + res, + '/workspace/untrusted', + ); + }); + + it.each([ + [{ kind: 'not_found' } as const, 404], + [ + { + kind: 'ambiguous', + runtimes: [] as WorkspaceRuntime[], + } as const, + 500, + ], + ])( + 'does not publish unresolved ownership for %o', + (resolution, statusCode) => { + const primary = runtime('/workspace/primary', { primary: true }); + const secondary = runtime('/workspace/secondary'); + const setup = registry({ + primary, + runtimes: [primary, secondary], + resolution, + }); + const res = response(); + + expect( + requireSessionRuntime({ + sessionId: 'missing', + route: 'POST /session/:id/rewind', + res, + workspaceRegistry: setup.registry, + }), + ).toBeUndefined(); + expect(res.statusCode).toBe(statusCode); + expect(telemetryMocks.setDaemonTelemetryWorkspace).not.toHaveBeenCalled(); + }, + ); +}); diff --git a/packages/cli/src/serve/routes/session-runtime.ts b/packages/cli/src/serve/routes/session-runtime.ts index 158e458195d..6128b91b799 100644 --- a/packages/cli/src/serve/routes/session-runtime.ts +++ b/packages/cli/src/serve/routes/session-runtime.ts @@ -11,6 +11,7 @@ import type { WorkspaceRuntime, } from '../workspace-registry.js'; import { sendUntrustedWorkspaceResponse } from '../workspace-route-runtime.js'; +import { setDaemonTelemetryWorkspace } from '../server/telemetry.js'; export function requireSessionRuntime(opts: { sessionId: string; @@ -29,12 +30,15 @@ export function requireSessionRuntime(opts: { details = {}, } = opts; if (workspaceRegistry.list().length === 1) { - return workspaceRegistry.primary; + const runtime = workspaceRegistry.primary; + setDaemonTelemetryWorkspace(res, runtime.workspaceCwd); + return runtime; } const resolution = workspaceRegistry.resolveLiveSessionOwner(sessionId); if (resolution.kind === 'found') { const runtime = resolution.runtime; + setDaemonTelemetryWorkspace(res, runtime.workspaceCwd); if (!runtime.primary && !runtime.trusted) { daemonLog?.warn('session routing failed', { route, diff --git a/packages/cli/src/serve/routes/session-telemetry.test.ts b/packages/cli/src/serve/routes/session-telemetry.test.ts new file mode 100644 index 00000000000..564b202fb93 --- /dev/null +++ b/packages/cli/src/serve/routes/session-telemetry.test.ts @@ -0,0 +1,185 @@ +/** + * @license + * Copyright 2026 Qwen Team + * SPDX-License-Identifier: Apache-2.0 + */ + +import path from 'node:path'; +import express, { type Response } from 'express'; +import request from 'supertest'; +import { beforeEach, describe, expect, it, vi } from 'vitest'; +import { + SessionNotFoundError, + type AcpSessionBridge, +} from '../acp-session-bridge.js'; +import { + createWorkspaceRegistry, + type WorkspaceRuntime, +} from '../workspace-registry.js'; + +const telemetryMocks = vi.hoisted(() => ({ + setDaemonTelemetryWorkspace: vi.fn(), +})); + +vi.mock('../server/telemetry.js', () => telemetryMocks); + +import { registerSessionRoutes } from './session.js'; + +function bridgeWithSessions(sessionIds: string[] = []): AcpSessionBridge { + return { + getSessionSummary: vi.fn((sessionId: string) => { + if (!sessionIds.includes(sessionId)) { + throw new SessionNotFoundError(sessionId); + } + return { sessionId }; + }), + } as unknown as AcpSessionBridge; +} + +function runtime(opts: { + workspaceId: string; + workspaceCwd: string; + primary: boolean; + bridge: AcpSessionBridge; + trusted?: boolean; +}): WorkspaceRuntime { + return { + ...opts, + trusted: opts.trusted !== false, + } as WorkspaceRuntime; +} + +function makeApp(runtimes: WorkspaceRuntime[]) { + const app = express(); + app.use(express.json()); + const registry = createWorkspaceRegistry(runtimes); + const primary = registry.primary; + registerSessionRoutes(app, { + boundWorkspace: primary.workspaceCwd, + bridge: primary.bridge, + workspaceRegistry: registry, + archiveCoordinator: { + runSharedMany: async (_sessionIds, fn) => await fn(), + } as Parameters[1]['archiveCoordinator'], + mutate: () => (_req, _res, next) => next(), + sendBridgeError: (res: Response) => { + res.status(500).json({ error: 'test bridge error' }); + }, + sessionShellCommandEnabled: true, + languageCodes: ['en'], + }); + return app; +} + +describe('special session resolver telemetry publication', () => { + const primaryCwd = path.resolve('/workspace/telemetry-primary'); + const secondaryCwd = path.resolve('/workspace/telemetry-secondary'); + + beforeEach(() => { + vi.clearAllMocks(); + }); + + it('publishes the runtime root for creation before later validation', async () => { + const primary = runtime({ + workspaceId: 'primary', + workspaceCwd: primaryCwd, + primary: true, + bridge: bridgeWithSessions(), + }); + + const res = await request(makeApp([primary])) + .post('/session') + .send({ sessionScope: 'invalid' }); + + expect(res.status).toBe(400); + expect(telemetryMocks.setDaemonTelemetryWorkspace).toHaveBeenCalledWith( + expect.anything(), + primaryCwd, + ); + }); + + it('publishes the restore target before later validation', async () => { + const primary = runtime({ + workspaceId: 'primary', + workspaceCwd: primaryCwd, + primary: true, + bridge: bridgeWithSessions(), + }); + + const res = await request(makeApp([primary])) + .post('/session/persisted/load') + .send({ approvalMode: 'invalid' }); + + expect(res.status).toBe(400); + expect(telemetryMocks.setDaemonTelemetryWorkspace).toHaveBeenCalledWith( + expect.anything(), + primaryCwd, + ); + }); + + it('retains the selected restore target on a live-owner conflict', async () => { + const primary = runtime({ + workspaceId: 'primary', + workspaceCwd: primaryCwd, + primary: true, + bridge: bridgeWithSessions(['shared-session']), + }); + const secondary = runtime({ + workspaceId: 'secondary', + workspaceCwd: secondaryCwd, + primary: false, + bridge: bridgeWithSessions(), + }); + + const res = await request(makeApp([primary, secondary])) + .post('/session/shared-session/load') + .send({ cwd: secondaryCwd }); + + expect(res.status).toBe(409); + expect(telemetryMocks.setDaemonTelemetryWorkspace).toHaveBeenCalledWith( + expect.anything(), + secondaryCwd, + ); + }); + + it('publishes the single transcript runtime before storage lookup', async () => { + const primary = runtime({ + workspaceId: 'primary', + workspaceCwd: primaryCwd, + primary: true, + bridge: bridgeWithSessions(), + }); + + const res = await request(makeApp([primary])).get( + '/session/missing/transcript', + ); + + expect(res.status).toBe(500); + expect(telemetryMocks.setDaemonTelemetryWorkspace).toHaveBeenCalledWith( + expect.anything(), + primaryCwd, + ); + }); + + it('does not publish a workspace for a creation workspace mismatch', async () => { + const primary = runtime({ + workspaceId: 'primary', + workspaceCwd: primaryCwd, + primary: true, + bridge: bridgeWithSessions(), + }); + const secondary = runtime({ + workspaceId: 'secondary', + workspaceCwd: secondaryCwd, + primary: false, + bridge: bridgeWithSessions(), + }); + + const res = await request(makeApp([primary, secondary])) + .post('/session') + .send({ cwd: '/workspace/not-registered' }); + + expect(res.status).toBe(400); + expect(telemetryMocks.setDaemonTelemetryWorkspace).not.toHaveBeenCalled(); + }); +}); diff --git a/packages/cli/src/serve/routes/session.ts b/packages/cli/src/serve/routes/session.ts index 583c4be30d8..4aa739c9c5a 100644 --- a/packages/cli/src/serve/routes/session.ts +++ b/packages/cli/src/serve/routes/session.ts @@ -73,6 +73,7 @@ import { parseSessionExportFormat, sessionExportFormatValues, } from '../server/session-export.js'; +import { setDaemonTelemetryWorkspace } from '../server/telemetry.js'; import { createSessionOrganizationService } from '../session-organization-helpers.js'; import { replayTranscriptRecordPage } from '../../acp-integration/session/history-replay-page.js'; import { GENERATION_MAX_PROMPT_BYTES } from '../../acp-integration/generation.js'; @@ -429,6 +430,7 @@ export function registerSessionRoutes( return undefined; } if (workspaceRegistry.list().length === 1) { + setDaemonTelemetryWorkspace(res, workspaceRegistry.primary.workspaceCwd); return { runtime: workspaceRegistry.primary, workspaceCwd: @@ -445,6 +447,7 @@ export function registerSessionRoutes( sendWorkspaceMismatch(res, key); return undefined; } + setDaemonTelemetryWorkspace(res, runtime.workspaceCwd); if (!runtime.primary && !runtime.trusted) { logSessionRoutingFailure('POST /session', 'untrusted_workspace', { workspaceId: runtime.workspaceId, @@ -787,6 +790,7 @@ export function registerSessionRoutes( sendWorkspaceMismatch(res, key); return undefined; } + setDaemonTelemetryWorkspace(res, runtime.workspaceCwd); if (!runtime.primary && !runtime.trusted) { logSessionRoutingFailure(route, 'untrusted_workspace', { workspaceId: runtime.workspaceId, @@ -920,6 +924,7 @@ export function registerSessionRoutes( if (workspaceRegistry.list().length === 1) { const runtime = workspaceRegistry.primary; + setDaemonTelemetryWorkspace(res, runtime.workspaceCwd); if (await activeInRuntime(runtime)) { return runtime; } @@ -932,6 +937,7 @@ export function registerSessionRoutes( return undefined; } if (liveOwner.kind === 'found') { + setDaemonTelemetryWorkspace(res, liveOwner.runtime.workspaceCwd); if ( !assertTrustedSessionOwner(res, route, sessionId, liveOwner.runtime) ) { @@ -981,6 +987,7 @@ export function registerSessionRoutes( } if (activeRuntimes.length === 1) { const runtime = activeRuntimes[0]!; + setDaemonTelemetryWorkspace(res, runtime.workspaceCwd); if (!assertTrustedSessionOwner(res, route, sessionId, runtime)) { return undefined; } diff --git a/packages/cli/src/serve/server.ts b/packages/cli/src/serve/server.ts index 22546b12619..5ebfe932201 100644 --- a/packages/cli/src/serve/server.ts +++ b/packages/cli/src/serve/server.ts @@ -1054,40 +1054,27 @@ export function createServeApp( }); app.use( - daemonTelemetryMiddleware( - (req) => { - const match = req.path.match(/^\/workspaces\/([^/]+)/); - const rawSelector = match?.[1]; - if (rawSelector) { - try { - const selector = decodeURIComponent(rawSelector); - const byId = workspaceRegistry.getByWorkspaceId(selector); - if (byId) return byId.workspaceCwd; - if (isPortableAbsolutePath(selector)) { - const runtime = resolveRegisteredWorkspaceRuntimeByPathSelector( - workspaceRegistry, - selector, - ); - if (runtime) return runtime.workspaceCwd; - } - } catch { - return primaryBoundWorkspace; - } - } - return primaryBoundWorkspace; - }, - deps.recordDaemonRequest, - (sessionId) => { + daemonTelemetryMiddleware((req) => { + const match = req.path.match(/^\/workspaces\/([^/]+)/); + const rawSelector = match?.[1]; + if (rawSelector) { try { - const owner = workspaceRegistry.resolveLiveSessionOwner(sessionId); - return owner.kind === 'found' - ? owner.runtime.workspaceCwd - : undefined; + const selector = decodeURIComponent(rawSelector); + const byId = workspaceRegistry.getByWorkspaceId(selector); + if (byId) return byId.workspaceCwd; + if (isPortableAbsolutePath(selector)) { + const runtime = resolveRegisteredWorkspaceRuntimeByPathSelector( + workspaceRegistry, + selector, + ); + if (runtime) return runtime.workspaceCwd; + } } catch { - return undefined; + return primaryBoundWorkspace; } - }, - ), + } + return primaryBoundWorkspace; + }, deps.recordDaemonRequest), ); const buildWorkspaceCtx = createBuildWorkspaceCtx(primaryBoundWorkspace); diff --git a/packages/cli/src/serve/server/telemetry-catalog.test.ts b/packages/cli/src/serve/server/telemetry-catalog.test.ts new file mode 100644 index 00000000000..90ebf05cb4e --- /dev/null +++ b/packages/cli/src/serve/server/telemetry-catalog.test.ts @@ -0,0 +1,104 @@ +/** + * @license + * Copyright 2026 Qwen Team + * SPDX-License-Identifier: Apache-2.0 + */ + +import express, { type RequestHandler } from 'express'; +import { describe, expect, it, vi } from 'vitest'; +import { registerA2uiActionRoutes } from '../routes/a2ui-action.js'; +import { registerPermissionRoutes } from '../routes/permission.js'; +import { registerSessionRoutes } from '../routes/session.js'; +import { registerSseEventsRoutes } from '../routes/sse-events.js'; +import { legacySessionTelemetryRoutes } from './telemetry.js'; + +interface RouterLayer { + route?: { + path?: unknown; + methods?: Record; + }; + handle?: { + stack?: RouterLayer[]; + }; +} + +function collectExplicitRoutes(layers: RouterLayer[]): string[] { + const routes: string[] = []; + for (const layer of layers) { + const routePath = layer.route?.path; + const paths = + typeof routePath === 'string' + ? [routePath] + : Array.isArray(routePath) + ? routePath.filter((path): path is string => typeof path === 'string') + : []; + for (const [method, enabled] of Object.entries( + layer.route?.methods ?? {}, + )) { + if (!enabled) continue; + for (const path of paths) { + if (/^\/(session|sessions|permission)(?:\/|$)/.test(path)) { + routes.push(`${method.toUpperCase()} ${path}`); + } + } + } + if (layer.handle?.stack) { + routes.push(...collectExplicitRoutes(layer.handle.stack)); + } + } + return routes; +} + +describe('legacy session telemetry route drift guard', () => { + it('matches the explicit Express route registrations in both directions', () => { + const app = express(); + const pass: RequestHandler = (_req, _res, next) => next(); + const mutate = () => pass; + + registerSessionRoutes(app, { + boundWorkspace: '/workspace/primary', + bridge: {} as Parameters[1]['bridge'], + workspaceRegistry: {} as Parameters< + typeof registerSessionRoutes + >[1]['workspaceRegistry'], + archiveCoordinator: {} as Parameters< + typeof registerSessionRoutes + >[1]['archiveCoordinator'], + mutate, + sendBridgeError: vi.fn(), + sessionShellCommandEnabled: true, + languageCodes: [], + }); + registerPermissionRoutes(app, { + bridge: {} as Parameters[1]['bridge'], + workspaceRegistry: {} as Parameters< + typeof registerPermissionRoutes + >[1]['workspaceRegistry'], + mutate, + sendPermissionVoteError: vi.fn(), + }); + registerSseEventsRoutes(app, { + bridge: {} as Parameters[1]['bridge'], + workspaceRegistry: {} as Parameters< + typeof registerSseEventsRoutes + >[1]['workspaceRegistry'], + sendBridgeError: vi.fn(), + }); + registerA2uiActionRoutes(app, { + boundWorkspace: '/workspace/primary', + mutate, + safeBody: () => ({}), + getMcpServers: async () => [], + }); + + const router = (app as unknown as { router: { stack: RouterLayer[] } }) + .router; + const registered = collectExplicitRoutes(router.stack).sort(); + const catalog = legacySessionTelemetryRoutes + .map(({ method, path }) => `${method} ${path}`) + .sort(); + + expect(registered).toHaveLength(48); + expect(registered).toEqual(catalog); + }); +}); diff --git a/packages/cli/src/serve/server/telemetry.test.ts b/packages/cli/src/serve/server/telemetry.test.ts index 4f8f402cd35..a7666dae3f6 100644 --- a/packages/cli/src/serve/server/telemetry.test.ts +++ b/packages/cli/src/serve/server/telemetry.test.ts @@ -13,8 +13,10 @@ const coreMocks = vi.hoisted(() => ({ recordDaemonError: vi.fn(), recordDaemonHttpRequest: vi.fn(), recordDaemonHttpResponse: vi.fn(), + spanSetAttribute: vi.fn(), withDaemonRequestSpan: vi.fn( - (_attrs: unknown, fn: (span: unknown) => Promise) => fn({}), + (_attrs: unknown, fn: (span: unknown) => Promise) => + fn({ setAttribute: coreMocks.spanSetAttribute }), ), })); @@ -25,7 +27,13 @@ vi.mock('@qwen-code/qwen-code-core', () => ({ ...coreMocks, })); -import { daemonTelemetryMiddleware } from './telemetry.js'; +import { + daemonTelemetryMiddleware, + legacySessionTelemetryRoutes, + resolveDaemonTelemetryRoute, + setDaemonTelemetryWorkspace, +} from './telemetry.js'; +import { MAX_CLIENT_ID_LENGTH } from './request-helpers.js'; function mockReq(method: string, path: string): Request { return { method, path, get: () => undefined } as unknown as Request; @@ -34,12 +42,17 @@ function mockReq(method: string, path: string): Request { function mockRes(statusCode: number): Response & EventEmitter { const res = new EventEmitter() as Response & EventEmitter; (res as { statusCode: number }).statusCode = statusCode; + Object.defineProperty(res, 'headersSent', { value: true, writable: true }); return res; } describe('daemonTelemetryMiddleware — recordRequest seam', () => { beforeEach(() => { vi.clearAllMocks(); + coreMocks.hashDaemonWorkspace.mockImplementation( + (workspace: string) => `hash:${workspace}`, + ); + coreMocks.spanSetAttribute.mockImplementation(() => undefined); }); it('calls recordRequest with (durationMs, statusCode) once the response finishes on a matched route', () => { @@ -83,6 +96,8 @@ describe('daemonTelemetryMiddleware — recordRequest seam', () => { res.emit('finish'); res.emit('close'); expect(recordRequest).toHaveBeenCalledTimes(1); + expect(coreMocks.recordDaemonHttpRequest).toHaveBeenCalledTimes(1); + expect(coreMocks.recordDaemonHttpResponse).toHaveBeenCalledTimes(1); }); it('does NOT call recordRequest for an unmatched route', () => { @@ -187,13 +202,8 @@ describe('daemonTelemetryMiddleware — recordRequest seam', () => { ); }); - it('attributes singular rewind and shell routes to the live session owner', () => { - const resolveSessionWorkspaceCwd = vi.fn(() => '/workspace/secondary'); - const mw = daemonTelemetryMiddleware( - () => '/workspace/primary', - undefined, - resolveSessionWorkspaceCwd, - ); + it('defers singular owner-routed workspace attribution until the handler selects a runtime', () => { + const mw = daemonTelemetryMiddleware(() => '/workspace/primary'); for (const [method, path, route] of [ [ @@ -206,31 +216,29 @@ describe('daemonTelemetryMiddleware — recordRequest seam', () => { ] as const) { const res = mockRes(200); mw(mockReq(method, path), res, vi.fn() as unknown as NextFunction); + expect(coreMocks.withDaemonRequestSpan).toHaveBeenLastCalledWith( + expect.not.objectContaining({ workspaceHash: expect.anything() }), + expect.any(Function), + ); + setDaemonTelemetryWorkspace(res, '/workspace/secondary'); res.emit('finish'); expect(coreMocks.withDaemonRequestSpan).toHaveBeenLastCalledWith( expect.objectContaining({ method, route, sessionId: 'secondary-session', - workspaceHash: 'hash:/workspace/secondary', }), expect.any(Function), ); + expect(coreMocks.spanSetAttribute).toHaveBeenLastCalledWith( + 'qwen-code.workspace.hash', + 'hash:/workspace/secondary', + ); } - - expect(resolveSessionWorkspaceCwd).toHaveBeenCalledTimes(3); - expect(resolveSessionWorkspaceCwd).toHaveBeenCalledWith( - 'secondary-session', - ); }); - it('decodes session ids before owner lookup and span attribution', () => { - const resolveSessionWorkspaceCwd = vi.fn(() => '/workspace/secondary'); - const mw = daemonTelemetryMiddleware( - () => '/workspace/primary', - undefined, - resolveSessionWorkspaceCwd, - ); + it('decodes session ids before span attribution', () => { + const mw = daemonTelemetryMiddleware(() => '/workspace/primary'); const res = mockRes(200); mw( @@ -240,9 +248,6 @@ describe('daemonTelemetryMiddleware — recordRequest seam', () => { ); res.emit('finish'); - expect(resolveSessionWorkspaceCwd).toHaveBeenCalledWith( - 'secondary/session', - ); expect(coreMocks.withDaemonRequestSpan).toHaveBeenCalledWith( expect.objectContaining({ sessionId: 'secondary/session' }), expect.any(Function), @@ -250,12 +255,7 @@ describe('daemonTelemetryMiddleware — recordRequest seam', () => { }); it('keeps malformed session id encodings without throwing', () => { - const resolveSessionWorkspaceCwd = vi.fn(() => undefined); - const mw = daemonTelemetryMiddleware( - () => '/workspace/primary', - undefined, - resolveSessionWorkspaceCwd, - ); + const mw = daemonTelemetryMiddleware(() => '/workspace/primary'); const res = mockRes(200); expect(() => { @@ -267,7 +267,6 @@ describe('daemonTelemetryMiddleware — recordRequest seam', () => { }).not.toThrow(); res.emit('finish'); - expect(resolveSessionWorkspaceCwd).toHaveBeenCalledWith('bad%ZZ'); expect(coreMocks.withDaemonRequestSpan).toHaveBeenCalledWith( expect.objectContaining({ sessionId: 'bad%ZZ' }), expect.any(Function), @@ -344,8 +343,116 @@ describe('daemonTelemetryMiddleware — recordRequest seam', () => { ); res.emit('finish'); expect(recordRequest).not.toHaveBeenCalled(); + expect(coreMocks.recordDaemonHttpRequest).toHaveBeenCalledTimes(1); + }); + + it('keeps heartbeat in OTel HTTP metrics but excludes it from the metrics ring', () => { + const recordRequest = vi.fn(); + const mw = daemonTelemetryMiddleware(() => '/ws', recordRequest); + const res = mockRes(200); + + mw( + mockReq('POST', '/session/abc/heartbeat'), + res, + vi.fn() as unknown as NextFunction, + ); + setDaemonTelemetryWorkspace(res, '/ws'); + res.emit('finish'); + + expect(coreMocks.recordDaemonHttpRequest).toHaveBeenCalledWith( + expect.any(Number), + 'POST /session/:id/heartbeat', + 200, + ); + expect(recordRequest).not.toHaveBeenCalled(); + }); + + it('does not record successful SSE connection lifetime as HTTP request latency', () => { + const recordRequest = vi.fn(); + const mw = daemonTelemetryMiddleware(() => '/ws', recordRequest); + const res = mockRes(200); + + mw( + mockReq('GET', '/session/abc/events'), + res, + vi.fn() as unknown as NextFunction, + ); + setDaemonTelemetryWorkspace(res, '/ws'); + res.emit('close'); + + expect(coreMocks.recordDaemonHttpResponse).toHaveBeenCalledTimes(1); + expect(coreMocks.recordDaemonHttpRequest).not.toHaveBeenCalled(); + expect(recordRequest).not.toHaveBeenCalled(); + }); + + it('records request-scoped generation SSE duration as ordinary HTTP latency', () => { + const recordRequest = vi.fn(); + const mw = daemonTelemetryMiddleware(() => '/ws', recordRequest); + const res = mockRes(200); + + mw( + mockReq('POST', '/session/abc/generate'), + res, + vi.fn() as unknown as NextFunction, + ); + setDaemonTelemetryWorkspace(res, '/ws'); + res.emit('finish'); + + expect(coreMocks.recordDaemonHttpRequest).toHaveBeenCalledWith( + expect.any(Number), + 'POST /session/:id/generate', + 200, + ); + expect(recordRequest).toHaveBeenCalledWith(expect.any(Number), 200); + }); + + it('counts a 200 SSE request that closes before response headers are sent', () => { + const recordRequest = vi.fn(); + const mw = daemonTelemetryMiddleware(() => '/ws', recordRequest); + const res = mockRes(200); + (res as unknown as { headersSent: boolean }).headersSent = false; + + mw( + mockReq('GET', '/session/abc/events'), + res, + vi.fn() as unknown as NextFunction, + ); + res.emit('close'); + + expect(coreMocks.recordDaemonHttpRequest).toHaveBeenCalledWith( + expect.any(Number), + 'GET /session/:id/events', + 200, + ); + expect(recordRequest).toHaveBeenCalledWith(expect.any(Number), 200); }); + it.each([400, 404, 429, 500])( + 'records an SSE handshake failure with status %s as an ordinary request', + (statusCode) => { + const recordRequest = vi.fn(); + const mw = daemonTelemetryMiddleware(() => '/ws', recordRequest); + const res = mockRes(statusCode); + + mw( + mockReq('GET', '/session/abc/events'), + res, + vi.fn() as unknown as NextFunction, + ); + res.emit('finish'); + + expect(coreMocks.recordDaemonHttpRequest).toHaveBeenCalledWith( + expect.any(Number), + 'GET /session/:id/events', + statusCode, + ); + expect(recordRequest).toHaveBeenCalledWith( + expect.any(Number), + statusCode, + ); + }, + ); + it('is a silent no-op when recordRequest is omitted (the optional-chaining path)', () => { const mw = daemonTelemetryMiddleware(() => '/ws'); const res = mockRes(200); @@ -359,9 +466,36 @@ describe('daemonTelemetryMiddleware — recordRequest seam', () => { }).not.toThrow(); }); + it('settles normally when telemetry is disabled and no span is created', () => { + const recordRequest = vi.fn(); + coreMocks.withDaemonRequestSpan.mockImplementationOnce( + (_attrs: unknown, fn: (span: unknown) => Promise) => fn(undefined), + ); + const mw = daemonTelemetryMiddleware( + () => '/workspace/primary', + recordRequest, + ); + const res = mockRes(200); + + mw( + mockReq('POST', '/session/abc/prompt'), + res, + vi.fn() as unknown as NextFunction, + ); + setDaemonTelemetryWorkspace(res, '/workspace/secondary'); + expect(() => res.emit('finish')).not.toThrow(); + + expect(coreMocks.hashDaemonWorkspace).not.toHaveBeenCalled(); + expect(coreMocks.recordDaemonHttpResponse).toHaveBeenCalledWith( + undefined, + 200, + ); + expect(coreMocks.recordDaemonHttpRequest).toHaveBeenCalledTimes(1); + expect(recordRequest).toHaveBeenCalledTimes(1); + }); + it('resolves workspace hash per request instead of closing over the primary workspace', () => { - let workspace = '/workspace/one'; - const mw = daemonTelemetryMiddleware(() => workspace); + const mw = daemonTelemetryMiddleware(() => '/workspace/primary'); const firstRes = mockRes(200); mw( @@ -369,15 +503,16 @@ describe('daemonTelemetryMiddleware — recordRequest seam', () => { firstRes, vi.fn() as unknown as NextFunction, ); + setDaemonTelemetryWorkspace(firstRes, '/workspace/one'); firstRes.emit('finish'); - workspace = '/workspace/two'; const secondRes = mockRes(200); mw( mockReq('POST', '/session/abc/prompt'), secondRes, vi.fn() as unknown as NextFunction, ); + setDaemonTelemetryWorkspace(secondRes, '/workspace/two'); secondRes.emit('finish'); expect(coreMocks.hashDaemonWorkspace).toHaveBeenNthCalledWith( @@ -388,12 +523,16 @@ describe('daemonTelemetryMiddleware — recordRequest seam', () => { 2, '/workspace/two', ); - expect(coreMocks.withDaemonRequestSpan.mock.calls[0]?.[0]).toMatchObject({ - workspaceHash: 'hash:/workspace/one', - }); - expect(coreMocks.withDaemonRequestSpan.mock.calls[1]?.[0]).toMatchObject({ - workspaceHash: 'hash:/workspace/two', - }); + expect(coreMocks.spanSetAttribute).toHaveBeenNthCalledWith( + 1, + 'qwen-code.workspace.hash', + 'hash:/workspace/one', + ); + expect(coreMocks.spanSetAttribute).toHaveBeenNthCalledWith( + 2, + 'qwen-code.workspace.hash', + 'hash:/workspace/two', + ); }); it('memoizes workspace hashes by resolved workspace cwd', () => { @@ -406,12 +545,14 @@ describe('daemonTelemetryMiddleware — recordRequest seam', () => { firstRes, vi.fn() as unknown as NextFunction, ); + setDaemonTelemetryWorkspace(firstRes, '/workspace/one'); firstRes.emit('finish'); mw( mockReq('POST', '/session/abc/prompt'), secondRes, vi.fn() as unknown as NextFunction, ); + setDaemonTelemetryWorkspace(secondRes, '/workspace/one'); secondRes.emit('finish'); expect(coreMocks.hashDaemonWorkspace).toHaveBeenCalledTimes(1); @@ -419,4 +560,275 @@ describe('daemonTelemetryMiddleware — recordRequest seam', () => { '/workspace/one', ); }); + + it('settles a published workspace after its runtime is removed', () => { + const resolveWorkspaceCwd = vi.fn(() => '/workspace/primary'); + const mw = daemonTelemetryMiddleware(resolveWorkspaceCwd); + const runtimes = new Map([ + ['secondary', { workspaceCwd: '/workspace/secondary' }], + ]); + const runtime = runtimes.get('secondary')!; + const res = mockRes(200); + + mw( + mockReq('POST', '/session/abc/prompt'), + res, + vi.fn() as unknown as NextFunction, + ); + setDaemonTelemetryWorkspace(res, runtime.workspaceCwd); + runtimes.delete('secondary'); + res.emit('finish'); + + expect(resolveWorkspaceCwd).not.toHaveBeenCalled(); + expect(coreMocks.spanSetAttribute).toHaveBeenCalledWith( + 'qwen-code.workspace.hash', + 'hash:/workspace/secondary', + ); + }); + + it('uses first-selection-wins and clears deferred context after settlement', () => { + const mw = daemonTelemetryMiddleware(() => '/workspace/primary'); + const res = mockRes(200); + mw( + mockReq('POST', '/session/abc/prompt'), + res, + vi.fn() as unknown as NextFunction, + ); + + setDaemonTelemetryWorkspace(res, '/workspace/first'); + setDaemonTelemetryWorkspace(res, '/workspace/first'); + setDaemonTelemetryWorkspace(res, '/workspace/second'); + res.emit('finish'); + setDaemonTelemetryWorkspace(res, '/workspace/after-finish'); + + expect(coreMocks.spanSetAttribute).toHaveBeenCalledTimes(1); + expect(coreMocks.spanSetAttribute).toHaveBeenCalledWith( + 'qwen-code.workspace.hash', + 'hash:/workspace/first', + ); + }); + + it('omits workspace hash when a dynamic target is never resolved', () => { + const resolveWorkspaceCwd = vi.fn(() => '/workspace/primary'); + const mw = daemonTelemetryMiddleware(resolveWorkspaceCwd); + const res = mockRes(404); + + mw( + mockReq('POST', '/session/missing/prompt'), + res, + vi.fn() as unknown as NextFunction, + ); + res.emit('finish'); + + expect(resolveWorkspaceCwd).not.toHaveBeenCalled(); + expect(coreMocks.hashDaemonWorkspace).not.toHaveBeenCalled(); + expect(coreMocks.spanSetAttribute).not.toHaveBeenCalled(); + }); + + it('keeps pre-resolved resolver failures from affecting request settlement', () => { + const recordRequest = vi.fn(); + const next = vi.fn() as unknown as NextFunction; + const mw = daemonTelemetryMiddleware(() => { + throw new Error('resolver failed'); + }, recordRequest); + const res = mockRes(200); + + expect(() => mw(mockReq('GET', '/daemon/status'), res, next)).not.toThrow(); + res.emit('finish'); + + expect(next).toHaveBeenCalledTimes(1); + expect(coreMocks.withDaemonRequestSpan).toHaveBeenCalledWith( + expect.not.objectContaining({ workspaceHash: expect.anything() }), + expect.any(Function), + ); + expect(coreMocks.recordDaemonHttpRequest).toHaveBeenCalledTimes(1); + }); + + it('keeps late hash and span attribute failures from affecting metrics', () => { + const recordRequest = vi.fn(); + const mw = daemonTelemetryMiddleware( + () => '/workspace/primary', + recordRequest, + ); + const hashFailureRes = mockRes(200); + coreMocks.hashDaemonWorkspace.mockImplementationOnce(() => { + throw new Error('hash failed'); + }); + + mw( + mockReq('POST', '/session/abc/prompt'), + hashFailureRes, + vi.fn() as unknown as NextFunction, + ); + setDaemonTelemetryWorkspace(hashFailureRes, '/workspace/secondary'); + expect(() => hashFailureRes.emit('finish')).not.toThrow(); + + const attributeFailureRes = mockRes(200); + coreMocks.spanSetAttribute.mockImplementationOnce(() => { + throw new Error('attribute failed'); + }); + mw( + mockReq('POST', '/session/def/prompt'), + attributeFailureRes, + vi.fn() as unknown as NextFunction, + ); + setDaemonTelemetryWorkspace(attributeFailureRes, '/workspace/secondary'); + expect(() => attributeFailureRes.emit('finish')).not.toThrow(); + + expect(coreMocks.recordDaemonHttpRequest).toHaveBeenCalledTimes(2); + expect(recordRequest).toHaveBeenCalledTimes(2); + }); + + it('is a safe no-op when workspace selection is published without middleware context', () => { + const res = mockRes(200); + expect(() => + setDaemonTelemetryWorkspace(res, '/workspace/secondary'), + ).not.toThrow(); + expect(coreMocks.spanSetAttribute).not.toHaveBeenCalled(); + }); + + it('continues a dynamic request when its Response cannot store telemetry context', () => { + const next = vi.fn() as unknown as NextFunction; + const mw = daemonTelemetryMiddleware(() => '/workspace/primary'); + const res = Object.preventExtensions(mockRes(200)); + + expect(() => + mw(mockReq('POST', '/session/abc/prompt'), res, next), + ).not.toThrow(); + expect(() => + setDaemonTelemetryWorkspace(res, '/workspace/secondary'), + ).not.toThrow(); + res.emit('finish'); + + expect(next).toHaveBeenCalledTimes(1); + expect(coreMocks.spanSetAttribute).not.toHaveBeenCalled(); + expect(coreMocks.recordDaemonHttpRequest).toHaveBeenCalledTimes(1); + }); +}); + +describe('legacy session telemetry route catalog', () => { + it('contains 48 unique routes with the audited 41/7 attribution split', () => { + const keys = legacySessionTelemetryRoutes.map( + ({ method, path }) => `${method} ${path}`, + ); + expect(keys).toHaveLength(48); + expect(new Set(keys).size).toBe(48); + expect( + legacySessionTelemetryRoutes.filter( + ({ attribution }) => attribution === 'handler_resolved', + ), + ).toHaveLength(41); + expect( + legacySessionTelemetryRoutes.filter( + ({ attribution }) => attribution === 'pre_resolved', + ), + ).toHaveLength(7); + expect( + legacySessionTelemetryRoutes + .filter(({ attribution }) => attribution === 'pre_resolved') + .map(({ method, path }) => `${method} ${path}`) + .sort(), + ).toEqual( + [ + 'GET /session/:id/export', + 'PATCH /session/:id/organization', + 'POST /permission/:requestId', + 'POST /session/:id/a2ui-action', + 'POST /sessions/archive', + 'POST /sessions/delete', + 'POST /sessions/unarchive', + ].sort(), + ); + for (const entry of legacySessionTelemetryRoutes) { + expect(entry.route).toBe(`${entry.method} ${entry.path}`); + } + }); + + it('matches every catalog entry with its declared canonical attribution', () => { + for (const entry of legacySessionTelemetryRoutes) { + const path = entry.path.replace( + /:([A-Za-z][A-Za-z0-9_]*)/g, + (_match, name: string) => { + if (name === 'id') return 'session-1'; + if (name === 'requestId') return 'request-1'; + return `${name}-1`; + }, + ); + + expect(resolveDaemonTelemetryRoute(mockReq(entry.method, path))).toEqual({ + route: entry.route, + attribution: entry.attribution, + ...(entry.path.includes('/:id') ? { sessionId: 'session-1' } : {}), + ...(entry.path.includes('/:requestId') + ? { permissionRequestId: 'request-1' } + : {}), + }); + } + }); + + it.each([ + ['POST', '/SeSsIoN/abc/PrOmPt/', 'POST /session/:id/prompt', 'abc'], + [ + 'POST', + '/session/session%2Fchild/prompt', + 'POST /session/:id/prompt', + 'session/child', + ], + [ + 'POST', + '/session/session%252Fchild/prompt', + 'POST /session/:id/prompt', + 'session%2Fchild', + ], + [ + 'GET', + '/session/%E4%BD%A0%E5%A5%BD/status', + 'GET /session/:id/status', + '你好', + ], + ['POST', '/session/bad%ZZ/rewind', 'POST /session/:id/rewind', 'bad%ZZ'], + ])( + 'matches %s %s with a canonical label', + (method, path, route, sessionId) => { + expect(resolveDaemonTelemetryRoute(mockReq(method, path))).toMatchObject({ + route, + sessionId, + }); + }, + ); + + it('decodes and validates permission request ids after segment matching', () => { + expect( + resolveDaemonTelemetryRoute( + mockReq('POST', '/session/abc/permission/%72eq-1'), + ), + ).toMatchObject({ + route: 'POST /session/:id/permission/:requestId', + sessionId: 'abc', + permissionRequestId: 'req-1', + }); + expect( + resolveDaemonTelemetryRoute( + mockReq('POST', '/session/abc/permission/req%2F1'), + ), + ).not.toHaveProperty('permissionRequestId'); + expect( + resolveDaemonTelemetryRoute(mockReq('POST', '/permission/bad%ZZ')), + ).not.toHaveProperty('permissionRequestId'); + expect( + resolveDaemonTelemetryRoute( + mockReq('POST', `/permission/${'a'.repeat(MAX_CLIENT_ID_LENGTH + 1)}`), + ), + ).not.toHaveProperty('permissionRequestId'); + }); + + it.each([ + ['GET', '/session/abc/prompt'], + ['POST', '/session/abc/prompt/extra'], + ['POST', '/session/abc/prompt//'], + ['POST', '/session//prompt'], + ['HEAD', '/session/abc/status'], + ])('does not match the wrong method or path: %s %s', (method, path) => { + expect(resolveDaemonTelemetryRoute(mockReq(method, path))).toBeUndefined(); + }); }); diff --git a/packages/cli/src/serve/server/telemetry.ts b/packages/cli/src/serve/server/telemetry.ts index 696f6a903f4..d0e2579a194 100644 --- a/packages/cli/src/serve/server/telemetry.ts +++ b/packages/cli/src/serve/server/telemetry.ts @@ -18,6 +18,323 @@ import { MAX_CLIENT_ID_LENGTH, } from './request-helpers.js'; +type LegacySessionTelemetryAttribution = 'handler_resolved' | 'pre_resolved'; + +interface LegacySessionTelemetryRoute { + method: 'DELETE' | 'GET' | 'PATCH' | 'POST'; + path: string; + attribution: LegacySessionTelemetryAttribution; + route: string; +} + +export const legacySessionTelemetryRoutes = [ + { + method: 'POST', + path: '/session', + attribution: 'handler_resolved', + route: 'POST /session', + }, + { + method: 'POST', + path: '/session/:id/load', + attribution: 'handler_resolved', + route: 'POST /session/:id/load', + }, + { + method: 'POST', + path: '/session/:id/resume', + attribution: 'handler_resolved', + route: 'POST /session/:id/resume', + }, + { + method: 'POST', + path: '/session/:id/branch', + attribution: 'handler_resolved', + route: 'POST /session/:id/branch', + }, + { + method: 'POST', + path: '/session/:id/fork', + attribution: 'handler_resolved', + route: 'POST /session/:id/fork', + }, + { + method: 'POST', + path: '/session/:id/cd', + attribution: 'handler_resolved', + route: 'POST /session/:id/cd', + }, + { + method: 'GET', + path: '/session/:id/status', + attribution: 'handler_resolved', + route: 'GET /session/:id/status', + }, + { + method: 'GET', + path: '/session/:id/export', + attribution: 'pre_resolved', + route: 'GET /session/:id/export', + }, + { + method: 'GET', + path: '/session/:id/transcript', + attribution: 'handler_resolved', + route: 'GET /session/:id/transcript', + }, + { + method: 'GET', + path: '/session/:id/context', + attribution: 'handler_resolved', + route: 'GET /session/:id/context', + }, + { + method: 'GET', + path: '/session/:id/context-usage', + attribution: 'handler_resolved', + route: 'GET /session/:id/context-usage', + }, + { + method: 'GET', + path: '/session/:id/stats', + attribution: 'handler_resolved', + route: 'GET /session/:id/stats', + }, + { + method: 'GET', + path: '/session/:id/supported-commands', + attribution: 'handler_resolved', + route: 'GET /session/:id/supported-commands', + }, + { + method: 'GET', + path: '/session/:id/tasks', + attribution: 'handler_resolved', + route: 'GET /session/:id/tasks', + }, + { + method: 'GET', + path: '/session/:id/lsp', + attribution: 'handler_resolved', + route: 'GET /session/:id/lsp', + }, + { + method: 'GET', + path: '/session/:id/hooks', + attribution: 'handler_resolved', + route: 'GET /session/:id/hooks', + }, + { + method: 'GET', + path: '/session/:id/artifacts', + attribution: 'handler_resolved', + route: 'GET /session/:id/artifacts', + }, + { + method: 'POST', + path: '/session/:id/artifacts', + attribution: 'handler_resolved', + route: 'POST /session/:id/artifacts', + }, + { + method: 'DELETE', + path: '/session/:id/artifacts/:artifactId', + attribution: 'handler_resolved', + route: 'DELETE /session/:id/artifacts/:artifactId', + }, + { + method: 'POST', + path: '/session/:id/tasks/:taskId/cancel', + attribution: 'handler_resolved', + route: 'POST /session/:id/tasks/:taskId/cancel', + }, + { + method: 'POST', + path: '/session/:id/goal/clear', + attribution: 'handler_resolved', + route: 'POST /session/:id/goal/clear', + }, + { + method: 'POST', + path: '/session/:id/continue', + attribution: 'handler_resolved', + route: 'POST /session/:id/continue', + }, + { + method: 'POST', + path: '/session/:id/prompt', + attribution: 'handler_resolved', + route: 'POST /session/:id/prompt', + }, + { + method: 'POST', + path: '/session/:id/generate', + attribution: 'handler_resolved', + route: 'POST /session/:id/generate', + }, + { + method: 'POST', + path: '/session/:id/heartbeat', + attribution: 'handler_resolved', + route: 'POST /session/:id/heartbeat', + }, + { + method: 'POST', + path: '/session/:id/detach', + attribution: 'handler_resolved', + route: 'POST /session/:id/detach', + }, + { + method: 'POST', + path: '/session/:id/cancel', + attribution: 'handler_resolved', + route: 'POST /session/:id/cancel', + }, + { + method: 'DELETE', + path: '/session/:id', + attribution: 'handler_resolved', + route: 'DELETE /session/:id', + }, + { + method: 'POST', + path: '/sessions/delete', + attribution: 'pre_resolved', + route: 'POST /sessions/delete', + }, + { + method: 'POST', + path: '/sessions/archive', + attribution: 'pre_resolved', + route: 'POST /sessions/archive', + }, + { + method: 'POST', + path: '/sessions/unarchive', + attribution: 'pre_resolved', + route: 'POST /sessions/unarchive', + }, + { + method: 'PATCH', + path: '/session/:id/metadata', + attribution: 'handler_resolved', + route: 'PATCH /session/:id/metadata', + }, + { + method: 'PATCH', + path: '/session/:id/organization', + attribution: 'pre_resolved', + route: 'PATCH /session/:id/organization', + }, + { + method: 'POST', + path: '/session/:id/model', + attribution: 'handler_resolved', + route: 'POST /session/:id/model', + }, + { + method: 'POST', + path: '/session/:id/recap', + attribution: 'handler_resolved', + route: 'POST /session/:id/recap', + }, + { + method: 'POST', + path: '/session/:id/btw', + attribution: 'handler_resolved', + route: 'POST /session/:id/btw', + }, + { + method: 'POST', + path: '/session/:id/mid-turn-message', + attribution: 'handler_resolved', + route: 'POST /session/:id/mid-turn-message', + }, + { + method: 'GET', + path: '/session/:id/pending-prompts', + attribution: 'handler_resolved', + route: 'GET /session/:id/pending-prompts', + }, + { + method: 'DELETE', + path: '/session/:id/pending-prompts/:promptId', + attribution: 'handler_resolved', + route: 'DELETE /session/:id/pending-prompts/:promptId', + }, + { + method: 'POST', + path: '/session/:id/shell', + attribution: 'handler_resolved', + route: 'POST /session/:id/shell', + }, + { + method: 'GET', + path: '/session/:id/rewind/snapshots', + attribution: 'handler_resolved', + route: 'GET /session/:id/rewind/snapshots', + }, + { + method: 'POST', + path: '/session/:id/rewind', + attribution: 'handler_resolved', + route: 'POST /session/:id/rewind', + }, + { + method: 'POST', + path: '/session/:id/approval-mode', + attribution: 'handler_resolved', + route: 'POST /session/:id/approval-mode', + }, + { + method: 'POST', + path: '/session/:id/language', + attribution: 'handler_resolved', + route: 'POST /session/:id/language', + }, + { + method: 'POST', + path: '/session/:id/permission/:requestId', + attribution: 'handler_resolved', + route: 'POST /session/:id/permission/:requestId', + }, + { + method: 'POST', + path: '/permission/:requestId', + attribution: 'pre_resolved', + route: 'POST /permission/:requestId', + }, + { + method: 'GET', + path: '/session/:id/events', + attribution: 'handler_resolved', + route: 'GET /session/:id/events', + }, + { + method: 'POST', + path: '/session/:id/a2ui-action', + attribution: 'pre_resolved', + route: 'POST /session/:id/a2ui-action', + }, +] as const satisfies readonly LegacySessionTelemetryRoute[]; + +interface ResolvedDaemonTelemetryRoute { + route: string; + sessionId?: string; + permissionRequestId?: string; + attribution?: LegacySessionTelemetryAttribution; +} + +interface DaemonTelemetryResponseContext { + workspaceCwd?: string; +} + +const daemonTelemetryResponseContext = Symbol('daemonTelemetryResponseContext'); + +type TelemetryResponse = Response & { + [daemonTelemetryResponseContext]?: DaemonTelemetryResponseContext; +}; + function decodePathSegment(value: string): string { try { return decodeURIComponent(value); @@ -26,105 +343,90 @@ function decodePathSegment(value: string): string { } } -// Route handlers are split across `routes/*.ts`; any added or renamed route -// that needs daemon telemetry must keep these patterns in sync. -export function resolveDaemonTelemetryRoute( - req: Request, -): - | { route: string; sessionId?: string; permissionRequestId?: string } - | undefined { - const path = req.path.replace(/\/$/, '') || '/'; - if (req.method === 'POST' && path === '/session') { - return { route: 'POST /session' }; - } - if (req.method === 'POST' && path === '/sessions/delete') { - return { route: 'POST /sessions/delete' }; - } - if (req.method === 'GET' && path === '/daemon/status') { - return { route: 'GET /daemon/status' }; - } - const rewindSnapshots = path.match(/^\/session\/([^/]+)\/rewind\/snapshots$/); - if (rewindSnapshots?.[1] && req.method === 'GET') { - return { - route: 'GET /session/:id/rewind/snapshots', - sessionId: rewindSnapshots[1], - }; - } - const sessionAction = path.match( - /^\/session\/([^/]+)\/(load|resume|prompt|cancel|recap|btw|mid-turn-message|model|shell|detach|rewind|approval-mode|language|a2ui-action)$/, - ); - const sessionActionId = sessionAction?.[1]; - const sessionActionName = sessionAction?.[2]; - if (sessionActionId && sessionActionName && req.method === 'POST') { - return { - route: `POST /session/:id/${sessionActionName}`, - sessionId: sessionActionId, - }; - } - const sessionMetadata = path.match(/^\/session\/([^/]+)\/metadata$/); - if (sessionMetadata?.[1] && req.method === 'PATCH') { - return { - route: 'PATCH /session/:id/metadata', - sessionId: sessionMetadata[1], - }; - } - const sessionArtifacts = path.match(/^\/session\/([^/]+)\/artifacts$/); - if (sessionArtifacts?.[1]) { - if (req.method === 'GET') { - return { - route: 'GET /session/:id/artifacts', - sessionId: sessionArtifacts[1], - }; - } - if (req.method === 'POST') { - return { - route: 'POST /session/:id/artifacts', - sessionId: sessionArtifacts[1], - }; - } - } - const sessionArtifact = path.match( - /^\/session\/([^/]+)\/artifacts\/([^/]+)$/, - ); - if (sessionArtifact?.[1] && req.method === 'DELETE') { - return { - route: 'DELETE /session/:id/artifacts/:artifactId', - sessionId: sessionArtifact[1], - }; - } - const sessionPermission = path.match( - /^\/session\/([^/]+)\/permission\/([^/]+)$/, - ); +function matchLegacySessionTelemetryRoute( + method: string, + requestPath: string, +): ResolvedDaemonTelemetryRoute | undefined { + const path = + requestPath.length > 1 && requestPath.endsWith('/') + ? requestPath.slice(0, -1) + : requestPath; + const requestSegments = path.split('/').slice(1); + const prefix = requestSegments[0]?.toLowerCase(); if ( - sessionPermission?.[1] && - sessionPermission?.[2] && - req.method === 'POST' + prefix !== 'session' && + prefix !== 'sessions' && + prefix !== 'permission' ) { - const rawRequestId = sessionPermission[2]; - return { - route: 'POST /session/:id/permission/:requestId', - sessionId: sessionPermission[1], - ...(rawRequestId.length <= MAX_CLIENT_ID_LENGTH && - CLIENT_ID_RE.test(rawRequestId) - ? { permissionRequestId: rawRequestId } - : {}), - }; + return undefined; } - const globalPermission = path.match(/^\/permission\/([^/]+)$/); - if (globalPermission?.[1] && req.method === 'POST') { - const rawRequestId = globalPermission[1]; + + for (const entry of legacySessionTelemetryRoutes) { + if (entry.method !== method) continue; + const templateSegments = entry.path.split('/').slice(1); + if (templateSegments.length !== requestSegments.length) continue; + const params = new Map(); + let matched = true; + for (let index = 0; index < templateSegments.length; index += 1) { + const templateSegment = templateSegments[index]!; + const requestSegment = requestSegments[index]!; + if (templateSegment.startsWith(':')) { + if (requestSegment === '') { + matched = false; + break; + } + params.set(templateSegment.slice(1), requestSegment); + } else if ( + templateSegment.toLowerCase() !== requestSegment.toLowerCase() + ) { + matched = false; + break; + } + } + if (!matched) continue; + + const rawSessionId = params.get('id'); + const rawRequestId = params.get('requestId'); + const requestId = + rawRequestId !== undefined ? decodePathSegment(rawRequestId) : undefined; return { - route: 'POST /permission/:requestId', - ...(rawRequestId.length <= MAX_CLIENT_ID_LENGTH && - CLIENT_ID_RE.test(rawRequestId) - ? { permissionRequestId: rawRequestId } + route: entry.route, + attribution: entry.attribution, + ...(rawSessionId ? { sessionId: decodePathSegment(rawSessionId) } : {}), + ...(requestId !== undefined && + requestId.length <= MAX_CLIENT_ID_LENGTH && + CLIENT_ID_RE.test(requestId) + ? { permissionRequestId: requestId } : {}), }; } - const deleteSession = path.match(/^\/session\/([^/]+)$/); - const deleteSessionId = deleteSession?.[1]; - if (deleteSessionId && req.method === 'DELETE') { - return { route: 'DELETE /session/:id', sessionId: deleteSessionId }; + return undefined; +} + +export function setDaemonTelemetryWorkspace( + res: Response, + workspaceCwd: string, +): void { + try { + const context = (res as TelemetryResponse)[daemonTelemetryResponseContext]; + if (context && context.workspaceCwd === undefined) { + context.workspaceCwd = workspaceCwd; + } + } catch { + // Telemetry must not affect request handling. + } +} + +// Route handlers are split across `routes/*.ts`; any added or renamed route +// that needs daemon telemetry must keep these patterns in sync. +export function resolveDaemonTelemetryRoute( + req: Request, +): ResolvedDaemonTelemetryRoute | undefined { + const legacyRoute = matchLegacySessionTelemetryRoute(req.method, req.path); + if (legacyRoute) return legacyRoute; + const path = req.path.replace(/\/$/, '') || '/'; + if (req.method === 'GET' && path === '/daemon/status') { + return { route: 'GET /daemon/status' }; } if (req.method === 'GET' && /^\/workspace\/[^/]+\/sessions$/.test(path)) { return { route: 'GET /workspace/:id/sessions' }; @@ -138,7 +440,7 @@ export function resolveDaemonTelemetryRoute( if (workspaceTranscript?.[1] && req.method === 'GET') { return { route: 'GET /workspaces/:workspace/session/:id/transcript', - sessionId: workspaceTranscript[1], + sessionId: decodePathSegment(workspaceTranscript[1]), }; } const workspaceExport = path.match( @@ -147,7 +449,7 @@ export function resolveDaemonTelemetryRoute( if (workspaceExport?.[1] && req.method === 'GET') { return { route: 'GET /workspaces/:workspace/session/:id/export', - sessionId: workspaceExport[1], + sessionId: decodePathSegment(workspaceExport[1]), }; } const workspaceArchivedExport = path.match( @@ -156,7 +458,7 @@ export function resolveDaemonTelemetryRoute( if (workspaceArchivedExport?.[1] && req.method === 'GET') { return { route: 'GET /workspaces/:workspace/session/:id/archive/export', - sessionId: workspaceArchivedExport[1], + sessionId: decodePathSegment(workspaceArchivedExport[1]), }; } const pluralWorkspacePrefix = /^\/workspaces\/[^/]+/; @@ -331,7 +633,6 @@ export function daemonTelemetryMiddleware( // the OTel counter's scope, so the "requests" line reflects daemon API // traffic rather than static-asset or unrouted noise. recordRequest?: (durationMs: number, statusCode: number) => void, - resolveSessionWorkspaceCwd?: (sessionId: string) => string | undefined, ): (req: Request, res: Response, next: NextFunction) => void { const workspaceHashByCwd = new Map(); const resolveWorkspaceHash = (workspaceCwd: string): string => { @@ -348,18 +649,15 @@ export function daemonTelemetryMiddleware( next(); return; } - const resolveOwnerWorkspace = - route.route === 'GET /session/:id/rewind/snapshots' || - route.route === 'POST /session/:id/rewind' || - route.route === 'POST /session/:id/shell'; - const sessionId = route.sessionId - ? decodePathSegment(route.sessionId) - : undefined; - const workspaceCwd = - (resolveOwnerWorkspace && sessionId - ? resolveSessionWorkspaceCwd?.(sessionId) - : undefined) ?? resolveWorkspaceCwd(req); - const workspaceHash = resolveWorkspaceHash(workspaceCwd); + const sessionId = route.sessionId; + let workspaceHash: string | undefined; + if (route.attribution !== 'handler_resolved') { + try { + workspaceHash = resolveWorkspaceHash(resolveWorkspaceCwd(req)); + } catch { + // Telemetry must not affect request handling. + } + } const rawClientId = req.get(CLIENT_ID_HEADER); const clientId = rawClientId !== undefined && @@ -369,11 +667,19 @@ export function daemonTelemetryMiddleware( ? rawClientId : undefined; const startMs = Date.now(); + const telemetryRes = res as TelemetryResponse; + if (route.attribution === 'handler_resolved') { + try { + telemetryRes[daemonTelemetryResponseContext] = {}; + } catch { + // Telemetry must not affect request handling. + } + } void withDaemonRequestSpan( { method: req.method, route: route.route, - workspaceHash, + ...(workspaceHash !== undefined ? { workspaceHash } : {}), ...(sessionId ? { sessionId } : {}), ...(route.permissionRequestId ? { permissionRequestId: route.permissionRequestId } @@ -386,14 +692,36 @@ export function daemonTelemetryMiddleware( const finish = () => { if (done) return; done = true; + try { + const context = telemetryRes[daemonTelemetryResponseContext]; + delete telemetryRes[daemonTelemetryResponseContext]; + if (context?.workspaceCwd !== undefined) { + span?.setAttribute( + 'qwen-code.workspace.hash', + resolveWorkspaceHash(context.workspaceCwd), + ); + } + } catch { + // Telemetry must not affect response or metrics settlement. + } recordDaemonHttpResponse(span, res.statusCode); const durationMs = Date.now() - startMs; - recordDaemonHttpRequest(durationMs, route.route, res.statusCode); + const successfulSse = + route.route === 'GET /session/:id/events' && + res.statusCode === 200 && + res.headersSent; + if (!successfulSse) { + recordDaemonHttpRequest(durationMs, route.route, res.statusCode); + } // Exclude the dashboard's own status poll from the metrics-ring // request rate/latency, or the Requests chart shows a baseline of // ≥1/window with no external traffic (the dashboard counting itself) // — misleading an operator investigating load. OTel still counts it. - if (route.route !== 'GET /daemon/status') { + if ( + !successfulSse && + route.route !== 'GET /daemon/status' && + route.route !== 'POST /session/:id/heartbeat' + ) { recordRequest?.(durationMs, res.statusCode); } resolve(); From 8167bf5bdcf03da927331bfddf2d94335230c943 Mon Sep 17 00:00:00 2001 From: doudouOUC Date: Thu, 16 Jul 2026 16:10:38 +0800 Subject: [PATCH 2/4] codex: address PR review feedback (#7003) Co-authored-by: Qwen-Coder --- .../serve/routes/session-telemetry.test.ts | 106 ++++++++++++++++++ 1 file changed, 106 insertions(+) diff --git a/packages/cli/src/serve/routes/session-telemetry.test.ts b/packages/cli/src/serve/routes/session-telemetry.test.ts index 564b202fb93..e36e5f789c0 100644 --- a/packages/cli/src/serve/routes/session-telemetry.test.ts +++ b/packages/cli/src/serve/routes/session-telemetry.test.ts @@ -20,8 +20,15 @@ import { const telemetryMocks = vi.hoisted(() => ({ setDaemonTelemetryWorkspace: vi.fn(), })); +const archiveMocks = vi.hoisted(() => ({ + assertSessionLoadable: vi.fn(), +})); vi.mock('../server/telemetry.js', () => telemetryMocks); +vi.mock('../server/session-archive.js', async (importOriginal) => ({ + ...(await importOriginal()), + assertSessionLoadable: archiveMocks.assertSessionLoadable, +})); import { registerSessionRoutes } from './session.js'; @@ -33,6 +40,7 @@ function bridgeWithSessions(sessionIds: string[] = []): AcpSessionBridge { } return { sessionId }; }), + getSessionTranscriptPage: vi.fn(async () => ({ records: [] })), } as unknown as AcpSessionBridge; } @@ -77,6 +85,7 @@ describe('special session resolver telemetry publication', () => { beforeEach(() => { vi.clearAllMocks(); + archiveMocks.assertSessionLoadable.mockResolvedValue(undefined); }); it('publishes the runtime root for creation before later validation', async () => { @@ -98,6 +107,32 @@ describe('special session resolver telemetry publication', () => { ); }); + it('publishes the secondary runtime selected for creation', async () => { + const primary = runtime({ + workspaceId: 'primary', + workspaceCwd: primaryCwd, + primary: true, + bridge: bridgeWithSessions(), + }); + const secondary = runtime({ + workspaceId: 'secondary', + workspaceCwd: secondaryCwd, + primary: false, + bridge: bridgeWithSessions(), + }); + + const res = await request(makeApp([primary, secondary])) + .post('/session') + .send({ cwd: secondaryCwd, sessionScope: 'invalid' }); + + expect(res.status).toBe(400); + expect(telemetryMocks.setDaemonTelemetryWorkspace).toHaveBeenCalledTimes(1); + expect(telemetryMocks.setDaemonTelemetryWorkspace).toHaveBeenCalledWith( + expect.anything(), + secondaryCwd, + ); + }); + it('publishes the restore target before later validation', async () => { const primary = runtime({ workspaceId: 'primary', @@ -161,6 +196,77 @@ describe('special session resolver telemetry publication', () => { ); }); + it('publishes the live transcript owner in a multi-workspace daemon', async () => { + archiveMocks.assertSessionLoadable.mockResolvedValue('active'); + const primary = runtime({ + workspaceId: 'primary', + workspaceCwd: primaryCwd, + primary: true, + bridge: bridgeWithSessions(), + }); + const secondary = runtime({ + workspaceId: 'secondary', + workspaceCwd: secondaryCwd, + primary: false, + bridge: bridgeWithSessions(['secondary-session']), + }); + + const res = await request(makeApp([primary, secondary])).get( + '/session/secondary-session/transcript', + ); + + expect(res.status).toBe(200); + expect(archiveMocks.assertSessionLoadable).toHaveBeenCalledTimes(1); + expect(archiveMocks.assertSessionLoadable).toHaveBeenCalledWith( + secondaryCwd, + 'secondary-session', + ); + expect(telemetryMocks.setDaemonTelemetryWorkspace).toHaveBeenCalledTimes(1); + expect(telemetryMocks.setDaemonTelemetryWorkspace).toHaveBeenCalledWith( + expect.anything(), + secondaryCwd, + ); + }); + + it('publishes the sole active transcript runtime after storage lookup', async () => { + archiveMocks.assertSessionLoadable.mockImplementation( + async (workspaceCwd: string) => + workspaceCwd === secondaryCwd ? 'active' : undefined, + ); + const primary = runtime({ + workspaceId: 'primary', + workspaceCwd: primaryCwd, + primary: true, + bridge: bridgeWithSessions(), + }); + const secondary = runtime({ + workspaceId: 'secondary', + workspaceCwd: secondaryCwd, + primary: false, + bridge: bridgeWithSessions(), + }); + + const res = await request(makeApp([primary, secondary])).get( + '/session/stored-secondary/transcript', + ); + + expect(res.status).toBe(200); + expect(archiveMocks.assertSessionLoadable).toHaveBeenCalledTimes(2); + expect(archiveMocks.assertSessionLoadable).toHaveBeenCalledWith( + primaryCwd, + 'stored-secondary', + ); + expect(archiveMocks.assertSessionLoadable).toHaveBeenCalledWith( + secondaryCwd, + 'stored-secondary', + ); + expect(telemetryMocks.setDaemonTelemetryWorkspace).toHaveBeenCalledTimes(1); + expect(telemetryMocks.setDaemonTelemetryWorkspace).toHaveBeenCalledWith( + expect.anything(), + secondaryCwd, + ); + }); + it('does not publish a workspace for a creation workspace mismatch', async () => { const primary = runtime({ workspaceId: 'primary', From cb7e92ce46e477285a34829920bface8b64aa4ca Mon Sep 17 00:00:00 2001 From: doudouOUC Date: Thu, 16 Jul 2026 17:50:37 +0800 Subject: [PATCH 3/4] codex: address PR review feedback (#7003) Co-authored-by: Qwen-Coder --- .../serve/routes/session-telemetry.test.ts | 25 +++++++++++++++++++ .../cli/src/serve/server/telemetry.test.ts | 4 +-- 2 files changed, 27 insertions(+), 2 deletions(-) diff --git a/packages/cli/src/serve/routes/session-telemetry.test.ts b/packages/cli/src/serve/routes/session-telemetry.test.ts index e36e5f789c0..a8f94e5695f 100644 --- a/packages/cli/src/serve/routes/session-telemetry.test.ts +++ b/packages/cli/src/serve/routes/session-telemetry.test.ts @@ -267,6 +267,31 @@ describe('special session resolver telemetry publication', () => { ); }); + it('does not publish a workspace for ambiguous transcript storage matches', async () => { + archiveMocks.assertSessionLoadable.mockResolvedValue('active'); + const primary = runtime({ + workspaceId: 'primary', + workspaceCwd: primaryCwd, + primary: true, + bridge: bridgeWithSessions(), + }); + const secondary = runtime({ + workspaceId: 'secondary', + workspaceCwd: secondaryCwd, + primary: false, + bridge: bridgeWithSessions(), + }); + + const res = await request(makeApp([primary, secondary])).get( + '/session/shared-storage/transcript', + ); + + expect(res.status).toBe(500); + expect(res.body.code).toBe('ambiguous_session_owner'); + expect(archiveMocks.assertSessionLoadable).toHaveBeenCalledTimes(2); + expect(telemetryMocks.setDaemonTelemetryWorkspace).not.toHaveBeenCalled(); + }); + it('does not publish a workspace for a creation workspace mismatch', async () => { const primary = runtime({ workspaceId: 'primary', diff --git a/packages/cli/src/serve/server/telemetry.test.ts b/packages/cli/src/serve/server/telemetry.test.ts index a7666dae3f6..293d7ef6908 100644 --- a/packages/cli/src/serve/server/telemetry.test.ts +++ b/packages/cli/src/serve/server/telemetry.test.ts @@ -138,7 +138,7 @@ describe('daemonTelemetryMiddleware — recordRequest seam', () => { const res = mockRes(200); mw( - mockReq('GET', '/workspaces/ws-secondary/session/session-1/transcript'), + mockReq('GET', '/workspaces/ws-secondary/session/session%2F1/transcript'), res, vi.fn() as unknown as NextFunction, ); @@ -148,7 +148,7 @@ describe('daemonTelemetryMiddleware — recordRequest seam', () => { expect.objectContaining({ method: 'GET', route: 'GET /workspaces/:workspace/session/:id/transcript', - sessionId: 'session-1', + sessionId: 'session/1', workspaceHash: 'hash:/workspace/secondary', }), expect.any(Function), From 1eefbedbf89034dd10b40b3dabd5d80d72b5595c Mon Sep 17 00:00:00 2001 From: doudouOUC Date: Fri, 17 Jul 2026 10:03:42 +0800 Subject: [PATCH 4/4] codex: address PR review feedback (#7003) Co-authored-by: Qwen-Coder --- packages/cli/src/serve/server/telemetry.test.ts | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/packages/cli/src/serve/server/telemetry.test.ts b/packages/cli/src/serve/server/telemetry.test.ts index 0d6e2f019b9..c1386761294 100644 --- a/packages/cli/src/serve/server/telemetry.test.ts +++ b/packages/cli/src/serve/server/telemetry.test.ts @@ -406,6 +406,7 @@ describe('daemonTelemetryMiddleware — recordRequest seam', () => { expect.any(Number), 'POST /session/:id/heartbeat', 200, + undefined, ); expect(recordRequest).not.toHaveBeenCalled(); }); @@ -445,6 +446,7 @@ describe('daemonTelemetryMiddleware — recordRequest seam', () => { expect.any(Number), 'POST /session/:id/generate', 200, + undefined, ); expect(recordRequest).toHaveBeenCalledWith(expect.any(Number), 200); }); @@ -466,6 +468,7 @@ describe('daemonTelemetryMiddleware — recordRequest seam', () => { expect.any(Number), 'GET /session/:id/events', 200, + undefined, ); expect(recordRequest).toHaveBeenCalledWith(expect.any(Number), 200); }); @@ -488,6 +491,7 @@ describe('daemonTelemetryMiddleware — recordRequest seam', () => { expect.any(Number), 'GET /session/:id/events', statusCode, + undefined, ); expect(recordRequest).toHaveBeenCalledWith( expect.any(Number),