From 4d47995333278f9a19491925bc555d30f8490956 Mon Sep 17 00:00:00 2001 From: Brian Love Date: Thu, 24 Sep 2026 11:24:29 -0700 Subject: [PATCH 1/2] test(runtime): verify installed checkpoint ownership --- fixtures/react-parity/runtime/README.md | 44 +++- .../runtime/angular-checkpoints.ts | 145 +++++++++++++ fixtures/react-parity/runtime/evidence.json | 119 ++++++++++- .../runtime/react-checkpoints.tsx | 156 ++++++++++++++ .../react-parity/runtime/runtime-entry.ts | 6 + fixtures/react-parity/runtime/scenarios.ts | 5 +- scripts/react-parity/checkpoint-execution.mjs | 190 ++++++++++++++++++ .../checkpoint-execution.spec.mjs | 88 ++++++++ scripts/react-parity/review-runtime.mjs | 7 +- scripts/react-parity/review-runtime.spec.mjs | 8 + scripts/react-parity/runtime-consumer.mjs | 28 ++- 11 files changed, 789 insertions(+), 7 deletions(-) create mode 100644 fixtures/react-parity/runtime/angular-checkpoints.ts create mode 100644 fixtures/react-parity/runtime/react-checkpoints.tsx create mode 100644 scripts/react-parity/checkpoint-execution.mjs create mode 100644 scripts/react-parity/checkpoint-execution.spec.mjs diff --git a/fixtures/react-parity/runtime/README.md b/fixtures/react-parity/runtime/README.md index 63fee1600..10e8066b9 100644 --- a/fixtures/react-parity/runtime/README.md +++ b/fixtures/react-parity/runtime/README.md @@ -1,5 +1,45 @@ # Installed native runtime consumers +## Checkpoint review + +Open `/?checkpoints` on either installed review URL for the separate fixed +`checkpoint-thread` workflow. One application-owned session lives outside +component lifetime. React observes it with `useAgent` under StrictMode; Angular +uses `observeAgent` in its component injection context. Selected checkpoint +references and command outcome/completion counts are application state. The +session's execution position is private and is not invented as a snapshot field. + +Follow the sequence displayed in the view: **Load → Select A → Select B → Select +A → Fork selected → Select B → Continue branch → Load → Select P → Fork selected +→ Drop branch → Reconnect branch → Dispose → Continue branch → Fork selected**. +Selection performs no I/O and changes no transcript or values. Fork A reads the +exact completed source and confirms A1; continuation uses A1 even while B is +selected. The subsequent Load reads exact A2 and retains the earlier B/A/P +history page. Pending P rejects after its source read without a creation POST or +optimistic state change. Drop from A2 ends with physical status running; explicit +reconnect joins that same run at its cursor and confirms A3. Commands after +disposal resolve aborted without I/O. + +`checkpoint-execution.mjs` is an independent, bounded wire oracle. The server's +global latest remains B. It checks all 15 requests in order: one history POST, +six exact checkpoint-read POSTs, three creation POSTs, four physical status GETs +and one cursor join GET. Root checkpoint maps, input, catalog and stream modes +are exact; the installed SDK serializes join modes as one JSON query parameter. +Unexpected requests fail verification. The same six browser scenario groups run +in both package verifiers and the interactive review command, alongside the main +and thread workflows; scenario totals are derived from completed assertions. +Node oracle tests also reject original-A/selected-B continuation, missing routing, +wrong root maps/catalog/modes, extra POSTs and wrong join cursors/modes. + +The development factory exports its existing fixture checkpoint vocabulary and +adds a narrow `fork(checkpoint, input, options?)` return signature. Installed type +probes pass an observed readonly history reference directly, reject malformed +inputs and reserved routing, and check `Promise`. Its emitted +declaration uses installed core and fixture data types only, with no private +source or SDK references. This fixture changes no public API or production +branch UI. The strict HTTP fixture complements the separate real-server tests; +it does not establish general backend compatibility or cross-client atomicity. + ## Durable tool claims The private runtime can use an application-supplied execution store. Only a newly @@ -299,7 +339,7 @@ are missing. `node scripts/react-parity/review-runtime.mjs --help` prints the prerequisite and review sequence. The runner packs those artifacts, installs and strictly type-checks isolated -React and Angular consumers, builds each app, and runs all twenty-one browser +React and Angular consumers, builds each app, and runs all main, thread and checkpoint browser scenarios on fresh fixture servers. Only after those checks pass does it print two new, untouched loopback URLs. Open each URL manually; no browser opens automatically. The review servers have made no SDK requests at that point. @@ -485,7 +525,7 @@ React uses a Vite production build. Angular uses the existing consumer template' installed Angular CLI application builder and real APF linking, with output in `dist/consumer/browser` and input evidence from `dist/consumer/stats.json`. -Both built apps run the same twenty-one browser scenarios in installed Playwright +Both built apps run the same main and thread browser scenarios in installed Playwright Chromium: inert mount, explicit history load, equal history refresh, empty history replacement, successful text, a real local tool handler and exact two-request result continuation, protected visible server error, held streaming diff --git a/fixtures/react-parity/runtime/angular-checkpoints.ts b/fixtures/react-parity/runtime/angular-checkpoints.ts new file mode 100644 index 000000000..4db12d02e --- /dev/null +++ b/fixtures/react-parity/runtime/angular-checkpoints.ts @@ -0,0 +1,145 @@ +import { Component, signal } from '@angular/core'; +import { bootstrapApplication } from '@angular/platform-browser'; +import { observeAgent } from '@threadplane/angular'; +import { createFixtureSession } from './runtime-entry.js'; +import { + checkpointInstructions, + display, + type FixtureCheckpoint, +} from './scenarios'; + +// Application ownership is independent of component creation and destruction. +let handlerCalls = 0; +const session = createFixtureSession('/api', 'checkpoint-thread', () => { + handlerCalls++; +}); + +@Component({ + selector: 'app-root', + template: ` +
+
+

Installed package review · Angular

+

Checkpoint review

+
+
+

Review sequence

+

{{ instructions }}

+
+
+

Application selection and session execution

+

+ The selected ID is local application state. The active execution + position is retained privately by the session after a confirmed + command; observed values below come from saved server state. +

+
+ + @for (id of ['A', 'B', 'P']; track id) { + + } + + + + + +
+
+ @for (field of fields(); track field[0]) { +
+

{{ field[1] }}

+ {{ + field[2] + }} +
+ } +
+
+
+ `, +}) +export class CheckpointApp { + readonly snapshot = observeAgent(session); + readonly instructions = checkpointInstructions; + readonly selected = signal(undefined); + readonly outcome = signal(''); + readonly finished = signal(0); + readonly owner = signal('active'); + readonly busy = signal(false); + reference(id: string) { + return this.snapshot().history?.find( + (entry) => entry.checkpoint.checkpoint_id === id + )?.checkpoint; + } + select(id: string) { + this.selected.set(this.reference(id)); + } + async command(action: () => Promise) { + this.busy.set(true); + this.outcome.set('running'); + try { + this.outcome.set((await action()) ?? 'loaded'); + } catch { + this.outcome.set('rejected'); + } finally { + this.finished.update((count) => count + 1); + this.busy.set(false); + } + } + load() { + return this.command(() => session.load!()); + } + fork() { + return this.command(() => session.fork(this.selected()!, 'Fork A')); + } + run(input: string) { + return this.command(() => session.submit(input)); + } + reconnect() { + return this.command(() => session.reconnect()); + } + async dispose() { + await session.dispose(); + this.owner.set('disposed'); + } + fields() { + const snapshot = this.snapshot(); + const view = display(snapshot); + return [ + [ + 'selected', + 'Selected checkpoint reference', + this.selected()?.checkpoint_id ?? 'none', + ], + ['owner', 'Application owner', this.owner()], + ['status', 'Session status', snapshot.status], + ['outcome', 'Last command outcome', this.outcome()], + ['finished', 'Completed commands', String(this.finished())], + ['text', 'Observed transcript', view.transcript], + ['values', 'Observed values', view.values], + ['history', 'Last loaded history page', view.history], + ['reconnect', 'Reconnect run', snapshot.reconnect?.runId ?? ''], + ['handlers', 'Tool handler calls', String(handlerCalls)], + ]; + } +} + +void bootstrapApplication(CheckpointApp); diff --git a/fixtures/react-parity/runtime/evidence.json b/fixtures/react-parity/runtime/evidence.json index efebd7da7..38cc1d980 100644 --- a/fixtures/react-parity/runtime/evidence.json +++ b/fixtures/react-parity/runtime/evidence.json @@ -360,5 +360,122 @@ "generatorsRun": [], "reason": "Only contributor fixture guidance, tests and private review infrastructure changed; no public docs/API/context generator inputs changed." }, - "verificationProvenance": "All listed commands ran in this increment. Runtime and native package production implementation/public exports remain unchanged. Historical foundation evidence remains unchanged. Parent audited spec, source, lifecycle, types, installed behavior and evidence. New independent subagent review was unavailable after earlier thread capacity exhaustion; hosted review must be inspected separately from job success." + "verificationProvenance": "All listed commands ran in this increment. Runtime and native package production implementation/public exports remain unchanged. Historical foundation evidence remains unchanged. Parent audited spec, source, lifecycle, types, installed behavior and evidence. New independent subagent review was unavailable after earlier thread capacity exhaustion; hosted review must be inspected separately from job success.", + "installedCheckpointReview": { + "recordScope": "This additive record covers only the September 24 installed checkpoint fixture milestone. The other top-level fields remain the historical September 22 thread-lifetime record, including its source fingerprint and manual browser claims.", + "observedOn": "2026-09-24", + "status": "local automated and manual verification passed; independent compliance and quality reviews approved; CI reported separately", + "verificationHead": "d3d8e018b31ca6ab71f238df38f8429e4c153765", + "integratedMain": "f4ccd582ccc0d36bcb3a6a0d51c0dfa8cb1ff33f", + "integrationEvidence": "PR #1156 required CI 36037710445 passed on verificationHead. Its main merge tree 51d158f52748566e8630bda74881ae9c66ecec0c is identical. Fixture changes were preserved byte-for-byte when moving onto that main commit.", + "workingTree": "Fixture-only changes on the reviewed production predecessor plus its recording-transport return-type follow-up; no production runtime, core, binding or dependency change in this milestone.", + "executed": [ + { + "log": "/tmp/installed-checkpoint-build.log", + "logSha256": "d301e3126fa75071ede57c077da114657206c2b396d986ac6c4f98bd85f322c4", + "command": "NX_DAEMON=false NX_TUI=false npx nx run-many -t build -p core,content,react,angular --skip-nx-cache", + "exitCode": 0 + }, + { + "log": "/tmp/installed-checkpoint-node-final.log", + "logSha256": "f8f5893229b03f595345ab7ff29a7723fbf9b29710f56a1f8a4e5c89a9cb89e3", + "command": "node --test scripts/react-parity/*.spec.mjs", + "exitCode": 0, + "testsPassed": 329 + }, + { + "log": "/tmp/installed-checkpoint-react-restored.log", + "logSha256": "9c30fcfda00d425a46a19dbb16f35c8f2e0afb2cb4da434d2c936be12079149e", + "command": "node scripts/react-parity/verify-packages.mjs", + "exitCode": 0, + "browserScenarios": 27, + "checkpointScenarioGroups": 6, + "installedTypeProbes": true, + "productionBuild": true + }, + { + "log": "/tmp/installed-checkpoint-angular.log", + "logSha256": "d444c583c5324e61e8dd2f1f01698baa456867fa10cde3a3eff212bb710b3a7d", + "command": "node scripts/react-parity/verify-angular-package.mjs", + "exitCode": 0, + "browserScenarios": 27, + "checkpointScenarioGroups": 6, + "installedTypeProbes": true, + "productionBuild": true + } + ], + "testFirstEvidence": [ + { + "log": "/tmp/installed-checkpoint-types-red.log", + "logSha256": "9f52b6058e55809923f7f4e72a174cb58b45ee6fa89df963e9017e41f88653af", + "expectedFailure": "Readonly observed history checkpoint cannot call missing factory fork method." + }, + { + "log": "/tmp/installed-checkpoint-oracle-red.log", + "logSha256": "192e8c826cb661b95f2c16b328f82d991c333acf665fd3146683076afe15f246", + "expectedFailure": "Checkpoint-thread HTTP routes are not implemented." + }, + { + "log": "/tmp/installed-checkpoint-browser-red.log", + "logSha256": "26977eeae155251e4f3bdfc771d4dc3af84931e1a4a09c7db9243d9e29f41fc4", + "expectedFailure": "Installed checkpoint view does not exist." + }, + { + "log": "/tmp/installed-checkpoint-review-evidence-red.log", + "logSha256": "a6339c88470a344fad5bb48e8341bdb83578337bb8a0a19da2ed189716b0603c", + "expectedFailure": "Review provenance has no input hashes." + }, + { + "log": "/tmp/installed-checkpoint-review-evidence-green.log", + "logSha256": "8f1891e83b6bdeba32a4cb15539f46a116e90460dd7ab1f191db99f0f20f0d93", + "exitCode": 0 + } + ], + "semanticNegativeControl": { + "log": "/tmp/installed-checkpoint-negative-control.log", + "logSha256": "158ea6ee10872775772a06e74c5d2fb67b33748a175bcc326598acb45ddb56ef", + "mutation": "Temporarily change actual create-session stream routing from A1 to original A (ID and root map) for continuation.", + "expectedFailure": "Continue branch: wire oracle errors; exact branch creation routing, input, catalog and modes", + "restoration": "Original production file bytes restored in finally, empty production diff checked; restored React installed verifier passed." + }, + "checkpointProtocol": { + "thread": "checkpoint-thread", + "globalLatest": "B", + "historyPosts": 1, + "exactCheckpointReadPosts": 6, + "creationPosts": 3, + "physicalStatusGets": 4, + "cursorJoinGets": 1, + "fullSequenceAsserted": true, + "selectionImplicitIO": false, + "postDisposalIO": false, + "handlerCalls": 0, + "joinModesEncoding": "One JSON array query parameter, matching the installed SDK BaseClient." + }, + "declarations": "Readonly native-observer history reference passes directly to fork; outcome is Promise; malformed input and reserved config routing fail type probes; emitted declaration rejects SDK/private-source references.", + "cleanup": "Both package verifiers close page contexts, browsers, HTTP connections and temporary installed consumers in finally. No persistent manual servers were started by implementation.", + "manualBrowserReview": { + "performedBy": "Parent using Chrome DevTools MCP for React and Codex in-app browser for Angular, on fresh loopback review servers.", + "observedOn": "2026-09-24", + "reactUrl": "http://127.0.0.1:55315/?checkpoints", + "angularUrl": "http://127.0.0.1:55316/?checkpoints", + "sequence": "Load; Select A/B/A; Fork selected; Select B; Continue branch; Load; Select P; Fork selected; Drop branch; Reconnect branch; Dispose; Continue branch; Fork selected.", + "observations": "Both views retained branch A1/A2 despite selected B, rejected pending P without replacing A2, recovered run-A3 into A3, and returned aborted for both commands after disposal. Final selected P, owner disposed, session idle, completed commands 9, handlers 0; history remained B/A/P.", + "consoleWarnings": 0, + "consoleErrors": 0, + "reactNetwork": "Chrome recorded exactly 15 HTTP requests: one history POST, six checkpoint reads, three run creations, four status GETs and one run-A3 cursor join GET; all returned 200.", + "limits": "Angular network sequence is asserted by the installed automated oracle; its manual review claims visible state and console observations only. No live deployment, SSR or performance-budget claim.", + "runnerLog": "/tmp/installed-checkpoint-manual-runner.log", + "cleanup": "Owned review tabs closed and loopback ports 55315/55316 confirmed closed; user-owned manual sessions preserved." + }, + "independentReview": "Fresh compliance and quality reviewers independently approved the fixture diff; each ran 28 focused Node tests. Parent additionally ran all 329 parity-script tests, 829 runtime tests plus runtime/public type targets, inventory, source boundaries and version consistency. A whole-workspace emitted-boundary scan was not completed in this fixture worktree because the unchanged ag-ui built artifact was absent; installed declaration checks passed for both exercised consumers.", + "documentation": { + "generatorsRun": [], + "reason": "Only contributor fixtures and review infrastructure changed; no public API/docs/context generator inputs." + }, + "limits": [ + "Strict deterministic HTTP evidence complements separate real-server tests; it does not certify general LangGraph compatibility or cross-client atomicity.", + "No new public API, branch UI library, dependency update, package-root migration or release." + ] + } } diff --git a/fixtures/react-parity/runtime/react-checkpoints.tsx b/fixtures/react-parity/runtime/react-checkpoints.tsx new file mode 100644 index 000000000..fc10ede7e --- /dev/null +++ b/fixtures/react-parity/runtime/react-checkpoints.tsx @@ -0,0 +1,156 @@ +import { StrictMode, useState } from 'react'; +import { createRoot } from 'react-dom/client'; +import { useAgent } from '@threadplane/react'; +import { createFixtureSession } from './runtime-entry.js'; +import { + checkpointInstructions, + display, + type FixtureCheckpoint, +} from './scenarios'; +import './review.css'; + +// Application ownership is independent of the observer and StrictMode lifecycle. +let handlerCalls = 0; +const session = createFixtureSession('/api', 'checkpoint-thread', () => { + handlerCalls++; +}); + +export function CheckpointApp() { + const snapshot = useAgent(session); + const [selected, setSelected] = useState(); + const [outcome, setOutcome] = useState(''); + const [finished, setFinished] = useState(0); + const [owner, setOwner] = useState('active'); + const [busy, setBusy] = useState(false); + const view = display(snapshot); + const command = async (action: () => Promise) => { + setBusy(true); + setOutcome('running'); + try { + setOutcome((await action()) ?? 'loaded'); + } catch { + setOutcome('rejected'); + } finally { + setFinished((count) => count + 1); + setBusy(false); + } + }; + const dispose = async () => { + await session.dispose(); + setOwner('disposed'); + }; + const fields = [ + [ + 'selected', + 'Selected checkpoint reference', + selected?.checkpoint_id ?? 'none', + ], + ['owner', 'Application owner', owner], + ['status', 'Session status', snapshot.status], + ['outcome', 'Last command outcome', outcome], + ['finished', 'Completed commands', String(finished)], + ['text', 'Observed transcript', view.transcript], + ['values', 'Observed values', view.values], + ['history', 'Last loaded history page', view.history], + ['reconnect', 'Reconnect run', snapshot.reconnect?.runId ?? ''], + ['handlers', 'Tool handler calls', String(handlerCalls)], + ]; + return ( +
+
+

Installed package review · React

+

Checkpoint review

+
+
+

Review sequence

+

{checkpointInstructions}

+
+
+

Application selection and session execution

+

+ The selected ID is local application state. The active execution + position is retained privately by the session after a confirmed + command; observed values below come from saved server state. +

+
+ + {['A', 'B', 'P'].map((id) => ( + + ))} + + + + + +
+
+ {fields.map(([id, label, value]) => ( +
+

{label}

+ {value} +
+ ))} +
+
+
+ ); +} + +const container = document.getElementById('root'); +if (!container) throw new Error('Missing fixture root'); +createRoot(container).render( + + + +); diff --git a/fixtures/react-parity/runtime/runtime-entry.ts b/fixtures/react-parity/runtime/runtime-entry.ts index 3cf00f09c..a75437e79 100644 --- a/fixtures/react-parity/runtime/runtime-entry.ts +++ b/fixtures/react-parity/runtime/runtime-entry.ts @@ -6,6 +6,7 @@ import type { // eslint-disable-next-line @nx/enforce-module-boundaries -- This development-only entry composes private source into a temporary fixture bundle, never a package export. import { createSession } from '../../../libs/langgraph/src/runtime/create-session'; import type { + FixtureCheckpoint, FixtureRunOptions, FixtureSnapshot, FixtureSubmitInput, @@ -23,6 +24,11 @@ export function createFixtureSession( input: FixtureSubmitInput, options?: FixtureRunOptions ): Promise; + fork( + checkpoint: FixtureCheckpoint, + input: FixtureSubmitInput, + options?: FixtureRunOptions + ): Promise; load?: (options?: { signal?: AbortSignal }) => Promise; resume( value?: PlainValue, diff --git a/fixtures/react-parity/runtime/scenarios.ts b/fixtures/react-parity/runtime/scenarios.ts index 9a1dbe39a..de0637415 100644 --- a/fixtures/react-parity/runtime/scenarios.ts +++ b/fixtures/react-parity/runtime/scenarios.ts @@ -10,6 +10,9 @@ import type { export const reviewInstructions = 'Click Load three times: saved history, equal refresh, then empty history. Continue with Send → Tool → Error → Hold → Stop → Pause → Stop → Resume → Resume → Drop → Reconnect → Send. Tool and Drop send model, reasoning effort, UI mode and itinerary state once; displayed values come from the server. Resume first sends both approval responses, then confirms the final action. Drop loses observation of a running run; Reconnect joins that same run without another submission. Finish with Unmount → Dispose → Send after dispose → Resume after dispose → Reconnect after dispose in the owner controls below. Only three Load requests and one Drop are available per server; restart the review command to reset. Reloading the page does not reset server state.'; +export const checkpointInstructions = + 'Load → Select A → Select B → Select A → Fork selected → Select B → Continue branch → Load → Select P → Fork selected (rejected) → Drop branch → Reconnect branch → Dispose → Continue branch → Fork selected. Selection only chooses a saved reference for Fork selected. Continue and Load follow the session’s confirmed branch position even while B is selected; the global latest remains B. One bounded sequence is available per server; restart the review command to reset.'; + /** Fixture-local input contract uses only the installed neutral data vocabulary. */ export type FixtureInputState = Readonly> & { readonly messages?: never; @@ -96,7 +99,7 @@ type FixtureInterrupt = { readonly ns?: readonly string[]; }; -type FixtureCheckpoint = { +export type FixtureCheckpoint = { readonly thread_id: string; readonly checkpoint_ns: string; readonly checkpoint_id: string | null | undefined; diff --git a/scripts/react-parity/checkpoint-execution.mjs b/scripts/react-parity/checkpoint-execution.mjs new file mode 100644 index 000000000..c5c3c3ae9 --- /dev/null +++ b/scripts/react-parity/checkpoint-execution.mjs @@ -0,0 +1,190 @@ +import assert from 'node:assert/strict'; +import { expect } from '@playwright/test'; + +const thread = 'checkpoint-thread'; +const root = `/api/threads/${thread}/`; +const modes = ['values', 'messages-tuple', 'updates', 'custom', 'checkpoints']; +const checkpoint = (id) => ({ thread_id: thread, checkpoint_ns: '', checkpoint_id: id, checkpoint_map: { '': id } }); +const human = (id, content) => ({ id, type: 'human', content }); +const assistant = (id, content) => ({ id, type: 'ai', content }); +const sse = (event, data, id) => `${id ? `id: ${id}\n` : ''}event: ${event}\ndata: ${JSON.stringify(data)}\n\n`; +const saved = (id, messages, parent = null, pending = false) => ({ + checkpoint: checkpoint(id), parent_checkpoint: parent ? checkpoint(parent) : null, + created_at: '2026-09-24T00:00:00Z', metadata: { run_id: `run-${id}` }, + values: { stage: id, messages }, next: pending ? ['approval'] : [], + tasks: pending ? [{ id: 'pending-approval', name: 'approval', error: null, result: null, interrupts: [{ id: 'approval', value: 'Approve P?' }] }] : [], +}); + +// Expectations are fixture-owned, never derived from a requested checkpoint, +// production helpers, selected UI state or the preceding request's routing. +const sequence = [ + ['POST', 'history', 'B/A/P'], + ['POST', 'state/checkpoint', 'A'], + ['POST', 'runs/stream', 'A', 'Fork A', 'A1'], + ['GET', 'runs/run-A1', 'success'], + ['POST', 'state/checkpoint', 'A1'], + ['POST', 'runs/stream', 'A1', 'Continue branch', 'A2'], + ['GET', 'runs/run-A2', 'success'], + ['POST', 'state/checkpoint', 'A2'], + ['POST', 'state/checkpoint', 'A2'], + ['POST', 'state/checkpoint', 'P'], + ['POST', 'runs/stream', 'A2', 'Drop branch', 'A3'], + ['GET', 'runs/run-A3', 'running'], + ['GET', 'runs/run-A3/stream', 'drop-cursor'], + ['GET', 'runs/run-A3', 'success'], + ['POST', 'state/checkpoint', 'A3'], +]; + +/** Independent bounded HTTP oracle; global latest B never follows the branch. */ +export function createCheckpointRoutes() { + const requests = []; + const a = saved('A', [human('a-user', 'Source A'), assistant('a-answer', 'Answer A')]); + const b = saved('B', [assistant('b-answer', 'Global B')], 'A'); + const p = saved('P', [assistant('p-answer', 'Pending P')], 'A', true); + const states = { A: a, B: b, P: p }; + const responses = new Set(); + const frame = (state) => sse('checkpoints', { + config: { configurable: { ...state.checkpoint, run_id: state.metadata.run_id } }, + values: state.values, next: state.next, tasks: state.tasks.map(({ id, name }) => ({ id, name })), + }, `cursor-${state.checkpoint.checkpoint_id}`); + return { + requests, + assertComplete() { assert.deepEqual(requests.map(({ method, path, target }) => [method, path, target]), sequence.map(step => step.slice(0, 3)), 'entire checkpoint request sequence'); }, + async handle(request, response, pathname) { + if (!pathname.startsWith(root)) return false; + const expected = sequence[requests.length]; + assert.ok(expected, 'no extra checkpoint operations'); + const [method, path, target, prompt, output] = expected; + const url = new URL(request.url, 'http://fixture'); + assert.equal(request.method, method, `checkpoint step ${requests.length}: method`); + assert.equal(pathname, root + path, `checkpoint step ${requests.length}: path`); + const chunks = []; + for await (const chunk of request) chunks.push(chunk); + const raw = Buffer.concat(chunks).toString(); + const body = raw ? JSON.parse(raw) : undefined; + let result; + let stream = false; + if (path.endsWith('/stream') && method === 'GET') { + assert.equal(body, undefined, 'join has no body'); + // The installed SDK BaseClient serializes array query values as JSON. + assert.deepEqual([...url.searchParams], [['cancel_on_disconnect', '0'], ['stream_mode', JSON.stringify(modes)]], 'exact checkpoint join modes'); + assert.equal(request.headers['last-event-id'], target, 'exact checkpoint join cursor'); + result = frame(states.A3); + stream = true; + } else { + assert.equal(url.search, '', 'no unexpected checkpoint query'); + if (path === 'history') { + assert.deepEqual(body, { limit: 10 }, 'bounded history request'); + result = [b, a, p]; + } else if (path === 'state/checkpoint') { + assert.deepEqual(body, { checkpoint: checkpoint(target) }, 'exact saved checkpoint read'); + assert.ok(states[target], 'saved checkpoint exists'); + result = states[target]; + } else if (path === 'runs/stream') { + const message = body?.input?.messages?.[0]; + assert.equal(typeof message?.id, 'string'); + assert.ok(message.id.length > 0); + assert.deepEqual(body, { + assistant_id: 'fixture-assistant', checkpoint: checkpoint(target), + input: { messages: [human(message.id, prompt)], client_tools: [{ name: 'weather', description: 'Current weather' }, { name: 'count', description: 'Count values' }] }, + stream_mode: modes, stream_subgraphs: true, stream_resumable: true, on_disconnect: 'continue', + }, 'exact branch creation routing, input, catalog and modes'); + const messages = [...states[target].values.messages, human(message.id, prompt), assistant(`answer-${output}`, output === 'A3' ? 'Branch partial recovered A3' : `Branch ${output}`)]; + states[output] = saved(output, messages, target); + result = output === 'A3' + ? sse('values', { stage: 'A2-running', messages: [...messages.slice(0, -1), assistant('answer-A3', 'Branch partial')] }, 'drop-cursor') + : frame(states[output]); + stream = true; + } else { + assert.equal(body, undefined, 'status has no body'); + result = { thread_id: thread, run_id: path.slice('runs/'.length), status: target }; + } + } + requests.push({ method, path, target, ...(body === undefined ? {} : { body }), ...(path.endsWith('/stream') && method === 'GET' ? { lastEventId: request.headers['last-event-id'], modes: url.searchParams.getAll('stream_mode') } : {}) }); + responses.add(response); + response.once('close', () => responses.delete(response)); + response.writeHead(200, stream ? { 'content-type': 'text/event-stream', 'cache-control': 'no-cache', ...(output ? { 'content-location': `/threads/${thread}/runs/run-${output}` } : {}) } : { 'content-type': 'application/json' }); + response.end(stream ? result : JSON.stringify(result)); + return true; + }, + close() { for (const response of responses) response.destroy(); }, + }; +} + +/** Identical assertions run against both installed native observers. */ +export async function runCheckpointScenarios(page, server) { + await page.goto(`${server.url}/?checkpoints`); + const field = (id) => page.getByTestId(`checkpoint-${id}`); + const click = (name) => page.getByRole('button', { name, exact: true }).click(); + const command = async (label, count, outcome) => { + await click(label); + await expect(field('finished')).toHaveText(String(count)); + assert.deepEqual(server.errors.map(String), [], `${label}: wire oracle errors`); + await expect(field('outcome')).toHaveText(outcome); + }; + const observed = async () => Promise.all(['text', 'values', 'history', 'status'].map(id => field(id).textContent())); + await expect(field('status')).toHaveText('idle'); + await expect(field('text')).toHaveText(''); + await expect(field('values')).toHaveText('unobserved'); + await expect(field('history')).toHaveText('unobserved'); + await expect(field('selected')).toHaveText('none'); + await expect(field('finished')).toHaveText('0'); + assert.deepEqual(server.checkpoints.requests, [], 'checkpoint mount is inert'); + await command('Load', 1, 'loaded'); + await expect(field('text')).toHaveText('Global B'); + const history = await field('history').textContent(); + assert.deepEqual(JSON.parse(history).map(entry => entry.checkpoint.checkpoint_id), ['B', 'A', 'P']); + const loaded = await observed(); + for (const id of ['A', 'B', 'A']) { + await click(`Select ${id}`); + await expect(field('selected')).toHaveText(id); + assert.deepEqual(await observed(), loaded, 'selection leaves observed state untouched'); + } + assert.equal(server.checkpoints.requests.length, 1, 'selection performs no I/O'); + await expect(field('handlers')).toHaveText('0'); + + await command('Fork selected', 2, 'success'); + await expect(field('text')).toHaveText('Source A\nAnswer A\nFork A\nBranch A1'); + await expect(field('values')).toHaveText('{"stage":"A1"}'); + await expect(field('status')).toHaveText('idle'); + assert.equal(server.checkpoints.requests.length, 5); + + await click('Select B'); + await command('Continue branch', 3, 'success'); + await expect(field('selected')).toHaveText('B'); + await expect(field('text')).toHaveText('Source A\nAnswer A\nFork A\nBranch A1\nContinue branch\nBranch A2'); + await expect(field('values')).toHaveText('{"stage":"A2"}'); + await command('Load', 4, 'loaded'); + await expect(field('history')).toHaveText(history); + await expect(field('values')).toHaveText('{"stage":"A2"}'); + assert.equal(server.checkpoints.requests.length, 9); + + await click('Select P'); + const beforeRejected = await observed(); + await command('Fork selected', 5, 'rejected'); + assert.deepEqual(await observed(), beforeRejected, 'pending fork publishes no optimistic state'); + assert.equal(server.checkpoints.requests.length, 10); + + await command('Drop branch', 6, 'interrupted'); + await expect(field('status')).toHaveText('error'); + await expect(field('reconnect')).toHaveText('run-A3'); + await expect(field('text')).toContainText('Branch partial'); + assert.equal(server.checkpoints.requests.length, 12, 'no automatic reconnect'); + await command('Reconnect branch', 7, 'success'); + await expect(field('status')).toHaveText('idle'); + await expect(field('reconnect')).toHaveText(''); + await expect(field('values')).toHaveText('{"stage":"A3"}'); + await expect(field('text')).toHaveText('Source A\nAnswer A\nFork A\nBranch A1\nContinue branch\nBranch A2\nDrop branch\nBranch partial recovered A3'); + assert.equal((await field('text').innerText()).split('Branch partial').length - 1, 1); + + await click('Dispose'); + await expect(field('owner')).toHaveText('disposed'); + const disposed = await observed(); + await command('Continue branch', 8, 'aborted'); + await command('Fork selected', 9, 'aborted'); + assert.deepEqual(await observed(), disposed); + await expect(field('handlers')).toHaveText('0'); + server.checkpoints.assertComplete(); + assert.deepEqual(server.errors, []); + return ['checkpoint inert mount, explicit history and local selection', 'completed A fork and exact A1 confirmation', 'branch continuation and exact A2 load despite selected B', 'pending P rejection without optimistic publication', 'branch clean EOF and explicit cursor reconnect to A3', 'disposed checkpoint commands abort without I/O']; +} diff --git a/scripts/react-parity/checkpoint-execution.spec.mjs b/scripts/react-parity/checkpoint-execution.spec.mjs new file mode 100644 index 000000000..671ef4513 --- /dev/null +++ b/scripts/react-parity/checkpoint-execution.spec.mjs @@ -0,0 +1,88 @@ +import assert from 'node:assert/strict'; +import test from 'node:test'; +import { serveRuntimeConsumer } from './runtime-consumer.mjs'; + +const ref = (id) => ({ thread_id: 'checkpoint-thread', checkpoint_ns: '', checkpoint_id: id, checkpoint_map: { '': id } }); +const modes = ['values', 'messages-tuple', 'updates', 'custom', 'checkpoints']; +const body = (id, content) => ({ + assistant_id: 'fixture-assistant', checkpoint: ref(id), + input: { messages: [{ id: `user-${content}`, type: 'human', content }], client_tools: [{ name: 'weather', description: 'Current weather' }, { name: 'count', description: 'Count values' }] }, + stream_mode: modes, stream_subgraphs: true, stream_resumable: true, on_disconnect: 'continue', +}); +const send = (server, path, method = 'GET', value, headers) => fetch(`${server.url}/api/threads/checkpoint-thread/${path}`, { method, headers, ...(value === undefined ? {} : { body: JSON.stringify(value) }) }); +const steps = [ + ['history', 'POST', { limit: 10 }], + ['state/checkpoint', 'POST', { checkpoint: ref('A') }], + ['runs/stream', 'POST', body('A', 'Fork A')], + ['runs/run-A1'], + ['state/checkpoint', 'POST', { checkpoint: ref('A1') }], + ['runs/stream', 'POST', body('A1', 'Continue branch')], + ['runs/run-A2'], + ['state/checkpoint', 'POST', { checkpoint: ref('A2') }], + ['state/checkpoint', 'POST', { checkpoint: ref('A2') }], + ['state/checkpoint', 'POST', { checkpoint: ref('P') }], + ['runs/stream', 'POST', body('A2', 'Drop branch')], + ['runs/run-A3'], + [`runs/run-A3/stream?cancel_on_disconnect=0&stream_mode=${encodeURIComponent(JSON.stringify(modes))}`, 'GET', undefined, { 'Last-Event-ID': 'drop-cursor' }], + ['runs/run-A3'], + ['state/checkpoint', 'POST', { checkpoint: ref('A3') }], +]; +async function advance(server, count) { + const responses = []; + for (const step of steps.slice(0, count)) { + const response = await send(server, ...step); + const text = await response.text(); + assert.equal(response.status, 200, `${step[0]}: ${server.errors.map(String).join('; ')}`); + responses.push(text); + } + return responses; +} + +test('checkpoint wire fixture independently bounds complete fork, continuation, rejection and reconnect sequence', async () => { + const server = await serveRuntimeConsumer('/unused'); + try { + const responses = await advance(server, steps.length); + assert.deepEqual(JSON.parse(responses[0]).map(state => state.checkpoint.checkpoint_id), ['B', 'A', 'P']); + assert.equal(JSON.parse(responses[0])[0].values.messages[0].content, 'Global B'); + assert.match(responses[2], /event: checkpoints/); + assert.equal(JSON.parse(responses[4]).values.stage, 'A1'); + assert.equal(JSON.parse(responses[7]).values.stage, 'A2'); + assert.deepEqual(JSON.parse(responses[9]).next, ['approval']); + assert.equal(JSON.parse(responses[11]).status, 'running'); + assert.equal(JSON.parse(responses[13]).status, 'success'); + assert.equal(JSON.parse(responses[14]).values.stage, 'A3'); + server.checkpoints.assertComplete(); + assert.deepEqual(server.errors, []); + assert.deepEqual(server.requests, []); + assert.deepEqual(server.historyRequests, []); + assert.equal((await send(server, 'runs/stream', 'POST', body('A3', 'Continue branch'))).status, 500); + assert.equal(server.errors.length, 1); + } finally { await server.close(); } +}); + +for (const [label, prefix, request] of [ + ['wrong source read', 1, ['state/checkpoint', 'POST', { checkpoint: ref('B') }]], + ['missing routing', 2, ['runs/stream', 'POST', { ...body('A', 'Fork A'), checkpoint: undefined }]], + ['wrong root map', 2, ['runs/stream', 'POST', { ...body('A', 'Fork A'), checkpoint: { ...ref('A'), checkpoint_map: {} } }]], + ['missing checkpoint mode', 2, ['runs/stream', 'POST', { ...body('A', 'Fork A'), stream_mode: modes.slice(0, -1) }]], + ['wrong catalog', 2, ['runs/stream', 'POST', { ...body('A', 'Fork A'), input: { ...body('A', 'Fork A').input, client_tools: [] } }]], + ['continue routed to original A', 5, ['runs/stream', 'POST', body('A', 'Continue branch')]], + ['continue routed to selected B', 5, ['runs/stream', 'POST', body('B', 'Continue branch')]], + ['extra creation POST', 3, ['runs/stream', 'POST', body('A', 'Fork A')]], + ['wrong join cursor', 12, [steps[12][0], 'GET', undefined, { 'Last-Event-ID': 'wrong' }]], + ['missing join modes', 12, ['runs/run-A3/stream?cancel_on_disconnect=0', 'GET', undefined, { 'Last-Event-ID': 'drop-cursor' }]], + ['wrong join modes', 12, [`runs/run-A3/stream?cancel_on_disconnect=0&stream_mode=${encodeURIComponent(JSON.stringify(modes.slice(0, -1)))}`, 'GET', undefined, { 'Last-Event-ID': 'drop-cursor' }]], + ['wrong method', 0, ['history']], +]) { + test(`checkpoint oracle rejects ${label} without advancing`, async () => { + const server = await serveRuntimeConsumer('/unused'); + try { + await advance(server, prefix); + const response = await send(server, ...request); + await response.text(); + assert.equal(response.status, 500); + assert.equal(server.errors.length, 1); + assert.equal(server.checkpoints.requests.length, prefix); + } finally { await server.close(); } + }); +} diff --git a/scripts/react-parity/review-runtime.mjs b/scripts/react-parity/review-runtime.mjs index 98c8b144c..e862cc659 100644 --- a/scripts/react-parity/review-runtime.mjs +++ b/scripts/react-parity/review-runtime.mjs @@ -17,7 +17,7 @@ const script = fileURLToPath(import.meta.url); const buildCommand = 'NX_DAEMON=false npx nx run-many -t build -p core,angular,react --skip-nx-cache'; const manualOrder = - 'Use three Load clicks (saved, equal refresh, empty); Send → Tool → Error → Hold → Stop → Pause → Stop → Resume → Resume → Drop → Reconnect → Send → Unmount → Dispose → Send after dispose → Resume after dispose → Reconnect after dispose. Resume answers both approvals, then the final confirmation. Reconnect joins the dropped run without resubmitting. Open /?threads on either review URL for application-owned conversation selection; follow its separate sequence. Only bounded requests are available per server; restart this command for a fresh review. Reloading the page does not reset server state.'; + 'Use three Load clicks (saved, equal refresh, empty); Send → Tool → Error → Hold → Stop → Pause → Stop → Resume → Resume → Drop → Reconnect → Send → Unmount → Dispose → Send after dispose → Resume after dispose → Reconnect after dispose. Resume answers both approvals, then the final confirmation. Reconnect joins the dropped run without resubmitting. Open /?threads for application-owned conversation selection or /?checkpoints for completed checkpoint fork, continued branch, pending rejection and cursor reconnect; follow each view’s separate sequence. Only bounded requests are available per server; restart this command for a fresh review. Reloading the page does not reset server state.'; function prerequisites(root) { const missing = ['core', 'angular', 'react'].filter( @@ -40,10 +40,15 @@ function sourceProvenance(root) { 'fixtures/react-parity/runtime', 'scripts/react-parity/runtime-consumer.mjs', 'scripts/react-parity/thread-lifetime.mjs', + 'scripts/react-parity/checkpoint-execution.mjs', 'scripts/react-parity/review-runtime.mjs', ]; return { head: git('rev-parse', 'HEAD'), + inputs: [...new Set(git('ls-files', '-z', '--cached', '--others', '--exclude-standard', '--', ...paths).split('\0').filter(Boolean))] + .filter(path => existsSync(join(root, path))) + .sort() + .map(path => ({ path, sha256: createHash('sha256').update(readFileSync(join(root, path))).digest('hex') })), trackedRuntimeFixtureStatus: git( 'status', '--short', diff --git a/scripts/react-parity/review-runtime.spec.mjs b/scripts/react-parity/review-runtime.spec.mjs index 4b168d70f..f873cb773 100644 --- a/scripts/react-parity/review-runtime.spec.mjs +++ b/scripts/react-parity/review-runtime.spec.mjs @@ -25,6 +25,10 @@ function fixtureRoot(t) { } const git = (...args) => execFileSync('git', args, { cwd: root, stdio: 'pipe' }); + for (const path of ['scripts/react-parity/checkpoint-execution.mjs', 'fixtures/react-parity/runtime/scenarios.ts']) { + mkdirSync(join(root, path, '..'), { recursive: true }); + writeFileSync(join(root, path), `fixture ${path}`); + } git('init'); git( '-c', @@ -155,6 +159,10 @@ test('workers have isolated unset build environments, e2e runs before fresh manu }).trim() ); assert.equal(review.artifacts.length, 2); + assert.ok(Array.isArray(review.provenance.inputs), 'review provenance includes input hashes'); + for (const path of ['scripts/react-parity/checkpoint-execution.mjs', 'fixtures/react-parity/runtime/scenarios.ts']) { + assert.match(review.provenance.inputs.find(input => input.path === path)?.sha256 ?? '', /^[a-f0-9]{64}$/, `${path} has review input evidence`); + } assert.match(review.artifacts[0].sha256, /^[a-f0-9]{64}$/); assert.ok( logs.some( diff --git a/scripts/react-parity/runtime-consumer.mjs b/scripts/react-parity/runtime-consumer.mjs index 959532a9f..efcf38653 100644 --- a/scripts/react-parity/runtime-consumer.mjs +++ b/scripts/react-parity/runtime-consumer.mjs @@ -8,6 +8,7 @@ import { build } from 'vite'; import ts from 'typescript'; import { chromium, expect } from '@playwright/test'; import { createThreadRoutes, runThreadScenarios } from './thread-lifetime.mjs'; +import { createCheckpointRoutes, runCheckpointScenarios } from './checkpoint-execution.mjs'; export function lockedReactManifest(lock) { const entries = (names) => Object.fromEntries(names.map((name) => { @@ -201,6 +202,21 @@ export function installedTypeSource(template, kind) { // @ts-expect-error Loaded pages remain readonly. history.pop(); const entry = history[0]; + const forked: Promise = session.fork(entry.checkpoint, 'Fork A', runOptions); + void forked; + void session.fork(entry.checkpoint, { message: 'Fork A', state: { color: 'blue' } } as const); + // @ts-expect-error Fork requires a checkpoint reference. + void session.fork('A', 'Fork A'); + // @ts-expect-error Fork input must be authored text/plain data. + void session.fork(entry.checkpoint, { message: 42 }); + // @ts-expect-error Fork cannot override thread routing. + void session.fork(entry.checkpoint, 'Fork A', { config: { configurable: { thread_id: 'other' } } }); + // @ts-expect-error Fork cannot override checkpoint routing. + void session.fork(entry.checkpoint, 'Fork A', { config: { configurable: { checkpoint_id: 'B' } } }); + // @ts-expect-error Fork cannot override root namespace. + void session.fork(entry.checkpoint, 'Fork A', { config: { configurable: { checkpoint_ns: 'child' } } }); + // @ts-expect-error Fork cannot override checkpoint map. + void session.fork(entry.checkpoint, 'Fork A', { config: { configurable: { checkpoint_map: {} } } }); const checkpointId: string | null | undefined = entry.checkpoint.checkpoint_id; const parentId: string | null | undefined = entry.parent_checkpoint?.checkpoint_id; const checkpointData: PlainValue = entry.checkpoint.checkpoint_map?.['branch']; @@ -331,15 +347,19 @@ export async function prepareRuntimeConsumer(root, consumer, kind) { const destination = kind === 'angular' ? join(consumer, 'src') : consumer; cpSync(join(temporary, 'bundle/runtime-entry.js'), join(destination, 'runtime-entry.js')); cpSync(join(temporary, 'types/fixtures/react-parity/runtime/runtime-entry.d.ts'), join(destination, 'runtime-entry.d.ts')); + const declaration = readFileSync(join(destination, 'runtime-entry.d.ts'), 'utf8'); + assert.doesNotMatch(declaration, /@langchain|libs\/|create-session/, 'emitted fixture declaration exposes no SDK/private references'); cpSync(join(fixture, 'scenarios.ts'), join(destination, 'scenarios.ts')); cpSync(join(fixture, 'thread-owner.ts'), join(destination, 'thread-owner.ts')); const threadView = `${kind}-threads.${kind === 'react' ? 'tsx' : 'ts'}`; cpSync(join(fixture, threadView), join(destination, threadView)); + const checkpointView = `${kind}-checkpoints.${kind === 'react' ? 'tsx' : 'ts'}`; + cpSync(join(fixture, checkpointView), join(destination, checkpointView)); cpSync(join(fixture, 'review.css'), join(destination, 'review.css')); const app = `${kind}-app.${kind === 'react' ? 'tsx' : 'ts'}`; cpSync(join(fixture, app), join(destination, app)); writeFileSync(join(destination, kind === 'react' ? 'main.tsx' : 'main.ts'), - `if (new URLSearchParams(location.search).has('threads')) {\n void import('./${kind}-threads');\n} else {\n void import('./${kind}-app');\n}\n`); + `if (new URLSearchParams(location.search).has('checkpoints')) {\n void import('./${kind}-checkpoints');\n} else if (new URLSearchParams(location.search).has('threads')) {\n void import('./${kind}-threads');\n} else {\n void import('./${kind}-app');\n}\n`); if (kind === 'angular') { const configPath = join(consumer, 'angular.json'); const config = JSON.parse(readFileSync(configPath, 'utf8')); @@ -359,6 +379,7 @@ export async function prepareRuntimeConsumer(root, consumer, kind) { /** Bounded fixture server: built files and deterministic history/run routes. */ export async function serveRuntimeConsumer(directory) { const threads = createThreadRoutes(); + const checkpoints = createCheckpointRoutes(); const requests = []; const historyRequests = []; const joinRequests = []; @@ -377,6 +398,7 @@ export async function serveRuntimeConsumer(directory) { const pathname = url.pathname; if (pathname.startsWith('/api/')) { if (await threads.handle(request, response, pathname)) return; + if (await checkpoints.handle(request, response, pathname)) return; const runPath = '/api/threads/fixture-thread/runs/drop-run'; if (pathname === runPath || pathname === `${runPath}/stream`) { assert.equal(request.method, 'GET', 'known-run recovery only performs GET'); @@ -447,9 +469,10 @@ export async function serveRuntimeConsumer(directory) { server.listen(0, '127.0.0.1'); await once(server, 'listening'); return { - url: `http://127.0.0.1:${server.address().port}`, requests, historyRequests, joinRequests, statusRequests, errors, holdStarted, holdAborted, threads, + url: `http://127.0.0.1:${server.address().port}`, requests, historyRequests, joinRequests, statusRequests, errors, holdStarted, holdAborted, threads, checkpoints, async close() { threads.close(); + checkpoints.close(); for (const response of held) response.destroy(); const closed = once(server, 'close'); server.close(); @@ -724,6 +747,7 @@ export async function runRuntimeScenarios(directory, kind) { assert.deepEqual(server.historyRequests, [{ limit: 10 }, { limit: 10 }, { limit: 10 }], 'only explicit loads read history'); completed.push('unmount and explicit disposal'); completed.push(...await runThreadScenarios(page, server)); + completed.push(...await runCheckpointScenarios(page, server)); assert.deepEqual(server.errors.map(String), []); assert.deepEqual(pageErrors, []); assert.deepEqual(unexpected, []); From 9ddeb1579323898cb5e470d658cfa53852931830 Mon Sep 17 00:00:00 2001 From: Brian Love Date: Thu, 24 Sep 2026 11:40:43 -0700 Subject: [PATCH 2/2] feat(runtime): preserve owned citation metadata --- fixtures/react-parity/runtime/README.md | 12 + fixtures/react-parity/runtime/angular-app.ts | 6 + fixtures/react-parity/runtime/evidence.json | 221 +++++++++ .../react-parity/runtime/installed-types.ts | 36 ++ fixtures/react-parity/runtime/react-app.tsx | 6 + fixtures/react-parity/runtime/scenarios.ts | 3 + libs/core/README.md | 9 +- libs/core/src/contracts/citation.ts | 15 + .../core/src/contracts/contracts.type-test.ts | 26 + libs/core/src/contracts/message.ts | 4 +- libs/core/src/index.ts | 1 + libs/langgraph/src/runtime/README.md | 22 + .../src/runtime/citation-projection.spec.ts | 130 +++++ .../src/runtime/citation-projection.ts | 55 +++ libs/langgraph/src/runtime/citations.spec.ts | 460 ++++++++++++++++++ .../src/runtime/history-projection.spec.ts | 29 ++ .../src/runtime/history-projection.ts | 7 +- .../src/runtime/message-reducer.spec.ts | 108 ++++ libs/langgraph/src/runtime/message-reducer.ts | 8 +- libs/langgraph/src/runtime/ownership.ts | 20 +- .../langgraph/src/runtime/publication.spec.ts | 54 ++ .../src/runtime/stream-projection.spec.ts | 32 +- .../src/runtime/stream-projection.ts | 2 + scripts/react-parity/baseline.json | 35 +- scripts/react-parity/dispositions.json | 13 +- scripts/react-parity/runtime-consumer.mjs | 10 +- 26 files changed, 1294 insertions(+), 30 deletions(-) create mode 100644 libs/core/src/contracts/citation.ts create mode 100644 libs/langgraph/src/runtime/citation-projection.spec.ts create mode 100644 libs/langgraph/src/runtime/citation-projection.ts create mode 100644 libs/langgraph/src/runtime/citations.spec.ts diff --git a/fixtures/react-parity/runtime/README.md b/fixtures/react-parity/runtime/README.md index 10e8066b9..4bb52b20a 100644 --- a/fixtures/react-parity/runtime/README.md +++ b/fixtures/react-parity/runtime/README.md @@ -1,5 +1,17 @@ # Installed native runtime consumers +## Citation projection + +The main view's plain-text Citations field proves the installed readonly core +`Citation` contract through both native bindings. The existing Load sequence +shows `saved-source: Saved reference`, preserves it on an equal refresh, and +clears it on empty history. The first Resume shows `approval-source: Approval +reference`; the next Resume canonically completes the same assistant message +without metadata and clears the field. These assertions add no requests, tool +executions or navigation. Installed declaration probes reject citation mutation, +Date timestamps and non-plain provider metadata. This is display projection only, +not a citation renderer or completion of rich-message parity. + ## Checkpoint review Open `/?checkpoints` on either installed review URL for the separate fixed diff --git a/fixtures/react-parity/runtime/angular-app.ts b/fixtures/react-parity/runtime/angular-app.ts index 19c3244ca..c4577e8a8 100644 --- a/fixtures/react-parity/runtime/angular-app.ts +++ b/fixtures/react-parity/runtime/angular-app.ts @@ -173,6 +173,12 @@ const submit = (input: string) => { view().transcript }} +
+

Citations

+ {{ + view().citations + }} +
diff --git a/fixtures/react-parity/runtime/evidence.json b/fixtures/react-parity/runtime/evidence.json index 38cc1d980..b3ad83516 100644 --- a/fixtures/react-parity/runtime/evidence.json +++ b/fixtures/react-parity/runtime/evidence.json @@ -477,5 +477,226 @@ "Strict deterministic HTTP evidence complements separate real-server tests; it does not certify general LangGraph compatibility or cross-client atomicity.", "No new public API, branch UI library, dependency update, package-root migration or release." ] + }, + "ownedCitationProjection": { + "observedOn": "2026-09-24", + "scope": "Private owned citation metadata in the shared message read model, observed through both installed native bindings; no public LangGraph root migration or citation renderer.", + "source": { + "baseCommit": "4d47995333278f9a19491925bc555d30f8490956", + "workingTree": "Uncommitted owned-citation implementation on the reviewed installed-checkpoint milestone, itself based on merged #1156. Hashes identify tested inputs, not a future commit.", + "algorithm": "SHA-256(file bytes); manifest hash is SHA-256 of sorted sha256 + two spaces + path + LF lines.", + "inputs": [ + { + "path": "fixtures/react-parity/runtime/angular-app.ts", + "sha256": "596f37317491248d7ec86d02fcc2138bc4231667c73d1e4bfa244a81c1970937" + }, + { + "path": "fixtures/react-parity/runtime/installed-types.ts", + "sha256": "3ef07efb0d9f3531dbe77ea464683cc4a7f0cfd6f10bad7b3e119978b1c7165d" + }, + { + "path": "fixtures/react-parity/runtime/react-app.tsx", + "sha256": "cde2df17216eb7ea132f495dd48b77ca38122f36122c924ffbe2c199588412bc" + }, + { + "path": "fixtures/react-parity/runtime/runtime-entry.ts", + "sha256": "72509ac95343b1e2d09c825d5017408fcd66891d5efb96e6e742270baf90aebd" + }, + { + "path": "fixtures/react-parity/runtime/scenarios.ts", + "sha256": "6a43712903b9529a2f16fead4e2164f3d1ba2f2da4556abccbb784d3ed660d4a" + }, + { + "path": "libs/core/src/contracts/citation.ts", + "sha256": "5b8b1acdbdfc6a4f74305f2cd3e507858727069582b3f9d81bc5a71498ceb1c1" + }, + { + "path": "libs/core/src/contracts/message.ts", + "sha256": "e8661f99972c20e22728566e6219105507e91ccd7166b1686d0a92d34c6d84ff" + }, + { + "path": "libs/core/src/index.ts", + "sha256": "90b04d322bb6ea03f41263a4756d7edae70e49794f1201fbdd6f7b5b99ba549d" + }, + { + "path": "libs/langgraph/src/runtime/checkpoint-authority.ts", + "sha256": "f56912f077e4ba891f9a55419e89743d62f573c435c3984fb64ac0803271c3dc" + }, + { + "path": "libs/langgraph/src/runtime/citation-projection.ts", + "sha256": "52d99b652713e964242f38ec9e81f74695ecac938d42f10255b70d3c7eab8201" + }, + { + "path": "libs/langgraph/src/runtime/create-session.ts", + "sha256": "f4c6a1e693d205cf48a7c9ec357ddf037bdcf4ceabe315c6529c6806fbdaefb1" + }, + { + "path": "libs/langgraph/src/runtime/history-projection.ts", + "sha256": "c9b18e186261d16d6b363ddc906793acbf3295e4a5c45e93bbf2d555bd2b2e92" + }, + { + "path": "libs/langgraph/src/runtime/message-reducer.ts", + "sha256": "5758caaae8c31361a9361b81a7995239c10a48bf670a135a9fe6a327da2e6531" + }, + { + "path": "libs/langgraph/src/runtime/ownership.ts", + "sha256": "0cfc4dabfedbc3aa9406b943257ab3df3bae240276ccaec6d50f1e51f11e5c69" + }, + { + "path": "libs/langgraph/src/runtime/publication.ts", + "sha256": "7d290721417d504313f893ee7e90154b85e9eee27d05f6cce816d592495791e5" + }, + { + "path": "libs/langgraph/src/runtime/stream-projection.ts", + "sha256": "076cf6eb83b633d67a42118d7f17931b3f5010ca9f00b7c15e97bf5214f42240" + }, + { + "path": "package-lock.json", + "sha256": "530c28574a832ab39c15c9607717de0eb8ccfae2593e09c7ae373bd8936d0b64" + }, + { + "path": "scripts/react-parity/runtime-consumer.mjs", + "sha256": "1620b2bea166aa3cb0ea9fd2bef835b8f60f298eeea55331a3e045f36068c8d8" + } + ], + "manifestSha256": "734886ae78fe391f149ca6a3fce7eedf08eddd03790bb1de78fbecde44ecba2f" + }, + "automated": { + "core": { + "targets": [ + "test", + "type-tests", + "lint", + "build" + ], + "path": "/tmp/owned-citations-core-final-restored.log", + "sha256": "e05e9066f318050de23cf28a6c6eadb71f0dabaacd5a747c0f186d1f125935ca" + }, + "runtime": { + "targets": [ + "runtime-quality", + "runtime-type-tests" + ], + "tests": 859, + "files": 46, + "path": "/tmp/owned-citations-runtime-final.log", + "sha256": "72e0e86b389deff226249eeb4bd3fcc31e8530c953c58b967a6d214a7ae6b960" + }, + "publicLangGraph": { + "targets": [ + "test", + "type-tests", + "build" + ], + "path": "/tmp/owned-citations-public-langgraph.log", + "sha256": "9f46bbff7513f767d905cc0dad955a1b7eaac1d307900d617d15d6bec289835f" + }, + "packageBuilds": { + "projects": [ + "core", + "content", + "react", + "angular" + ], + "path": "/tmp/owned-citations-build-consumers.log", + "sha256": "35d5810824fba54bc69b0cf54ab9b13a1ab147e3d62ddcb60b6fbc07f13e96e8" + }, + "parityScripts": { + "tests": 329, + "path": "/tmp/owned-citations-parity-scripts.log", + "sha256": "36923eccbb0ff34d967909d8bbd7f928c2167d57f9e71808739e80d55d89142b" + }, + "react": { + "scenarios": 27, + "path": "/tmp/owned-citations-installed-react-restored.log", + "sha256": "2b11a4267ff6e40cfaae745f805118d25916f7a8276787a8bf814d0056e39e9f" + }, + "angular": { + "scenarios": 27, + "path": "/tmp/owned-citations-installed-angular-restored.log", + "sha256": "ea4619deb6f8a1a487a5cd09733256f94c31a78820af1e9c385649f7276c1618" + }, + "boundaries": { + "path": "/tmp/owned-citations-boundaries.log", + "sha256": "995db896b38cf7de5ca9db6590dba47e082d26abdf69bab04b92eee484696063" + }, + "versions": { + "path": "/tmp/owned-citations-versions.log", + "sha256": "1d9f9440bfffad48cca3d74ddec850cf803b8f449dcee2216776392713bf9374" + }, + "inventory": { + "path": "/tmp/owned-citations-inventory-check.log", + "sha256": "ee3072f80303358bc95e35f0d8f71b478dc82b1049814291166802e581948719" + } + }, + "redGreen": [ + { + "path": "/tmp/owned-citations-contract-red.log", + "sha256": "4d1b4438ae3840b3b25139d69d37336dddd380f2b3acb9cda48d984da39bd880", + "expectedFailure": "Core has no exported Citation type or Message.citations field." + }, + { + "path": "/tmp/owned-citations-normalizer-red.log", + "sha256": "f34bb7c33368482aec4e736d0323535aef1b43147356607e91ccfa37e26cd774", + "expectedFailure": "Neutral citation projection module missing." + }, + { + "path": "/tmp/owned-citations-projection-red.log", + "sha256": "3b41e28eae7a61e86610c440f4a6c6201ff81138738bee4a12cc4541834db320", + "expectedFailure": "Six ownership/reduction/history/stream/publication citation tests fail while existing cases pass." + } + ], + "semanticNegativeControl": { + "mutation": "Temporarily omit only the production history-projection citations assignment, keeping the native fields and wire citations in place.", + "react": { + "path": "/tmp/owned-citations-installed-react-negative.log", + "sha256": "8fa5d26b97a748146996ddde49597b97814095c05178110f4c70fbc984968472" + }, + "angular": { + "path": "/tmp/owned-citations-installed-angular-negative.log", + "sha256": "48301b5348f2f86d61b37e3557a41211f6e171daceccb946bcc096631cbc7c74" + }, + "expectedFailure": "Both native installed consumers expected saved-source: Saved reference and received empty citation text.", + "restoration": "Original history-projection source restored in finally and byte-compared with cmp; both installed consumers then passed all 27 scenarios.", + "restoredHistorySha256": "c9b18e186261d16d6b363ddc906793acbf3295e4a5c45e93bbf2d555bd2b2e92" + }, + "installedAssertions": [ + "Readonly Citation and nested plain metadata through both native bindings with skipLibCheck:false.", + "Saved source alias normalizes on history load and equal refresh, and empty history clears it.", + "First Resume publishes approval-source; second Resume canonically clears citations on the same assistant message.", + "Existing main/thread/checkpoint request counts, tool runs and lifecycle assertions stay unchanged." + ], + "manualReview": { + "date": "2026-09-24", + "react": "Chrome MCP, owned port 56833", + "angular": "In-app browser, owned port 56834", + "sequence": "Load, Load, Load, Pause, Resume, Resume, Dispose", + "observations": [ + "Saved reference appeared on the first history load, survived equal history, and cleared with empty history.", + "Approval reference appeared during resumed approval and cleared at canonical completion in both native views.", + "Both views finished three loads and two resumes, recorded zero local tool handlers, and disposed their owner.", + "Both browser consoles had zero warnings and errors.", + "Chrome independently recorded exactly six successful POST requests: three histories and three runs. Angular network behavior was checked by the automated oracle." + ], + "runner": "/tmp/owned-citations-manual-runner.log", + "freshAutomatedScenarios": 54, + "cleanup": "Only the two owned tabs and runner were closed; ports 56833 and 56834 were confirmed closed. Existing manual sessions were preserved.", + "scope": "Citation lifecycle walkthrough only; not an additional manual review of every main, thread or checkpoint scenario." + }, + "independentReview": { + "compliance": "Approved; 106 focused runtime tests passed. A misleading test title was clarified and all 19 session citation tests rerun.", + "quality": "Approved with no must-fix findings; 39 citation/metadata tests passed. All 18 input hashes and 15 recorded log hashes matched.", + "parent": "Inspected production and fixture changes, reran all 859 runtime tests and runtime/public types, and performed the scoped browser walkthrough." + }, + "documentation": { + "generatorsRun": [], + "reason": "Only private core/runtime and contributor fixtures changed. API generator registers langgraph/chat/render/ag-ui/a2ui/middleware public roots, with no core entry. No public API, narrative docs or public agent-context input changed." + }, + "limits": [ + "Display strings do not certify source truth or safe HTML/navigation.", + "Strict checkpoint ingress still rejects Date and cyclic/non-plain values before citation normalization.", + "Legacy Angular arbitrary-extra semantics remain separate until backend cutover.", + "No full T09/T15 parity, rich blocks, reasoning/timing, citation UI, dependency upgrade, release, live-server, SSR or performance claim." + ] } } diff --git a/fixtures/react-parity/runtime/installed-types.ts b/fixtures/react-parity/runtime/installed-types.ts index 26ab15c4e..58da0540f 100644 --- a/fixtures/react-parity/runtime/installed-types.ts +++ b/fixtures/react-parity/runtime/installed-types.ts @@ -2,6 +2,7 @@ import type { AgentError, AgentSession, AgentSnapshot, + Citation, CompleteOutcome, Message, PlainValue, @@ -14,6 +15,41 @@ import type { FixtureTools } from './scenarios'; export function assertSnapshot(snapshot: AgentSnapshot) { /* BACKEND_VALUES */ + const citations: readonly Citation[] | undefined = + snapshot.messages[0]?.citations; + if (citations) { + const citation: Citation = citations[0]; + const publishedAt: string | number | undefined = citation.publishedAt; + const extra: Readonly> | undefined = + citation.extra; + // @ts-expect-error Citation collections remain readonly through the native binding. + citations.push(citation); + // @ts-expect-error Citation fields remain readonly. + citation.title = 'changed'; + if (extra) { + // @ts-expect-error Provider records remain deeply readonly plain data. + extra['changed'] = true; + const nested = extra['nested']; + if (nested && typeof nested === 'object') { + // @ts-expect-error Nested provider records remain readonly. + nested['value'] = true; + } + } + void publishedAt; + } + const badTimestamp: Citation = { + id: 'c', + index: 1, + // @ts-expect-error Date timestamps are not in the neutral contract. + publishedAt: new Date(), + }; + const badExtra: Citation = { + id: 'c', + index: 1, + // @ts-expect-error SDK instances and callbacks cannot enter plain extras. + extra: { date: new Date(), callback: () => true }, + }; + void [badTimestamp, badExtra]; for (const call of snapshot.toolCalls) { if (call.name === 'weather') { const city: string = call.args.city; diff --git a/fixtures/react-parity/runtime/react-app.tsx b/fixtures/react-parity/runtime/react-app.tsx index 9ef3af43b..154b19057 100644 --- a/fixtures/react-parity/runtime/react-app.tsx +++ b/fixtures/react-parity/runtime/react-app.tsx @@ -221,6 +221,12 @@ function App() { {view.transcript} +
+

Citations

+ + {view.citations} + +
diff --git a/fixtures/react-parity/runtime/scenarios.ts b/fixtures/react-parity/runtime/scenarios.ts index de0637415..8469548a3 100644 --- a/fixtures/react-parity/runtime/scenarios.ts +++ b/fixtures/react-parity/runtime/scenarios.ts @@ -139,6 +139,9 @@ export function display(snapshot: FixtureSnapshot) { return { text: assistant.map((message) => message.content).join('\n'), transcript: snapshot.messages.map((message) => message.content).join('\n'), + citations: snapshot.messages.flatMap((message) => + (message.citations ?? []).map((citation) => `${citation.id}: ${citation.title ?? ''}`) + ).join('\n'), values: JSON.stringify(snapshot.values) ?? 'unobserved', history: JSON.stringify(snapshot.history) ?? 'unobserved', interrupts: JSON.stringify(snapshot.interrupts), diff --git a/libs/core/README.md b/libs/core/README.md index 03b73aa48..b85c408b5 100644 --- a/libs/core/README.md +++ b/libs/core/README.md @@ -1,7 +1,7 @@ # @threadplane/core Private, unpublished agent contracts for the shared runtime work. The root exports -readonly text messages, delivery outcomes, plain error projections, authored tool +readonly text messages and citations, delivery outcomes, plain error projections, authored tool call types, snapshots and the minimal session interface. This package implements delivery constructors and error projection, not an execution owner or public store. @@ -16,6 +16,13 @@ causes, controllers, SDK instances and arbitrary extras belong to effect boundar `projectAgentError` copies already classified display fields; it does not classify transport failures. Tool failures carry a plain error string instead. +`Message.citations` optionally contains readonly `Citation` source metadata: stable +IDs/order, title, URL, snippet, source type, icon URL, primitive publication time, +and deeply readonly plain provider extras. These strings are authored display data; +they do not certify source accuracy or safe HTML/navigation. No URL is fetched or +footnote resolved. Date instances, SDK objects and functions are outside this +portable contract. + Tool contracts pair authored TypeScript argument/result types by name, including ordinary interfaces. Pending calls have decoded, finalized arguments awaiting execution; partial streamed arguments are not exposed as authored types. `void` diff --git a/libs/core/src/contracts/citation.ts b/libs/core/src/contracts/citation.ts new file mode 100644 index 000000000..d61fb2284 --- /dev/null +++ b/libs/core/src/contracts/citation.ts @@ -0,0 +1,15 @@ +import type { PlainValue } from './tool.js'; + +/** Authored source metadata, owned by a snapshot. Strings are display data, + * not a certification of safe navigation, HTML or source accuracy. */ +export interface Citation { + readonly id: string; + readonly index: number; + readonly title?: string; + readonly url?: string; + readonly snippet?: string; + readonly sourceType?: string; + readonly iconUrl?: string; + readonly publishedAt?: string | number; + readonly extra?: Readonly>; +} diff --git a/libs/core/src/contracts/contracts.type-test.ts b/libs/core/src/contracts/contracts.type-test.ts index 2528b37d4..8ca507b52 100644 --- a/libs/core/src/contracts/contracts.type-test.ts +++ b/libs/core/src/contracts/contracts.type-test.ts @@ -2,6 +2,7 @@ import type { AgentError, AgentSession, AgentSnapshot, + Citation, Message, MessageDelivery, PlainValue, @@ -10,6 +11,31 @@ import type { declare const snapshot: AgentSnapshot; declare const session: AgentSession; +declare const citation: Citation; +const citationList: Message['citations'] = [citation]; +// @ts-expect-error citation collections are readonly +citationList.push(citation); +// @ts-expect-error citation fields are readonly +citation.title = 'changed'; +if (citation.extra) { + // @ts-expect-error provider metadata is readonly + citation.extra['changed'] = true; +} +// @ts-expect-error timestamps are portable primitives +const datedCitation: Citation = { id: 'c', index: 1, publishedAt: new Date() }; +const opaqueCitation: Citation = { + id: 'c', + index: 1, + // @ts-expect-error SDK instances are not portable metadata + extra: { date: new Date() }, +}; +const callableCitation: Citation = { + id: 'c', + index: 1, + // @ts-expect-error functions are not portable metadata + extra: { run: () => 1 }, +}; +void [datedCitation, opaqueCitation, callableCitation]; // @ts-expect-error snapshot fields are readonly snapshot.status = 'running'; // @ts-expect-error snapshot collections are readonly diff --git a/libs/core/src/contracts/message.ts b/libs/core/src/contracts/message.ts index 09d74e2b8..19ca2857a 100644 --- a/libs/core/src/contracts/message.ts +++ b/libs/core/src/contracts/message.ts @@ -1,12 +1,14 @@ import type { MessageDelivery } from './delivery.js'; +import type { Citation } from './citation.js'; export type Role = 'user' | 'assistant' | 'system' | 'tool'; -/** Owned text projection; rich content and arbitrary SDK extras are not in this slice. */ +/** Owned text and source metadata; rich content and SDK objects remain outside this slice. */ export interface Message { readonly id: string; readonly role: Role; readonly content: string; + readonly citations?: readonly Citation[]; readonly delivery: MessageDelivery; readonly toolCallId?: string; readonly toolCallIds?: readonly string[]; diff --git a/libs/core/src/index.ts b/libs/core/src/index.ts index e40710c22..d0ef69c91 100644 --- a/libs/core/src/index.ts +++ b/libs/core/src/index.ts @@ -5,6 +5,7 @@ export type { } from './contracts/agent-session.js'; export type { AgentSnapshot, AgentStatus } from './contracts/agent-snapshot.js'; export type { Message, Role } from './contracts/message.js'; +export type { Citation } from './contracts/citation.js'; export { completeDelivery, staticDelivery, diff --git a/libs/langgraph/src/runtime/README.md b/libs/langgraph/src/runtime/README.md index 9980787de..17f796e41 100644 --- a/libs/langgraph/src/runtime/README.md +++ b/libs/langgraph/src/runtime/README.md @@ -4,6 +4,28 @@ This directory stages the framework-independent session owner. It is not the published LangGraph package entry point. Core and framework bindings do not own backend execution positions. +## Owned citations + +History, root streams and read-only child streams project citations from +`additional_kwargs.citations ?? additional_kwargs.sources`. Supported arrays replace +the list; an empty array clears it. Missing metadata retains the current list on +interim events within one generation, while canonical messages and authoritative +history reads define the complete state and clear omissions. Later interim events +cannot undo committed canonical metadata. Equal metadata shares owned references, +including when only text changes. Terminal candidates and queued publications +capture metadata before callers can mutate it. + +The private normalizer stages the existing Angular extractor's aliases without an +Angular dependency. Consolidate this bounded duplication at the backend cutover. +The adapters do not promise identical arbitrary-extra semantics: neutral snapshots +accept plain provider records and string/finite-number timestamps, omit unsupported +optional fields, and reject instances/cycles in extras at the ownership boundary. +Exact checkpoint ingress remains stricter: it captures the complete plain state +before projection and rejects Date instances even in otherwise optional metadata. +Citation changes never authorize tools, change execution positions, or reopen +delivery. Rich blocks, reasoning/timing, event render state and citation components +remain separate future capabilities. + ## Completed checkpoint forks `session.fork(checkpoint, input, options?)` submits new input in the session's diff --git a/libs/langgraph/src/runtime/citation-projection.spec.ts b/libs/langgraph/src/runtime/citation-projection.spec.ts new file mode 100644 index 000000000..79904241a --- /dev/null +++ b/libs/langgraph/src/runtime/citation-projection.spec.ts @@ -0,0 +1,130 @@ +import { describe, expect, it } from 'vitest'; +import { projectCitations } from './citation-projection'; + +const project = (citations: unknown) => + projectCitations({ additional_kwargs: { citations } }); + +describe('owned citation normalization', () => { + it('normalizes aliases and deterministic positions without inventing source facts', () => { + expect( + project([ + 'https://example.test/one', + { + refId: 'ref', + index: 4.5, + name: 'Name', + href: '/two', + content: 'Excerpt', + sourceType: 'web', + iconUrl: '/icon', + publishedAt: '2026-09-24', + }, + { + id: 'primary', + refId: 'ignored', + title: '', + name: 'ignored', + url: '', + href: 'ignored', + snippet: '', + excerpt: 'ignored', + publishedAt: 42, + }, + { source: '/four', excerpt: 'Four' }, + null, + ]) + ).toEqual([ + { id: 'c1', index: 1, url: 'https://example.test/one' }, + { + id: 'ref', + index: 4.5, + title: 'Name', + url: '/two', + snippet: 'Excerpt', + sourceType: 'web', + iconUrl: '/icon', + publishedAt: '2026-09-24', + }, + { + id: 'primary', + index: 3, + title: '', + url: '', + snippet: '', + publishedAt: 42, + }, + { id: 'c4', index: 4, url: '/four', snippet: 'Four' }, + { id: 'c5', index: 5 }, + ]); + }); + + it('preserves absent versus empty and nullish precedence over sources', () => { + expect(projectCitations({})).toBeUndefined(); + expect( + projectCitations({ + additional_kwargs: { citations: [], sources: ['ignored'] }, + }) + ).toEqual([]); + expect( + projectCitations({ + additional_kwargs: { citations: null, sources: ['source'] }, + }) + ).toEqual([{ id: 'c1', index: 1, url: 'source' }]); + expect( + projectCitations({ + additional_kwargs: { citations: 'unsupported', sources: ['ignored'] }, + }) + ).toBeUndefined(); + }); + + it('omits unsupported optional fields and accepts only finite numeric indexes and timestamps', () => { + expect( + project([ + { index: Infinity, publishedAt: NaN, title: 4, extra: 'ignored' }, + { index: NaN, publishedAt: new Date() }, + { index: -2, publishedAt: 0 }, + ]) + ).toEqual([ + { id: 'c1', index: 1 }, + { id: 'c2', index: 2 }, + { id: 'c3', index: -2, publishedAt: 0 }, + ]); + }); + + it('owns arrays, entries and nested provider metadata without freezing callers', () => { + const extra = { nested: { tags: ['original'] } }; + const entry = { title: 'Original', extra }; + const source = [entry]; + const result = project(source)!; + source.push({ title: 'Another', extra }); + entry.title = 'Mutated'; + extra.nested.tags.push('mutated'); + expect(result).toEqual([ + { + id: 'c1', + index: 1, + title: 'Original', + extra: { nested: { tags: ['original'] } }, + }, + ]); + expect(Object.isFrozen(result)).toBe(true); + expect(Object.isFrozen(result[0])).toBe(true); + expect(Object.isFrozen(result[0].extra?.['nested'])).toBe(true); + expect(Object.isFrozen(entry)).toBe(false); + expect(Object.isFrozen(extra)).toBe(false); + }); + + it('rejects opaque and cyclic extras at the existing ownership boundary', () => { + const cycle: Record = {}; + cycle['self'] = cycle; + for (const extra of [ + new Date(), + { nested: new Map() }, + cycle, + { run: () => 1 }, + ]) { + expect(() => project([{ extra }])).toThrow(/plain data/); + expect(Object.isFrozen(extra)).toBe(false); + } + }); +}); diff --git a/libs/langgraph/src/runtime/citation-projection.ts b/libs/langgraph/src/runtime/citation-projection.ts new file mode 100644 index 000000000..d30d4cfd4 --- /dev/null +++ b/libs/langgraph/src/runtime/citation-projection.ts @@ -0,0 +1,55 @@ +import type { Citation, PlainValue } from '@threadplane/core'; +import { ownValue } from './ownership'; +import { record } from './wire-message'; + +/** Bounded staging copy of the legacy extractor's aliases. Consolidate at the + * backend cutover: this neutral contract deliberately requires plain extras and + * primitive timestamps. Undefined means no update; a supported list replaces. */ +export function projectCitations( + message: Record +): readonly Citation[] | undefined { + const kwargs = record(message['additional_kwargs']); + const raw = kwargs?.['citations'] ?? kwargs?.['sources']; + if (!Array.isArray(raw)) return undefined; + return ownValue( + raw.map((entry, position) => { + const index = position + 1; + if (typeof entry === 'string') + return { id: `c${index}`, index, url: entry }; + const source = record(entry); + const firstString = (...keys: string[]) => { + for (const key of keys) { + const value = source?.[key]; + if (typeof value === 'string') return value; + } + return undefined; + }; + const explicitIndex = source?.['index']; + const publishedAt = source?.['publishedAt']; + const extra = record(source?.['extra']); + const result: Record = { + id: firstString('id', 'refId') ?? `c${index}`, + index: + typeof explicitIndex === 'number' && Number.isFinite(explicitIndex) + ? explicitIndex + : index, + }; + const fields = { + title: firstString('title', 'name'), + url: firstString('url', 'href', 'source'), + snippet: firstString('snippet', 'content', 'excerpt'), + sourceType: firstString('sourceType'), + iconUrl: firstString('iconUrl'), + publishedAt: + typeof publishedAt === 'string' || + (typeof publishedAt === 'number' && Number.isFinite(publishedAt)) + ? publishedAt + : undefined, + extra: extra as PlainValue, + }; + for (const [key, value] of Object.entries(fields)) + if (value !== undefined) result[key] = value; + return result; + }) + ) as unknown as readonly Citation[]; +} diff --git a/libs/langgraph/src/runtime/citations.spec.ts b/libs/langgraph/src/runtime/citations.spec.ts new file mode 100644 index 000000000..af1a95f2e --- /dev/null +++ b/libs/langgraph/src/runtime/citations.spec.ts @@ -0,0 +1,460 @@ +import { afterEach, describe, expect, it, vi } from 'vitest'; +import { createSession } from './create-session'; +import { + captureCheckpointState, + captureCheckpointEvent, + confirmCheckpoint, +} from './checkpoint-authority'; +import type { AgentTransport, StreamEvent } from './transport.types'; +import { controlledTransport } from './testing/controlled-transport'; +import { deferred } from './testing/deferred'; +import { + checkpointEvent, + fixture as checkpointFixture, + position, + saved, +} from './testing/checkpoint-fixture'; + +const ai = (citations?: unknown, content = 'Answer') => ({ + type: 'ai', + id: 'answer', + content, + ...(citations === undefined ? {} : { additional_kwargs: { citations } }), +}); +const citation = (title = 'Source') => ({ + id: 'source', + title, + extra: { tags: ['original'] }, +}); +const frame = (citations?: unknown, content = 'Answer'): StreamEvent => ({ + type: 'values', + data: { messages: [ai(citations, content)] }, +}); +const cleanups: (() => Promise)[] = []; +function fixture() { + const streams: ReturnType>[] = []; + const starts = Array.from({ length: 4 }, () => deferred()); + const acknowledgments: ReturnType>[] = []; + const stream = vi.fn((_a, _t, _input, signal) => { + const index = streams.length; + const controlled = controlledTransport({ + signal, + ignoreAbort: true, + }); + streams.push(controlled); + starts[index].resolve(); + const iterator: AsyncIterableIterator = { + [Symbol.asyncIterator]() { + return iterator; + }, + next() { + acknowledgments[index]?.resolve(); + return controlled.stream.next(); + }, + return() { + acknowledgments[index]?.resolve(); + return controlled.stream.return(); + }, + }; + return iterator; + }); + const getHistory = vi.fn>( + async () => [] + ); + const session = createSession({ + assistantId: 'agent', + threadId: 'thread', + transport: { stream, getHistory }, + }); + cleanups.push(() => session.dispose()); + return { + session, + stream, + streams, + getHistory, + started: (index = 0) => starts[index].promise, + async emit(event: StreamEvent, index = 0) { + acknowledgments[index] = deferred(); + streams[index].release(event); + await acknowledgments[index].promise; + }, + }; +} +afterEach(async () => { + await Promise.all(cleanups.splice(0).map((cleanup) => cleanup())); +}); + +describe('session citation ownership', () => { + it('captures terminal candidates, publishes metadata-only updates and clears canonical omissions without reopening delivery', async () => { + const f = fixture(); + const run = f.session.submit('Go'); + await f.started(); + await f.emit(frame([citation()])); + const first = f.session.getSnapshot().messages.at(-1)!; + const raw = [citation('Corrected')]; + await f.emit(frame(raw)); + const corrected = f.session.getSnapshot().messages.at(-1)!; + expect(corrected.content).toBe(first.content); + expect(corrected.citations?.[0].title).toBe('Corrected'); + expect(corrected).not.toBe(first); + raw[0].title = 'Mutated'; + raw[0].extra.tags.push('mutated'); + f.streams[0].finish(); + expect(await run).toBe('success'); + expect(f.session.getSnapshot().messages.at(-1)).toMatchObject({ + citations: [{ title: 'Corrected', extra: { tags: ['original'] } }], + delivery: { phase: 'complete', outcome: 'success' }, + }); + const beforeLoad = f.session.getSnapshot(); + f.getHistory.mockResolvedValue([saved('clean', [ai()])]); + await f.session.load!(); + expect(f.session.getSnapshot().messages[0].citations).toBeUndefined(); + expect(beforeLoad.messages.at(-1)?.citations?.[0].title).toBe('Corrected'); + expect(f.stream).toHaveBeenCalledOnce(); + expect(f.getHistory).toHaveBeenCalledOnce(); + }); + + it('replaces repeated delta metadata, retains omitted interim data and commits canonical removal', async () => { + const f = fixture(); + const run = f.session.submit('Go'); + await f.started(); + const delta = (citations?: unknown, content = ''): StreamEvent => ({ + type: 'messages', + messageMetadata: {}, + messages: [{ ...ai(citations, content), type: 'AIMessageChunk' }], + }); + await f.emit(delta([citation()], 'A')); + const first = f.session.getSnapshot().messages.at(-1)!; + await f.emit(delta([citation()])); + expect(f.session.getSnapshot().messages.at(-1)).toBe(first); + await f.emit(delta(undefined, 'B')); + expect(f.session.getSnapshot().messages.at(-1)?.citations).toBe( + first.citations + ); + await f.emit(delta([citation('Next')])); + expect(f.session.getSnapshot().messages.at(-1)?.citations).toHaveLength(1); + await f.emit(frame(undefined, 'Final')); + expect(f.session.getSnapshot().messages.at(-1)?.citations?.[0].title).toBe( + 'Next' + ); + f.streams[0].finish(); + expect(await run).toBe('success'); + expect(f.session.getSnapshot().messages.at(-1)?.citations).toBeUndefined(); + }); + + it('captures history metadata before later transport getters mutate its source and preserves equal reload identity', async () => { + const f = fixture(); + const raw = [citation()]; + const checkpoint = saved('history', [], { + values: { + messages: [ai(raw)], + get application() { + raw[0].title = 'Mutated'; + raw[0].extra.tags.push('mutated'); + return true; + }, + }, + }); + f.getHistory.mockResolvedValueOnce([checkpoint]); + await f.session.load!(); + const before = f.session.getSnapshot(); + expect(before.messages[0].citations?.[0]).toMatchObject({ + title: 'Source', + extra: { tags: ['original'] }, + }); + f.getHistory.mockResolvedValue([ + saved('history', [ai([citation()])], { + values: { messages: [ai([citation()])], application: true }, + }), + ]); + await f.session.load!(); + expect(f.session.getSnapshot()).toBe(before); + expect(Object.isFrozen(raw[0])).toBe(false); + }); + + it.each(['history', 'stream'] as const)( + 'rejects stale citation ownership caused by a reentrant %s getter', + async (mode) => { + const f = fixture(); + let nested: Promise | undefined; + let fired = false; + const raw = [ + { + get title() { + if (!fired) { + fired = true; + nested = f.session.submit('Replacement'); + } + return 'Stale secret'; + }, + }, + ]; + if (mode === 'history') { + f.getHistory.mockResolvedValue([saved('stale', [ai(raw)])]); + await f.session.load!(); + await f.started(); + } else { + const run = f.session.submit('Original'); + await f.started(); + await f.emit(frame(raw)); + await run; + await f.started(1); + } + expect(fired).toBe(true); + expect( + f.session + .getSnapshot() + .messages.some((message) => + message.citations?.some((item) => item.title === 'Stale secret') + ) + ).toBe(false); + await f.session.stop(); + await nested; + } + ); + + it('retains stopped metadata and isolates stale attempts during explicit citation replacement', async () => { + const f = fixture(); + const oldRun = f.session.submit('Old'); + await f.started(); + await f.emit(frame([citation('Old')])); + await f.session.stop(); + expect(await oldRun).toBe('aborted'); + const stopped = f.session.getSnapshot(); + const nextRun = f.session.submit('New'); + await f.started(1); + f.streams[0].release(frame([citation('Stale')])); + await f.emit(frame([citation('New')]), 1); + f.streams[1].finish(); + expect(await nextRun).toBe('success'); + expect( + f.session + .getSnapshot() + .messages.find((message) => message.id === 'answer')?.citations?.[0] + .title + ).toBe('New'); + expect(stopped.messages.at(-1)?.citations?.[0].title).toBe('Old'); + expect(stopped.messages.at(-1)?.delivery).toMatchObject({ + phase: 'complete', + outcome: 'aborted', + }); + }); + + it('isolates equal message IDs in child projections and never executes child citations or tool metadata', async () => { + const f = fixture(); + const run = f.session.submit('Go'); + await f.started(); + await f.emit(frame([citation('Root')])); + await f.emit({ ...frame([citation('Child')]), namespace: ['child'] }); + await f.emit({ ...frame([citation('Sibling')]), namespace: ['sibling'] }); + const retained = f.session.getSnapshot(); + expect(retained.messages.at(-1)?.citations?.[0].title).toBe('Root'); + expect( + retained.subgraphs.map((child) => child.messages[0].citations?.[0].title) + ).toEqual(['Child', 'Sibling']); + await f.emit({ ...frame([], 'Changed'), namespace: ['child'] }); + expect(f.session.getSnapshot().subgraphs[1]).toBe(retained.subgraphs[1]); + expect(retained.subgraphs[0].messages[0].citations?.[0].title).toBe( + 'Child' + ); + f.streams[0].finish(); + await run; + expect(f.stream).toHaveBeenCalledOnce(); + expect(f.getHistory).not.toHaveBeenCalled(); + }); + + it.each( + (['history', 'stream'] as const).flatMap((mode) => [ + { + mode, + kind: 'throwing getter', + extra: () => ({ + get secret() { + throw new Error('PRIVATE_PROVIDER_DATA'); + }, + }), + }, + { mode, kind: 'SDK instance', extra: () => ({ instance: new Date() }) }, + { + mode, + kind: 'callback', + extra: () => ({ callback: () => 'PRIVATE_PROVIDER_DATA' }), + }, + { + mode, + kind: 'cycle', + extra: () => { + const value: Record = {}; + value['self'] = value; + return value; + }, + }, + ]) + )( + 'contains $kind citation extras on $mode without partial publication or raw error data', + async ({ mode, extra }) => { + const f = fixture(); + f.getHistory.mockResolvedValueOnce([ + saved('safe', [ai([citation('Safe')])]), + ]); + await f.session.load!(); + const raw = extra(); + if (mode === 'history') { + const before = f.session.getSnapshot(); + f.getHistory.mockResolvedValueOnce([ + saved('bad', [ai([{ extra: raw }])]), + ]); + await expect(f.session.load!()).rejects.not.toThrow( + 'PRIVATE_PROVIDER_DATA' + ); + expect(f.session.getSnapshot()).toBe(before); + } else { + const run = f.session.submit('Go'); + await f.started(); + await f.emit(frame([{ extra: raw }])); + await run; + expect(JSON.stringify(f.session.getSnapshot())).not.toContain( + 'PRIVATE_PROVIDER_DATA' + ); + expect( + f.session + .getSnapshot() + .messages.find((message) => message.id === 'answer')?.citations?.[0] + .title + ).toBe('Safe'); + } + } + ); +}); + +describe('citation checkpoint and reconnect evidence', () => { + it('rejects a Date-bearing exact checkpoint before normalization without admitting a fork or publishing metadata', async () => { + const f = checkpointFixture(); + f.source.values = { messages: [ai([{ publishedAt: new Date() }])] }; + const session = createSession({ + assistantId: 'agent', + threadId: 'thread', + transport: f.transport, + }); + cleanups.push(() => session.dispose()); + const before = session.getSnapshot(); + await expect(session.fork(position('a'), 'Rejected')).rejects.toThrow(); + expect(session.getSnapshot()).toBe(before); + expect(f.transport.getState).toHaveBeenCalledOnce(); + expect(f.transport.stream).not.toHaveBeenCalled(); + expect(f.transport.getHistory).not.toHaveBeenCalled(); + }); + + it('retains supported primitive metadata through strict checkpoint capture and preserves evidence comparison', () => { + const state = saved('a', [ + ai([{ publishedAt: '2026-09-24' }, { publishedAt: 0 }]), + ]); + const captured = captureCheckpointState(state, position('a')); + expect(captured.values).toEqual(state.values); + const candidate = captureCheckpointEvent(checkpointEvent(state), 'thread'); + expect(confirmCheckpoint(candidate, state, 'run-a').state.values).toEqual( + state.values + ); + expect(() => + confirmCheckpoint( + candidate, + saved('a', [ai([{ publishedAt: 1 }])]), + 'run-a' + ) + ).toThrow(); + const cycle: Record = {}; + cycle['self'] = cycle; + for (const unsupported of [{ publishedAt: new Date() }, { extra: cycle }]) { + const bad = saved('a', [ai([unsupported])]); + expect(() => captureCheckpointState(bad, position('a'))).toThrow( + /plain data/ + ); + expect(() => + captureCheckpointEvent(checkpointEvent(bad), 'thread') + ).toThrow(/plain data/); + expect(Object.isFrozen(unsupported)).toBe(false); + } + }); + + it('loads the exact branch with citations while historical tool evidence stays inert and request counts stay fixed', async () => { + const f = checkpointFixture(); + const historical = [ + { + ...ai([citation('Historical')]), + tool_calls: [{ id: 'old', name: 'work', args: {} }], + }, + { type: 'tool', id: 'result', tool_call_id: 'old', content: 'Saved' }, + ]; + f.source.values = { messages: historical }; + const handler = vi.fn(() => 'must not execute'); + const session = createSession({ + assistantId: 'agent', + threadId: 'thread', + transport: f.transport, + tools: { work: { description: 'Work', handler } }, + }); + cleanups.push(() => session.dispose()); + expect(await session.fork(position('a'), 'Continue')).toBe('success'); + const branch = f.states.get('result-1')!; + (branch.values as { messages: unknown[] }).messages[0] = { + ...ai([citation('Reloaded')]), + tool_calls: [{ id: 'old', name: 'work', args: {} }], + }; + await session.load!(); + expect(session.getSnapshot().messages[0].citations?.[0].title).toBe( + 'Reloaded' + ); + expect(f.transport.getState).toHaveBeenCalledTimes(3); + expect(vi.mocked(f.transport.getState!).mock.calls[2][1]).toEqual( + position('result-1') + ); + expect(f.transport.getHistory).not.toHaveBeenCalled(); + expect(f.transport.stream).toHaveBeenCalledOnce(); + expect(handler).not.toHaveBeenCalled(); + }); + + it('lets an explicit reconnect canonical correction remove earlier citation metadata without a second run', async () => { + const stream = vi.fn(async function* ( + _a, + _t, + _input, + _signal, + options + ) { + options?.onRunCreated?.({ run_id: 'run', thread_id: 'thread' }); + yield { + type: 'messages', + messageMetadata: {}, + sseId: '1', + messages: [ai([citation()], 'Partial')], + }; + }); + const joinStream = vi.fn>( + async function* () { + yield { ...frame(undefined, 'Final'), sseId: '2' }; + } + ); + const getRunStatus = vi.fn>( + async () => 'running' + ); + const session = createSession({ + assistantId: 'agent', + threadId: 'thread', + transport: { stream, joinStream, getRunStatus }, + }); + cleanups.push(() => session.dispose()); + await session.submit('Go'); + const interrupted = session.getSnapshot(); + expect(interrupted.messages.at(-1)?.citations?.[0].title).toBe('Source'); + getRunStatus.mockResolvedValue('success'); + expect(await session.reconnect()).toBe('success'); + expect(session.getSnapshot().messages.at(-1)).toMatchObject({ + content: 'Final', + delivery: { phase: 'complete', outcome: 'success' }, + }); + expect(session.getSnapshot().messages.at(-1)?.citations).toBeUndefined(); + expect(interrupted.messages.at(-1)?.citations?.[0].title).toBe('Source'); + expect(stream).toHaveBeenCalledOnce(); + expect(joinStream).toHaveBeenCalledOnce(); + }); +}); diff --git a/libs/langgraph/src/runtime/history-projection.spec.ts b/libs/langgraph/src/runtime/history-projection.spec.ts index ef9190dde..c0bbe83f3 100644 --- a/libs/langgraph/src/runtime/history-projection.spec.ts +++ b/libs/langgraph/src/runtime/history-projection.spec.ts @@ -49,6 +49,35 @@ function checkpoint( const initial = () => initialMessageState(); describe('pure authoritative history projection', () => { + it('owns citations, shares equal histories and equal metadata across text changes, and clears authoritative omissions', () => { + const citations = [{ title: 'Source', extra: { tags: ['original'] } }]; + const input = () => + checkpoint([{ ...ai('a'), additional_kwargs: { citations } }]); + const first = projectHistory(initial(), [input()]); + expect(first.messages[0].citations?.[0]).toEqual({ + id: 'c1', + index: 1, + title: 'Source', + extra: { tags: ['original'] }, + }); + expect(projectHistory(first, [input()])).toBe(first); + const changed = projectHistory(first, [ + checkpoint([{ ...ai('a', 'Changed'), additional_kwargs: { citations } }]), + ]); + expect(changed.messages[0].citations).toBe(first.messages[0].citations); + citations[0].title = 'Mutated'; + citations[0].extra.tags.push('mutated'); + citations.push({ title: 'Another', extra: { tags: [] } }); + expect(first.messages[0].citations).toHaveLength(1); + expect(first.messages[0].citations?.[0].title).toBe('Source'); + expect(first.messages[0].citations?.[0].extra).toEqual({ + tags: ['original'], + }); + expect( + projectHistory(first, [checkpoint([ai('a')])]).messages[0].citations + ).toBeUndefined(); + }); + it('uses an explicit empty values control as breakpoint evidence for the latest assistant', () => { const messages = [human('old-u'), ai('old-a'), human('new-u'), ai('new-a')]; const state = projectHistory(initial(), [ diff --git a/libs/langgraph/src/runtime/history-projection.ts b/libs/langgraph/src/runtime/history-projection.ts index cd1d6fe6b..ca3b26361 100644 --- a/libs/langgraph/src/runtime/history-projection.ts +++ b/libs/langgraph/src/runtime/history-projection.ts @@ -17,6 +17,7 @@ import { record, roleOf, textContent } from './wire-message'; import { projectHistoryInterrupts } from './interrupt-projection'; import type { LangGraphInterrupt } from './langgraph-snapshot'; import { observeInvocation, type ToolInvocation } from './tool-invocations'; +import { projectCitations } from './citation-projection'; export interface HistoryProjectionOptions { /** Omit for broad wire observation. A supplied catalog exposes only its @@ -162,6 +163,7 @@ export function projectHistory( id, role, content: textContent(message['content']), + citations: projectCitations(message), delivery: staticDelivery(id), ...(typeof message['name'] === 'string' ? { name: message['name'] } : {}), ...(typeof message['tool_call_id'] === 'string' @@ -222,7 +224,10 @@ export function projectHistory( } const ownedMessages = projectedMessages.map((message) => { const prior = previousMessages.get(message.id); - return ownMessage(prior && sameMessage(prior, message) ? prior : message); + return ownMessage( + prior && sameMessage(prior, message) ? prior : message, + prior + ); }); const messages = sameEntries(ownedMessages, previous.messages) ? previous.messages diff --git a/libs/langgraph/src/runtime/message-reducer.spec.ts b/libs/langgraph/src/runtime/message-reducer.spec.ts index fba7a9879..47220a861 100644 --- a/libs/langgraph/src/runtime/message-reducer.spec.ts +++ b/libs/langgraph/src/runtime/message-reducer.spec.ts @@ -20,6 +20,114 @@ function message( } describe('pure text and tool transitions', () => { + it('replaces citation lists, preserves omitted interim metadata and shares equal metadata across text changes', () => { + const citations = [{ id: 'c', index: 1, extra: { nested: ['original'] } }]; + const first = reduceMessages(initialMessageState(), { + type: 'message', + mode: 'snapshot', + message: { ...message('m', 'A'), citations }, + }); + const owned = first.messages[0].citations; + citations[0].extra.nested.push('mutated'); + expect(owned?.[0].extra).toEqual({ nested: ['original'] }); + const delta = reduceMessages(first, { + type: 'message', + mode: 'delta', + message: message('m', 'B'), + }); + expect(delta.messages[0].content).toBe('AB'); + expect(delta.messages[0].citations).toBe(owned); + const equal = reduceMessages(delta, { + type: 'message', + mode: 'snapshot', + message: { + ...message('m', 'ABC'), + citations: [{ id: 'c', index: 1, extra: { nested: ['original'] } }], + }, + }); + expect(equal.messages[0].citations).toBe(owned); + const updatedTitle = reduceMessages(equal, { + type: 'message', + mode: 'snapshot', + message: { + ...message('m', 'ABC'), + citations: [ + { + id: 'c', + index: 1, + title: 'Updated', + extra: { nested: ['original'] }, + }, + ], + }, + }); + expect(updatedTitle.messages[0].citations).not.toBe(owned); + expect(updatedTitle.messages[0].citations?.[0].extra).toBe(owned?.[0].extra); + const changed = reduceMessages(equal, { + type: 'message', + mode: 'delta', + message: { ...message('m', ''), citations: [{ id: 'next', index: 2 }] }, + }); + expect(changed.messages[0].citations).toEqual([{ id: 'next', index: 2 }]); + expect( + reduceMessages(changed, { + type: 'message', + mode: 'delta', + message: { ...message('m', ''), citations: [{ id: 'next', index: 2 }] }, + }) + ).toBe(changed); + expect( + reduceMessages(changed, { + type: 'message', + mode: 'snapshot', + message: { ...message('m', ''), citations: [] }, + }).messages[0].citations + ).toEqual([]); + }); + + it.each([undefined, []])( + 'canonical citation removal %j bars late resurrection while ordered corrections remain authoritative', + (removed) => { + const withCitation = { + ...message('m', 'Text'), + citations: [{ id: 'c', index: 1 }], + }; + const first = reduceMessages(initialMessageState(), { + type: 'message', + mode: 'snapshot', + message: withCitation, + }); + const canonical = reduceMessages(first, { + type: 'message', + mode: 'canonical', + message: { ...message('m', 'Final'), citations: removed }, + }); + expect(canonical.messages[0].citations).toEqual(removed); + for (const mode of ['delta', 'snapshot'] as const) { + const late = reduceMessages(canonical, { + type: 'message', + mode, + message: withCitation, + }); + expect(late.messages[0].content).toBe('Final'); + expect(late.messages[0].citations).toBe(canonical.messages[0].citations); + } + expect( + reduceMessages(canonical, { + type: 'message', + mode: 'canonical', + message: withCitation, + }).messages[0].citations + ).toEqual(withCitation.citations); + const nextGeneration = reduceMessages(first, { + type: 'message', + mode: 'snapshot', + message: { ...message('m', 'New'), delivery: streamingDelivery('run-2') }, + }); + expect(nextGeneration.messages[0].citations).toBeUndefined(); + } + ); + it.each([ { label: 'empty to holes', from: 0, to: 2, populated: false }, { label: 'shorter trailing holes', from: 3, to: 1, populated: true }, diff --git a/libs/langgraph/src/runtime/message-reducer.ts b/libs/langgraph/src/runtime/message-reducer.ts index 8527a0a4c..60182270a 100644 --- a/libs/langgraph/src/runtime/message-reducer.ts +++ b/libs/langgraph/src/runtime/message-reducer.ts @@ -194,6 +194,12 @@ export function reduceMessages( ...incoming, id, content, + citations: + sameGeneration && event.mode !== 'canonical' + ? canonical + ? previous.citations + : incoming.citations ?? previous.citations + : incoming.citations, // A late snapshot cannot reopen a finalized generation. delivery: sameGeneration && previous.delivery.phase === 'complete' @@ -207,7 +213,7 @@ export function reduceMessages( if (!previous || !sameMessage(previous, candidate)) { const next = [...messages]; if (index < 0) next.push(ownMessage(candidate)); - else next[index] = ownMessage(candidate); + else next[index] = ownMessage(candidate, previous); messages = Object.freeze(next); } const aliases = diff --git a/libs/langgraph/src/runtime/ownership.ts b/libs/langgraph/src/runtime/ownership.ts index 69004c694..b746d01ef 100644 --- a/libs/langgraph/src/runtime/ownership.ts +++ b/libs/langgraph/src/runtime/ownership.ts @@ -170,6 +170,10 @@ export function sameMessage(a: Message, b: Message): boolean { (a.id === b.id && a.role === b.role && a.content === b.content && + sameOwnedValue( + a.citations as unknown as PlainValue, + b.citations as unknown as PlainValue + ) && a.name === b.name && a.toolCallId === b.toolCallId && sameOwnedValue(a.toolCallIds, b.toolCallIds) && @@ -181,12 +185,17 @@ export function sameMessage(a: Message, b: Message): boolean { ); } -export function ownMessage(message: Message): Message { - if (owned.has(message)) return message; +export function ownMessage(message: Message, previous?: Message): Message { + const citations = ownValueWithSharing( + message.citations as unknown as PlainValue, + previous?.citations as unknown as PlainValue + ) as unknown as Message['citations']; + if (owned.has(message) && citations === message.citations) return message; return freeze({ id: message.id, role: message.role, content: message.content, + citations, delivery: message.delivery.phase === 'complete' ? completeDelivery( @@ -233,7 +242,7 @@ export function ownToolCall(call: ToolCall): ToolCall { function ownArray( values: readonly T[], previous: readonly T[] | undefined, - project: (value: T) => T, + project: (value: T, previous?: T) => T, equal: (a: T, b: T) => boolean ): readonly T[] { if (values === previous) return values; @@ -242,7 +251,10 @@ function ownArray( const next = values.map((value, index) => { // Ownership permits reuse, not a change notification. A distinct owned // array may still equal the current one (including queued publications). - const projected = alreadyOwned ? value : project(value); + const projected = + alreadyOwned && previous?.[index] === undefined + ? value + : project(value, previous?.[index]); return previous?.[index] !== undefined && equal(projected, previous[index]) ? previous[index] : projected; diff --git a/libs/langgraph/src/runtime/publication.spec.ts b/libs/langgraph/src/runtime/publication.spec.ts index 7cea0d4a2..659acdedb 100644 --- a/libs/langgraph/src/runtime/publication.spec.ts +++ b/libs/langgraph/src/runtime/publication.spec.ts @@ -33,6 +33,60 @@ function snapshot(content = 'one'): LangGraphSnapshot { } describe('private snapshot publication', () => { + it('captures citation ingress during synchronous reentrant publication and suppresses an equal queued value', () => { + const publication = createPublication(snapshot()); + const citations = [ + { id: 'c', index: 1, title: 'Original', extra: { tags: ['original'] } }, + ]; + const input = (): LangGraphSnapshot => ({ + ...snapshot('two'), + messages: [{ ...snapshot('two').messages[0], citations }], + }); + const seen: LangGraphSnapshot[] = []; + publication.subscribe(() => { + seen.push(publication.getSnapshot()); + if (seen.length === 1) { + publication.publish(input()); + publication.publish(input()); + citations[0].title = 'Mutated'; + citations[0].extra.tags.push('mutated'); + citations.push({ + id: 'other', + index: 2, + title: 'Other', + extra: { tags: [] }, + }); + } + }); + publication.publish(snapshot('two')); + expect(seen).toHaveLength(2); + expect(seen[1].messages[0].citations).toEqual([ + { id: 'c', index: 1, title: 'Original', extra: { tags: ['original'] } }, + ]); + expect(Object.isFrozen(citations[0])).toBe(false); + const prior = publication.getSnapshot(); + publication.publish({ + ...prior, + messages: [ + { + ...prior.messages[0], + content: 'three', + citations: [ + { + id: 'c', + index: 1, + title: 'Original', + extra: { tags: ['original'] }, + }, + ], + }, + ], + }); + expect(publication.getSnapshot().messages[0].citations).toBe( + prior.messages[0].citations + ); + }); + it.each([ { label: 'empty to holes', from: 0, to: 2, populated: false }, { label: 'longer trailing holes', from: 1, to: 3, populated: true }, diff --git a/libs/langgraph/src/runtime/stream-projection.spec.ts b/libs/langgraph/src/runtime/stream-projection.spec.ts index 5af7d1096..5c53f85cf 100644 --- a/libs/langgraph/src/runtime/stream-projection.spec.ts +++ b/libs/langgraph/src/runtime/stream-projection.spec.ts @@ -24,7 +24,7 @@ const projection = (): StreamProjection => ({ paused: false, canonical: [], }); -function start(messages = [assistant('step', ['c1'])]) { +function start(messages: unknown[] = [assistant('step', ['c1'])]) { return projectStream(initialMessageState(), projection(), { type: 'values', data: { messages: [user, ...messages] }, @@ -32,6 +32,36 @@ function start(messages = [assistant('step', ['c1'])]) { } describe('authoritative tool ownership', () => { + it('captures terminal citations before finalization and clears omitted canonical metadata', () => { + const citations = [ + { id: 'source', title: 'Original', extra: { tags: ['original'] } }, + ]; + const before = start([ + { ...assistant('step'), additional_kwargs: { citations } }, + ]); + expect(before.state.messages[1].citations?.[0].title).toBe('Original'); + citations[0].title = 'Mutated'; + citations[0].extra.tags.push('mutated'); + const final = finalizeProjection(before.state, before.projection); + expect(final.messages[1].citations?.[0]).toEqual({ + id: 'source', + index: 1, + title: 'Original', + extra: { tags: ['original'] }, + }); + const correction = projectStream(final, before.projection, { + type: 'values', + data: { messages: [user, assistant('step', undefined, 'Corrected')] }, + }); + expect(correction.state.messages[1].citations).toBe( + final.messages[1].citations + ); + expect( + finalizeProjection(correction.state, correction.projection).messages[1] + .citations + ).toBeUndefined(); + }); + it('leaves pause classification to the aggregate interrupt projector without reading control getters', () => { const before = start(); const next = projectStream(before.state, before.projection, { diff --git a/libs/langgraph/src/runtime/stream-projection.ts b/libs/langgraph/src/runtime/stream-projection.ts index 57b57b173..b2a97611d 100644 --- a/libs/langgraph/src/runtime/stream-projection.ts +++ b/libs/langgraph/src/runtime/stream-projection.ts @@ -13,6 +13,7 @@ import { import type { StreamEvent } from './transport.types'; import { ownMessage, ownToolCall } from './ownership'; import { record, roleOf, textContent } from './wire-message'; +import { projectCitations } from './citation-projection'; export { record } from './wire-message'; @@ -170,6 +171,7 @@ export function projectStream( id: wireId, role, content, + citations: projectCitations(raw), delivery: !current && previous ? previous.delivery diff --git a/scripts/react-parity/baseline.json b/scripts/react-parity/baseline.json index a53fa14ca..027950de4 100644 --- a/scripts/react-parity/baseline.json +++ b/scripts/react-parity/baseline.json @@ -1,21 +1,16 @@ { "schemaVersion": 1, - "baselineHead": "f8de78635270a3158de156aa060dce2bdd875fd6", + "baselineHead": "4d47995333278f9a19491925bc555d30f8490956", "sourceState": { "modified": [ - "libs/langgraph/src/lib/transport/fetch-stream.transport.ts", - "libs/langgraph/src/runtime/create-session.ts", - "libs/langgraph/src/runtime/stream-projection.ts", - "libs/langgraph/src/runtime/tool-persistence.ts", - "libs/langgraph/src/runtime/transport.types.ts" + "libs/langgraph/src/runtime/README.md", + "libs/langgraph/src/runtime/history-projection.ts", + "libs/langgraph/src/runtime/message-reducer.ts", + "libs/langgraph/src/runtime/ownership.ts", + "libs/langgraph/src/runtime/stream-projection.ts" ], "untracked": [ - "libs/langgraph/src/lib/transport/checkpoint-position.ts", - "libs/langgraph/src/runtime/README.md", - "libs/langgraph/src/runtime/checkpoint-authority.ts", - "libs/langgraph/src/runtime/checkpoint-state.ts", - "libs/langgraph/src/runtime/checkpoint-tool-evidence.ts", - "libs/langgraph/src/runtime/testing/checkpoint-fixture.ts" + "libs/langgraph/src/runtime/citation-projection.ts" ] }, "scope": { @@ -605,7 +600,7 @@ "id": "asset:libs/langgraph/src/runtime/README.md", "kind": "asset", "path": "libs/langgraph/src/runtime/README.md", - "sha256": "45d9aad401d27f6e53dbe2aa5f495e4665a89a69bdf4208525b7e7ca619c818b" + "sha256": "289b901753126e898d40f9a84c093afeb630a7f78692ef6bf867e955e1c77f48" }, { "id": "asset:libs/langgraph/test/fixtures/streaming-reasoning-puzzle.json", @@ -12505,6 +12500,12 @@ "path": "libs/langgraph/src/runtime/checkpoint-tool-evidence.ts", "sha256": "833f052dee8017dd6bdf1968457d85e0efa7a204dbc0d75b490b5ec82e6d2db2" }, + { + "id": "source:libs/langgraph/src/runtime/citation-projection.ts", + "kind": "source", + "path": "libs/langgraph/src/runtime/citation-projection.ts", + "sha256": "52d99b652713e964242f38ec9e81f74695ecac938d42f10255b70d3c7eab8201" + }, { "id": "source:libs/langgraph/src/runtime/create-session.ts", "kind": "source", @@ -12521,7 +12522,7 @@ "id": "source:libs/langgraph/src/runtime/history-projection.ts", "kind": "source", "path": "libs/langgraph/src/runtime/history-projection.ts", - "sha256": "9a12195d4b74caab2776cf5cbb4da72949f60db5abab8ed762c75608dbbe5239" + "sha256": "c9b18e186261d16d6b363ddc906793acbf3295e4a5c45e93bbf2d555bd2b2e92" }, { "id": "source:libs/langgraph/src/runtime/interrupt-projection.ts", @@ -12539,7 +12540,7 @@ "id": "source:libs/langgraph/src/runtime/message-reducer.ts", "kind": "source", "path": "libs/langgraph/src/runtime/message-reducer.ts", - "sha256": "b4842acc1b3c957b7606fcdf3c9c42ab4f569bc858abab46f722e8d2a98b55b2" + "sha256": "5758caaae8c31361a9361b81a7995239c10a48bf670a135a9fe6a327da2e6531" }, { "id": "source:libs/langgraph/src/runtime/operation-errors.ts", @@ -12551,7 +12552,7 @@ "id": "source:libs/langgraph/src/runtime/ownership.ts", "kind": "source", "path": "libs/langgraph/src/runtime/ownership.ts", - "sha256": "f7f3d00e1f1b7bd196a09dc907d06618a85c8513b69a5ffa229d7567451b7876" + "sha256": "0cfc4dabfedbc3aa9406b943257ab3df3bae240276ccaec6d50f1e51f11e5c69" }, { "id": "source:libs/langgraph/src/runtime/publication.ts", @@ -12575,7 +12576,7 @@ "id": "source:libs/langgraph/src/runtime/stream-projection.ts", "kind": "source", "path": "libs/langgraph/src/runtime/stream-projection.ts", - "sha256": "85ac624f5921d63ba5bff724c14612279934447c5178654b6d5db2d257f99635" + "sha256": "076cf6eb83b633d67a42118d7f17931b3f5010ca9f00b7c15e97bf5214f42240" }, { "id": "source:libs/langgraph/src/runtime/subgraph-projection.ts", diff --git a/scripts/react-parity/dispositions.json b/scripts/react-parity/dispositions.json index 4741082a2..bf2c4f188 100644 --- a/scripts/react-parity/dispositions.json +++ b/scripts/react-parity/dispositions.json @@ -11215,7 +11215,8 @@ "T08" ], "treatment": "shared", - "status": "planned" + "status": "planned", + "note": "Private runtime citation projection now stages the aliases with owned plain metadata. Consolidation/removal at backend cutover remains planned; legacy arbitrary-extra and Date semantics are not claimed equivalent." }, { "id": "source:libs/langgraph/src/lib/internals/stream-manager.bridge.ts", @@ -11346,6 +11347,14 @@ "treatment": "shared", "status": "planned" }, + { + "id": "source:libs/langgraph/src/runtime/citation-projection.ts", + "taskIds": ["T08", "T09", "T15"], + "treatment": "internal", + "status": "in-progress", + "reason": "Private owned citation metadata projection, verified through both installed native bindings; not a public backend cutover.", + "note": "Normalizes existing aliases with absent/empty distinction, finite primitive timestamps and deeply owned plain extras. Rich blocks, reasoning/timing, citation UI and full T09/T15 parity remain open." + }, { "id": "source:libs/langgraph/src/runtime/run-options.ts", "taskIds": [ @@ -11478,7 +11487,7 @@ "treatment": "internal", "status": "in-progress", "reason": "Private staged LangGraph runtime subset; fixture-only and not a neutral public LangGraph root or tarball.", - "note": "Owns immutable interrupt batches alongside values/messages, preserving unchanged nested identity. Broader execution/publication migration remains open." + "note": "Owns immutable interrupt batches and citation arrays/entries/plain extras alongside values/messages, preserving equal metadata identity across changed text and queued publications. Broader execution/publication migration remains open." }, { "id": "source:libs/langgraph/src/runtime/publication.ts", diff --git a/scripts/react-parity/runtime-consumer.mjs b/scripts/react-parity/runtime-consumer.mjs index efcf38653..5fe706d36 100644 --- a/scripts/react-parity/runtime-consumer.mjs +++ b/scripts/react-parity/runtime-consumer.mjs @@ -45,7 +45,7 @@ const savedHistory = [{ { id: 'saved-count', name: 'count', args: { values: ['saved'] }, type: 'tool_call' }, ] }, { id: 'saved-result', type: 'tool', tool_call_id: 'saved-weather', content: 'Raw historical weather result' }, - { id: 'saved-final', type: 'ai', content: [{ type: 'text', text: 'Saved final answer' }] }, + { id: 'saved-final', type: 'ai', content: [{ type: 'text', text: 'Saved final answer' }], additional_kwargs: { sources: [{ refId: 'saved-source', name: 'Saved reference', publishedAt: '2026-09-24', extra: { provider: { labels: ['history'] } } }] } }, ] }, next: ['review', 'confirmation'], tasks: savedInterrupts.map((interrupt, index) => ({ id: `saved-task-${index}`, name: index === 0 ? 'review' : 'confirmation', error: null, checkpoint: null, state: null, interrupts: [interrupt] })), @@ -73,7 +73,7 @@ export function runtimeResponse(body) { assert.deepEqual(body, { assistant_id: 'fixture-assistant', input: null, command: { resume: response }, stream_mode: ['values', 'messages-tuple', 'updates', 'custom'], stream_subgraphs: true, ...resumableRunFields, ...configuredRunFields }, 'exact resume run fields'); if (Object.hasOwn(response ?? {}, 'live-approval')) { assert.deepEqual(response, { 'live-approval': 'yes', 'live-confirmation': false }, 'exact initial response map'); - return sse('values', { stage: 'final-approval', messages: [{ type: 'ai', id: 'resume-answer', content: 'One final approval' }] }) + return sse('values', { stage: 'final-approval', messages: [{ type: 'ai', id: 'resume-answer', content: 'One final approval', additional_kwargs: { citations: [{ id: 'approval-source', title: 'Approval reference', publishedAt: 0 }] } }] }) + sse('values|review:child', { stage: 'child-final-approval', messages: [{ type: 'ai', id: 'child-review', content: 'Child final approval' }] }) + sse('updates|review:child', { __interrupt__: [{ id: 'child-final', value: 'Child confirmation' }] }) + sse('updates', { __interrupt__: [{ id: 'final-approval', value: { question: 'Confirm final action?' } }] }); @@ -529,6 +529,7 @@ export async function runRuntimeScenarios(directory, kind) { await expectInterrupts([]); assert.deepEqual(await children(), []); await expect(page.getByTestId('history')).toHaveText('unobserved'); + await expect(page.getByTestId('citations')).toHaveText(''); assert.equal(server.requests.length, 0, 'mount/observation performs no I/O'); assert.equal(server.historyRequests.length, 0, 'mount/observation performs no history reads'); completed.push('inert mount'); @@ -542,6 +543,7 @@ export async function runRuntimeScenarios(directory, kind) { await expectValues({ stage: 'saved', profile: { name: 'Saved user' } }); await expectInterrupts(savedInterrupts); await expect(page.getByTestId('text')).toHaveText('Saved tool request\nSaved final answer'); + await expect(page.getByTestId('citations')).toHaveText('saved-source: Saved reference'); await expect(page.getByTestId('transcript')).toContainText('Saved question'); await expect(page.getByTestId('transcript')).toContainText('Raw historical weather result'); await expect(page.getByTestId('delivery')).toHaveText('complete:paused'); @@ -564,6 +566,7 @@ export async function runRuntimeScenarios(directory, kind) { assert.equal(server.historyRequests.length, 2); assert.equal(server.requests.length, 0); completed.push('equal history refresh'); + await expect(page.getByTestId('citations')).toHaveText('saved-source: Saved reference'); await page.getByRole('button', { name: 'Load', exact: true }).click(); await expect(page.getByTestId('loads-finished')).toHaveText('3'); @@ -573,6 +576,7 @@ export async function runRuntimeScenarios(directory, kind) { await expectInterrupts([]); await expect(page.getByTestId('text')).toHaveText(''); await expect(page.getByTestId('transcript')).toHaveText(''); + await expect(page.getByTestId('citations')).toHaveText(''); await expect(page.getByTestId('tool')).toHaveText('[]'); await expect(page.getByTestId('handler-calls')).toHaveText('0'); assert.equal(server.historyRequests.length, 3); @@ -653,6 +657,7 @@ export async function runRuntimeScenarios(directory, kind) { await expect(page.getByTestId('resume-outcome')).toHaveText('paused'); await expect(page.getByTestId('delivery')).toHaveText('complete:paused'); await expect(page.getByTestId('text')).toContainText('One final approval'); + await expect(page.getByTestId('citations')).toHaveText('approval-source: Approval reference'); await expectInterrupts([{ id: 'final-approval', value: { question: 'Confirm final action?' } }]); await expectValues({ stage: 'final-approval' }); await expectChild(['review:child'], 'Child final approval', 'complete:paused', { stage: 'child-final-approval' }, [{ id: 'child-final', value: 'Child confirmation' }]); @@ -668,6 +673,7 @@ export async function runRuntimeScenarios(directory, kind) { await expect(page.getByTestId('delivery')).toHaveText('complete:success'); await expect(page.getByTestId('text')).toContainText('Approvals complete'); await expect(page.getByTestId('text')).not.toContainText('One final approval'); + await expect(page.getByTestId('citations')).toHaveText(''); await expectInterrupts([]); await expectValues({ stage: 'approved' }); await expectChild(['review:child'], 'Child approved', 'complete:success', { stage: 'child-approved' });