From 0381fc9e860eea5cf6808d6f904dcc637b61e553 Mon Sep 17 00:00:00 2001 From: iamtoruk Date: Sun, 26 Jul 2026 14:25:48 -0700 Subject: [PATCH] feat(core): dual-registry bridge + qwen exemplar migration (phase 8.1) --- docs/core/phase8-recipe.md | 336 ++++++++++++++++++ packages/cli/src/providers/bridge.ts | 96 +++++ packages/cli/src/providers/qwen.ts | 176 +++------ .../-Users-t-alpha-repo/chats/session-1.jsonl | 6 + packages/cli/tests/provider-registry.test.ts | 54 +++ .../cli/tests/providers/qwen-bridge.test.ts | 118 ++++++ packages/core/package.json | 4 + packages/core/src/providers/qwen/decode.ts | 165 +++++++++ packages/core/src/providers/qwen/index.ts | 30 ++ .../core/src/providers/qwen/observations.ts | 96 +++++ packages/core/src/providers/qwen/types.ts | 73 ++++ packages/core/tests/architecture-gate.test.ts | 2 + packages/core/tests/content-smuggling.test.ts | 72 ++++ .../core/tests/providers/qwen-decode.test.ts | 111 ++++++ packages/core/tsup.config.ts | 1 + 15 files changed, 1206 insertions(+), 134 deletions(-) create mode 100644 docs/core/phase8-recipe.md create mode 100644 packages/cli/src/providers/bridge.ts create mode 100644 packages/cli/tests/fixtures/qwen-parity/-Users-t-alpha-repo/chats/session-1.jsonl create mode 100644 packages/cli/tests/providers/qwen-bridge.test.ts create mode 100644 packages/core/src/providers/qwen/decode.ts create mode 100644 packages/core/src/providers/qwen/index.ts create mode 100644 packages/core/src/providers/qwen/observations.ts create mode 100644 packages/core/src/providers/qwen/types.ts create mode 100644 packages/core/tests/providers/qwen-decode.test.ts diff --git a/docs/core/phase8-recipe.md b/docs/core/phase8-recipe.md new file mode 100644 index 00000000..22e04547 --- /dev/null +++ b/docs/core/phase8-recipe.md @@ -0,0 +1,336 @@ +# Phase 8 fan-out recipe: bridge a provider's pure decode into `@codeburn/core` + +Phase 8 of the `@codeburn/core` extraction (issue #809) moves each provider's +**pure record decode** into `packages/core/src/providers//` while the CLI +keeps discovery, file/DB I/O, caching, and pricing. A thin **bridge** wraps the +core decoder so every existing consumer — the `parser.ts` scan loop, `doctor` +probeRoots, `--provider` filters, tool/model display names — sees the exact same +`Provider` interface as before. + +The seam and one exemplar (`qwen`) shipped in **phase 8.1**. This document is the +mechanical recipe for the remaining providers. **Follow it per-provider; each +provider is one small PR.** It mirrors the structure of +[`phase0-recipe.md`](./phase0-recipe.md); read that first if you have not. + +> **Prerequisite:** every provider you bridge here must already be **Phase 0** +> converted (emits `costBasis` instead of calling `calculateCost`). If it still +> prices itself, do the Phase 0 conversion first — the bridge assumes cost has +> already left the decoder. + +--- + +## The seam (already in place — do not re-add) + +`packages/cli/src/providers/bridge.ts` exposes two things: + +- **`BridgedProviderSpec`** — how you express a core-migrated provider: + `{ name, displayName, modelDisplayName, toolDisplayName, discoverSessions, + probeRoots?, readRecords, decode, toProviderCall }`. +- **`createBridgedProvider(spec)`** — wraps that spec into the `Provider` + interface. The generated `createSessionParser` calls `readRecords` (CLI I/O), + runs `decode` (core, threading the shared `seenKeys` dedup set), and yields + each `toProviderCall(rich)`. + +**You never edit `bridge.ts`, `parser.ts`, `pricing-pass.ts`, or +`session-cache.ts`.** The registry (`providers/index.ts`) already resolves +bridged and legacy providers uniformly — a bridged provider is registered exactly +like any other (`import { foo } from './foo.js'`, add to `coreProviders`), because +`createBridgedProvider` returns a plain `Provider`. + +This **generalizes** the bespoke adapters claude/codex wrote in phases 3/4 (each +reads its file, calls its core decoder, and maps the rich result) for the +**simple** case: one whole-file pass, a live cross-file dedup set, **no +incremental cache and therefore no serializable resume state**. Providers that +cache/resume (codex) or hold a stateful multi-store cache (antigravity/kiro) keep +their bespoke adapters — the bridge is deliberately not built to cover them (see +DO-NOT list). + +### What the two layers look like in core (the qwen shape) + +`packages/core/src/providers//`: + +- **`types.ts`** — raw record types (moved verbatim from the CLI) + the rich + `…DecodedCall` (mirrors `ParsedProviderCall` **minus** cost fields) + any + `…ToolCall` used for `toolSequence`. +- **`decode.ts`** — `decode({ records, context, seenKeys? })` returning + `{ calls, diagnostics }`. **Pure**: no `fs`/env/clock, no pricing, no + `strip-ansi`/bash base-name extraction (that stays CLI-side). Dedup threads the + live `seenKeys` set. The tool-name map moves here too. +- **`observations.ts`** — `toObservations(decode, ctx)` mapping the rich decode + into the strict `SessionObservation[]`. This is where the content-smuggling + guarantees bind: only fingerprints, enums, numbers, timestamps, dedup keys, and + canonical tool names cross the boundary. Fingerprint file paths through + `extractResourceRefs` from `../resource-refs.js`. +- **`index.ts`** — barrel re-exporting decode + observations + types. + +Then wire the build/exports (three edits, all add-only): + +1. `packages/core/tsup.config.ts` — add `'src/providers//index.ts'` to `entry`. +2. `packages/core/package.json` — add the `"./providers/"` exports entry + (`types` + `import` pointing at `dist/providers//index.*`). The + import-smoke guardrail auto-discovers this subpath from `package.json`. +3. `packages/core/tests/architecture-gate.test.ts` — if the rich decode carries + `userMessage`, add `src/providers//decode.ts` and `…/types.ts` to + `USER_MESSAGE_ALLOWLIST` (they are legitimate host-side carriers, exactly like + claude/codex/qwen). + +--- + +## Before you start — classify the provider + +Each provider falls into exactly one bucket. **Only Category A is the cheap +tail.** The full classification of the current tail is the table at the end. + +| Bucket | Signal | This recipe? | +|--------|--------|--------------| +| **A. Simple file-based (JSONL/JSON)** | `readSessionFile`/`readFile`, one file (or a few) per session, no sqlite, no cache resume | **Yes — the cheap tail.** | +| **B. sqlite-rowed** | imports `better-sqlite3`/`node:sqlite` or `sqlite-session-parser` | Yes, but the **sqlite driver + query stay CLI-side**; only the row→observation decode moves. Bigger PR — see "sqlite variant". | +| **C. shared-module consumers** | imports `vscode-cline-parser` / `opencode-file-parser` / `session-message` / `sqlite-session-parser` | Migrate the **shared decode module once** into core, then wire every consumer to it. One PR for the module, then trivial per-consumer. | +| **D. stateful multi-store** | a session cache with `byteOffset`/resume, or several on-disk stores stitched with in-memory `Map`s | **No — judgment-tier (Opus).** Do not force it into the simple bridge. | + +Run this to place a provider yourself: + +``` +grep -n "better-sqlite3\|node:sqlite\|sqlite-session-parser" packages/cli/src/providers/.ts # B +grep -n "byteOffset\|-cache.js\|CacheEntry\|new Map(" packages/cli/src/providers/.ts # D signals +grep -n "vscode-cline-parser\|opencode-file-parser\|session-message" packages/cli/src/providers/.ts # C +``` + +--- + +## Category A — the move checklist (simple file-based) + +1. **Create the core module** (`types.ts`, `decode.ts`, `observations.ts`, + `index.ts`) as above. Move the raw record types, the tool-name map, the + line-parse, and the per-record decode **verbatim**. Strip anything host-side: + pricing (already gone post-Phase 0), and bash base-name extraction — the + decoder emits **raw** command strings (`rawBashCommands`), the CLI reduces them. +2. **Wire build/exports** (tsup entry, package.json exports, arch-gate allowlist). +3. **Rebuild core** so the CLI resolves the new dist: `npm run build -w @codeburn/core`. +4. **Convert the CLI provider** to `createBridgedProvider`: + - Keep `discoverSessions` / `probeRoots` unchanged (discovery stays CLI-side). + - Add `readRecords(source)` — the file read + line split that used to live + inside `createSessionParser`. + - `decode: decode` (the core ref). + - `toProviderCall(rich)` — map the rich call to `ParsedProviderCall`, add + `costBasis: 'estimated'` (or `'measured'` per Phase 0 Pattern B), and run any + CLI-only reduction (`extractBashCommands`). +5. **Add tests** (below). +6. **Verify parity**, then revert `src/data`. + +### Parity is the gate + +- **Fixture golden (always).** Capture the pre-migration `parse()` output for a + representative fixture **before** you touch the provider, commit it as the + expected array, and assert the bridged provider reproduces it **byte-for-byte** + (`toEqual`). This is exactly what `packages/cli/tests/providers/qwen-bridge.test.ts` + does. Because `priceProviderCall` is deterministic, byte-identical raw output + guarantees byte-identical priced output. +- **Frozen-corpus compare (conditional — DO NOT SKIP when it applies).** qwen has + **no** data in the frozen corpus, so a fixture golden was sufficient. **Any + provider that IS present in the frozen corpus (e.g. `cursor`, `gemini`, and the + others listed in the corpus manifest) MUST run the frozen-corpus compare** as + part of its migration PR — the corpus is the byte-level parity oracle for real + data and a fixture cannot substitute for it. Check membership first: + + ``` + grep -rl "" packages/cli/tests/fixtures/**/frozen* 2>/dev/null # adjust to the corpus path + ``` + + If the provider appears, the frozen-corpus test is a required, non-negotiable + gate alongside the fixture golden. + +### Test wiring (Category A) + +Add two test files: + +- `packages/core/tests/providers/-decode.test.ts` — unit-tests the rich + decode (dedup, skips, token buckets, tool mapping, model fallback) **and** that + `toObservations` produces a schema-valid, secret-free envelope. +- extend `packages/core/tests/content-smuggling.test.ts` — one hostile-transcript + block planting every `SECRETS.*` value in the free-text fields the decode + captures (user prompt, command, file path, a tool NAME carrying a command line) + and asserting the minimized envelope surfaces none of them. Copy the qwen block. +- `packages/cli/tests/providers/-bridge.test.ts` — the fixture-golden parity + test above. + +Import-smoke needs **no** manual edit — it enumerates `package.json` exports, so +the new subpath is covered the moment you add it there. + +--- + +## The sqlite variant (Category B) + +Same shape, one rule: **the sqlite driver and the SQL query never leave the CLI.** +`readRecords` runs the query host-side and returns the **rows** (plain objects); +`decode` is pure over those rows. Do not move `better-sqlite3` into core +(it would break the import-smoke I/O guardrail instantly). Everything else — +types, row→call decode, tool map, `toObservations` — moves exactly as Category A. + +`readRecords` returns `Row[]`; the core decoder's `records: unknown[]` are those +rows. That is the only difference from the JSONL case. + +--- + +## The shared-module variant (Category C) + +`vscode-cline-parser`, `opencode-file-parser`, `sqlite-session-parser`, and +`session-message` are decode logic shared by several providers. **Migrate the +shared module into core once** (as its own `core/src/providers//` or a +shared decode under the consuming provider), then point each consumer's `decode` +at it. Do the module PR first; the per-consumer PRs after are trivial rewires. +Keep the sqlite driver (for `sqlite-session-parser`) CLI-side per Category B. + +--- + +## Per-provider acceptance checklist + +1. `grep -n "decode logic markers" packages/cli/src/providers/.ts` → the CLI + file has **no** record parsing / token extraction / tool-map left (discovery + + I/O + `toProviderCall` only). The decode logic exists in **exactly one** + package (core). +2. Core: `npm run -w @codeburn/core typecheck` clean; `npm test -w @codeburn/core` + green (new `-decode` + smuggling tests pass; import-smoke covers the new + subpath). +3. Rebuild core dist, then root `npm test` green — count unchanged except the new + bridge test(s). The provider's own CLI tests pass **unchanged**. +4. `npx tsc --noEmit` clean in both packages; `npm run build` ok. +5. **Fixture golden byte-identical.** If the provider is in the frozen corpus, + **frozen-corpus compare green too.** +6. No `PROVIDER_PARSE_VERSIONS` / `parse-versions` change (grep the diff). +7. `git checkout -- packages/cli/src/data` before committing (the build refreshes + the litellm/pricing snapshots). + +--- + +## The DO-NOT list + +- **Never touch `claude`, `codex`, `session-cache`, or `parse-versions`.** claude + and codex already have bespoke adapters (phases 3/4); the bridge generalizes + their shape for simple providers, it does **not** rework them. `session-cache.ts` + and `parse-versions` are host infrastructure — off limits. +- **sqlite driver + DB reads stay CLI-side.** Only the row→observation decode + moves. Moving `better-sqlite3`/`node:sqlite` into core breaks the import-smoke + I/O guardrail. +- **Stateful providers (antigravity/kiro-style caches) are judgment-tier, not the + cheap tail.** They stitch multiple stores with resume offsets / in-memory maps; + hand them to Opus, do not force them through `createBridgedProvider`. +- **No pricing in core.** The decoder emits token buckets + (host-side) + `costBasis`; `priceProviderCall` fills `costUSD`. No `calculateCost` in core. +- **No free text in the envelope.** `toObservations` emits only fingerprints, + enums, numbers, timestamps, dedup keys, and canonical tool names. The + architecture gate + content-smuggling tests enforce this. + +--- + +## Worked diff — the `qwen` exemplar (phase 8.1) + +**Core, new file `packages/core/src/providers/qwen/decode.ts`** (the moved decode, +pure; tool map + line parse + per-record decode; raw commands, no base-name +extraction): + +```ts +export const qwenToolNameMap: Record = { read_file: 'Read', execute_command: 'Bash', /* … */ } + +export function decodeQwen({ records, seenKeys: liveSeen }: QwenDecodeInput): QwenDecodeResult { + const seen = liveSeen ?? new Set() + const calls: QwenDecodedCall[] = [] + let pendingUserMessage = '' + for (const rawLine of records) { + const entry = typeof rawLine === 'string' ? parseQwenLine(rawLine) : (rawLine as QwenEntry | null) + if (!entry) continue + if (entry.type === 'user' && entry.message) { /* set pendingUserMessage (thought parts filtered) */; continue } + if (entry.type !== 'assistant' || !entry.usageMetadata) continue + // …zero-token skip, `qwen::` dedup against `seen`… + calls.push({ provider: 'qwen', model, inputTokens, /* buckets */, tools: [...new Set(tools)], + rawBashCommands, timestamp, speed: 'standard', deduplicationKey, userMessage: pendingUserMessage, sessionId, ...(toolSequence.length ? { toolSequence } : {}) }) + pendingUserMessage = '' + } + return { calls, diagnostics: [] } +} +``` + +**CLI, `packages/cli/src/providers/qwen.ts`** — the whole file becomes discovery + +I/O + map, wrapped by the bridge: + +```diff +-import { readSessionFile } from '../fs-utils.js' +-import { extractBashCommands } from '../bash-utils.js' +-import type { Provider, SessionSource, SessionParser, ParsedProviderCall } from './types.js' +- +-const toolNameMap: Record = { read_file: 'Read', /* … */ } +-type QwenPart = { /* … */ } +-type QwenEntry = { /* … */ } +-function extractTools(parts) { /* … */ } +-function createParser(source, seenKeys): SessionParser { /* JSON.parse loop, token extraction, yield */ } ++import { decodeQwen, qwenToolNameMap } from '@codeburn/core/providers/qwen' ++import type { QwenDecodedCall } from '@codeburn/core/providers/qwen' ++import { readSessionFile } from '../fs-utils.js' ++import { extractBashCommands } from '../bash-utils.js' ++import { createBridgedProvider } from './bridge.js' ++import type { Provider, SessionSource, ParsedProviderCall } from './types.js' ++ ++function toProviderCall(rich: QwenDecodedCall): ParsedProviderCall { ++ return { provider: 'qwen', model: rich.model, /* buckets */, ++ costBasis: 'estimated', ++ tools: rich.tools, ++ bashCommands: [...new Set(rich.rawBashCommands.flatMap(c => extractBashCommands(c)))], ++ /* timestamp, speed, deduplicationKey, userMessage, sessionId */ } ++} + + export function createQwenProvider(overrideDir?: string): Provider { + const projectsDir = overrideDir ?? getQwenProjectsDir() +- return { name: 'qwen', displayName: 'Qwen', +- modelDisplayName(m) { return m }, +- toolDisplayName(t) { return toolNameMap[t] ?? t }, +- async discoverSessions() { /* … unchanged … */ }, +- createSessionParser(source, seenKeys) { return createParser(source, seenKeys) }, +- } ++ return createBridgedProvider({ ++ name: 'qwen', displayName: 'Qwen', ++ modelDisplayName(m) { return m }, ++ toolDisplayName(t) { return qwenToolNameMap[t] ?? t }, ++ async discoverSessions() { /* … unchanged … */ }, ++ async readRecords(source) { ++ const raw = await readSessionFile(source.path) ++ if (raw === null) return null ++ return raw.split('\n').filter(l => l.trim()) ++ }, ++ decode: decodeQwen, ++ toProviderCall, ++ }) + } +``` + +Discovery is untouched. The `JSON.parse` loop + token extraction moved to core. +Bash base-name extraction (`extractBashCommands`, with its `strip-ansi` dep) stays +CLI-side and runs in `toProviderCall`. `costBasis: 'estimated'` (from Phase 0) is +priced downstream by `parser.ts`, byte-identical. + +--- + +## Tail classification (as of phase 8.1) + +35 provider identities remain after `claude`/`codex`/`qwen` (done). Counts: + +| Category | Count | Providers | +|----------|-------|-----------| +| **A. Simple file-based (cheap tail)** | **15** | codebuff, codewhale, droid, gemini, grok, kimi, kimicode, lingtai-tui, mistral-vibe¹, mux, open-design, openclaw, pi, omp, zerostack | +| **B. sqlite-rowed** | **11** | copilot, crush, cursor-agent, devin, forge, goose, hermes, quickdesk, warp, zcode, zed | +| **C. shared-module consumers** | **5** | cline, ibm-bob, roo-code (`vscode-cline-parser`); kilo-code, opencode (`sqlite-session-parser`/`opencode-file-parser`) — plus the shared modules themselves: `vscode-cline-parser`, `sqlite-session-parser`, `opencode-file-parser`, `session-message` | +| **D. stateful multi-store (judgment-tier — NOT the cheap tail)** | **3** | antigravity, cursor, kiro | +| **Network (special: no on-disk file, re-fetch each run)** | **1** | vercel-gateway | + +¹ `mistral-vibe` is decode-simple (JSONL, file-based) but a **Phase 0 pricing +misfit** (session-cost allocation, whitelisted). Its decode move is Category A; +its pricing was handled separately in Phase 0. Bridge the decode; do not touch its +cost path. + +`pi` and `omp` share `pi.ts` (one file, two registered providers) — one PR covers +both. `antigravity`/`cursor` are sqlite **and** stateful; classified D because the +stateful cache dominates. `kilo-code`/`opencode` are sqlite **and** shared-module; +classified C because the shared decode module should move first. + +**Recommended order:** Category A first (highest confidence, smallest PRs), then +the Category C shared modules (one module PR unblocks several consumers), then +Category B (sqlite), and finally hand Category D to Opus. diff --git a/packages/cli/src/providers/bridge.ts b/packages/cli/src/providers/bridge.ts new file mode 100644 index 00000000..b03b353f --- /dev/null +++ b/packages/cli/src/providers/bridge.ts @@ -0,0 +1,96 @@ +import type { DecodeContext } from '@codeburn/core' + +import type { Provider, ProbeRoot, SessionSource, SessionParser, ParsedProviderCall } from './types.js' + +// ── The dual-registry bridge ──────────────────────────────────────────────── +// +// A provider can be expressed in one of two shapes and the registry resolves +// both uniformly (issue #809, phase 8): +// +// - LEGACY: a hand-written `Provider` whose `createSessionParser` decodes the +// session in-CLI (the ~32 providers not yet migrated). +// - CORE-MIGRATED: the pure record decode lives in `@codeburn/core/providers/ +// `; the CLI keeps discovery, file I/O, and pricing and merely feeds +// records to the core decoder. Expressed as a `BridgedProviderSpec` and +// wrapped by `createBridgedProvider` into the SAME `Provider` interface, so +// every existing consumer — the parser.ts scan loop, `doctor` probeRoots, +// `--provider` filters, tool/model display names — sees no difference. +// +// This GENERALIZES the shape claude/codex established in phases 3/4 (each wrote +// its own bespoke adapter that reads the file, calls its core decoder, and maps +// the rich result to a `ParsedProviderCall`) for the SIMPLE providers: one whole- +// file pass, a live cross-file dedup set, no incremental cache and therefore no +// serializable resume state. Providers that DO cache/resume (codex) or carry a +// stateful session cache (antigravity/kiro) keep their bespoke adapters — the +// bridge is deliberately not built to cover them. + +/** + * A core-migrated provider. Discovery + I/O stay CLI-side; pure record decode + * delegates to a core decoder that returns rich, cost-free calls. + * + * `TRich` is the core decoder's per-call output type (e.g. `QwenDecodedCall`). + */ +export interface BridgedProviderSpec { + name: string + displayName: string + network?: boolean + durableSources?: boolean + modelDisplayName: (model: string) => string + toolDisplayName: (rawTool: string) => string + + // Discovery stays CLI-side (fs walk, env overrides, source shaping). + discoverSessions: () => Promise + probeRoots?: () => Promise + + // I/O adapter (CLI-side): read one discovered source into the raw records the + // core decoder consumes (e.g. the JSONL lines of a file). Return `null` to + // skip the source entirely (unreadable / oversized / empty) so a transient + // read failure yields nothing rather than an empty-but-cached result. + readRecords: (source: SessionSource) => Promise + + // Pure core decode: records -> rich cost-free calls, threading the host's + // shared cross-file dedup set (mutated in place). NO fs/env/clock, NO pricing. + decode: (input: { records: unknown[]; context: DecodeContext; seenKeys: Set }) => { calls: TRich[] } + + // Host-side map of one rich call -> the host's `ParsedProviderCall`. This is + // where cost re-enters (via `costBasis: 'estimated'` + the parser.ts pricing + // pass) and where CLI-only concerns (bash base-name extraction) are applied. + toProviderCall: (rich: TRich) => ParsedProviderCall +} + +/** + * Wrap a `BridgedProviderSpec` into the `Provider` interface every consumer + * already speaks. The generated `createSessionParser` reads the source (CLI + * I/O), runs the core decoder against the shared dedup set, and yields each + * mapped call. Dedup lives inside the core decoder (it threads `seenKeys`), so + * the wrapper simply relays what the decoder returns. + */ +export function createBridgedProvider(spec: BridgedProviderSpec): Provider { + return { + name: spec.name, + displayName: spec.displayName, + ...(spec.network ? { network: spec.network } : {}), + ...(spec.durableSources ? { durableSources: spec.durableSources } : {}), + modelDisplayName: spec.modelDisplayName, + toolDisplayName: spec.toolDisplayName, + discoverSessions: spec.discoverSessions, + ...(spec.probeRoots ? { probeRoots: spec.probeRoots } : {}), + + createSessionParser(source: SessionSource, seenKeys: Set): SessionParser { + return { + async *parse(): AsyncGenerator { + const records = await spec.readRecords(source) + if (records === null) return + // The CLI holds the rich decode only; minimization / fingerprinting + // happens on the sync path, so an empty privacy key is correct here + // (the rich decoder never consumes it), matching claude/codex. + const context: DecodeContext = { privacyKey: '', providerId: spec.name, sourceRef: source.path } + const { calls } = spec.decode({ records, context, seenKeys }) + for (const rich of calls) { + yield spec.toProviderCall(rich) + } + }, + } + }, + } +} diff --git a/packages/cli/src/providers/qwen.ts b/packages/cli/src/providers/qwen.ts index ea10d2fe..87cf3a48 100644 --- a/packages/cli/src/providers/qwen.ts +++ b/packages/cli/src/providers/qwen.ts @@ -1,52 +1,14 @@ import { readdir, stat } from 'fs/promises' -import { basename, join } from 'path' +import { join } from 'path' import { homedir } from 'os' +import { decodeQwen, qwenToolNameMap } from '@codeburn/core/providers/qwen' +import type { QwenDecodedCall } from '@codeburn/core/providers/qwen' + import { readSessionFile } from '../fs-utils.js' import { extractBashCommands } from '../bash-utils.js' -import type { Provider, SessionSource, SessionParser, ParsedProviderCall } from './types.js' - -const toolNameMap: Record = { - read_file: 'Read', - write_to_file: 'Write', - edit_file: 'Edit', - execute_command: 'Bash', - search_files: 'Grep', - list_files: 'LS', - list_directory: 'LS', - browser_action: 'WebFetch', - web_search: 'WebSearch', - ask_followup_question: 'AskUser', - attempt_completion: 'Complete', -} - -type QwenPart = { - text?: string - thought?: boolean - functionCall?: { name?: string; args?: Record } - functionResponse?: unknown -} - -type QwenEntry = { - uuid: string - sessionId: string - timestamp: string - type: string - subtype?: string - cwd?: string - model?: string - message?: { - role: string - parts: QwenPart[] - } - usageMetadata?: { - promptTokenCount: number - candidatesTokenCount: number - thoughtsTokenCount: number - totalTokenCount: number - cachedContentTokenCount: number - } -} +import { createBridgedProvider } from './bridge.js' +import type { Provider, SessionSource, ParsedProviderCall } from './types.js' function getQwenProjectsDir(): string { return process.env['QWEN_DATA_DIR'] ?? join(homedir(), '.qwen', 'projects') @@ -57,99 +19,38 @@ function projectNameFromDirName(dirName: string): string { return parts[parts.length - 1] || dirName } -function extractTools(parts: QwenPart[]): { tools: string[]; bashCommands: string[] } { - const tools: string[] = [] - const bashCommands: string[] = [] - - for (const part of parts) { - if (part.functionCall?.name) { - const mapped = toolNameMap[part.functionCall.name] ?? part.functionCall.name - tools.push(mapped) - if (mapped === 'Bash' && part.functionCall.args && typeof part.functionCall.args['command'] === 'string') { - bashCommands.push(...extractBashCommands(part.functionCall.args['command'] as string)) - } - } - } - - return { tools, bashCommands } -} - -function createParser(source: SessionSource, seenKeys: Set): SessionParser { +// Map one rich, cost-free decoder call into the host's ParsedProviderCall. Cost +// re-enters here: `costBasis: 'estimated'` marks the call so the parser.ts +// pricing pass fills `costUSD` from the token buckets (byte-identical to the +// pre-migration in-decoder `calculateCost`, retired in Phase 0). Bash base-name +// extraction (and its `strip-ansi` dependency) stays CLI-side: the core decoder +// carries the raw command strings; the host reduces them to base names here. +function toProviderCall(rich: QwenDecodedCall): ParsedProviderCall { return { - async *parse(): AsyncGenerator { - const raw = await readSessionFile(source.path) - if (raw === null) return - - const lines = raw.split('\n').filter(l => l.trim()) - let pendingUserMessage = '' - - for (const line of lines) { - let entry: QwenEntry - try { - entry = JSON.parse(line) - } catch { - continue - } - - if (entry.type === 'user' && entry.message) { - const texts = (entry.message.parts ?? []) - .filter(p => p.text && !p.thought) - .map(p => p.text!) - if (texts.length > 0) { - pendingUserMessage = texts.join(' ').slice(0, 500) - } - continue - } - - if (entry.type !== 'assistant' || !entry.usageMetadata) continue - - const usage = entry.usageMetadata - if (usage.promptTokenCount === 0 && usage.candidatesTokenCount === 0) continue - - const dedupKey = `qwen:${entry.sessionId}:${entry.uuid}` - if (seenKeys.has(dedupKey)) continue - seenKeys.add(dedupKey) - - const model = entry.model || 'qwen-auto' - const { tools, bashCommands } = extractTools(entry.message?.parts ?? []) - - const inputTokens = usage.promptTokenCount - const outputTokens = usage.candidatesTokenCount - const reasoningTokens = usage.thoughtsTokenCount ?? 0 - const cachedTokens = usage.cachedContentTokenCount ?? 0 - - yield { - provider: 'qwen', - model, - inputTokens, - outputTokens, - cacheCreationInputTokens: 0, - cacheReadInputTokens: cachedTokens, - cachedInputTokens: cachedTokens, - reasoningTokens, - webSearchRequests: 0, - // Priced host-side from these buckets: reasoning tokens are billed at - // the output rate and cachedTokens as cache-read (see pricing-pass.ts). - costBasis: 'estimated', - tools: [...new Set(tools)], - bashCommands: [...new Set(bashCommands)], - timestamp: entry.timestamp || '', - speed: 'standard', - deduplicationKey: dedupKey, - userMessage: pendingUserMessage, - sessionId: entry.sessionId, - } - - pendingUserMessage = '' - } - }, + provider: 'qwen', + model: rich.model, + inputTokens: rich.inputTokens, + outputTokens: rich.outputTokens, + cacheCreationInputTokens: rich.cacheCreationInputTokens, + cacheReadInputTokens: rich.cacheReadInputTokens, + cachedInputTokens: rich.cachedInputTokens, + reasoningTokens: rich.reasoningTokens, + webSearchRequests: rich.webSearchRequests, + costBasis: 'estimated', + tools: rich.tools, + bashCommands: [...new Set(rich.rawBashCommands.flatMap(c => extractBashCommands(c)))], + timestamp: rich.timestamp, + speed: rich.speed, + deduplicationKey: rich.deduplicationKey, + userMessage: rich.userMessage, + sessionId: rich.sessionId, } } export function createQwenProvider(overrideDir?: string): Provider { const projectsDir = overrideDir ?? getQwenProjectsDir() - return { + return createBridgedProvider({ name: 'qwen', displayName: 'Qwen', @@ -158,7 +59,7 @@ export function createQwenProvider(overrideDir?: string): Provider { }, toolDisplayName(rawTool: string): string { - return toolNameMap[rawTool] ?? rawTool + return qwenToolNameMap[rawTool] ?? rawTool }, async discoverSessions(): Promise { @@ -194,10 +95,17 @@ export function createQwenProvider(overrideDir?: string): Provider { return sources }, - createSessionParser(source: SessionSource, seenKeys: Set): SessionParser { - return createParser(source, seenKeys) + // I/O adapter: read the chat file and split into JSONL records. The core + // decoder parses each line and threads the shared dedup set. + async readRecords(source: SessionSource): Promise { + const raw = await readSessionFile(source.path) + if (raw === null) return null + return raw.split('\n').filter(l => l.trim()) }, - } + + decode: decodeQwen, + toProviderCall, + }) } export const qwen = createQwenProvider() diff --git a/packages/cli/tests/fixtures/qwen-parity/-Users-t-alpha-repo/chats/session-1.jsonl b/packages/cli/tests/fixtures/qwen-parity/-Users-t-alpha-repo/chats/session-1.jsonl new file mode 100644 index 00000000..b92fae26 --- /dev/null +++ b/packages/cli/tests/fixtures/qwen-parity/-Users-t-alpha-repo/chats/session-1.jsonl @@ -0,0 +1,6 @@ +{"uuid":"u-1","sessionId":"sess-alpha","timestamp":"2026-05-16T10:00:00.000Z","type":"user","message":{"role":"user","parts":[{"text":"please fix the parser"},{"text":"and run the tests","thought":false},{"text":"internal reasoning here","thought":true}]}} +{"uuid":"a-1","sessionId":"sess-alpha","timestamp":"2026-05-16T10:00:05.000Z","type":"assistant","model":"qwen3-coder-plus","message":{"role":"assistant","parts":[{"functionCall":{"name":"read_file","args":{"path":"src/parser.ts"}}},{"functionCall":{"name":"execute_command","args":{"command":"npm test && npx vitest run"}}},{"text":"done"}]},"usageMetadata":{"promptTokenCount":1200,"candidatesTokenCount":340,"thoughtsTokenCount":90,"totalTokenCount":1630,"cachedContentTokenCount":800}} +{"uuid":"a-1","sessionId":"sess-alpha","timestamp":"2026-05-16T10:00:06.000Z","type":"assistant","model":"qwen3-coder-plus","message":{"role":"assistant","parts":[{"functionCall":{"name":"execute_command","args":{"command":"echo dup"}}}]},"usageMetadata":{"promptTokenCount":10,"candidatesTokenCount":10,"thoughtsTokenCount":0,"totalTokenCount":20,"cachedContentTokenCount":0}} +{"uuid":"a-zero","sessionId":"sess-alpha","timestamp":"2026-05-16T10:00:07.000Z","type":"assistant","model":"qwen3-coder-plus","message":{"role":"assistant","parts":[{"text":"nothing"}]},"usageMetadata":{"promptTokenCount":0,"candidatesTokenCount":0,"thoughtsTokenCount":0,"totalTokenCount":0,"cachedContentTokenCount":0}} +{"uuid":"u-2","sessionId":"sess-alpha","timestamp":"2026-05-16T10:01:00.000Z","type":"user","message":{"role":"user","parts":[{"text":"now write the changelog"}]}} +{"uuid":"a-2","sessionId":"sess-alpha","timestamp":"2026-05-16T10:01:05.000Z","type":"assistant","message":{"role":"assistant","parts":[{"functionCall":{"name":"write_to_file","args":{"path":"CHANGELOG.md"}}},{"functionCall":{"name":"some_native_tool","args":{}}}]},"usageMetadata":{"promptTokenCount":500,"candidatesTokenCount":120,"thoughtsTokenCount":0,"totalTokenCount":620,"cachedContentTokenCount":0}} diff --git a/packages/cli/tests/provider-registry.test.ts b/packages/cli/tests/provider-registry.test.ts index 5f5f2012..515efa7e 100644 --- a/packages/cli/tests/provider-registry.test.ts +++ b/packages/cli/tests/provider-registry.test.ts @@ -129,6 +129,60 @@ describe('provider registry', () => { expect(cursor.modelDisplayName('unknown-model')).toBe('unknown-model') }) + describe('dual-registry bridge (legacy + core-migrated resolve identically)', () => { + // qwen is a CORE-MIGRATED provider (pure decode lives in + // @codeburn/core/providers/qwen, wrapped by the bridge); cline is a LEGACY + // in-CLI provider. Both must be indistinguishable through the public + // registry API — the whole point of the bridge seam. + const REGISTRY_SHAPE = ['name', 'displayName', 'modelDisplayName', 'toolDisplayName', 'discoverSessions', 'createSessionParser'] as const + + it('a bridged (qwen) and a legacy (cline) provider expose the same Provider surface', () => { + const bridged = providers.find(p => p.name === 'qwen')! + const legacy = providers.find(p => p.name === 'cline')! + for (const provider of [bridged, legacy]) { + expect(typeof provider.name).toBe('string') + expect(typeof provider.displayName).toBe('string') + expect(typeof provider.modelDisplayName).toBe('function') + expect(typeof provider.toolDisplayName).toBe('function') + expect(typeof provider.discoverSessions).toBe('function') + expect(typeof provider.createSessionParser).toBe('function') + } + // The public shape a consumer relies on is identical for both. + for (const key of REGISTRY_SHAPE) { + expect(bridged).toHaveProperty(key) + expect(legacy).toHaveProperty(key) + } + }) + + it('both resolve through getProvider and back the same Provider object', async () => { + const bridged = await getProvider('qwen') + const legacy = await getProvider('cline') + expect(bridged).toBeDefined() + expect(legacy).toBeDefined() + expect(bridged).toBe(providers.find(p => p.name === 'qwen')) + expect(legacy).toBe(providers.find(p => p.name === 'cline')) + }) + + it('the bridged provider participates in discoverAllSessions like any other', async () => { + // An injected legacy fake and the real bridged qwen both flow through the + // same discovery loop; qwen's own discovery finds nothing here (no data + // dir), which is exactly how a legacy provider with no data behaves. + const fake = fakeProvider('legacy-fake', async () => [{ path: '/x.jsonl', project: 'p', provider: 'legacy-fake' }]) + const qwen = providers.find(p => p.name === 'qwen')! + const sources = await discoverAllSessions('all', [qwen, fake]) + expect(sources.map(s => s.provider)).toEqual(['legacy-fake']) + }) + + it('the bridged provider normalizes tool display names through the registry', () => { + const qwen = providers.find(p => p.name === 'qwen')! + // Delegates to the core qwenToolNameMap, resolved identically to how a + // legacy provider resolves its own map. + expect(qwen.toolDisplayName('execute_command')).toBe('Bash') + expect(qwen.toolDisplayName('read_file')).toBe('Read') + expect(qwen.toolDisplayName('unknown_tool')).toBe('unknown_tool') + }) + }) + describe('provider-discovery isolation', () => { it('safeDiscoverSessions returns [] and warns once instead of propagating', async () => { const warn = vi.spyOn(process.stderr, 'write').mockReturnValue(true) diff --git a/packages/cli/tests/providers/qwen-bridge.test.ts b/packages/cli/tests/providers/qwen-bridge.test.ts new file mode 100644 index 00000000..e4df9b39 --- /dev/null +++ b/packages/cli/tests/providers/qwen-bridge.test.ts @@ -0,0 +1,118 @@ +import { dirname, resolve } from 'path' +import { fileURLToPath } from 'url' + +import { describe, it, expect } from 'vitest' + +import { createQwenProvider } from '../../src/providers/qwen.js' +import { priceProviderCall } from '../../src/pricing-pass.js' +import type { ParsedProviderCall, SessionSource } from '../../src/providers/types.js' + +// Byte-identical parity gate for the qwen bridge migration (phase 8.1). Qwen is +// not present in the frozen corpus (no qwen data), so a committed fixture golden +// is THE parity gate: the bridged provider (discovery + I/O CLI-side, pure decode +// delegated to @codeburn/core/providers/qwen) must reproduce exactly what the +// pre-migration in-CLI decode produced. The GOLDEN below was captured from the +// legacy provider before the migration. + +const here = dirname(fileURLToPath(import.meta.url)) +const FIXTURE_DIR = resolve(here, '../fixtures/qwen-parity') + +// The exact ParsedProviderCall stream the pre-migration qwen decode produced for +// the fixture: user-message attribution (thought parts filtered), tool mapping +// (read_file->Read, execute_command->Bash, write_to_file->Write, unknown passes +// through), bash base-name extraction, uuid dedup (the second a-1 drops), the +// zero-token skip, cached/reasoning buckets, and the 'qwen-auto' model fallback. +const GOLDEN: ParsedProviderCall[] = [ + { + provider: 'qwen', + model: 'qwen3-coder-plus', + inputTokens: 1200, + outputTokens: 340, + cacheCreationInputTokens: 0, + cacheReadInputTokens: 800, + cachedInputTokens: 800, + reasoningTokens: 90, + webSearchRequests: 0, + costBasis: 'estimated', + tools: ['Read', 'Bash'], + bashCommands: ['npm', 'vitest'], + timestamp: '2026-05-16T10:00:05.000Z', + speed: 'standard', + deduplicationKey: 'qwen:sess-alpha:a-1', + userMessage: 'please fix the parser and run the tests', + sessionId: 'sess-alpha', + }, + { + provider: 'qwen', + model: 'qwen-auto', + inputTokens: 500, + outputTokens: 120, + cacheCreationInputTokens: 0, + cacheReadInputTokens: 0, + cachedInputTokens: 0, + reasoningTokens: 0, + webSearchRequests: 0, + costBasis: 'estimated', + tools: ['Write', 'some_native_tool'], + bashCommands: [], + timestamp: '2026-05-16T10:01:05.000Z', + speed: 'standard', + deduplicationKey: 'qwen:sess-alpha:a-2', + userMessage: 'now write the changelog', + sessionId: 'sess-alpha', + }, +] + +async function collect(): Promise { + const provider = createQwenProvider(FIXTURE_DIR) + const sources: SessionSource[] = await provider.discoverSessions() + sources.sort((a, b) => a.path.localeCompare(b.path)) + const seen = new Set() + const calls: ParsedProviderCall[] = [] + for (const source of sources) { + for await (const call of provider.createSessionParser(source, seen).parse()) { + calls.push(call) + } + } + return calls +} + +describe('qwen bridge — fixture parity', () => { + it('the bridged provider reproduces the pre-migration decode byte-for-byte', async () => { + expect(await collect()).toEqual(GOLDEN) + }) + + it('the priced output survives the pricing pass with only costUSD added', async () => { + const raw = await collect() + const priced = raw.map(priceProviderCall) + // Raw parity is byte-identical (test above) and priceProviderCall is + // deterministic, so priced parity follows. The pass adds exactly costUSD + // (a finite number) and leaves every other field — including the + // 'estimated' marker — untouched. + priced.forEach((call, i) => { + expect(typeof call.costUSD).toBe('number') + expect(Number.isFinite(call.costUSD)).toBe(true) + expect(call.costBasis).toBe('estimated') + const { costUSD, ...rest } = call + expect(rest).toEqual(raw[i]) + }) + }) + + it('discovery, I/O, and dedup stay CLI-side; the shared seenKeys set dedups', async () => { + // A second scan sharing the SAME seenKeys set must yield nothing (every call + // already seen), proving dedup threads through the host-owned set. + const provider = createQwenProvider(FIXTURE_DIR) + const sources = await provider.discoverSessions() + const seen = new Set() + const first: ParsedProviderCall[] = [] + for (const source of sources) { + for await (const call of provider.createSessionParser(source, seen).parse()) first.push(call) + } + const second: ParsedProviderCall[] = [] + for (const source of sources) { + for await (const call of provider.createSessionParser(source, seen).parse()) second.push(call) + } + expect(first.length).toBe(2) + expect(second).toEqual([]) + }) +}) diff --git a/packages/core/package.json b/packages/core/package.json index 5168feaa..165bc53b 100644 --- a/packages/core/package.json +++ b/packages/core/package.json @@ -42,6 +42,10 @@ "./providers/codex": { "types": "./dist/providers/codex/index.d.ts", "import": "./dist/providers/codex/index.js" + }, + "./providers/qwen": { + "types": "./dist/providers/qwen/index.d.ts", + "import": "./dist/providers/qwen/index.js" } }, "files": [ diff --git a/packages/core/src/providers/qwen/decode.ts b/packages/core/src/providers/qwen/decode.ts new file mode 100644 index 00000000..9c0d6dc5 --- /dev/null +++ b/packages/core/src/providers/qwen/decode.ts @@ -0,0 +1,165 @@ +// @codeburn/core Qwen decoder: pure decode over supplied chat JSONL records. No +// fs / env / clock — the host reads the file and hands lines straight through. +// The rich output carries token buckets but NO pricing (cost leaves the decoder; +// the host prices via its estimated-cost seam) and NO bash base-name extraction +// (that, with its `strip-ansi` dependency, stays host-side). +// +// Qwen is a "simple JSONL" provider: one chat file is one logical session, the +// host re-reads the whole file every run (no incremental cache), so the decoder +// is a single pass with no serializable resume state. The only cross-record +// memory it needs is the pending user message (threaded within the one pass) and +// the cross-file dedup set (threaded live by the host, exactly like codex). + +import type { DecodeContext } from '../../contracts.js' +import type { RecordDiagnostic } from '../../diagnostics.js' +import type { QwenDecodedCall, QwenEntry, QwenPart, QwenToolCall } from './types.js' + +// Qwen (Gemini-lineage) tool ids mapped to the canonical vocabulary. An id with +// no mapping passes through unchanged so a provider-native tool still shows up. +export const qwenToolNameMap: Record = { + read_file: 'Read', + write_to_file: 'Write', + edit_file: 'Edit', + execute_command: 'Bash', + search_files: 'Grep', + list_files: 'LS', + list_directory: 'LS', + browser_action: 'WebFetch', + web_search: 'WebSearch', + ask_followup_question: 'AskUser', + attempt_completion: 'Complete', +} + +// Parse one chat line. Qwen chat lines are ordinary JSONL; a blank or malformed +// line is skipped (returns null) exactly as the pre-migration decode did. +export function parseQwenLine(line: string): QwenEntry | null { + if (!line.trim()) return null + try { + return JSON.parse(line) as QwenEntry + } catch { + return null + } +} + +// Pull the file path a tool call targets, honoring both Qwen arg spellings. +function fileFromArgs(args: Record | undefined): string | undefined { + const fp = args?.['file_path'] ?? args?.['path'] + return typeof fp === 'string' ? fp : undefined +} + +function extractCalls(parts: QwenPart[]): { + tools: string[] + rawBashCommands: string[] + toolSequence: QwenToolCall[][] +} { + const tools: string[] = [] + const rawBashCommands: string[] = [] + const toolSequence: QwenToolCall[][] = [] + + for (const part of parts) { + const name = part.functionCall?.name + if (!name) continue + const mapped = qwenToolNameMap[name] ?? name + tools.push(mapped) + + const args = part.functionCall?.args + const call: QwenToolCall = { tool: mapped } + const file = fileFromArgs(args) + if (file) call.file = file + if (mapped === 'Bash' && args && typeof args['command'] === 'string') { + const command = args['command'] as string + call.command = command + rawBashCommands.push(command) + } + toolSequence.push([call]) + } + + return { tools, rawBashCommands, toolSequence } +} + +export type QwenDecodeInput = { + records: unknown[] + context: DecodeContext + // Optional live dedup set the host mutates in place (its shared cross-file + // seenKeys). Threaded exactly like codex's live set. Simple JSONL providers + // never persist resume state, so there is no serialized `seenKeys` fallback. + seenKeys?: Set +} + +export type QwenDecodeResult = { + calls: QwenDecodedCall[] + diagnostics: RecordDiagnostic[] +} + +/** + * Decode a Qwen chat file's records into rich, cost-free calls. A single pass: + * user messages set the pending prompt for the next assistant call; assistant + * messages that carry token usage flush into a call. Dedup is keyed on + * `qwen::` against the live `seenKeys` set (host-owned), so a + * fork/replay that repeats a uuid collides and drops. + */ +// `context` is part of the decode contract but the rich layer never consumes it: +// minimization / fingerprinting happens in toObservations. +export function decodeQwen({ records, seenKeys: liveSeen }: QwenDecodeInput): QwenDecodeResult { + const seen = liveSeen ?? new Set() + const calls: QwenDecodedCall[] = [] + const diagnostics: RecordDiagnostic[] = [] + + let pendingUserMessage = '' + + for (const rawLine of records) { + const entry = typeof rawLine === 'string' ? parseQwenLine(rawLine) : (rawLine as QwenEntry | null) + if (!entry) continue + + if (entry.type === 'user' && entry.message) { + const texts = (entry.message.parts ?? []) + .filter(p => p.text && !p.thought) + .map(p => p.text!) + if (texts.length > 0) { + pendingUserMessage = texts.join(' ').slice(0, 500) + } + continue + } + + if (entry.type !== 'assistant' || !entry.usageMetadata) continue + + const usage = entry.usageMetadata + const promptTokenCount = usage.promptTokenCount ?? 0 + const candidatesTokenCount = usage.candidatesTokenCount ?? 0 + if (promptTokenCount === 0 && candidatesTokenCount === 0) continue + + const dedupKey = `qwen:${entry.sessionId ?? ''}:${entry.uuid ?? ''}` + if (seen.has(dedupKey)) continue + seen.add(dedupKey) + + const model = entry.model || 'qwen-auto' + const { tools, rawBashCommands, toolSequence } = extractCalls(entry.message?.parts ?? []) + + const reasoningTokens = usage.thoughtsTokenCount ?? 0 + const cachedTokens = usage.cachedContentTokenCount ?? 0 + + calls.push({ + provider: 'qwen', + model, + inputTokens: promptTokenCount, + outputTokens: candidatesTokenCount, + cacheCreationInputTokens: 0, + cacheReadInputTokens: cachedTokens, + cachedInputTokens: cachedTokens, + reasoningTokens, + webSearchRequests: 0, + tools: [...new Set(tools)], + rawBashCommands, + timestamp: entry.timestamp || '', + speed: 'standard', + deduplicationKey: dedupKey, + userMessage: pendingUserMessage, + sessionId: entry.sessionId ?? '', + ...(toolSequence.length > 0 ? { toolSequence } : {}), + }) + + pendingUserMessage = '' + } + + return { calls, diagnostics } +} diff --git a/packages/core/src/providers/qwen/index.ts b/packages/core/src/providers/qwen/index.ts new file mode 100644 index 00000000..4d87de55 --- /dev/null +++ b/packages/core/src/providers/qwen/index.ts @@ -0,0 +1,30 @@ +// @codeburn/core Qwen provider. +// +// Two layers: +// - Rich pure decode (`decodeQwen`): host-facing, NOT part of the stable +// minimized surface. Pure over supplied records; carries content in-memory +// but no pricing (cost leaves the decoder) and no bash base-name extraction +// (that stays host-side with its `strip-ansi` dependency). +// - Minimizing transform (`toObservations`): maps the rich decode into the +// strict observation envelope; the content-smuggling guarantees bind here. + +export { + decodeQwen, + parseQwenLine, + qwenToolNameMap, + type QwenDecodeInput, + type QwenDecodeResult, +} from './decode.js' + +export { + toObservations, + type RichQwenSessionDecode, + type QwenToObservationsContext, +} from './observations.js' + +export type { + QwenDecodedCall, + QwenEntry, + QwenPart, + QwenToolCall, +} from './types.js' diff --git a/packages/core/src/providers/qwen/observations.ts b/packages/core/src/providers/qwen/observations.ts new file mode 100644 index 00000000..0eb67394 --- /dev/null +++ b/packages/core/src/providers/qwen/observations.ts @@ -0,0 +1,96 @@ +// Minimizing transform: rich Qwen decode -> the strict observation envelope. +// This is where the content-smuggling guarantees bind. Only opaque ids, +// fingerprints, enums, numbers, timestamps, and CANONICAL tool names cross into +// the output — never the user message, the cwd/project path, a shell command, an +// edited/read file path, or a tool argument. `.strict()` on the schemas rejects +// any extra field; this transform simply never emits one. + +import { projectRef, sessionRef } from '../../fingerprint.js' +import type { RecordDiagnostic } from '../../diagnostics.js' +import type { CallObservation, SessionObservation } from '../../observations.js' +import { extractResourceRefs } from '../resource-refs.js' +import type { QwenDecodedCall } from './types.js' + +/** One Qwen session's rich decode, as the host holds it before minimization. */ +export interface RichQwenSessionDecode { + sessionId: string + /** Absolute project path (the session cwd); fingerprinted, never emitted raw. */ + projectPath: string + /** Rich, cost-free calls in decode order (as decodeQwen emits them). */ + calls: QwenDecodedCall[] +} + +export interface QwenToObservationsContext { + /** HMAC key that scopes every fingerprint. */ + privacyKey: string + /** Provider id stamped onto sessions/calls and folded into sessionRef. */ + provider?: string +} + +// Canonical tool-name charset, mirroring core's CanonicalToolName schema. A name +// that does not match (a provider-native id with a slash, an argument blob) is +// dropped rather than emitted. +const CANONICAL_TOOL_NAME = /^[A-Za-z0-9_.-]{1,64}$/ + +function toCallObservation(call: QwenDecodedCall, turnIndex: number, privacyKey: string): CallObservation { + return { + provider: call.provider, + model: call.model, + tokens: { + input: call.inputTokens, + output: call.outputTokens, + reasoning: call.reasoningTokens, + cacheRead: call.cacheReadInputTokens, + cacheCreate: call.cacheCreationInputTokens, + }, + webSearchRequests: call.webSearchRequests, + speed: call.speed, + // Qwen calls are priced from token buckets by the host's pricing table; they + // carry no provider-reported dollar figure. + costBasis: 'estimated', + timestamp: call.timestamp, + dedupKey: call.deduplicationKey, + toolNames: call.tools.filter(t => CANONICAL_TOOL_NAME.test(t)), + turnIndex, + ...extractResourceRefs(privacyKey, call.toolSequence), + } +} + +function toSessionObservation(decode: RichQwenSessionDecode, ctx: QwenToObservationsContext): SessionObservation { + const provider = ctx.provider ?? 'qwen' + // Each assistant call flushes a distinct turn; Qwen has no turn id to group on. + const calls: CallObservation[] = decode.calls.map((call, i) => toCallObservation(call, i, ctx.privacyKey)) + + const timestamps = calls.map(c => c.timestamp).filter(t => t.length > 0).sort() + const startedAt = timestamps[0] ?? '' + const endedAt = timestamps.length > 0 ? timestamps[timestamps.length - 1]! : '' + + const session: SessionObservation = { + sessionRef: sessionRef(ctx.privacyKey, provider, decode.sessionId), + projectRef: projectRef(ctx.privacyKey, decode.projectPath), + providerId: provider, + startedAt, + ...(endedAt ? { endedAt } : {}), + calls, + turnCount: calls.length, + } + return session +} + +/** + * Map a rich Qwen decode (one or many sessions) into the minimized observation + * layer. Returns the `sessions` array plus any per-record `diagnostics`. + * + * Content-smuggling guarantee: no free text (user message, cwd, project path, + * command, read/edited file path, tool argument) is ever copied into the result. + * Only fingerprints, enums, numbers, timestamps, dedup keys, and canonical tool + * names cross the boundary. + */ +export function toObservations( + decode: RichQwenSessionDecode | RichQwenSessionDecode[], + ctx: QwenToObservationsContext, +): { sessions: SessionObservation[]; diagnostics: RecordDiagnostic[] } { + const decodes = Array.isArray(decode) ? decode : [decode] + const sessions = decodes.map(d => toSessionObservation(d, ctx)) + return { sessions, diagnostics: [] } +} diff --git a/packages/core/src/providers/qwen/types.ts b/packages/core/src/providers/qwen/types.ts new file mode 100644 index 00000000..5f971de2 --- /dev/null +++ b/packages/core/src/providers/qwen/types.ts @@ -0,0 +1,73 @@ +// Raw record + rich-decode types for the Qwen provider. +// +// The record types (QwenPart, QwenEntry) describe the shape of a Qwen Code CLI +// chat JSONL line. The Decoded* types are the rich decode layer's output: pure +// over supplied records, carrying content in-memory but NO pricing (the host +// prices them). The CLI adapter maps QwenDecodedCall into its own +// ParsedProviderCall by adding `costBasis: 'estimated'`, extracting bash base +// commands, and running the pricing pass. + +export type QwenPart = { + text?: string + thought?: boolean + functionCall?: { name?: string; args?: Record } + functionResponse?: unknown +} + +export type QwenEntry = { + uuid?: string + sessionId?: string + timestamp?: string + type?: string + subtype?: string + cwd?: string + model?: string + message?: { + role?: string + parts?: QwenPart[] + } + usageMetadata?: { + promptTokenCount?: number + candidatesTokenCount?: number + thoughtsTokenCount?: number + totalTokenCount?: number + cachedContentTokenCount?: number + } +} + +// A single tool invocation captured in a turn's tool sequence. Mirrors the +// CLI's ToolCall so the host can consume it without a shape conversion; `file` +// is host-side only (fingerprinted before it can reach an observation). +export type QwenToolCall = { + tool: string + file?: string + command?: string +} + +// The rich decode of one Qwen call (one assistant message with usage), pre- +// pricing. Mirrors the host's ParsedProviderCall minus cost fields (the host +// adds those). `rawBashCommands` are the un-split shell command strings from +// Bash-mapped tool calls; the CLI adapter runs its own base-name extraction on +// them to build the `bashCommands` field (that extraction, and its `strip-ansi` +// dependency, stay CLI-side). `toolSequence` carries raw file paths host-side +// so the observation transform can fingerprint them; it never leaves the host +// as-is. +export type QwenDecodedCall = { + provider: 'qwen' + model: string + inputTokens: number + outputTokens: number + cacheCreationInputTokens: number + cacheReadInputTokens: number + cachedInputTokens: number + reasoningTokens: number + webSearchRequests: number + tools: string[] + rawBashCommands: string[] + timestamp: string + speed: 'standard' | 'fast' + deduplicationKey: string + userMessage: string + sessionId: string + toolSequence?: QwenToolCall[][] +} diff --git a/packages/core/tests/architecture-gate.test.ts b/packages/core/tests/architecture-gate.test.ts index 1d51f50f..da295f85 100644 --- a/packages/core/tests/architecture-gate.test.ts +++ b/packages/core/tests/architecture-gate.test.ts @@ -120,6 +120,8 @@ const USER_MESSAGE_ALLOWLIST = new Set([ 'src/providers/claude/types.ts', 'src/providers/codex/decode.ts', 'src/providers/codex/types.ts', + 'src/providers/qwen/decode.ts', + 'src/providers/qwen/types.ts', ]) describe('architecture gate: no classification or free text in @codeburn/core source', () => { diff --git a/packages/core/tests/content-smuggling.test.ts b/packages/core/tests/content-smuggling.test.ts index 89dbbadd..a86b2132 100644 --- a/packages/core/tests/content-smuggling.test.ts +++ b/packages/core/tests/content-smuggling.test.ts @@ -19,6 +19,7 @@ import { } from '../src/providers/claude/index.js' import type { JournalEntry, ToolResultMeta } from '../src/providers/claude/index.js' import { decodeCodex, toObservations as toCodexObservations } from '../src/providers/codex/index.js' +import { decodeQwen, toObservations as toQwenObservations } from '../src/providers/qwen/index.js' import type { DecodeContext } from '../src/contracts.js' const here = dirname(fileURLToPath(import.meta.url)) @@ -285,6 +286,77 @@ describe('content-smuggling guardrail: real codex decode -> toObservations is se }) }) +describe('content-smuggling guardrail: real qwen decode -> toObservations is secret-free', () => { + // A hostile Qwen chat planting every secret in the free-text fields a real + // decode captures: the user prompt, an execute_command shell line, and a + // read_file path — plus a tool NAME carrying a command line. Decoding it fully + // and minimizing MUST surface none of them. + const qwenContext: DecodeContext = { privacyKey: 'test-privacy-key', providerId: 'qwen', sourceRef: 'ref' } + + function decodeAndMinimize() { + const records = [ + JSON.stringify({ + uuid: 'u-1', sessionId: 'sess-hostile', timestamp: '2026-07-17T10:00:00.000Z', type: 'user', + message: { role: 'user', parts: [{ text: `${SECRETS.prompt} ${SECRETS.apiKey} ${SECRETS.fileContent}` }] }, + }), + JSON.stringify({ + uuid: 'a-1', sessionId: 'sess-hostile', timestamp: '2026-07-17T10:00:05.000Z', type: 'assistant', model: 'qwen3-coder-plus', + message: { + role: 'assistant', + parts: [ + { functionCall: { name: 'execute_command', args: { command: SECRETS.commandLine } } }, + { functionCall: { name: 'read_file', args: { path: SECRETS.absPath } } }, + // A hostile tool NAME carrying a command line (spaces + slashes): it + // fails the canonical charset and must be dropped, not emitted. + { functionCall: { name: SECRETS.commandLine, args: {} } }, + ], + }, + usageMetadata: { promptTokenCount: 500, candidatesTokenCount: 200, thoughtsTokenCount: 0, totalTokenCount: 700, cachedContentTokenCount: 0 }, + }), + ] + const { calls } = decodeQwen({ records, context: qwenContext }) + const { sessions } = toQwenObservations( + { sessionId: 'sess-hostile', projectPath: SECRETS.absPath, calls }, + { privacyKey: 'test-privacy-key', provider: 'qwen' }, + ) + return { + schemaVersion: OBSERVATION_SCHEMA_VERSION, + generator: { name: '@codeburn/core', version: '0.0.0-test' }, + sessions, + } + } + + it('produces a schema-valid envelope from the hostile chat', () => { + expect(ObservationEnvelope.safeParse(decodeAndMinimize()).success).toBe(true) + }) + + it('the serialized envelope contains none of the planted secrets', () => { + const serialized = JSON.stringify(decodeAndMinimize()) + for (const secret of ALL_SECRETS) { + expect(serialized).not.toContain(secret) + } + }) + + it('keeps canonical tool names (Bash/Read) and drops the argument-carrying name', () => { + const env = decodeAndMinimize() + const allToolNames = env.sessions.flatMap(s => s.calls.flatMap(c => c.toolNames)) + expect(allToolNames).toContain('Bash') + expect(allToolNames).toContain('Read') + expect(allToolNames).not.toContain(SECRETS.commandLine) + }) + + it('fingerprints the read_file path into a 16-hex resourceRead, never the raw path', () => { + const env = decodeAndMinimize() + const reads = env.sessions.flatMap(s => s.calls.flatMap(c => c.resourceReads ?? [])) + expect(reads.length).toBeGreaterThan(0) + for (const ref of reads) { + expect(ref.resourceId).toMatch(/^[0-9a-f]{16}$/) + expect(typeof ref.resourceClass).toBe('string') + } + expect(allStrings(reads)).not.toContain(SECRETS.absPath) + }) +}) + describe('content-smuggling guardrail: diagnostic detail rejects paths', () => { it('rejects an absolute path', () => { expect(DiagnosticDetail.safeParse(SECRETS.absPath).success).toBe(false) diff --git a/packages/core/tests/providers/qwen-decode.test.ts b/packages/core/tests/providers/qwen-decode.test.ts new file mode 100644 index 00000000..92ad7878 --- /dev/null +++ b/packages/core/tests/providers/qwen-decode.test.ts @@ -0,0 +1,111 @@ +import { describe, expect, it } from 'vitest' + +import { decodeQwen, toObservations } from '../../src/providers/qwen/index.js' +import { ObservationEnvelope } from '../../src/observations.js' +import { OBSERVATION_SCHEMA_VERSION } from '../../src/schema.js' +import type { DecodeContext } from '../../src/contracts.js' + +const context: DecodeContext = { privacyKey: 'k', providerId: 'qwen', sourceRef: 'ref' } + +function userMessage(uuid: string, ts: string, parts: unknown[]) { + return JSON.stringify({ uuid, sessionId: 'sess-a', timestamp: ts, type: 'user', message: { role: 'user', parts } }) +} +function assistant(uuid: string, ts: string, opts: { model?: string; parts?: unknown[]; prompt?: number; candidates?: number; thoughts?: number; cached?: number }) { + return JSON.stringify({ + uuid, + sessionId: 'sess-a', + timestamp: ts, + type: 'assistant', + ...(opts.model ? { model: opts.model } : {}), + message: { role: 'assistant', parts: opts.parts ?? [] }, + usageMetadata: { + promptTokenCount: opts.prompt ?? 0, + candidatesTokenCount: opts.candidates ?? 0, + thoughtsTokenCount: opts.thoughts ?? 0, + totalTokenCount: (opts.prompt ?? 0) + (opts.candidates ?? 0), + cachedContentTokenCount: opts.cached ?? 0, + }, + }) +} + +const RECORDS: string[] = [ + userMessage('u-1', '2026-05-16T10:00:00Z', [ + { text: 'fix the parser' }, + { text: 'and run tests', thought: false }, + { text: 'secret reasoning', thought: true }, + ]), + assistant('a-1', '2026-05-16T10:00:05Z', { + model: 'qwen3-coder-plus', + prompt: 1200, candidates: 340, thoughts: 90, cached: 800, + parts: [ + { functionCall: { name: 'read_file', args: { path: 'src/parser.ts' } } }, + { functionCall: { name: 'execute_command', args: { command: 'npm test && npx vitest run' } } }, + { text: 'done' }, + ], + }), + // Same uuid replayed with different content: must drop (dedup). + assistant('a-1', '2026-05-16T10:00:06Z', { model: 'qwen3-coder-plus', prompt: 10, candidates: 10 }), + // Zero-token assistant: must skip. + assistant('a-zero', '2026-05-16T10:00:07Z', { model: 'qwen3-coder-plus' }), + // No model on the entry: falls back to 'qwen-auto'. + assistant('a-2', '2026-05-16T10:01:05Z', { + prompt: 500, candidates: 120, + parts: [{ functionCall: { name: 'some_native_tool', args: {} } }], + }), +] + +describe('qwen rich decode (moved to @codeburn/core)', () => { + it('decodes assistant calls into cost-free rich calls, deduped and zero-skipped', () => { + const { calls } = decodeQwen({ records: RECORDS, context }) + expect(calls).toHaveLength(2) + + const [first, second] = calls + // No pricing crosses into the decode layer. + expect(first).not.toHaveProperty('costUSD') + expect(first).not.toHaveProperty('costBasis') + + expect(first!.model).toBe('qwen3-coder-plus') + expect(first!.inputTokens).toBe(1200) + expect(first!.outputTokens).toBe(340) + expect(first!.reasoningTokens).toBe(90) + expect(first!.cacheReadInputTokens).toBe(800) + expect(first!.cachedInputTokens).toBe(800) + expect(first!.tools).toEqual(['Read', 'Bash']) + // Raw command strings survive host-side; base-name extraction is the CLI's job. + expect(first!.rawBashCommands).toEqual(['npm test && npx vitest run']) + // The thought part is filtered out of the pending user message. + expect(first!.userMessage).toBe('fix the parser and run tests') + expect(first!.deduplicationKey).toBe('qwen:sess-a:a-1') + + // Unknown model -> 'qwen-auto'; unknown tool id passes through unchanged. + expect(second!.model).toBe('qwen-auto') + expect(second!.tools).toEqual(['some_native_tool']) + }) + + it('threads a live seenKeys set so a repeated uuid across passes drops', () => { + const seen = new Set() + const first = decodeQwen({ records: RECORDS.slice(0, 2), context, seenKeys: seen }).calls + expect(first).toHaveLength(1) + // Re-decoding the same records with the shared set yields nothing. + const again = decodeQwen({ records: RECORDS.slice(0, 2), context, seenKeys: seen }).calls + expect(again).toEqual([]) + }) + + it('toObservations produces a schema-valid, content-free envelope', () => { + const { calls } = decodeQwen({ records: RECORDS, context }) + const { sessions } = toObservations( + { sessionId: 'sess-a', projectPath: '/Users/t/alpha', calls }, + { privacyKey: 'test-privacy-key', provider: 'qwen' }, + ) + const envelope = { + schemaVersion: OBSERVATION_SCHEMA_VERSION, + generator: { name: '@codeburn/core', version: '0.0.0-test' }, + sessions, + } + expect(ObservationEnvelope.safeParse(envelope).success).toBe(true) + // The read_file path is fingerprinted into a resourceRead, never emitted raw. + const reads = sessions.flatMap(s => s.calls.flatMap(c => c.resourceReads ?? [])) + expect(reads.length).toBeGreaterThan(0) + for (const ref of reads) expect(ref.resourceId).toMatch(/^[0-9a-f]{16}$/) + }) +}) diff --git a/packages/core/tsup.config.ts b/packages/core/tsup.config.ts index 7a6acb5d..0faec29c 100644 --- a/packages/core/tsup.config.ts +++ b/packages/core/tsup.config.ts @@ -13,6 +13,7 @@ export default defineConfig({ 'src/detectors/index.ts', 'src/providers/claude/index.ts', 'src/providers/codex/index.ts', + 'src/providers/qwen/index.ts', ], format: ['esm'], target: 'node20',