From 45232e8350ab18632544c5448e749713f39265c9 Mon Sep 17 00:00:00 2001 From: Brian Love Date: Thu, 24 Sep 2026 10:43:38 -0700 Subject: [PATCH 1/2] feat(langgraph): retain completed checkpoint fork authority --- .../content/docs/langgraph/api/api-docs.json | 72 ++- .../angular/e2e/checkpoint-session.spec.ts | 213 +++++++++ .../angular/e2e/checkpoint-session.spec.ts | 196 ++++++++ .../interrupts/angular/e2e/tsconfig.json | 6 + .../src/lib/transport/checkpoint-position.ts | 59 +++ .../lib/transport/fetch-stream.transport.ts | 41 +- libs/langgraph/src/runtime/README.md | 65 +++ .../src/runtime/checkpoint-admission.spec.ts | 268 +++++++++++ .../src/runtime/checkpoint-authority.spec.ts | 198 ++++++++ .../src/runtime/checkpoint-authority.ts | 226 +++++++++ .../src/runtime/checkpoint-execution.spec.ts | 348 ++++++++++++++ .../runtime/checkpoint-execution.type-test.ts | 21 + .../runtime/checkpoint-persistence.spec.ts | 154 ++++++ .../src/runtime/checkpoint-recovery.spec.ts | 93 ++++ .../src/runtime/checkpoint-state.spec.ts | 34 ++ .../langgraph/src/runtime/checkpoint-state.ts | 35 ++ .../runtime/checkpoint-tool-evidence.spec.ts | 369 +++++++++++++++ .../src/runtime/checkpoint-tool-evidence.ts | 80 ++++ .../src/runtime/checkpoint-transport.spec.ts | 153 ++++++ libs/langgraph/src/runtime/create-session.ts | 444 +++++++++++++++++- .../src/runtime/stream-projection.ts | 43 +- .../src/runtime/testing/checkpoint-fixture.ts | 87 ++++ .../langgraph/src/runtime/tool-persistence.ts | 9 +- libs/langgraph/src/runtime/transport.types.ts | 22 +- scripts/react-parity/baseline.json | 74 ++- scripts/react-parity/dispositions.json | 52 +- 26 files changed, 3291 insertions(+), 71 deletions(-) create mode 100644 cockpit/langgraph/client-tools/angular/e2e/checkpoint-session.spec.ts create mode 100644 cockpit/langgraph/interrupts/angular/e2e/checkpoint-session.spec.ts create mode 100644 libs/langgraph/src/lib/transport/checkpoint-position.ts create mode 100644 libs/langgraph/src/runtime/README.md create mode 100644 libs/langgraph/src/runtime/checkpoint-admission.spec.ts create mode 100644 libs/langgraph/src/runtime/checkpoint-authority.spec.ts create mode 100644 libs/langgraph/src/runtime/checkpoint-authority.ts create mode 100644 libs/langgraph/src/runtime/checkpoint-execution.spec.ts create mode 100644 libs/langgraph/src/runtime/checkpoint-execution.type-test.ts create mode 100644 libs/langgraph/src/runtime/checkpoint-persistence.spec.ts create mode 100644 libs/langgraph/src/runtime/checkpoint-recovery.spec.ts create mode 100644 libs/langgraph/src/runtime/checkpoint-state.spec.ts create mode 100644 libs/langgraph/src/runtime/checkpoint-state.ts create mode 100644 libs/langgraph/src/runtime/checkpoint-tool-evidence.spec.ts create mode 100644 libs/langgraph/src/runtime/checkpoint-tool-evidence.ts create mode 100644 libs/langgraph/src/runtime/checkpoint-transport.spec.ts create mode 100644 libs/langgraph/src/runtime/testing/checkpoint-fixture.ts diff --git a/apps/website/content/docs/langgraph/api/api-docs.json b/apps/website/content/docs/langgraph/api/api-docs.json index cf5b0820a..0e9760335 100644 --- a/apps/website/content/docs/langgraph/api/api-docs.json +++ b/apps/website/content/docs/langgraph/api/api-docs.json @@ -293,9 +293,34 @@ } ] }, + { + "name": "getState", + "signature": "getState(threadId: string, checkpoint: OwnedCheckpointPosition, signal: AbortSignal): Promise>", + "description": "Read one exact saved checkpoint with the command's cancellation signal.", + "params": [ + { + "name": "threadId", + "type": "string", + "description": "", + "optional": false + }, + { + "name": "checkpoint", + "type": "OwnedCheckpointPosition", + "description": "", + "optional": false + }, + { + "name": "signal", + "type": "AbortSignal", + "description": "", + "optional": false + } + ] + }, { "name": "joinStream", - "signature": "joinStream(threadId: string, runId: string, lastEventId: string | undefined, signal: AbortSignal): AsyncIterable", + "signature": "joinStream(threadId: string, runId: string, lastEventId: string | undefined, signal: AbortSignal, options: object): AsyncIterable", "description": "Join an already-started run without creating a new thread.", "params": [ { @@ -321,6 +346,12 @@ "type": "AbortSignal", "description": "", "optional": false + }, + { + "name": "options", + "type": "object", + "description": "", + "optional": true } ] }, @@ -363,8 +394,8 @@ }, { "name": "updateState", - "signature": "updateState(threadId: string, values: Record, signal: AbortSignal, options: object): Promise", - "description": "Update server-side thread state, e.g. to remove messages for regenerate rollback.", + "signature": "updateState(threadId: string, values: Record, signal: AbortSignal, options: object): Promise", + "description": "Update state once and return only usable root routing, when supplied.", "params": [ { "name": "threadId", @@ -1273,9 +1304,34 @@ } ] }, + { + "name": "getState", + "signature": "getState(threadId: string, checkpoint: OwnedCheckpointPosition, signal: AbortSignal): Promise>", + "description": "Exact saved root read. Branch effects must not infer authority from latest.", + "params": [ + { + "name": "threadId", + "type": "string", + "description": "", + "optional": false + }, + { + "name": "checkpoint", + "type": "OwnedCheckpointPosition", + "description": "", + "optional": false + }, + { + "name": "signal", + "type": "AbortSignal", + "description": "", + "optional": false + } + ] + }, { "name": "joinStream", - "signature": "joinStream(threadId: string, runId: string, lastEventId: string | undefined, signal: AbortSignal): AsyncIterable", + "signature": "joinStream(threadId: string, runId: string, lastEventId: string | undefined, signal: AbortSignal, options: object): AsyncIterable", "description": "Optional: join an already-started run without creating a new one.", "params": [ { @@ -1301,6 +1357,12 @@ "type": "AbortSignal", "description": "", "optional": false + }, + { + "name": "options", + "type": "object", + "description": "", + "optional": true } ] }, @@ -1343,7 +1405,7 @@ }, { "name": "updateState", - "signature": "updateState(threadId: string, values: Record, signal: AbortSignal, options: object): Promise", + "signature": "updateState(threadId: string, values: Record, signal: AbortSignal, options: object): Promise", "description": "Optional: update server-side thread state (e.g. to emit RemoveMessage\nentries for regenerate rollback). Forwards to the LangGraph\n`threads.updateState` API.\n\n`options.asNode` corresponds to LangGraph's `as_node` parameter — the\nserver treats the update as if that node had just produced the values,\nwhich determines what the next pull resumes. `regenerate()` passes\n`asNode: '__start__'` so the next `submit(null)` resumes at the entry\nnode and re-runs `generate` against the rolled-back state.", "params": [ { diff --git a/cockpit/langgraph/client-tools/angular/e2e/checkpoint-session.spec.ts b/cockpit/langgraph/client-tools/angular/e2e/checkpoint-session.spec.ts new file mode 100644 index 000000000..d45cbb5ae --- /dev/null +++ b/cockpit/langgraph/client-tools/angular/e2e/checkpoint-session.spec.ts @@ -0,0 +1,213 @@ +import { expect, test } from '@playwright/test'; +import { createSession } from '../../../../../libs/langgraph/src/runtime/create-session'; +import { FetchStreamTransport } from '../../../../../libs/langgraph/src/lib/transport/fetch-stream.transport'; +import { + backendUrl, client, modelJournal, position as wirePosition, prompt, run, +} from './checkpoint-protocol.helpers'; + +const sourcePrompt = 'Checkpoint branch original second turn'; +const competitorPrompt = 'Checkpoint branch competing later turn'; +const forkPrompt = 'Checkpoint branch fork from first turn'; +const advancePrompt = 'Checkpoint continuity advance the competing branch'; +const followUpPrompt = 'Checkpoint continuity continue the completed fork'; + +function position(config: Parameters[0]) { + const owned = wirePosition(config); + // The server omits an empty map on saved states but can echo {} in frames + // when the request supplied it. Both represent the same root routing. + return { ...owned, checkpoint_map: owned.checkpoint_map ?? {} }; +} + +// Observe the real adapter without supplying synthetic events, state or receipts. +class ObservedTransport extends FetchStreamTransport { + readonly creations: Parameters[] = []; + readonly checkpoints: ReturnType[] = []; + readonly writes: { args: Parameters; result: unknown }[] = []; + + override async *stream(...args: Parameters) { + this.creations.push(args); + for await (const event of super.stream(...args)) { + if (event.type === 'checkpoints') { + const data = event.data as { config: { configurable?: Record } }; + this.checkpoints.push(position(data.config.configurable)); + } + yield event; + } + } + + override async updateState(...args: Parameters) { + const result = await super.updateState(...args); + this.writes.push({ args, result }); + return result; + } + + finalPosition() { + const checkpoint = this.checkpoints.at(-1); + if (!checkpoint) throw new Error('Session must obtain a real root checkpoint'); + return checkpoint; + } +} + +function deferred() { + let resolve!: () => void; + const promise = new Promise((done) => { resolve = done; }); + return { promise, resolve }; +} + +async function descendsFrom( + api: ReturnType, threadId: string, + child: Awaited>, parent: ReturnType, + excluded: readonly string[], +) { + let ancestor = child; + const visited = new Set(); + for (let depth = 0; depth < 8; depth++) { + const current = position(ancestor.checkpoint); + expect(current.thread_id).toBe(threadId); + expect(current.checkpoint_ns).toBe(''); + expect(excluded).not.toContain(current.checkpoint_id); + expect(visited.has(current.checkpoint_id)).toBe(false); + visited.add(current.checkpoint_id); + if (current.checkpoint_id === parent.checkpoint_id) break; + ancestor = await api.threads.getState(threadId, position(ancestor.parent_checkpoint ?? undefined)); + } + expect(position(ancestor.checkpoint)).toEqual(parent); +} + +for (const followUp of [true, false]) { + test(`checkpoint session: retains its fork through reads, submissions and ${followUp ? 'tool continuation' : 'terminal tool persistence'}`, async () => { + const api = client(); + const { thread_id: threadId } = await api.threads.create(); + const started = deferred(); + const release = deferred(); + let handlers = 0; + const transport = new ObservedTransport(backendUrl(), undefined, { maxRetries: 0 }); + const session = createSession({ + assistantId: 'client-tools', threadId, transport, + tools: { + get_weather: { + description: 'Get weather for a location', + parameters: { type: 'object', properties: { location: { type: 'string' } } }, + followUp, + handler: async (args: { location: string }) => { + handlers++; + expect(args).toEqual({ location: 'Paris' }); + started.resolve(); + await release.promise; + return '68°F'; + }, + }, + }, + }); + try { + const source = await run(api, threadId, undefined, { + messages: [{ id: 'source-question', type: 'human', content: sourcePrompt }], + }); + const sourcePosition = position(source.checkpoint); + const competitor = await run(api, threadId, sourcePosition, { + messages: [{ id: 'competitor-question', type: 'human', content: competitorPrompt }], + }); + await session.load?.(); + expect(session.getSnapshot().messages.at(-1)?.content).toBe('Competing later answer.'); + const beforeFork = (await modelJournal(forkPrompt)).length; + expect(await session.fork(sourcePosition, forkPrompt)).toBe('success'); + expect(session.getSnapshot().messages.map((message) => message.content)).toEqual([ + ...source.values.messages.map((message) => message.content), forkPrompt, 'Fork answer from first turn.', + ]); + const fork = transport.finalPosition(); + const savedFork = await api.threads.getState(threadId, fork); + expect(savedFork.values.messages).toEqual([ + ...source.values.messages, + expect.objectContaining({ type: 'human', content: forkPrompt }), + expect.objectContaining({ type: 'ai', content: 'Fork answer from first turn.' }), + ]); + await descendsFrom(api, threadId, savedFork, sourcePosition, [position(competitor.checkpoint).checkpoint_id]); + const forkRequests = (await modelJournal(forkPrompt)).slice(beforeFork); + expect(forkRequests).toHaveLength(1); + expect(forkRequests[0].body?.messages?.filter((message) => message.role !== 'system').map((message) => message.content ?? '')) + .toEqual([...source.values.messages.map((message) => message.content), forkPrompt]); + + const advanced = await run(api, threadId, position(competitor.checkpoint), { + messages: [{ id: 'advanced-competitor', type: 'human', content: advancePrompt }], + }); + expect(position((await api.threads.getState(threadId)).checkpoint)).toEqual(position(advanced.checkpoint)); + await session.load?.(); + expect(session.getSnapshot().messages.map((message) => message.content)) + .toEqual(savedFork.values.messages.map((message) => message.content)); + const beforeFollowUp = (await modelJournal(followUpPrompt)).length; + expect(await session.submit(followUpPrompt)).toBe('success'); + const followed = transport.finalPosition(); + const savedFollowed = await api.threads.getState(threadId, followed); + expect(savedFollowed.values.messages.map((message) => message.content)).toEqual([ + ...savedFork.values.messages.map((message) => message.content), followUpPrompt, 'Completed fork follow-up answer.', + ]); + await descendsFrom(api, threadId, savedFollowed, fork, [competitor, advanced].map((state) => position(state.checkpoint).checkpoint_id)); + const followUpRequests = (await modelJournal(followUpPrompt)).slice(beforeFollowUp); + expect(followUpRequests).toHaveLength(1); + expect(followUpRequests[0].body?.messages?.filter((message) => message.role !== 'system').map((message) => message.content ?? '')) + .toEqual([...savedFork.values.messages.map((message) => message.content), followUpPrompt]); + + const beforeTool = (await modelJournal()).length; + const beforeCreations = transport.creations.length; + const outcome = session.submit(prompt); + // Await actual handler entry, not a timer or an outcome label. Its result + // is held while another writer moves the global tip. + await Promise.race([ + started.promise, + outcome.then((result) => { throw new Error(`Tool never started: ${result}; ${JSON.stringify(session.getSnapshot().error)}`); }), + ]); + const pending = transport.finalPosition(); + const pendingState = await api.threads.getState(threadId, pending); + const call = pendingState.values.messages.flatMap((message) => message.tool_calls ?? []).at(-1); + expect(call).toMatchObject({ name: 'get_weather', args: { location: 'Paris' } }); + expect(pendingState.values.messages).toEqual([ + ...savedFollowed.values.messages, + expect.objectContaining({ type: 'human', content: prompt }), + expect.objectContaining({ type: 'ai', tool_calls: [call] }), + ]); + const toolRequests = (await modelJournal()).slice(beforeTool); + expect(toolRequests).toHaveLength(1); + expect(toolRequests[0].body?.messages?.filter((message) => message.role !== 'system').map((message) => message.content ?? '')) + .toEqual([...savedFollowed.values.messages.map((message) => message.content), prompt]); + await descendsFrom(api, threadId, pendingState, followed, [competitor, advanced].map((state) => position(state.checkpoint).checkpoint_id)); + const other = await run(api, threadId, position(advanced.checkpoint), { + messages: [{ id: 'during-handler-competitor', type: 'human', content: competitorPrompt }], + }); + expect(position((await api.threads.getState(threadId)).checkpoint)).toEqual(position(other.checkpoint)); + release.resolve(); + expect(await outcome).toBe('success'); + expect(handlers).toBe(1); + expect(transport.creations.length - beforeCreations).toBe(followUp ? 2 : 1); + expect(transport.writes).toHaveLength(followUp ? 0 : 1); + const exact = followUp ? transport.finalPosition() : position(transport.writes[0].result as Record); + const saved = await api.threads.getState(threadId, exact); + const toolMessage = expect.objectContaining({ + type: 'tool', tool_call_id: call?.id, content: '68°F', + }); + expect(saved.values.messages).toEqual([ + ...pendingState.values.messages, toolMessage, + ...(followUp ? [expect.objectContaining({ type: 'ai', content: 'Protocol weather complete: Paris is 68°F.' })] : []), + ]); + expect(saved.next).toEqual([]); + const requests = (await modelJournal()).slice(beforeTool); + expect(requests).toHaveLength(followUp ? 2 : 1); + if (followUp) { + expect(requests[1].body?.messages?.filter((message) => message.role === 'tool')) + .toEqual([expect.objectContaining({ tool_call_id: call?.id, content: '68°F' })]); + } + + // Read only exact parents. A fixture's answer alone cannot prove routing. + await descendsFrom(api, threadId, saved, pending, [competitor, advanced, other].map((state) => position(state.checkpoint).checkpoint_id)); + await session.load?.(); + expect(session.getSnapshot().messages.map((message) => message.content)) + .toEqual(saved.values.messages.map((message) => message.content)); + for (const original of [source, competitor, advanced, other]) { + expect((await api.threads.getState(threadId, position(original.checkpoint))).values).toEqual(original.values); + } + } finally { + release.resolve(); + session.dispose(); + await api.threads.delete(threadId); + } + }); +} diff --git a/cockpit/langgraph/interrupts/angular/e2e/checkpoint-session.spec.ts b/cockpit/langgraph/interrupts/angular/e2e/checkpoint-session.spec.ts new file mode 100644 index 000000000..926bd198d --- /dev/null +++ b/cockpit/langgraph/interrupts/angular/e2e/checkpoint-session.spec.ts @@ -0,0 +1,196 @@ +import { expect, test } from '@playwright/test'; +import { Client, type Checkpoint } from '@langchain/langgraph-sdk'; +import { createSession } from '../../../../../libs/langgraph/src/runtime/create-session'; +import { FetchStreamTransport } from '../../../../../libs/langgraph/src/lib/transport/fetch-stream.transport'; + +const seedPrompt = 'Refund $129.00 to customer cus_z19fp for an unrecognized charge.'; +const prompt = 'Refund $47.50 to customer cus_a8x2k — they were charged twice for the same order.'; +const acknowledgment = 'Understood — a $47.50 refund to cus_a8x2k for a duplicate charge. Pausing for operator approval; no refund is issued until a human reviews and approves.'; +const cancellation = 'Refund cancelled by operator. No charge issued.'; +type State = { + messages: { id: string; type: string; content: unknown }[]; + customer_id?: string; + amount?: number; + reason?: string; + decision_approved?: boolean; + refund_id?: string; +}; +type Position = Checkpoint & { checkpoint_id: string }; + +function environment(name: string) { + const value = process.env[name]; + if (!value) throw new Error(`Global setup must expose ${name}`); + return value; +} + +function position(config: Record | undefined): Position { + if (!config) throw new Error('Expected a full root checkpoint reference'); + expect(config['thread_id']).toEqual(expect.stringMatching(/\S/)); + expect(config['checkpoint_ns']).toBe(''); + expect(config['checkpoint_id']).toEqual(expect.stringMatching(/\S/)); + const map = config['checkpoint_map']; + if (map !== undefined) { + expect(map).not.toBeNull(); + expect(typeof map).toBe('object'); + expect(Array.isArray(map)).toBe(false); + for (const value of Object.values(map as object)) expect(value).toEqual(expect.any(String)); + } + return { + thread_id: config['thread_id'] as string, checkpoint_ns: '', checkpoint_id: config['checkpoint_id'] as string, + // Saved state may omit the empty map echoed by checkpoint frames. + checkpoint_map: map === undefined ? {} : { ...map as Record }, + }; +} + +class ObservedTransport extends FetchStreamTransport { + creations = 0; + readonly checkpoints: Position[] = []; + override async *stream(...args: Parameters) { + this.creations++; + for await (const event of super.stream(...args)) { + if (event.type === 'checkpoints') { + const data = event.data as { config: { configurable?: Record } }; + this.checkpoints.push(position(data.config.configurable)); + } + yield event; + } + } + finalPosition() { + const exact = this.checkpoints.at(-1); + if (!exact) throw new Error('Session must observe a real root checkpoint'); + return exact; + } +} + +async function run(api: Client, threadId: string, checkpoint?: Position, response?: { approved: boolean; amount?: number }) { + let exact: Position | undefined; + let runId: string | undefined; + for await (const event of api.runs.stream(threadId, 'interrupts', { + input: checkpoint ? null : { messages: [{ id: 'completed-source', type: 'human', content: seedPrompt }] }, + checkpoint, command: response ? { resume: response } : undefined, + streamMode: ['values', 'checkpoints'], signal: AbortSignal.timeout(20_000), + })) { + expect(event.event).not.toBe('error'); + if (event.event === 'metadata') runId = event.data.run_id; + if (event.event === 'checkpoints') exact = position((event.data as { config: { configurable?: Record } }).config.configurable); + } + if (!exact || !runId) throw new Error('Expected physical run and saved root checkpoint'); + const saved = await api.threads.getState(threadId, exact); + expect(position(saved.checkpoint)).toEqual(exact); + expect(saved.metadata?.['run_id']).toBe(runId); + expect((await api.runs.get(threadId, runId)).status).toBe('success'); + return saved; +} + +async function requests() { + const response = await fetch(`${environment('INTERRUPTS_AIMOCK_URL')}/__aimock/journal`, { signal: AbortSignal.timeout(5_000) }); + expect(response.ok).toBe(true); + const entries = await response.json() as { body?: { messages?: { role: string; content: unknown }[] } }[]; + return entries.filter((entry) => entry.body?.messages?.some((message) => message.role === 'user' && message.content === prompt)); +} + +for (const consumed of [false, true]) { + test(`checkpoint session: ${consumed ? 'rejects a task consumed by another client before resume' : 'resumes its own unconsumed pause after a distinct branch completes'}`, async () => { + const api = new Client({ + apiUrl: environment('INTERRUPTS_API_URL'), apiKey: null, callerOptions: { maxRetries: 0 }, timeoutMs: 20_000, + }); + const { thread_id: threadId } = await api.threads.create(); + const transport = new ObservedTransport(environment('INTERRUPTS_API_URL'), undefined, { maxRetries: 0 }); + const session = createSession({ assistantId: 'interrupts', threadId, transport }); + try { + const firstPause = await run(api, threadId); + const source = await run(api, threadId, position(firstPause.checkpoint), { approved: false }); + expect(source.next).toEqual([]); + expect(source.tasks).toEqual([]); + expect(source.values.refund_id).toBeUndefined(); + const before = (await requests()).length; + expect(await session.fork(source.checkpoint, prompt)).toBe('paused'); + const pausedPosition = transport.finalPosition(); + const paused = await api.threads.getState(threadId, pausedPosition); + expect(paused.values.messages.map((message) => message.content)).toEqual([ + ...source.values.messages.map((message) => message.content), prompt, acknowledgment, + ]); + expect(paused.next).toEqual(['request_approval']); + expect(paused.tasks).toHaveLength(1); + const task = paused.tasks[0]; + expect(task).toMatchObject({ id: expect.stringMatching(/\S/), name: 'request_approval', result: null }); + expect(task.interrupts).toEqual([expect.objectContaining({ + id: expect.stringMatching(/\S/), + value: { kind: 'refund_approval', customer_id: 'cus_a8x2k', amount: 47.5, reason: 'Customer was charged twice for the same order.' }, + })]); + expect(await requests()).toHaveLength(before + 2); + const snapshot = session.getSnapshot(); + const beforeCreations = transport.creations; + + if (consumed) { + const approved = await run(api, threadId, pausedPosition, { approved: true, amount: 31.25 }); + expect(approved.values.decision_approved).toBe(true); + const reread = await api.threads.getState(threadId, pausedPosition); + // Checkpoint values/interrupts alone conceal consumption on this backend. + expect(reread.values).toEqual(paused.values); + expect(reread.tasks[0]).toMatchObject({ id: task.id, interrupts: task.interrupts, result: expect.objectContaining({ decision_approved: true }) }); + const publications: ReturnType[] = []; + const unsubscribe = session.subscribe(() => publications.push(session.getSnapshot())); + let rejected = false; + try { + await session.resume({ approved: false }); + } catch { + rejected = true; + } finally { + unsubscribe(); + } + expect(transport.creations).toBe(beforeCreations); + expect(rejected).toBe(true); + for (const published of publications) { + expect(published.messages).toEqual(snapshot.messages); + expect(published.values).toEqual(snapshot.values); + } + expect(session.getSnapshot().messages).toEqual(snapshot.messages); + expect(session.getSnapshot().values).toEqual(snapshot.values); + expect(position((await api.threads.getState(threadId)).checkpoint)).toEqual(position(approved.checkpoint)); + } else { + const fork = position((await api.threads.updateState(threadId, { + checkpoint: pausedPosition, values: { amount: 31.25 }, asNode: 'draft', signal: AbortSignal.timeout(10_000), + })).configurable); + const forked = await api.threads.getState(threadId, fork); + expect(forked.tasks[0].id).not.toBe(task.id); + const approved = await run(api, threadId, fork, { approved: true, amount: 31.25 }); + expect(approved.values).toMatchObject({ amount: 31.25, decision_approved: true, refund_id: 're_demo__a8x2k' }); + expect(position((await api.threads.getState(threadId)).checkpoint)).toEqual(position(approved.checkpoint)); + const stillPaused = await api.threads.getState(threadId, pausedPosition); + expect(stillPaused.tasks).toEqual([expect.objectContaining({ id: task.id, result: null, interrupts: task.interrupts })]); + expect(await session.resume({ approved: false })).toBe('success'); + expect(transport.creations).toBe(beforeCreations + 1); + const exact = transport.finalPosition(); + const rejected = await api.threads.getState(threadId, exact); + expect(rejected.values).toEqual({ + ...paused.values, decision_approved: false, + messages: [...paused.values.messages, expect.objectContaining({ type: 'ai', content: cancellation })], + }); + expect(rejected.values.refund_id).toBeUndefined(); + expect(rejected.next).toEqual([]); + expect(rejected.tasks).toEqual([]); + let ancestor = rejected; + const visited = new Set(); + for (let depth = 0; depth < 6; depth++) { + const current = position(ancestor.checkpoint); + expect(current.thread_id).toBe(threadId); + expect([fork.checkpoint_id, approved.checkpoint.checkpoint_id]).not.toContain(current.checkpoint_id); + expect(visited.has(current.checkpoint_id)).toBe(false); + visited.add(current.checkpoint_id); + if (current.checkpoint_id === pausedPosition.checkpoint_id) break; + ancestor = await api.threads.getState(threadId, position(ancestor.parent_checkpoint ?? undefined)); + } + expect(position(ancestor.checkpoint)).toEqual(pausedPosition); + await session.load?.(); + expect(session.getSnapshot().messages.map((message) => message.content)).toEqual(rejected.values.messages.map((message) => message.content)); + expect((await api.threads.getState(threadId, position(approved.checkpoint))).values).toEqual(approved.values); + } + expect(await requests()).toHaveLength(before + 2); + expect((await api.threads.getState(threadId, position(source.checkpoint))).values).toEqual(source.values); + } finally { + session.dispose(); + await api.threads.delete(threadId); + } + }); +} diff --git a/cockpit/langgraph/interrupts/angular/e2e/tsconfig.json b/cockpit/langgraph/interrupts/angular/e2e/tsconfig.json index 0fc9befb1..e09b301a5 100644 --- a/cockpit/langgraph/interrupts/angular/e2e/tsconfig.json +++ b/cockpit/langgraph/interrupts/angular/e2e/tsconfig.json @@ -12,6 +12,12 @@ ], "baseUrl": "../../../../..", "paths": { + "@threadplane/core": [ + "libs/core/src/index.ts" + ], + "@threadplane/core/tools": [ + "libs/core/src/tools/index.ts" + ], "@threadplane-internal/e2e-harness": [ "libs/e2e-harness/src/index.ts" ], diff --git a/libs/langgraph/src/lib/transport/checkpoint-position.ts b/libs/langgraph/src/lib/transport/checkpoint-position.ts new file mode 100644 index 000000000..6b74f22b6 --- /dev/null +++ b/libs/langgraph/src/lib/transport/checkpoint-position.ts @@ -0,0 +1,59 @@ +import type { Checkpoint } from '@langchain/langgraph-sdk'; + +/** SDK/history input shape; complete root routing is checked at capture. */ +export type CheckpointReference = Readonly< + Omit +> & { + readonly checkpoint_map?: Readonly> | null; +}; +export type OwnedCheckpointPosition = { + readonly thread_id: string; + readonly checkpoint_ns: ''; + readonly checkpoint_id: string; + readonly checkpoint_map: Readonly>; +}; + +/** Shared wire routing only. Never copy arbitrary configurable metadata. */ +export function captureCheckpoint( + value: unknown, + threadId: string +): OwnedCheckpointPosition { + const invalid = () => + new Error('Checkpoint execution authority is unavailable.'); + if (!value || typeof value !== 'object' || Array.isArray(value)) + throw invalid(); + const input = value as Record; + const thread = input['thread_id']; + const namespace = input['checkpoint_ns']; + const id = input['checkpoint_id']; + const map = input['checkpoint_map']; + if ( + thread !== threadId || + !threadId.trim() || + namespace !== '' || + typeof id !== 'string' || + !id.trim() + ) + throw invalid(); + let entries: [string, string][] = []; + if (map !== undefined) { + if ( + !map || + typeof map !== 'object' || + Array.isArray(map) || + (Object.getPrototypeOf(map) !== Object.prototype && + Object.getPrototypeOf(map) !== null) + ) + throw invalid(); + entries = Object.entries(map).map(([key, value]) => { + if (typeof value !== 'string') throw invalid(); + return [key, value]; + }); + } + return Object.freeze({ + thread_id: threadId, + checkpoint_ns: '', + checkpoint_id: id, + checkpoint_map: Object.freeze(Object.fromEntries(entries)), + }); +} diff --git a/libs/langgraph/src/lib/transport/fetch-stream.transport.ts b/libs/langgraph/src/lib/transport/fetch-stream.transport.ts index b88b976cf..6d8fc577c 100644 --- a/libs/langgraph/src/lib/transport/fetch-stream.transport.ts +++ b/libs/langgraph/src/lib/transport/fetch-stream.transport.ts @@ -1,4 +1,5 @@ import type { Client, Run, StreamMode, ThreadState } from '@langchain/langgraph-sdk'; +import { captureCheckpoint, type OwnedCheckpointPosition } from './checkpoint-position'; import type { AgentQueueEntry, AgentTransport, LangGraphClientOptions, LangGraphSubmitOptions, StreamEvent } from '../../runtime/transport.types'; import { createLangGraphClient, @@ -108,12 +109,14 @@ export class FetchStreamTransport implements AgentTransport { runId: string, lastEventId: string | undefined, signal: AbortSignal, + options?: { streamMode?: StreamMode[] } ): AsyncIterable { // SDK joinStream: joins an already-started run without creating a new one. let run: ReturnType; try { run = this.client.runs.joinStream(threadId, runId, { signal, + ...(options?.streamMode ? { streamMode: options.streamMode } : {}), ...(lastEventId !== undefined ? { lastEventId } : {}), }); } catch (error) { @@ -189,19 +192,47 @@ export class FetchStreamTransport implements AgentTransport { } } - /** Update server-side thread state, e.g. to remove messages for regenerate rollback. */ + /** Read one exact saved checkpoint with the command's cancellation signal. */ + async getState( + threadId: string, + checkpoint: OwnedCheckpointPosition, + signal: AbortSignal + ): Promise { + try { + return await this.client.threads.getState(threadId, checkpoint, { + signal, + }); + } catch (error) { + return this.rethrowOperationError(error, signal); + } + } + + /** Update state once and return only usable root routing, when supplied. */ async updateState( threadId: string, values: Record, signal: AbortSignal, - options?: { asNode?: string }, - ): Promise { - const body: { values: Record; signal: AbortSignal; asNode?: string } = { values, signal }; + options?: { asNode?: string; checkpoint?: OwnedCheckpointPosition } + ): Promise { + const body: { + values: Record; + signal: AbortSignal; + asNode?: string; + checkpoint?: OwnedCheckpointPosition; + } = { values, signal }; if (options?.asNode !== undefined) { body.asNode = options.asNode; } + if (options?.checkpoint) body.checkpoint = options.checkpoint; try { - await this.client.threads.updateState(threadId, body); + const result = await this.client.threads.updateState(threadId, body); + // Legacy non-branch callers need no routing acknowledgment. The branch + // effect owner rejects absence without replaying the successful request. + try { + return captureCheckpoint(result.configurable, threadId); + } catch { + return undefined; + } } catch (error) { this.rethrowOperationError(error, signal); } diff --git a/libs/langgraph/src/runtime/README.md b/libs/langgraph/src/runtime/README.md new file mode 100644 index 000000000..9980787de --- /dev/null +++ b/libs/langgraph/src/runtime/README.md @@ -0,0 +1,65 @@ +# Private LangGraph session development + +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. + +## Completed checkpoint forks + +`session.fork(checkpoint, input, options?)` submits new input in the session's +fixed thread from an exact, completed root checkpoint. Pass the full SDK/history +checkpoint reference; capture checks the thread, root namespace, ID and optional +map. The source must have no remaining graph work, interrupts or unanswered tool +calls. Selecting history in a UI does not change this authority. + +After activation, subsequent submit, tool follow-up, result persistence and load +use the session's resulting checkpoint. Root checkpoint frames are candidates; +the owner confirms the physical run has ended and reads its exact saved state +before permitting another effect. A custom transport must perform single requests, +return exact checkpoint routing after writes, and deliver checkpoint frames during +both creation and joining. An owned SDK client configured to retry requests cannot +activate this capability. Retrying an ambiguous write could create a sibling. + +New tool calls retain the existing execution-store and follow-up policies. +Tool candidates for new effects must match the invocation and resolution evidence +in the confirmed saved messages. Transient stream calls or results cannot +independently authorize or suppress execution. +Locally settled invocations retain identity-bound completion evidence for their +branch owner across loads. Graphs may consume a result and remove its wire message; +that does not authorize another handler invocation when a pause resumes. History +and wire results cannot create this completion evidence, and a new fork starts +with none. +Historical ToolMessages establish entry eligibility only; they do not restore +typed results or prove historical guard coverage. A new invocation cannot reuse a +baseline call ID. Serialized result writes retain their originating owner and +advance through each returned checkpoint, including already-authorized cleanup +after local stop or disposal. + +Missing final evidence and ambiguous write acknowledgments leave authority +unavailable. Commands cannot silently use an older checkpoint or adopt the global +tip. `checkStatus()` does not repair a branch: ready is a no-op; active, disposed +or uncertain states reject. Explicit reconnect can join a retained physical run +when sufficient run/cursor evidence exists. It never submits another run. + +## Dynamic pauses + +Branch resume reads the exact retained pause before sending a decision. Task and +interrupt identities, payloads and unconsumed task results must match. Undefined +does not mean replay; pass an explicit defined response, including null where +appropriate. Static breakpoints and child execution are outside this capability. + +On the locked interrupts backend (LangGraph 1.1.6 / API 0.7.96), a physically +successful run can be paused. Resuming the same consumed task again can reuse its +first decision even though checkpoint values and interrupt metadata still look +paused. The preflight detects already-observed consumption, but it is not an atomic +server claim. Applications must coordinate responders to the same task. Distinct +checkpoint branches do not provide independent opposing decisions on one task. + +## Verification + +Run `langgraph:runtime-quality`, `langgraph:runtime-type-tests` and +`langgraph:type-tests`. Actual server coverage lives in the client-tools and +interrupts cockpit `checkpoint-session.spec.ts` files; adjacent protocol tests +characterize the locked SDK/backend behavior. Installed Angular/React consumers +verify shared ownership and lifecycle regressions separately. No replay API, +branch-tree UI, cross-client lease, public package migration or release is implied. diff --git a/libs/langgraph/src/runtime/checkpoint-admission.spec.ts b/libs/langgraph/src/runtime/checkpoint-admission.spec.ts new file mode 100644 index 000000000..b91db8816 --- /dev/null +++ b/libs/langgraph/src/runtime/checkpoint-admission.spec.ts @@ -0,0 +1,268 @@ +import { afterEach, describe, expect, it, vi } from 'vitest'; +import { createSession, type SessionOptions } from './create-session'; +import type { ThreadState } from '@langchain/langgraph-sdk'; +import { + checkpointEvent, + fixture, + position, + saved, +} from './testing/checkpoint-fixture'; +import { deferred } from './testing/deferred'; + +afterEach(() => vi.unstubAllGlobals()); +describe('checkpoint command admission', () => { + it('rejects unsupported pre-branch load without retaining ownership, then loads an activated branch exactly', async () => { + const f = fixture(); + delete f.transport.getHistory; + const session = createSession({ + assistantId: 'agent', + threadId: 'thread', + transport: f.transport, + }); + let outcome: 'pending' | 'resolved' | 'rejected' = 'pending'; + const read = session.load?.(); + void read?.then( + () => { + outcome = 'resolved'; + }, + () => { + outcome = 'rejected'; + } + ); + try { + // Drain already-queued promise work; no backend request can settle this read. + await new Promise((resolve) => setImmediate(resolve)); + expect(outcome).toBe('rejected'); + expect(f.transport.getState).not.toHaveBeenCalled(); + expect(await session.fork(position('a'), 'Fork')).toBe('success'); + await session.load?.(); + expect(f.transport.getState).toHaveBeenLastCalledWith( + 'thread', + position('result-1'), + expect.any(AbortSignal) + ); + } finally { + await session.dispose(); + await read?.catch(() => undefined); + } + }); + it.each(['checkpoint', 'input', 'options', 'signal', 'state'] as const)( + 'prevents stale commits after a reentrant %s getter stops preparation', + async (field) => { + const f = fixture(); + const session = createSession({ + assistantId: 'agent', + threadId: 'thread', + transport: f.transport, + }); + const stop = () => { + void session.stop(); + }; + const checkpoint = + field === 'checkpoint' + ? { + ...position('a'), + get checkpoint_id() { + stop(); + return 'a'; + }, + } + : position('a'); + const input = + field === 'input' + ? { + get message() { + stop(); + return 'Fork'; + }, + } + : 'Fork'; + const options = + field === 'signal' + ? { + get signal() { + stop(); + return new AbortController().signal; + }, + } + : field === 'options' + ? { + get context() { + stop(); + return {}; + }, + } + : undefined; + if (field === 'state') + f.states.set('a', { + ...f.source, + get values() { + stop(); + return f.source.values; + }, + }); + const snapshot = session.getSnapshot(); + expect(await session.fork(checkpoint, input, options)).toBe('aborted'); + expect(f.transport.stream).not.toHaveBeenCalled(); + expect(session.getSnapshot()).toBe(snapshot); + expect(f.transport.getState).toHaveBeenCalledTimes( + field === 'state' ? 1 : 0 + ); + } + ); + it('lets an observer stop an atomically installed baseline before the creation POST', async () => { + const f = fixture(); + const session = createSession({ + assistantId: 'agent', + threadId: 'thread', + transport: f.transport, + }); + const observed: string[][] = []; + session.subscribe(() => { + observed.push( + session.getSnapshot().messages.map((message) => message.content) + ); + if (session.getSnapshot().status === 'running') void session.stop(); + }); + expect(await session.fork(position('a'), 'Fork')).toBe('aborted'); + expect(observed[0]).toEqual(['Source A', 'Answer A', 'Fork']); + expect(f.transport.stream).not.toHaveBeenCalled(); + await session.checkStatus?.(); + }); + it('captures owned SDK retries once, rejects positive retries before any I/O, and ignores supplied transport client options', async () => { + const fetch = vi.fn(); + vi.stubGlobal('fetch', fetch); + let retryReads = 0; + let retries = 1; + const session = createSession({ + assistantId: 'agent', + threadId: 'thread', + apiUrl: 'https://runtime.example', + clientOptions: { + get maxRetries() { + retryReads++; + return retries; + }, + }, + }); + retries = 0; + await expect(session.fork(position('a'), 'Fork')).rejects.toThrow( + 'maxRetries' + ); + expect(retryReads).toBe(1); + expect(fetch).not.toHaveBeenCalled(); + const f = fixture(); + const supplied = createSession({ + assistantId: 'agent', + threadId: 'thread', + transport: f.transport, + get clientOptions(): SessionOptions['clientOptions'] { + throw new Error('Must be ignored'); + }, + }); + expect(await supplied.fork(position('a'), 'Fork')).toBe('success'); + }); + it('keeps source failures safe and leaves prior observations untouched', async () => { + const f = fixture(); + const session = createSession({ + assistantId: 'agent', + threadId: 'thread', + transport: f.transport, + }); + await session.load?.(); + const before = session.getSnapshot(); + f.states.set( + 'a', + saved('a', [ + { + type: 'ai', + id: 'hidden', + tool_calls: [{ id: 'pending', name: 'unregistered', args: {} }], + }, + ]) + ); + await expect(session.fork(position('a'), 'Fork')).rejects.not.toThrow( + 'secret' + ); + expect(session.getSnapshot()).toBe(before); + expect(f.transport.stream).not.toHaveBeenCalled(); + }); + it('rejects fork during history reads and branch commands during exact load', async () => { + const f = fixture(); + const history = deferred(); + f.transport.getHistory = vi.fn(() => history.promise); + const session = createSession({ + assistantId: 'agent', + threadId: 'thread', + transport: f.transport, + }); + const reading = session.load?.(); + await expect(session.fork(position('a'), 'Fork')).rejects.toThrow(); + history.resolve([]); + await reading; + expect(await session.fork(position('a'), 'Fork')).toBe('success'); + }); + it('blocks baseline ID reuse even when the new turn reuses its historical assistant message ID', async () => { + const f = fixture(); + const call = { id: 'old', name: 'work', args: {} }; + const historical = { + type: 'ai', + id: 'same-assistant', + content: '', + tool_calls: [call], + }; + f.states.set( + 'a', + saved('a', [ + historical, + { type: 'tool', tool_call_id: 'old', content: 'old result' }, + ]) + ); + f.transport.stream = vi.fn(async function* (_a, _t, input, _s, options) { + options?.onRunCreated?.({ run_id: 'run-collision' }); + yield checkpointEvent( + saved('collision', [ + ...(input as { messages: unknown[] }).messages, + historical, + ]) + ); + }); + const session = createSession({ + assistantId: 'agent', + threadId: 'thread', + transport: f.transport, + }); + expect(await session.fork(position('a'), 'New turn')).toBe('interrupted'); + expect(session.getSnapshot().error?.message).toContain('identity conflict'); + await expect(session.submit('Again')).rejects.toThrow(); + }); + it('rejects changed baseline assistant content with a repeated call even when the stream omits the user anchor', async () => { + const f = fixture(); + const historical = { + type: 'ai', + id: 'same', + content: 'Old', + tool_calls: [{ id: 'old', name: 'work', args: {} }], + }; + f.states.set( + 'a', + saved('a', [ + historical, + { type: 'tool', tool_call_id: 'old', content: 'done' }, + ]) + ); + f.transport.stream = vi.fn(async function* (_a, _t, _i, _s, options) { + options?.onRunCreated?.({ run_id: 'run-b' }); + const result = saved('b', [{ ...historical, content: 'New' }]); + f.states.set('b', result); + yield checkpointEvent(result); + }); + const session = createSession({ + assistantId: 'agent', + threadId: 'thread', + transport: f.transport, + }); + expect(await session.fork(position('a'), 'Run')).toBe('interrupted'); + expect(session.getSnapshot().error?.message).toContain('identity conflict'); + }); +}); diff --git a/libs/langgraph/src/runtime/checkpoint-authority.spec.ts b/libs/langgraph/src/runtime/checkpoint-authority.spec.ts new file mode 100644 index 000000000..3848dd700 --- /dev/null +++ b/libs/langgraph/src/runtime/checkpoint-authority.spec.ts @@ -0,0 +1,198 @@ +import { describe, expect, it } from 'vitest'; +import { + captureCheckpoint, + captureCompletedCheckpoint, + captureCheckpointEvent, + confirmCheckpoint, + assertResumeCheckpoint, +} from './checkpoint-authority'; + +const checkpoint = { + thread_id: 'thread', + checkpoint_ns: '' as const, + checkpoint_id: 'a', + checkpoint_map: { '': 'a' }, +}; +const saved = (changes: Record = {}) => ({ + checkpoint, + values: { messages: [] }, + next: [], + tasks: [], + metadata: { run_id: 'run' }, + ...changes, +}); +const call = { + type: 'ai', + id: 'assistant', + tool_calls: [{ id: 'old', name: 'unknown', args: { x: 1 } }], +}; + +describe('checkpoint authority capture', () => { + it.each([{}, false, 0, null, '', [{ id: 'interrupt', value: 'pending' }]])( + 'rejects completed source root interrupts %j', + (interrupts) => { + expect(() => + captureCompletedCheckpoint(saved({ interrupts }), checkpoint) + ).toThrow(); + } + ); + it.each([undefined, []])( + 'accepts optional empty completed source root interrupts %j', + (interrupts) => { + expect(() => + captureCompletedCheckpoint(saved({ interrupts }), checkpoint) + ).not.toThrow(); + } + ); + it('owns only routing fields and deep freezes maps', () => { + const input = { + ...checkpoint, + checkpoint_map: { '': 'a' }, + secret: 'omit', + }; + const captured = captureCheckpoint(input, 'thread'); + input.checkpoint_map[''] = 'changed'; + expect(captured).toEqual(checkpoint); + expect(Object.isFrozen(captured.checkpoint_map)).toBe(true); + }); + it.each([ + { thread_id: 'other' }, + { checkpoint_ns: 'child' }, + { checkpoint_id: undefined }, + { checkpoint_id: '' }, + { checkpoint_map: { '': 3 } }, + { checkpoint_map: null }, + ])('rejects unsupported routing %j', (change) => { + expect(() => + captureCheckpoint({ ...checkpoint, ...change }, 'thread') + ).toThrow(); + }); + it.each([ + { next: ['work'] }, + { tasks: [{}] }, + { values: { __interrupt__: [], messages: [] } }, + { values: { messages: [call] } }, + ])('rejects incomplete or hidden calls %j', (change) => { + expect(() => + captureCompletedCheckpoint(saved(change), checkpoint) + ).toThrow(); + }); + it('accepts authoritative historical tool results without reconstructing typed provenance', () => { + const source = captureCompletedCheckpoint( + saved({ + values: { + messages: [ + call, + { type: 'tool', tool_call_id: 'old', content: '{"x":1}' }, + ], + }, + }), + checkpoint + ); + expect(source.calls.map((entry) => entry.id)).toEqual(['old']); + expect(source.state.values).toEqual({ + messages: [ + call, + { type: 'tool', tool_call_id: 'old', content: '{"x":1}' }, + ], + }); + }); +}); + +describe('final checkpoint evidence', () => { + const event = (changes: Record = {}) => ({ + type: 'checkpoints' as const, + data: { + config: { configurable: { ...checkpoint, run_id: 'run' } }, + values: { messages: [] }, + next: [], + tasks: [], + ...changes, + }, + }); + it.each([{}, false, 0, null, '', [{ id: 'interrupt', value: 'pending' }]])( + 'rejects completed confirmation root interrupts %j', + (interrupts) => { + const candidate = captureCheckpointEvent(event(), 'thread'); + expect(() => + confirmCheckpoint(candidate, saved({ interrupts }), 'run') + ).toThrow(); + } + ); + it.each([undefined, []])( + 'accepts optional empty completed confirmation root interrupts %j', + (interrupts) => { + const candidate = captureCheckpointEvent(event(), 'thread'); + expect( + confirmCheckpoint(candidate, saved({ interrupts }), 'run').paused + ).toBe(false); + } + ); + it('requires matching saved root identity, physical run, and final state', () => { + const candidate = captureCheckpointEvent(event(), 'thread'); + expect(confirmCheckpoint(candidate, saved(), 'run').position).toEqual( + checkpoint + ); + for (const change of [ + { checkpoint: { ...checkpoint, checkpoint_id: 'other' } }, + { metadata: { run_id: 'other' } }, + { values: { messages: ['other'] } }, + { next: ['work'] }, + ]) + expect(() => + confirmCheckpoint(candidate, saved(change), 'run') + ).toThrow(); + expect(() => confirmCheckpoint(undefined, saved(), 'run')).toThrow(); + expect( + captureCheckpointEvent({ ...event(), namespace: ['child'] }, 'thread') + ).toBeUndefined(); + }); + it('supports only unconsumed dynamic root tasks and detects consumption before resume', () => { + const tasks = [ + { + id: 'task', + name: 'approval', + error: null, + result: null, + interrupts: [{ id: 'interrupt', value: { amount: 10 } }], + }, + ]; + const pause = saved({ next: ['approval'], tasks }); + const candidate = captureCheckpointEvent( + event({ next: ['approval'], tasks: [{ id: 'task', name: 'approval' }] }), + 'thread' + ); + const confirmed = confirmCheckpoint(candidate, pause, 'run'); + expect(confirmed.paused).toBe(true); + expect(() => assertResumeCheckpoint(confirmed, pause)).not.toThrow(); + expect(() => + assertResumeCheckpoint(confirmed, { + ...pause, + values: { messages: [], amount: 99 }, + }) + ).toThrow(); + expect(() => + assertResumeCheckpoint( + confirmed, + saved({ + next: ['approval'], + tasks: [{ ...tasks[0], result: { decision: true } }], + }) + ) + ).toThrow(); + expect(() => + assertResumeCheckpoint( + confirmed, + saved({ + next: ['approval'], + tasks: [ + { + ...tasks[0], + interrupts: [{ id: 'different', value: { amount: 10 } }], + }, + ], + }) + ) + ).toThrow(); + }); +}); diff --git a/libs/langgraph/src/runtime/checkpoint-authority.ts b/libs/langgraph/src/runtime/checkpoint-authority.ts new file mode 100644 index 000000000..4726d4240 --- /dev/null +++ b/libs/langgraph/src/runtime/checkpoint-authority.ts @@ -0,0 +1,226 @@ +import type { ThreadState } from '@langchain/langgraph-sdk'; +import type { PlainValue } from '@threadplane/core'; +import { ownValue, sameOwnedValue, sameToolInvocation } from './ownership'; +import { record, roleOf } from './wire-message'; +import type { ToolInvocation } from './tool-invocations'; +import type { StreamEvent } from './transport.types'; + +import { + captureCheckpoint, + type OwnedCheckpointPosition, +} from '../lib/transport/checkpoint-position'; +export { + captureCheckpoint, + type OwnedCheckpointPosition, + type CheckpointReference, +} from '../lib/transport/checkpoint-position'; + +const unavailable = () => + new Error('Checkpoint execution authority is unavailable.'); +const nonempty = (value: unknown): value is string => + typeof value === 'string' && value.trim().length > 0; + +type CapturedCheckpointState = ThreadState & { + readonly interrupts?: readonly unknown[]; +}; + +/** Select and own once before any projection can traverse transport getters. */ +export function captureCheckpointState( + value: unknown, + position: OwnedCheckpointPosition +): CapturedCheckpointState { + const raw = record(value); + if (!raw) throw unavailable(); + const state = ownValue({ + checkpoint: raw['checkpoint'], + values: raw['values'], + next: raw['next'], + tasks: raw['tasks'], + metadata: raw['metadata'], + interrupts: raw['interrupts'], + } as PlainValue) as unknown as CapturedCheckpointState; + if ( + !sameOwnedValue( + captureCheckpoint(state.checkpoint, position.thread_id), + position + ) || + !record(state.values) || + !Array.isArray(state.next) || + !state.next.every((next) => typeof next === 'string') || + !Array.isArray(state.tasks) || + (state.interrupts !== undefined && !Array.isArray(state.interrupts)) + ) + throw unavailable(); + return state; +} + +export function captureCompletedCheckpoint( + value: unknown, + position: OwnedCheckpointPosition +) { + const state = captureCheckpointState(value, position); + const values = record(state.values)!; + if ( + state.next.length || + state.tasks.length || + Object.hasOwn(values, '__interrupt__') || + state.interrupts?.length + ) + throw unavailable(); + const messages = values['messages']; + if (messages !== undefined && !Array.isArray(messages)) throw unavailable(); + const calls = new Map(); + const results = new Set(); + for (const raw of (messages as unknown[] | undefined) ?? []) { + const message = record(raw); + if (!message) throw unavailable(); + if (roleOf(message) === 'tool') { + if (!nonempty(message['tool_call_id'])) throw unavailable(); + results.add(message['tool_call_id']); + } + if (roleOf(message) !== 'assistant') continue; + const entries = message['tool_calls']; + if (entries === undefined) continue; + if (!Array.isArray(entries) || message['type'] === 'AIMessageChunk') + throw unavailable(); + for (const rawCall of entries) { + const entry = record(rawCall); + const id = entry?.['id']; + const name = entry?.['name']; + if (!nonempty(id) || !nonempty(name)) throw unavailable(); + const call = Object.freeze({ + id, + name, + args: entry?.['args'] as PlainValue, + }); + const previous = calls.get(id); + if (previous && !sameToolInvocation(previous, call)) throw unavailable(); + calls.set(id, call); + } + } + if ([...calls.keys()].some((id) => !results.has(id))) throw unavailable(); + return { state, calls: Object.freeze([...calls.values()]) }; +} + +export interface CheckpointCandidate { + readonly position: OwnedCheckpointPosition; + readonly runId: string; + readonly values: PlainValue; + readonly next: readonly string[]; + readonly tasks: readonly { readonly id: string; readonly name: string }[]; +} +export interface ConfirmedCheckpoint { + readonly position: OwnedCheckpointPosition; + readonly state: ThreadState; + readonly paused: boolean; +} + +export function captureCheckpointEvent( + event: StreamEvent, + threadId: string +): CheckpointCandidate | undefined { + if (event.type !== 'checkpoints' || event.namespace?.length) return undefined; + const raw = record(event['data']); + const config = record(record(raw?.['config'])?.['configurable']); + const position = captureCheckpoint(config, threadId); + const runId = config?.['run_id']; + const values = ownValue(raw?.['values'] as PlainValue); + const next = ownValue(raw?.['next'] as PlainValue); + const tasks = ownValue(raw?.['tasks'] as PlainValue); + if ( + !nonempty(runId) || + !record(values) || + !Array.isArray(next) || + !next.every((entry) => typeof entry === 'string') || + !Array.isArray(tasks) + ) + throw unavailable(); + const identities = tasks.map((value) => { + const task = record(value); + if (!nonempty(task?.['id']) || !nonempty(task?.['name'])) + throw unavailable(); + return Object.freeze({ id: task['id'], name: task['name'] }); + }); + return Object.freeze({ + position, + runId, + values, + next: next as readonly string[], + tasks: Object.freeze(identities), + }); +} + +function pauseEvidence(state: ThreadState) { + if (!state.next.length || !state.tasks.length) throw unavailable(); + const ids = new Set(); + return state.tasks.map((raw) => { + const task = record(raw); + if ( + !task || + !nonempty(task['id']) || + !nonempty(task['name']) || + task['result'] !== null || + task['error'] != null || + task['state'] != null || + !state.next.includes(task['name']) || + ids.has(task['id']) + ) + throw unavailable(); + ids.add(task['id']); + const interrupts = task['interrupts']; + if ( + !Array.isArray(interrupts) || + !interrupts.length || + interrupts.some( + (entry) => + !nonempty(record(entry)?.['id']) || !Object.hasOwn(entry, 'value') + ) + ) + throw unavailable(); + return { id: task['id'], name: task['name'], interrupts }; + }); +} + +export function confirmCheckpoint( + candidate: CheckpointCandidate | undefined, + value: unknown, + runId: string +): ConfirmedCheckpoint { + if (!candidate || candidate.runId !== runId) throw unavailable(); + const state = captureCheckpointState(value, candidate.position); + if ( + state.metadata?.['run_id'] !== runId || + !sameOwnedValue(state.values as PlainValue, candidate.values) || + !sameOwnedValue(state.next, candidate.next) || + !sameOwnedValue( + state.tasks.map((task) => ({ id: task.id, name: task.name })), + candidate.tasks + ) + ) + throw unavailable(); + const paused = !!(state.next.length || state.tasks.length); + if (paused) pauseEvidence(state); + else if ( + Object.hasOwn(record(state.values)!, '__interrupt__') || + state.interrupts?.length + ) + throw unavailable(); + return Object.freeze({ position: candidate.position, state, paused }); +} + +export function assertResumeCheckpoint( + previous: ConfirmedCheckpoint, + value: unknown +) { + if (!previous.paused) throw unavailable(); + const state = captureCheckpointState(value, previous.position); + if ( + !sameOwnedValue(pauseEvidence(previous.state), pauseEvidence(state)) || + !sameOwnedValue(previous.state.next, state.next) || + !sameOwnedValue( + previous.state.values as PlainValue, + state.values as PlainValue + ) + ) + throw unavailable(); +} diff --git a/libs/langgraph/src/runtime/checkpoint-execution.spec.ts b/libs/langgraph/src/runtime/checkpoint-execution.spec.ts new file mode 100644 index 000000000..4b3d904ab --- /dev/null +++ b/libs/langgraph/src/runtime/checkpoint-execution.spec.ts @@ -0,0 +1,348 @@ +import type { ThreadState } from '@langchain/langgraph-sdk'; +import { describe, expect, it, vi } from 'vitest'; +import { createSession } from './create-session'; +import type { StreamEvent } from './transport.types'; +import { deferred } from './testing/deferred'; +import { controlledTransport } from './testing/controlled-transport'; + +import { + checkpointEvent, + fixture, + position, + saved, +} from './testing/checkpoint-fixture'; + +describe('checkpoint execution', () => { + it('acknowledges a confirmed tool-result follow-up that pauses so its owned interrupt can resume', async () => { + const f = fixture(); + const tools = saved('tools', [ + { + type: 'ai', + id: 'tool-step', + content: '', + tool_calls: [{ id: 'fresh', name: 'work', args: {} }], + }, + ]); + const pause = saved('pause', [], { + next: ['approval'], + tasks: [ + { + id: 'task', + name: 'approval', + result: null, + error: null, + interrupts: [{ id: 'decision', value: 'Approve?' }], + }, + ], + }); + const done = saved('done', [{ type: 'ai', id: 'done', content: 'Done' }]); + let calls = 0; + f.transport.stream = vi.fn(async function* (_a, _t, input, _s, options) { + const result = [tools, pause, done][calls++]; + if (result === pause) result.values = input as ThreadState['values']; + f.states.set(result.checkpoint.checkpoint_id!, result); + options?.onRunCreated?.({ run_id: String(result.metadata?.['run_id']) }); + yield checkpointEvent(result); + }); + const handler = vi.fn(() => 'Tool result'); + const session = createSession({ + assistantId: 'agent', + threadId: 'thread', + transport: f.transport, + tools: { work: { description: 'Work', followUp: true, handler } }, + }); + expect(await session.fork(position('a'), 'Run')).toBe('paused'); + expect(handler).toHaveBeenCalledOnce(); + expect(await session.resume(null)).toBe('success'); + expect(vi.mocked(f.transport.stream).mock.calls[1][4]?.checkpoint).toEqual( + position('tools') + ); + expect(vi.mocked(f.transport.stream).mock.calls[2][4]?.checkpoint).toEqual( + position('pause') + ); + }); + it('keeps historical echoes inert but blocks identical new-turn baseline call ID reuse monotonically', async () => { + const f = fixture(); + const call = { id: 'old', name: 'work', args: {} }; + const history = [ + { type: 'ai', id: 'historic', content: '', tool_calls: [call] }, + { type: 'tool', id: 'result', tool_call_id: 'old', content: 'done' }, + ]; + f.states.set('a', saved('a', history)); + let collide = false; + f.transport.stream = vi.fn(async function* (_a, _t, input, _s, options) { + options?.onRunCreated?.({ run_id: 'run-echo' }); + const result = saved('echo', [ + ...history, + ...(input as { messages: unknown[] }).messages, + { + type: 'ai', + id: collide ? 'new' : 'answer', + content: 'New answer', + ...(collide ? { tool_calls: [call] } : {}), + }, + ]); + f.states.set('echo', result); + yield checkpointEvent(result); + }); + const handler = vi.fn(() => 'should not run'); + const session = createSession({ + assistantId: 'agent', + threadId: 'thread', + transport: f.transport, + tools: { work: { description: 'Work', handler } }, + }); + expect(await session.fork(position('a'), 'Echo')).toBe('success'); + expect(handler).not.toHaveBeenCalled(); + collide = true; + expect(await session.submit('Collide')).toBe('interrupted'); + expect(session.getSnapshot().error?.message).toContain('identity conflict'); + await expect(session.load?.()).rejects.toThrow(); + await expect(session.submit('Again')).rejects.toThrow(); + expect(handler).not.toHaveBeenCalled(); + }); + it('is inert until commanded, adopts A atomically, and retains its position across submit/load/check', async () => { + const f = fixture(); + const session = createSession({ + assistantId: 'agent', + threadId: 'thread', + transport: f.transport, + }); + const views: string[][] = []; + session.subscribe(() => + views.push(session.getSnapshot().messages.map((m) => m.content)) + ); + expect(f.transport.getState).not.toHaveBeenCalled(); + expect(await session.fork(position('a'), 'Fork')).toBe('success'); + expect(views[0]).toEqual(['Source A', 'Answer A', 'Fork']); + expect(f.requests[0].options?.checkpoint).toEqual(position('a')); + expect(f.requests[0].options?.streamMode).toContain('checkpoints'); + expect(await session.submit('Next')).toBe('success'); + expect(f.requests[1].options?.checkpoint).toEqual(position('result-1')); + await session.load?.(); + expect(f.transport.getState).toHaveBeenLastCalledWith( + 'thread', + position('result-2'), + expect.any(AbortSignal) + ); + const reads = vi.mocked(f.transport.getState!).mock.calls.length; + await session.checkStatus?.(); + expect(f.transport.getState).toHaveBeenCalledTimes(reads); + expect(f.transport.getHistory).not.toHaveBeenCalled(); + expect(session.getSnapshot().history).toBeUndefined(); + }); + it('does not replace the prior snapshot or POST after stop during source preparation', async () => { + const f = fixture(); + const waiting = deferred(); + f.transport.getState = vi.fn(() => waiting.promise); + const session = createSession({ + assistantId: 'agent', + threadId: 'thread', + transport: f.transport, + }); + const before = session.getSnapshot(); + const run = session.fork(position('a'), 'Fork'); + await vi.waitFor(() => expect(f.transport.getState).toHaveBeenCalled()); + const signal = vi.mocked(f.transport.getState).mock.calls[0][2]; + await session.stop(); + expect(await run).toBe('aborted'); + expect(signal.aborted).toBe(true); + waiting.resolve(f.source); + await Promise.resolve(); + expect(session.getSnapshot()).toBe(before); + expect(f.transport.stream).not.toHaveBeenCalled(); + }); + it('owns input before awaiting and rejects overlapping branch consumers', async () => { + const f = fixture(); + const waiting = deferred(); + f.transport.getState = vi + .fn() + .mockReturnValueOnce(waiting.promise) + .mockImplementation(async (_thread, cp) => + f.states.get(cp.checkpoint_id) + ); + const session = createSession({ + assistantId: 'agent', + threadId: 'thread', + transport: f.transport, + }); + const input = { message: 'Fork', state: { preference: { color: 'blue' } } }; + const run = session.fork(position('a'), input); + await vi.waitFor(() => expect(f.transport.getState).toHaveBeenCalled()); + input.state.preference.color = 'red'; + await expect(session.submit('Overlap')).rejects.toThrow(); + await expect(session.load?.()).rejects.toThrow(); + await expect(session.fork(position('a'), 'Overlap')).rejects.toThrow(); + waiting.resolve(f.source); + expect(await run).toBe('success'); + expect(f.requests[0].input).toMatchObject({ + preference: { color: 'blue' }, + }); + }); + it('retains physical run evidence when stopped between terminal status and exact confirmation', async () => { + const f = fixture(); + const waiting = deferred(); + f.transport.getState = vi + .fn() + .mockImplementationOnce(async () => f.source) + .mockImplementationOnce(() => waiting.promise) + .mockImplementation(async (_thread, cp) => + f.states.get(cp.checkpoint_id) + ); + const session = createSession({ + assistantId: 'agent', + threadId: 'thread', + transport: f.transport, + }); + const run = session.fork(position('a'), 'Fork'); + await vi.waitFor(() => + expect(f.transport.getState).toHaveBeenCalledTimes(2) + ); + await session.stop(); + await session.stop(); + expect(await run).toBe('aborted'); + await expect(session.submit('Wrong')).rejects.toThrow(); + await expect(session.checkStatus?.()).rejects.toThrow(); + expect(session.getSnapshot().reconnect?.runId).toBe('run-result-1'); + expect(await session.reconnect()).toBe('success'); + expect(f.transport.joinStream).toHaveBeenCalledWith( + 'thread', + 'run-result-1', + 'cursor-result-1', + expect.any(AbortSignal), + { streamMode: expect.arrayContaining(['checkpoints']) } + ); + expect(f.transport.stream).toHaveBeenCalledTimes(1); + waiting.resolve(f.states.get('result-1')!); + }); + it('does not run a handler until the physical position is confirmed', async () => { + const f = fixture(); + const wire = controlledTransport(); + const confirmed = deferred(); + const handler = vi.fn(() => 'done'); + const result = saved('tools', [ + { + id: 'assistant-tools', + type: 'ai', + content: '', + tool_calls: [{ id: 'fresh', name: 'work', args: {} }], + }, + ]); + f.transport.stream = vi.fn((_a, _t, _p, _s, options) => { + options?.onRunCreated?.({ run_id: 'run-tools' }); + return wire.stream; + }); + f.transport.getState = vi + .fn() + .mockResolvedValueOnce(f.source) + .mockImplementation(() => confirmed.promise); + f.transport.updateState = vi.fn(async () => position('written')); + const session = createSession({ + assistantId: 'agent', + threadId: 'thread', + transport: f.transport, + tools: { work: { description: 'Work', handler, followUp: false } }, + }); + const run = session.fork(position('a'), 'Fork'); + await vi.waitFor(() => expect(f.transport.stream).toHaveBeenCalled()); + wire.release(checkpointEvent(result)); + await Promise.resolve(); + expect(handler).not.toHaveBeenCalled(); + wire.finish(); + await vi.waitFor(() => + expect(f.transport.getState).toHaveBeenCalledTimes(2) + ); + expect(handler).not.toHaveBeenCalled(); + confirmed.resolve(result); + expect(await run).toBe('success'); + expect(handler).toHaveBeenCalledOnce(); + expect(f.transport.updateState).toHaveBeenCalledWith( + 'thread', + expect.any(Object), + expect.any(AbortSignal), + { checkpoint: position('tools') } + ); + }); + it('rejects undefined branch resume and an already consumed paused task without mutation', async () => { + const f = fixture(); + const tasks = [ + { + id: 'task', + name: 'approval', + error: null, + result: null, + interrupts: [{ id: 'decision', value: { amount: 10 } }], + }, + ]; + const pause = saved('pause', [], { next: ['approval'], tasks }); + f.states.set('pause', pause); + f.transport.stream = vi.fn(async function* (_a, _t, _p, _s, options) { + options?.onRunCreated?.({ run_id: 'run-pause' }); + yield checkpointEvent(pause); + }); + const session = createSession({ + assistantId: 'agent', + threadId: 'thread', + transport: f.transport, + }); + expect(await session.fork(position('a'), 'Pause')).toBe('paused'); + const before = session.getSnapshot(); + const reads = vi.mocked(f.transport.getState!).mock.calls.length; + await expect(session.resume()).rejects.toThrow(); + expect(f.transport.getState).toHaveBeenCalledTimes(reads); + f.states.set( + 'pause', + saved('pause', [], { + next: ['approval'], + tasks: [{ ...tasks[0], result: { decision: true } }], + }) + ); + await expect(session.resume(null)).rejects.toThrow(); + expect(session.getSnapshot()).toBe(before); + expect(f.transport.stream).toHaveBeenCalledTimes(1); + }); + it('resumes an unconsumed exact paused position with explicit null and the command signal', async () => { + const f = fixture(); + const tasks = [ + { + id: 'task', + name: 'approval', + error: null, + result: null, + interrupts: [{ id: 'decision', value: { amount: 10 } }], + }, + ]; + const pause = saved('pause', [], { next: ['approval'], tasks }); + const complete = saved('done', [ + { type: 'ai', id: 'done', content: 'Decided' }, + ]); + f.states.set('pause', pause); + f.states.set('done', complete); + let calls = 0; + f.transport.stream = vi.fn(async function* (_a, _t, _p, _s, options) { + const result = calls++ ? complete : pause; + options?.onRunCreated?.({ run_id: String(result.metadata?.['run_id']) }); + yield checkpointEvent(result); + }); + const session = createSession({ + assistantId: 'agent', + threadId: 'thread', + transport: f.transport, + }); + expect(await session.fork(position('a'), 'Pause')).toBe('paused'); + const external = new AbortController(); + expect(await session.resume(null, { signal: external.signal })).toBe( + 'success' + ); + const stream = vi.mocked(f.transport.stream).mock.calls[1]; + expect(stream[2]).toBeNull(); + expect(stream[4]).toMatchObject({ + checkpoint: position('pause'), + command: { resume: null }, + }); + const preflight = vi.mocked(f.transport.getState!).mock.calls[2]; + expect(preflight[1]).toEqual(position('pause')); + expect(preflight[2]).not.toBe(external.signal); + expect(preflight[2].aborted).toBe(false); + }); +}); diff --git a/libs/langgraph/src/runtime/checkpoint-execution.type-test.ts b/libs/langgraph/src/runtime/checkpoint-execution.type-test.ts new file mode 100644 index 000000000..320532616 --- /dev/null +++ b/libs/langgraph/src/runtime/checkpoint-execution.type-test.ts @@ -0,0 +1,21 @@ +import type { Checkpoint, ThreadState } from '@langchain/langgraph-sdk'; +import type { LangGraphSession } from './create-session'; +import type { LangGraphHistoryEntry } from './langgraph-snapshot'; + +declare const session: LangGraphSession; +declare const checkpoint: Checkpoint; +declare const state: ThreadState; +declare const history: LangGraphHistoryEntry; + +void session.fork(checkpoint, 'Fork'); +void session.fork(state.checkpoint, { + message: 'Fork', + state: { locale: 'fr' }, +}); +void session.fork(history.checkpoint, 'Fork', { + signal: new AbortController().signal, +}); +// @ts-expect-error A complete reference is required; IDs are not aliases. +void session.fork('id', 'Fork'); +// @ts-expect-error Null input historical replay is not the fork capability. +void session.fork(checkpoint, null); diff --git a/libs/langgraph/src/runtime/checkpoint-persistence.spec.ts b/libs/langgraph/src/runtime/checkpoint-persistence.spec.ts new file mode 100644 index 000000000..4545d2646 --- /dev/null +++ b/libs/langgraph/src/runtime/checkpoint-persistence.spec.ts @@ -0,0 +1,154 @@ +import { describe, expect, it, vi } from 'vitest'; +import type { ToolExecutionStore } from '@threadplane/core/tools'; +import { createSession } from './create-session'; +import type { + AgentTransport, + OwnedCheckpointPosition, +} from './transport.types'; +import { deferred } from './testing/deferred'; +import { + checkpointEvent, + fixture, + position, + saved, +} from './testing/checkpoint-fixture'; + +const drain = () => new Promise((resolve) => setImmediate(resolve)); +function persistenceFixture() { + const f = fixture(); + const recording = [deferred(), deferred()]; + const writes = [ + deferred(), + deferred(), + ]; + const result = saved('tools', [ + { + type: 'ai', + id: 'tools', + content: '', + tool_calls: ['first', 'second'].map((id) => ({ + id, + name: 'work', + args: { id }, + })), + }, + ]); + f.states.set('tools', result); + f.transport.stream = vi.fn(async function* (_a, _t, _i, _s, options) { + options?.onRunCreated?.({ run_id: 'run-tools' }); + yield checkpointEvent(result); + }); + let writesStarted = 0; + const update = vi.fn>( + () => writes[writesStarted++].promise + ); + f.transport.updateState = update; + const executionStore: ToolExecutionStore = { + acquire: vi.fn(async () => ({ + status: 'acquired', + token: 'owner', + })), + settle: vi.fn(async ({ toolCallId }) => { + await recording[toolCallId === 'first' ? 0 : 1].promise; + return 'accepted'; + }), + }; + const handler = vi.fn(({ id }: { id: string }) => `Result ${id}`); + const session = createSession({ + assistantId: 'agent', + threadId: 'thread', + transport: f.transport, + executionStore, + tools: { work: { description: 'Work', handler, followUp: false } }, + }); + return { ...f, session, executionStore, handler, recording, writes, update }; +} + +describe('checkpoint persistence', () => { + it.each(['stop', 'dispose'] as const)( + 'chains two late writes through one captured owner after %s without publication', + async (command) => { + const f = persistenceFixture(); + const run = f.session.fork(position('a'), 'Tools'); + await vi.waitFor(() => + expect(f.executionStore.settle).toHaveBeenCalledTimes(2) + ); + await f.session[command](); + expect(await run).toBe('aborted'); + const snapshot = f.session.getSnapshot(); + const notify = vi.fn(); + f.session.subscribe(notify); + f.recording[0].resolve(); + await vi.waitFor(() => expect(f.update).toHaveBeenCalledTimes(1)); + expect(f.update.mock.calls[0][3]).toEqual({ + checkpoint: position('tools'), + }); + f.recording[1].resolve(); + await drain(); + expect(f.update).toHaveBeenCalledTimes(1); + if (command === 'stop') + await expect( + f.session.fork(position('a'), 'Overtake') + ).rejects.toThrow(); + f.writes[0].resolve(position('write-one')); + await vi.waitFor(() => expect(f.update).toHaveBeenCalledTimes(2)); + expect(f.update.mock.calls[1][3]).toEqual({ + checkpoint: position('write-one'), + }); + expect(f.update.mock.calls[1][1]).toMatchObject({ + messages: [{ tool_call_id: 'second' }], + }); + f.writes[1].resolve(position('write-two')); + await drain(); + expect(f.session.getSnapshot()).toBe(snapshot); + expect(notify).not.toHaveBeenCalled(); + expect(f.transport.stream).toHaveBeenCalledTimes(1); + if (command === 'stop') { + await f.session.submit('Next'); + expect( + vi.mocked(f.transport.stream).mock.calls[1][4]?.checkpoint + ).toEqual(position('write-two')); + } + await f.session.dispose(); + } + ); + it.each(['lost', 'void', 'wrong-thread'] as const)( + 'keeps a %s acknowledgment uncertain through late arrivals and all commands', + async (failure) => { + const f = persistenceFixture(); + const run = f.session.fork(position('a'), 'Tools'); + await vi.waitFor(() => + expect(f.executionStore.settle).toHaveBeenCalledTimes(2) + ); + await f.session.stop(); + expect(await run).toBe('aborted'); + f.recording[0].resolve(); + await vi.waitFor(() => expect(f.update).toHaveBeenCalledTimes(1)); + if (failure === 'lost') + f.writes[0].reject(new Error('Committed but acknowledgment lost')); + else + f.writes[0].resolve( + failure === 'void' + ? undefined + : { ...position('written'), thread_id: 'other' } + ); + await drain(); + f.recording[1].resolve(); + await drain(); + await f.session.stop(); + for (const command of [ + () => f.session.submit('Next'), + () => f.session.fork(position('a'), 'Fork'), + () => f.session.resume(null), + () => f.session.load?.(), + () => f.session.checkStatus?.(), + () => f.session.reconnect(), + ]) + await expect(command()).rejects.toThrow(); + expect(f.update).toHaveBeenCalledTimes(1); + expect(f.transport.stream).toHaveBeenCalledTimes(1); + expect(f.transport.getHistory).not.toHaveBeenCalled(); + await f.session.dispose(); + } + ); +}); diff --git a/libs/langgraph/src/runtime/checkpoint-recovery.spec.ts b/libs/langgraph/src/runtime/checkpoint-recovery.spec.ts new file mode 100644 index 000000000..5c551d386 --- /dev/null +++ b/libs/langgraph/src/runtime/checkpoint-recovery.spec.ts @@ -0,0 +1,93 @@ +import { describe, expect, it, vi } from 'vitest'; +import { createSession } from './create-session'; +import type { StreamEvent } from './transport.types'; +import { controlledTransport } from './testing/controlled-transport'; +import { + checkpointEvent, + fixture, + position, + saved, +} from './testing/checkpoint-fixture'; + +describe('checkpoint retained authority', () => { + it.each(['running', 'pending', 'interrupted', 'error', 'timeout'] as const)( + 'does not authorize an intermediate checkpoint under physical status %s', + async (status) => { + const f = fixture(); + f.transport.getRunStatus = vi.fn(async () => status); + const handler = vi.fn(() => 'done'); + const intermediate = saved('intermediate', [ + { + type: 'ai', + id: 'tool', + tool_calls: [{ id: 'work', name: 'work', args: {} }], + }, + ]); + f.states.set('intermediate', intermediate); + f.transport.stream = vi.fn(async function* (_a, _t, _i, _s, options) { + options?.onRunCreated?.({ run_id: 'run-intermediate' }); + yield checkpointEvent(intermediate); + }); + const session = createSession({ + assistantId: 'agent', + threadId: 'thread', + transport: f.transport, + tools: { work: { description: 'Work', handler, followUp: false } }, + }); + expect(await session.fork(position('a'), 'Run')).toBe('interrupted'); + expect(handler).not.toHaveBeenCalled(); + expect(f.transport.getState).toHaveBeenCalledTimes(1); + await expect(session.load?.()).rejects.toThrow(); + await expect(session.checkStatus?.()).rejects.toThrow(); + expect(f.transport.getHistory).not.toHaveBeenCalled(); + } + ); + it.each(['identity', 'cursor'] as const)( + 'keeps missing %s uncertain without advertising recovery', + async (missing) => { + const f = fixture(); + const stream = controlledTransport(); + f.transport.stream = vi.fn((_a, _t, _i, _s, options) => { + if (missing !== 'identity') + options?.onRunCreated?.({ run_id: 'run-b' }); + return stream.stream; + }); + const session = createSession({ + assistantId: 'agent', + threadId: 'thread', + transport: f.transport, + }); + const run = session.fork(position('a'), 'Run'); + await vi.waitFor(() => expect(f.transport.stream).toHaveBeenCalled()); + const event = checkpointEvent(saved('b')); + stream.release( + missing === 'cursor' ? { ...event, sseId: undefined } : event + ); + await new Promise((resolve) => setImmediate(resolve)); + await session.stop(); + expect(await run).toBe('aborted'); + expect(session.getSnapshot().reconnect).toBeUndefined(); + await expect(session.reconnect()).rejects.toThrow(); + await expect(session.submit('Wrong')).rejects.toThrow(); + } + ); + it('retains a candidate across a root stream error for explicit same-run reconnect', async () => { + const f = fixture(); + const result = saved('b'); + f.states.set('b', result); + f.transport.stream = vi.fn(async function* (_a, _t, _i, _s, options) { + options?.onRunCreated?.({ run_id: 'run-b' }); + yield checkpointEvent(result); + yield { type: 'error' as const, data: { message: 'connection failed' } }; + }); + const session = createSession({ + assistantId: 'agent', + threadId: 'thread', + transport: f.transport, + }); + expect(await session.fork(position('a'), 'Run')).toBe('error'); + expect(session.getSnapshot().reconnect?.runId).toBe('run-b'); + expect(await session.reconnect()).toBe('success'); + expect(f.transport.stream).toHaveBeenCalledTimes(1); + }); +}); diff --git a/libs/langgraph/src/runtime/checkpoint-state.spec.ts b/libs/langgraph/src/runtime/checkpoint-state.spec.ts new file mode 100644 index 000000000..1da21d0e6 --- /dev/null +++ b/libs/langgraph/src/runtime/checkpoint-state.spec.ts @@ -0,0 +1,34 @@ +import { expect, it } from 'vitest'; +import { + beginCheckpointEffect, + readyCheckpoint, + uncertainCheckpoint, +} from './checkpoint-state'; +it('invalidates old authority during every effect and never restores it after uncertainty', () => { + const position = { + thread_id: 'thread', + checkpoint_ns: '' as const, + checkpoint_id: 'a', + checkpoint_map: {}, + }; + const confirmed = { + position, + state: { + checkpoint: position, + values: {}, + next: [], + tasks: [], + created_at: '', + parent_checkpoint: null, + metadata: {}, + }, + paused: false, + }; + const ready = { kind: 'ready' as const, confirmed }; + expect(readyCheckpoint(ready).position.checkpoint_id).toBe('a'); + const inflight = beginCheckpointEffect(ready, 'run'); + expect(() => readyCheckpoint(inflight)).toThrow(); + const uncertain = uncertainCheckpoint(inflight); + expect(() => readyCheckpoint(uncertain)).toThrow(); + expect(uncertainCheckpoint(uncertain)).toBe(uncertain); +}); diff --git a/libs/langgraph/src/runtime/checkpoint-state.ts b/libs/langgraph/src/runtime/checkpoint-state.ts new file mode 100644 index 000000000..576672e17 --- /dev/null +++ b/libs/langgraph/src/runtime/checkpoint-state.ts @@ -0,0 +1,35 @@ +import type { ConfirmedCheckpoint } from './checkpoint-authority'; + +/** Position belongs to one session branch owner; physical run evidence remains + * in the existing attempt. An in-flight checkpoint can authorize no new effect. */ +export type CheckpointState = + | { readonly kind: 'ready'; readonly confirmed: ConfirmedCheckpoint } + | { readonly kind: 'run-inflight' | 'write-inflight' | 'uncertain' }; + +export interface CheckpointOwner { + state: CheckpointState; + readonly baselineCallIds: readonly string[]; + /** Only executeTool's settled outcomes in this owner establish this proof. + * Wire history, presentation, and provisional stop cancellation do not. */ + readonly settledToolIds: Set; +} + +export function readyCheckpoint(state: CheckpointState): ConfirmedCheckpoint { + if (state.kind !== 'ready') + throw new Error('Checkpoint execution authority is unavailable.'); + return state.confirmed; +} + +export function beginCheckpointEffect( + state: CheckpointState, + effect: 'run' | 'write' +): CheckpointState { + readyCheckpoint(state); + return { kind: effect === 'run' ? 'run-inflight' : 'write-inflight' }; +} + +export function uncertainCheckpoint(state: CheckpointState): CheckpointState { + return state.kind === 'uncertain' || state.kind === 'ready' + ? state + : { kind: 'uncertain' }; +} diff --git a/libs/langgraph/src/runtime/checkpoint-tool-evidence.spec.ts b/libs/langgraph/src/runtime/checkpoint-tool-evidence.spec.ts new file mode 100644 index 000000000..b16180e01 --- /dev/null +++ b/libs/langgraph/src/runtime/checkpoint-tool-evidence.spec.ts @@ -0,0 +1,369 @@ +import { describe, expect, it, vi } from 'vitest'; +import { createSession } from './create-session'; +import type { AgentTransport, StreamEvent } from './transport.types'; +import { + checkpointEvent, + fixture, + position, + saved, +} from './testing/checkpoint-fixture'; + +const call = { + type: 'ai', + id: 'tool-step', + content: '', + tool_calls: [{ id: 'fresh', name: 'work', args: { value: 'saved' } }], +}; +const toolMessage = { + type: 'tool', + id: 'tool-result', + tool_call_id: 'fresh', + content: 'Result', +}; +const values = (messages: unknown[], cursor: string): StreamEvent => ({ + type: 'values', + sseId: cursor, + data: { messages }, +}); + +describe('confirmed checkpoint tool evidence', () => { + it.each([ + 'missing-before', + 'missing-after', + 'changed-name', + 'changed-args', + 'phantom-resolution-before', + 'phantom-resolution-after', + 'omitted-authoritative-call', + ] as const)( + 'blocks effects and retains unavailable authority for %s', + async (mismatch) => { + const f = fixture(); + const missing = mismatch.startsWith('missing'); + const final = saved( + 'final', + missing ? [{ type: 'ai', id: 'final', content: 'Done' }] : [call] + ); + f.states.set('final', final); + f.transport.stream = vi.fn(async function* ( + _a, + _t, + input, + _s, + options + ) { + options?.onRunCreated?.({ run_id: 'run-final' }); + const user = (input as { messages: unknown[] }).messages; + if (mismatch === 'missing-before') + yield values([...user, call], 'before'); + if (mismatch === 'phantom-resolution-before') + yield values([...user, call, toolMessage], 'before'); + yield checkpointEvent(final); + if (mismatch === 'missing-after') + yield values([...user, call], 'after'); + if (mismatch === 'phantom-resolution-after') + yield values([toolMessage], 'after'); + if (mismatch === 'changed-name' || mismatch === 'changed-args') + yield values( + [ + { + ...call, + tool_calls: [ + { + ...call.tool_calls[0], + ...(mismatch === 'changed-name' + ? { name: 'other' } + : { args: { value: 'unsaved' } }), + }, + ], + }, + ], + 'after' + ); + if (mismatch === 'omitted-authoritative-call') + yield values([{ ...call, tool_calls: [] }], 'after'); + }); + f.transport.updateState = vi.fn(async () => position('written')); + const handler = vi.fn(() => 'Side effect'); + const session = createSession({ + assistantId: 'agent', + threadId: 'thread', + transport: f.transport, + tools: { + work: { description: 'Work', handler, followUp: false }, + other: { description: 'Other work', handler, followUp: false }, + }, + }); + const outcome = await session.fork(position('a'), 'Run'); + expect(handler).not.toHaveBeenCalled(); + expect(f.transport.updateState).not.toHaveBeenCalled(); + expect(outcome).not.toBe('success'); + expect(f.transport.getState).toHaveBeenCalledTimes(2); + await expect(session.submit('Next')).rejects.toThrow(); + await expect(session.load?.()).rejects.toThrow(); + await expect(session.checkStatus?.()).rejects.toThrow(); + expect(f.transport.stream).toHaveBeenCalledTimes(1); + await session.dispose(); + } + ); + + it.each(['unresolved', 'resolved', 'follow-up', 'partial-stream'] as const)( + 'accepts matching saved %s invocation evidence', + async (kind) => { + const f = fixture(); + const final = saved( + 'final', + kind === 'resolved' ? [call, toolMessage] : [call] + ); + f.states.set('final', final); + let streams = 0; + f.transport.stream = vi.fn(async function* ( + _a, + _t, + input, + _s, + options + ) { + const result = streams++ + ? saved('follow-up', [ + ...(input as { messages: unknown[] }).messages, + { type: 'ai', id: 'done', content: 'Done' }, + ]) + : final; + f.states.set(result.checkpoint.checkpoint_id!, result); + options?.onRunCreated?.({ + run_id: String(result.metadata?.['run_id']), + }); + if (kind === 'partial-stream') + yield { + type: 'messages', + sseId: 'partial', + messages: [ + { + type: 'AIMessageChunk', + id: 'tool-step', + content: 'Preparing', + tool_call_chunks: [ + { index: 0, id: 'fresh', name: 'work', args: '{' }, + ], + }, + ], + }; + yield checkpointEvent(result); + }); + f.transport.updateState = vi.fn(async () => position('written')); + const handler = vi.fn((args: { value: string }) => args.value); + const session = createSession({ + assistantId: 'agent', + threadId: 'thread', + transport: f.transport, + tools: { + work: { + description: 'Work', + handler, + followUp: kind === 'follow-up', + }, + }, + }); + expect(await session.fork(position('a'), 'Run')).toBe('success'); + expect(handler).toHaveBeenCalledTimes(kind === 'resolved' ? 0 : 1); + if (kind !== 'resolved') + expect(handler.mock.calls[0][0]).toEqual({ value: 'saved' }); + expect(f.transport.updateState).toHaveBeenCalledTimes( + kind === 'unresolved' || kind === 'partial-stream' ? 1 : 0 + ); + expect(f.transport.stream).toHaveBeenCalledTimes( + kind === 'follow-up' ? 2 : 1 + ); + if (kind === 'follow-up') + expect( + vi.mocked(f.transport.stream).mock.calls[1][4]?.checkpoint + ).toEqual(position('final')); + await session.dispose(); + } + ); +}); + +describe('branch-owned local tool settlement', () => { + it.each([ + { load: false, remove: false }, + { load: true, remove: false }, + { load: false, remove: true }, + { load: true, remove: true }, + ])( + 'resumes without repeating a consumed local result (load=$load, remove=$remove)', + async ({ load, remove }) => { + const f = fixture(); + let streams = 0; + let user: unknown; + let receipt: unknown; + f.transport.stream = vi.fn(async function* ( + _a, + _t, + input, + _s, + options + ) { + streams++; + if (streams === 1) + user = (input as { messages: unknown[] }).messages[0]; + if (streams === 2) + receipt = (input as { messages: unknown[] }).messages[0]; + const result = saved( + String(streams), + [ + user, + ...(streams === 3 && remove ? [] : [call]), + { + type: 'ai', + id: `answer-${streams}`, + content: streams === 3 ? 'Done' : 'Waiting', + }, + ], + streams === 2 + ? { + next: ['approval'], + tasks: [ + { + id: 'task', + name: 'approval', + error: null, + result: null, + interrupts: [{ id: 'decision', value: 'Approve?' }], + }, + ], + } + : {} + ); + f.states.set(String(streams), result); + options?.onRunCreated?.({ run_id: `run-${streams}` }); + if (streams === 3 && remove) + yield values([user, call], 'old-call-echo'); + yield checkpointEvent(result); + }); + const handler = vi.fn(() => 'ACTUAL RESULT'); + const session = createSession({ + assistantId: 'agent', + threadId: 'thread', + transport: f.transport, + tools: { work: { description: 'Work', handler, followUp: true } }, + }); + expect(await session.fork(position('a'), 'Run')).toBe('paused'); + expect(receipt).toMatchObject({ + tool_call_id: 'fresh', + content: 'ACTUAL RESULT', + }); + if (load) await session.load?.(); + expect(await session.resume(null)).toBe('success'); + expect(handler).toHaveBeenCalledOnce(); + expect(f.transport.stream).toHaveBeenCalledTimes(3); + await session.dispose(); + } + ); + + it.each(['name', 'args'] as const)( + 'keeps changed settled %s identities as sticky conflicts after load', + async (changed) => { + const f = fixture(); + let streams = 0; + let user: unknown; + f.transport.stream = vi.fn(async function* ( + _a, + _t, + input, + _s, + options + ) { + streams++; + if (streams === 1) + user = (input as { messages: unknown[] }).messages[0]; + const current = + streams === 3 + ? { + ...call, + tool_calls: [ + { + ...call.tool_calls[0], + ...(changed === 'name' + ? { name: 'other' } + : { args: { value: 'changed' } }), + }, + ], + } + : call; + const result = saved( + String(streams), + [user, current], + streams === 2 + ? { + next: ['approval'], + tasks: [ + { + id: 'task', + name: 'approval', + error: null, + result: null, + interrupts: [{ id: 'decision', value: 'Approve?' }], + }, + ], + } + : {} + ); + f.states.set(String(streams), result); + options?.onRunCreated?.({ run_id: `run-${streams}` }); + yield checkpointEvent(result); + }); + const handler = vi.fn(() => 'ACTUAL RESULT'); + const session = createSession({ + assistantId: 'agent', + threadId: 'thread', + transport: f.transport, + tools: { work: { description: 'Work', handler, followUp: true } }, + }); + expect(await session.fork(position('a'), 'Run')).toBe('paused'); + await session.load?.(); + expect(await session.resume(null)).toBe('interrupted'); + expect(session.getSnapshot().error?.message).toContain( + 'identity conflict' + ); + await expect(session.fork(position('a'), 'Again')).rejects.toThrow(); + expect(handler).toHaveBeenCalledOnce(); + await session.dispose(); + } + ); + + it('does not carry local settlement proof into a newly selected branch owner', async () => { + const f = fixture(); + let streams = 0; + let user: unknown; + f.transport.stream = vi.fn(async function* ( + _a, + _t, + input, + _s, + options + ) { + streams++; + if (streams !== 2) user = (input as { messages: unknown[] }).messages[0]; + const result = saved(String(streams), [user, call]); + f.states.set(String(streams), result); + options?.onRunCreated?.({ run_id: `run-${streams}` }); + yield checkpointEvent(result); + }); + const handler = vi.fn(() => 'ACTUAL RESULT'); + const session = createSession({ + assistantId: 'agent', + threadId: 'thread', + transport: f.transport, + tools: { work: { description: 'Work', handler, followUp: true } }, + }); + expect(await session.fork(position('a'), 'First branch')).toBe('success'); + expect(await session.fork(position('a'), 'Second branch')).not.toBe( + 'success' + ); + expect(handler).toHaveBeenCalledOnce(); + expect(f.transport.stream).toHaveBeenCalledTimes(3); + await expect(session.submit('Again')).rejects.toThrow(); + await session.dispose(); + }); +}); diff --git a/libs/langgraph/src/runtime/checkpoint-tool-evidence.ts b/libs/langgraph/src/runtime/checkpoint-tool-evidence.ts new file mode 100644 index 000000000..de9c36dab --- /dev/null +++ b/libs/langgraph/src/runtime/checkpoint-tool-evidence.ts @@ -0,0 +1,80 @@ +import type { ThreadState } from '@langchain/langgraph-sdk'; +import { initialMessageState, type MessageState } from './message-reducer'; +import { sameToolInvocation } from './ownership'; +import { projectStream, type StreamProjection } from './stream-projection'; + +/** Confirm execution eligibility against only the exact saved messages. Reuse + * turn selection, but none of the stream's accumulated calls, completions, + * canonical messages, or transient turn exclusions may become saved evidence. + * This projection is never installed: authored results and monotone invocation + * conflicts remain owned by the session's existing message state. */ +export function assertCheckpointToolEvidence( + saved: ThreadState, + state: MessageState, + projection: StreamProjection, + resolvedTools: ReadonlySet, + locallySettledTools: ReadonlySet +): void { + const authoritative = projectStream( + initialMessageState(), + { + generation: projection.generation, + messageIdPrefix: projection.messageIdPrefix, + userId: projection.userId, + baselineIds: projection.baselineIds, + sawAssistant: false, + terminal: false, + paused: false, + canonical: [], + ...(projection.resume + ? { resume: { turnIds: projection.resume.turnIds } } + : {}), + }, + { type: 'values', data: saved.values } + ); + const observedIds = new Set(projection.toolCallIds ?? []); + const savedIds = new Set(authoritative.projection.toolCallIds ?? []); + const observedCalls = new Map(state.toolCalls.map((call) => [call.id, call])); + const savedCalls = new Map( + authoritative.state.toolCalls.map((call) => [call.id, call]) + ); + const mismatch = () => + new Error('Saved checkpoint tool evidence does not match the owned run.'); + const invocations = new Map(state.invocations.map((call) => [call.id, call])); + // A graph may transform or remove an already consumed local result or call. + // Authenticate those echoes against this owner's settlement proof and the + // monotone invocation ledger before excluding them from new effect admission. + // The broader resolvedTools set also contains wire history and cannot do this. + for (const id of new Set([...observedIds, ...savedIds])) { + if (!locallySettledTools.has(id)) continue; + const invocation = invocations.get(id); + const observed = observedCalls.get(id); + const persisted = savedCalls.get(id); + if ( + !invocation || + invocation.conflicted || + (observedIds.has(id) && !observed) || + (savedIds.has(id) && !persisted) || + (observed && !sameToolInvocation(observed, invocation)) || + (persisted && !sameToolInvocation(persisted, invocation)) + ) + throw mismatch(); + observedIds.delete(id); + savedIds.delete(id); + } + if (observedIds.size !== savedIds.size) throw mismatch(); + for (const id of observedIds) { + const observed = observedCalls.get(id); + const persisted = savedCalls.get(id); + if ( + !savedIds.has(id) || + !observed || + !persisted || + !sameToolInvocation(observed, persisted) + ) + throw mismatch(); + const observedResolved = + resolvedTools.has(id) || observed.status !== 'pending'; + if (observedResolved !== (persisted.status !== 'pending')) throw mismatch(); + } +} diff --git a/libs/langgraph/src/runtime/checkpoint-transport.spec.ts b/libs/langgraph/src/runtime/checkpoint-transport.spec.ts new file mode 100644 index 000000000..9cf33314b --- /dev/null +++ b/libs/langgraph/src/runtime/checkpoint-transport.spec.ts @@ -0,0 +1,153 @@ +import { afterEach, describe, expect, it, vi } from 'vitest'; +import { FetchStreamTransport } from '../lib/transport/fetch-stream.transport'; + +afterEach(() => vi.unstubAllGlobals()); +const checkpoint = { + thread_id: 'thread', + checkpoint_ns: '' as const, + checkpoint_id: 'a', + checkpoint_map: { '': 'a' }, +}; +describe('exact checkpoint transport', () => { + it('reads the exact checkpoint and normalizes a write response without leaking configuration', async () => { + const request = vi.fn(async (url) => + String(url).endsWith('/checkpoint') + ? Response.json({ checkpoint, values: {}, next: [], tasks: [] }) + : Response.json({ + configurable: { + ...checkpoint, + checkpoint_id: 'written', + secret: 'omit', + }, + }) + ); + vi.stubGlobal('fetch', request); + const transport = new FetchStreamTransport( + 'https://runtime.example', + undefined, + { maxRetries: 0 } + ); + const signal = new AbortController().signal; + expect( + (await transport.getState('thread', checkpoint, signal)).checkpoint + ).toEqual(checkpoint); + expect( + await transport.updateState('thread', { messages: [] }, signal, { + asNode: '__start__', + checkpoint, + }) + ).toEqual({ ...checkpoint, checkpoint_id: 'written' }); + expect( + request.mock.calls.map(([url, init]) => ({ + url: String(url), + body: JSON.parse(String(init?.body)), + })) + ).toEqual([ + { + url: 'https://runtime.example/threads/thread/state/checkpoint', + body: { checkpoint }, + }, + { + url: 'https://runtime.example/threads/thread/state', + body: { values: { messages: [] }, checkpoint, as_node: '__start__' }, + }, + ]); + expect( + request.mock.calls.every( + ([, init]) => init?.signal instanceof AbortSignal + ) + ).toBe(true); + }); + it('permits a legacy successful write without routing and requests checkpoint delivery when joining', async () => { + const request = vi.fn(async (url) => + String(url).includes('/stream') + ? new Response('', { headers: { 'content-type': 'text/event-stream' } }) + : Response.json({}) + ); + vi.stubGlobal('fetch', request); + const transport = new FetchStreamTransport( + 'https://runtime.example', + undefined, + { maxRetries: 0 } + ); + const signal = new AbortController().signal; + expect(await transport.updateState('thread', {}, signal)).toBeUndefined(); + for await (const event of transport.joinStream( + 'thread', + 'run', + 'cursor', + signal, + { streamMode: ['values', 'checkpoints'] } + )) + void event; + expect( + JSON.parse( + new URL(String(request.mock.calls[1][0])).searchParams.get( + 'stream_mode' + )! + ) + ).toEqual(['values', 'checkpoints']); + expect( + new Headers(request.mock.calls[1][1]?.headers).get('last-event-id') + ).toBe('cursor'); + }); + it('protects exact-state errors and does not retry a lost write acknowledgment', async () => { + const request = vi.fn(async () => { + throw new Error('NetworkError secret'); + }); + vi.stubGlobal('fetch', request); + const transport = new FetchStreamTransport( + 'https://runtime.example', + undefined, + { maxRetries: 0, defaultHeaders: { authorization: 'secret' } } + ); + const signal = new AbortController().signal; + await expect( + transport.getState('thread', checkpoint, signal) + ).rejects.not.toThrow('secret'); + await expect( + transport.updateState('thread', {}, signal, { checkpoint }) + ).rejects.not.toThrow('secret'); + expect(request).toHaveBeenCalledTimes(2); + }); + it.each(['read', 'write'] as const)( + 'forwards cancellation to an outstanding exact %s request', + async (operation) => { + const aborted = vi.fn(); + const request = vi.fn( + async (_url, init) => + new Promise((_resolve, reject) => { + init?.signal?.addEventListener( + 'abort', + () => { + aborted(); + reject(init.signal?.reason); + }, + { once: true } + ); + }) + ); + vi.stubGlobal('fetch', request); + const transport = new FetchStreamTransport( + 'https://runtime.example', + undefined, + { maxRetries: 0, defaultHeaders: {} } + ); + const controller = new AbortController(); + const result = + operation === 'read' + ? transport.getState('thread', checkpoint, controller.signal) + : transport.updateState('thread', {}, controller.signal, { + checkpoint, + }); + const rejection = expect(result).rejects.toMatchObject({ + name: 'AbortError', + }); + await vi.waitFor(() => expect(request).toHaveBeenCalledTimes(1)); + controller.abort(); + await rejection; + expect(aborted).toHaveBeenCalledOnce(); + expect(request).toHaveBeenCalledTimes(1); + } + ); +}); diff --git a/libs/langgraph/src/runtime/create-session.ts b/libs/langgraph/src/runtime/create-session.ts index 0fb7a2b89..b671cf9f7 100644 --- a/libs/langgraph/src/runtime/create-session.ts +++ b/libs/langgraph/src/runtime/create-session.ts @@ -63,6 +63,25 @@ import { import { createSafeRequestError } from './operation-errors'; import { createToolPersistence } from './tool-persistence'; import { hasInvocationConflict } from './tool-invocations'; +import { observeInvocation } from './tool-invocations'; +import { + assertResumeCheckpoint, + captureCheckpoint, + captureCheckpointEvent, + captureCheckpointState, + captureCompletedCheckpoint, + confirmCheckpoint, + type CheckpointCandidate, + type CheckpointReference, + type OwnedCheckpointPosition, +} from './checkpoint-authority'; +import { + beginCheckpointEffect, + readyCheckpoint, + uncertainCheckpoint, + type CheckpointOwner, +} from './checkpoint-state'; +import { assertCheckpointToolEvidence } from './checkpoint-tool-evidence'; import { failureProjection, finalizeProjection, @@ -109,6 +128,11 @@ export type LangGraphSession< input: LangGraphSubmitInput, options?: LangGraphRunOptions ): Promise; + fork( + checkpoint: CheckpointReference, + input: LangGraphSubmitInput, + options?: LangGraphRunOptions + ): Promise; reconnect(options?: { readonly signal?: AbortSignal; }): Promise; @@ -127,6 +151,13 @@ interface HistoryRead { unlink?: () => void; } +interface Preparation extends Omit { + readonly result: Promise; + readonly resolve: (outcome: CompleteOutcome) => void; + readonly checkpoint: OwnedCheckpointPosition; + readonly adopt: (saved: ThreadState) => () => Attempt; +} + type AttemptInput = | { readonly kind: 'submit'; @@ -141,6 +172,8 @@ interface PhysicalRun { captureOpen: boolean; received: boolean; confirmed: boolean; + checkpoint?: CheckpointCandidate; + positionConfirmed?: boolean; } interface Attempt { @@ -150,6 +183,7 @@ interface Attempt { readonly resolve: (outcome: CompleteOutcome) => void; readonly input: AttemptInput; readonly runOptions: CapturedRunOptions | undefined; + readonly branch?: CheckpointOwner; projection: StreamProjection; subgraphs: SubgraphObservation; readonly calls: Map; @@ -205,18 +239,23 @@ export function createSession( // Execution dedupe survives transcript replacement. Authored result provenance // belongs only to the current transcript; a wire string cannot restore it. const authoredTools = new Set(); + const suppliedTransport = options.transport; + const ownedClientOptions = suppliedTransport + ? undefined + : { ...options.clientOptions }; + const ownedRetries = ownedClientOptions?.maxRetries ?? 0; const transport = - options.transport ?? + suppliedTransport ?? new FetchStreamTransport(options.apiUrl ?? '', undefined, { - ...options.clientOptions, - maxRetries: options.clientOptions?.maxRetries ?? 0, + ...ownedClientOptions, + maxRetries: ownedRetries, }); const protectedTransport = transport instanceof FetchStreamTransport && transport.protectsOperationErrors; - const persistence = createToolPersistence( + const persistence = createToolPersistence( buffer, - async (messages, signal) => { + async (messages, signal, branch) => { // A queued flush can start after history/stream observation found a conflict. // Durable claim/record ownership may finish, but its results stay quarantined. assertInvocationIdentity(); @@ -224,7 +263,26 @@ export function createSession( throw new Error( 'Persisting terminal tool results requires transport.updateState().' ); - await transport.updateState(threadId, { messages }, signal); + const confirmed = branch && readyCheckpoint(branch.state); + if (branch) branch.state = beginCheckpointEffect(branch.state, 'write'); + try { + const result = await transport.updateState( + threadId, + { messages }, + signal, + confirmed ? { checkpoint: confirmed.position } : undefined + ); + if (branch && confirmed) { + const position = captureCheckpoint(result, threadId); + branch.state = { + kind: 'ready', + confirmed: { ...confirmed, position }, + }; + } + } catch (error) { + if (branch) branch.state = uncertainCheckpoint(branch.state); + throw error; + } } ); const getHistory = @@ -234,6 +292,7 @@ export function createSession( const canCheck = !!getHistory; const joinStream = transport.joinStream?.bind(transport); const getRunStatus = transport.getRunStatus?.bind(transport); + const getState = transport.getState?.bind(transport); const canReconnect = !!joinStream && !!getRunStatus; const publication = createPublication({ status: 'idle', @@ -258,6 +317,15 @@ export function createSession( let checkController: AbortController | undefined; let loading: HistoryRead | undefined; let pendingToolSettlements = 0; + let branch: CheckpointOwner | undefined; + let preparing: Preparation | undefined; + const branchStreamModes = [ + 'values', + 'messages-tuple', + 'updates', + 'custom', + 'checkpoints', + ] as const; function invocationConflictError(): AgentError { return { @@ -287,12 +355,33 @@ export function createSession( } function admitSubmission() { assertInvocationIdentity(); + if (branch || preparing) admitBranchCommand(); if (unsettledTools.size) throw new Error('Submission cannot replace unsettled tool execution.'); if (persistence.pending) throw new Error('Submission cannot replace pending tool persistence.'); } + function admitBranchCommand() { + assertInvocationIdentity(); + if ( + owner || + loading || + checkController || + preparing || + recoveryAttempt || + retained || + pendingToolSettlements || + persistence.pending || + unsettledTools.size || + buffer.snapshot().messages.length + ) + throw new Error( + 'Checkpoint commands cannot replace active or unsettled work.' + ); + if (branch) readyCheckpoint(branch.state); + } + const owns = (attempt: Attempt) => owner === attempt && !disposed; function publish(status: 'idle' | 'running' | 'error', error?: AgentError) { if (hasInvocationConflict(state.invocations)) { @@ -345,6 +434,17 @@ export function createSession( if (abort) read.controller.abort(); unlink?.(); } + function detachPreparation() { + const previous = preparing; + preparing = undefined; + previous?.resolve('aborted'); + return previous; + } + function closePreparation(preparation: Preparation | undefined) { + if (!preparation) return; + preparation.controller.abort(); + preparation.unlink?.(); + } function close(attempt: Attempt, abort = false, final = true) { const unlink = final ? attempt.unlink : undefined; if (final) attempt.unlink = undefined; @@ -369,6 +469,7 @@ export function createSession( retainRun = false ) { if (!owns(attempt)) return; + retainRun ||= !!attempt.branch; if (hasInvocationConflict(state.invocations)) { outcome = 'interrupted'; error = invocationConflictError(); @@ -379,7 +480,7 @@ export function createSession( retainRun && canReconnect && run && - !run.confirmed && + !(attempt.branch ? run.positionConfirmed : run.confirmed) && !run.evidence.unsafe && run.evidence.runId && run.evidence.cursor @@ -387,6 +488,8 @@ export function createSession( : undefined; detach(outcome); retained = candidate; + if (attempt.branch) + attempt.branch.state = uncertainCheckpoint(attempt.branch.state); recoveryAttempt = error?.recovery === 'check' ? attempt : undefined; publish(error ? 'error' : 'idle', error); // An error can end local consumption before the HTTP body closes. SDK @@ -489,6 +592,7 @@ export function createSession( return; } const result = outcome.result; + attempt.branch?.settledToolIds.add(call.id); resolvedTools.add(call.id); authoredTools.add(call.id); buffer.stage(call.id, result); @@ -504,7 +608,10 @@ export function createSession( // Required durable cleanup may finish after stop/dispose. It can only // persist results; it has no route back to publication or run creation. try { - await persistence.flush(new AbortController().signal); + await persistence.flush( + new AbortController().signal, + attempt.branch + ); } catch { /* The staged result remains available for explicit handoff. */ } @@ -521,15 +628,33 @@ export function createSession( : 'complete'; } function stopExecution() { + const preparation = detachPreparation(); const hadRetained = !!retained; + const previousRetained = retained; + const previousOwner = owner; + const run = previousOwner?.physical; + const retainBranch = + previousOwner?.branch && + run && + !run.positionConfirmed && + !run.evidence.unsafe && + run.evidence.runId && + run.evidence.cursor && + canReconnect + ? { attempt: previousOwner, run } + : undefined; retained = undefined; const reading = detachLoad(); const checking = invalidateCheck(); const attempt = detach('aborted'); + if (attempt?.branch) + attempt.branch.state = uncertainCheckpoint(attempt.branch.state); + retained = retainBranch ?? (branch ? previousRetained : undefined); if (attempt || hadRetained) publish('idle'); checking?.abort(); if (attempt) close(attempt, true); closeLoad(reading); + closePreparation(preparation); } function reconcile(attempt: Attempt, history: ThreadState[]) { @@ -625,6 +750,7 @@ export function createSession( async function execute(attempt: Attempt): Promise { if (!owns(attempt)) return; const attemptCanCheck = () => + !attempt.branch && attempt.input.kind === 'submit' && canCheck && !attempt.physical?.evidence.runId && @@ -659,6 +785,15 @@ export function createSession( ? attempt.input.state : undefined ); + const from = + attempt.branch && !joining + ? readyCheckpoint(attempt.branch.state).position + : undefined; + if (attempt.branch && !joining) + attempt.branch.state = beginCheckpointEffect( + attempt.branch.state, + 'run' + ); // Ownership is captured before the first effect. The signal always belongs // to us, even when the caller also supplied an external AbortSignal. const capture = (metadata: { run_id: string; thread_id?: string }) => { @@ -683,7 +818,10 @@ export function createSession( threadId, run.evidence.runId, joining, - attempt.controller.signal + attempt.controller.signal, + ...(attempt.branch + ? [{ streamMode: [...branchStreamModes] }] + : []) ) : transport.stream( assistantId, @@ -697,6 +835,12 @@ export function createSession( attempt.input.value !== undefined) ? { ...attempt.runOptions, + ...(from + ? { + checkpoint: from, + streamMode: [...branchStreamModes], + } + : {}), ...(canReconnect ? { streamResumable: true, @@ -741,6 +885,9 @@ export function createSession( return; } const previousEvidence = run.evidence; + const checkpoint = attempt.branch + ? captureCheckpointEvent(event, threadId) + : undefined; const cursor = advanceCursor(previousEvidence, event, joining); if (!owns(attempt)) return; if (cursor.replay) { @@ -772,6 +919,7 @@ export function createSession( run.evidence === previousEvidence ? cursor.evidence : { unsafe: true }; + if (checkpoint) run.checkpoint = checkpoint; state = projected.state; values = projectedValues; interrupts = projectedInterrupts; @@ -809,7 +957,43 @@ export function createSession( /* Exact-run inspection failed; history cannot replace it. */ } if (!owns(attempt)) return; - if ( + if (attempt.branch) { + if (status !== 'success' || !run.checkpoint || !getState) { + outcome = 'interrupted'; + } else { + run.confirmed = true; + const saved = await getState( + threadId, + run.checkpoint.position, + attempt.controller.signal + ); + if (!owns(attempt)) return; + const confirmed = confirmCheckpoint( + run.checkpoint, + saved, + run.evidence.runId + ); + assertCheckpointToolEvidence( + confirmed.state, + state, + attempt.projection, + resolvedTools, + attempt.branch.settledToolIds + ); + const projectedInterrupts = projectHistoryInterrupts(interrupts, [ + confirmed.state, + ]); + if (!owns(attempt)) return; + attempt.branch.state = { kind: 'ready', confirmed }; + run.positionConfirmed = true; + interrupts = projectedInterrupts; + attempt.projection = { + ...attempt.projection, + paused: confirmed.paused, + }; + outcome = confirmed.paused ? 'paused' : 'success'; + } + } else if ( (status === 'success' || status === 'interrupted') && attempt.projection.paused ) { @@ -836,7 +1020,8 @@ export function createSession( ); return; } else outcome = 'interrupted'; - } else if (run.evidence.unsafe) outcome = 'interrupted'; + } else if (run.evidence.unsafe || attempt.branch) + outcome = 'interrupted'; if ( outcome === 'interrupted' && attemptCanCheck() && @@ -870,6 +1055,10 @@ export function createSession( state = finalizeProjection(state, attempt.projection); subgraphs = settleSubgraphs(subgraphs, outcome, true); attempt.subgraphs = subgraphs; + // The exact saved pause confirms the submitted handoff just as a + // completed run does. Keeping that batch would block its own resume. + if (outcome === 'paused' && attempt.branch && run.positionConfirmed) + persistence.acknowledge(batch); } if (outcome === 'success') { state = reduceMessages(state, { @@ -884,7 +1073,10 @@ export function createSession( // Persist legitimate mixed-group results without continuing past // an unavailable call. Failed writes keep the exact staged buffer. try { - await persistence.flush(attempt.controller.signal); + await persistence.flush( + attempt.controller.signal, + attempt.branch + ); } catch { /* retained */ } @@ -912,6 +1104,9 @@ export function createSession( terminal: false, paused: false, canonical: [], + ...(attempt.branch + ? { baselineCallIds: attempt.branch.baselineCallIds } + : {}), ...(attempt.input.kind === 'resume' ? { resume: { @@ -925,7 +1120,7 @@ export function createSession( continue; } try { - await persistence.flush(attempt.controller.signal); + await persistence.flush(attempt.controller.signal, attempt.branch); } catch { if (owns(attempt)) settle(attempt, 'error', { @@ -1007,6 +1202,7 @@ export function createSession( handoffIds: [], input, runOptions, + branch, subgraphs, projection: { generation, @@ -1016,6 +1212,7 @@ export function createSession( terminal: false, paused: false, canonical: [], + ...(branch ? { baselineCallIds: branch.baselineCallIds } : {}), ...(input.kind === 'resume' ? { resume: { turnIds: turnIds(userId) } } : {}), @@ -1069,6 +1266,150 @@ export function createSession( }); } + function prepare( + checkpoint: OwnedCheckpointPosition, + external: AbortSignal | undefined, + adopt: Preparation['adopt'] + ) { + let resolve!: Preparation['resolve']; + let reject!: Preparation['reject']; + const result = new Promise((done, failed) => { + resolve = done; + reject = failed; + }); + const created: Preparation = { + controller: new AbortController(), + result, + resolve, + reject, + checkpoint, + adopt, + }; + preparing = created; + if (external) { + const abort = () => { + void publication.command(() => { + if (preparing === created) stopExecution(); + }); + }; + created.unlink = () => external.removeEventListener('abort', abort); + external.addEventListener('abort', abort, { once: true }); + if (external.aborted) abort(); + } + return created; + } + + async function runPreparation(preparation: Preparation) { + if (preparing !== preparation || disposed || !getState) return; + try { + const saved = await getState( + threadId, + preparation.checkpoint, + preparation.controller.signal + ); + let attempt: Attempt | undefined; + await publication.command(() => { + if (preparing !== preparation || disposed) return; + const commit = preparation.adopt(saved); + if (preparing !== preparation || disposed) return; + preparing = undefined; + attempt = commit(); + preparation.unlink?.(); + void attempt.result.then(preparation.resolve); + }); + if (attempt) void execute(attempt); + } catch { + await publication.command(() => { + if (preparing !== preparation || disposed) return; + preparing = undefined; + preparation.reject(createSafeRequestError()); + closePreparation(preparation); + }); + } + } + + function fork( + checkpoint: CheckpointReference, + input: LangGraphSubmitInput, + options?: LangGraphRunOptions + ): Promise { + let preparation: Preparation | undefined; + const beginning = publication.command(() => { + if (disposed) return; + admitBranchCommand(); + if (!getState || !getRunStatus || !joinStream) + throw new Error( + 'Fork requires exact checkpoint and physical-run transport capabilities.' + ); + if (!suppliedTransport && ownedRetries > 0) + throw new Error( + 'Fork requires an SDK transport with maxRetries set to zero.' + ); + const capturedRevision = revision; + const external = options?.signal; + const current = () => + !disposed && !external?.aborted && capturedRevision === revision; + if (!current()) return; + const position = captureCheckpoint(checkpoint, threadId); + if (!current()) return; + const captured = captureSubmitInput(input); + if (!current()) return; + const runOptions = captureRunOptions(options); + if (!current()) return; + admitBranchCommand(); + preparation = prepare(position, external, (raw) => { + const source = captureCompletedCheckpoint(raw, position); + let projected = projectHistory(state, [source.state], { + interrupts: [], + ...(typedTools + ? { registeredTools: new Set(definitions.keys()) } + : {}), + }); + let invocations = projected.invocations; + for (const call of source.calls) + invocations = observeInvocation(invocations, call, true); + projected = Object.freeze({ ...projected, invocations }); + if (hasInvocationConflict(invocations)) + throw new Error(invocationConflictError().message); + const projectedValues = projectHistoryValues(values, [source.state]); + return () => { + state = projected; + values = projectedValues; + authoredTools.clear(); + for (const call of source.calls) resolvedTools.add(call.id); + branch = { + state: { + kind: 'ready', + confirmed: { position, state: source.state, paused: false }, + }, + baselineCallIds: source.calls.map((call) => call.id), + settledToolIds: new Set(), + }; + return beginAttempt( + { + kind: 'submit', + state: captured.state, + messages: [ + { + type: 'human', + id: crypto.randomUUID(), + content: captured.message, + }, + ], + }, + external, + runOptions + ); + }; + }); + }); + return beginning.then(() => { + if (!preparation) return 'aborted'; + void runPreparation(preparation); + return preparation.result; + }); + } + function submit( input: LangGraphSubmitInput, options?: LangGraphRunOptions @@ -1129,6 +1470,7 @@ export function createSession( options?: LangGraphRunOptions ): Promise { let attempt: Attempt | undefined; + let preparation: Preparation | undefined; const beginning = publication.command(() => { const capturedRevision = revision; const capturedLoad = loading; @@ -1142,6 +1484,7 @@ export function createSession( return; const admit = () => { assertInvocationIdentity(); + if (branch || preparing) admitBranchCommand(); if ( owner || loading || @@ -1159,6 +1502,8 @@ export function createSession( throw new Error('Resume requires an observed interrupt.'); }; admit(); + if (branch && value === undefined) + throw new Error('Branch resume requires a defined response value.'); const captured = ownValue(value); if ( disposed || @@ -1177,13 +1522,36 @@ export function createSession( ) return; admit(); + if (branch) { + if (!getState) + throw new Error('Branch resume requires an exact checkpoint read.'); + const confirmed = readyCheckpoint(branch.state); + preparation = prepare(confirmed.position, external, (saved) => { + assertResumeCheckpoint(confirmed, saved); + return () => + beginAttempt( + { kind: 'resume', value: captured }, + external, + runOptions + ); + }); + return; + } attempt = beginAttempt( { kind: 'resume', value: captured }, external, runOptions ); }); - return dispatch(beginning, () => attempt); + return beginning.then(() => { + if (preparation) { + void runPreparation(preparation); + return preparation.result; + } + if (!attempt) return 'aborted'; + void execute(attempt); + return attempt.result; + }); } function reconnect(options?: { @@ -1199,6 +1567,7 @@ export function createSession( if ( owner || loading || + preparing || checkController || pendingToolSettlements || persistence.pending || @@ -1246,6 +1615,7 @@ export function createSession( resolve, input: candidate.attempt.input, runOptions: candidate.attempt.runOptions, + branch: candidate.attempt.branch, calls: new Map(), groups: candidate.attempt.groups, handoffIds: candidate.attempt.handoffIds, @@ -1277,9 +1647,23 @@ export function createSession( } async function readHistory(read: HistoryRead) { - if (!ownsLoad(read) || !getHistory) return; + if (!ownsLoad(read)) return; try { - const history = await getHistory(threadId, read.controller.signal); + const capturedBranch = branch; + const confirmed = capturedBranch && readyCheckpoint(capturedBranch.state); + const history = + confirmed && getState + ? [ + captureCheckpointState( + await getState( + threadId, + confirmed.position, + read.controller.signal + ), + confirmed.position + ), + ] + : await getHistory!(threadId, read.controller.signal); await publication.command(() => { if (!ownsLoad(read)) return; const previousState = state; @@ -1313,7 +1697,7 @@ export function createSession( state = projected; values = projectedValues; interrupts = projectedInterrupts; - historyPage = projectedHistory; + if (!capturedBranch) historyPage = projectedHistory; subgraphs = initialSubgraphs(); authoredTools.clear(); loading = undefined; @@ -1338,6 +1722,11 @@ export function createSession( let read: HistoryRead | undefined; const beginning = publication.command(() => { if (disposed || options?.signal?.aborted) return; + if (!branch && !getHistory) + throw new Error( + 'Loading before checkpoint activation requires transport.getHistory().' + ); + if (branch || preparing) admitBranchCommand(); if ( owner || recoveryAttempt || @@ -1392,6 +1781,22 @@ export function createSession( | undefined; await publication.command(() => { if (disposed) throw new Error('Agent has been disposed'); + if (branch || preparing) { + if ( + owner || + loading || + checkController || + preparing || + pendingToolSettlements || + persistence.pending + ) + throw new Error( + 'Checkpoint status cannot be checked during an active operation.' + ); + if (!branch) throw new Error('Checkpoint preparation is active.'); + readyCheckpoint(branch.state); + return; + } if (owner) throw new Error('Stop the active request before checking status'); if (!recoveryAttempt) return; @@ -1460,6 +1865,7 @@ export function createSession( subscribe: (notify) => disposed ? () => undefined : publication.subscribe(notify), submit, + fork, resume, reconnect, stop: () => @@ -1467,11 +1873,12 @@ export function createSession( if (disposed) return; stopExecution(); }), - ...(canCheck ? { checkStatus, load } : {}), + ...(canCheck || getState ? { checkStatus, load } : {}), dispose: () => publication.command(() => { if (disposed) return; disposed = true; + const preparation = detachPreparation(); retained = undefined; const reading = detachLoad(); const checking = invalidateCheck(); @@ -1482,6 +1889,7 @@ export function createSession( checking?.abort(); if (attempt) close(attempt, true); closeLoad(reading); + closePreparation(preparation); }), }; } diff --git a/libs/langgraph/src/runtime/stream-projection.ts b/libs/langgraph/src/runtime/stream-projection.ts index 0037016bb..57b57b173 100644 --- a/libs/langgraph/src/runtime/stream-projection.ts +++ b/libs/langgraph/src/runtime/stream-projection.ts @@ -32,6 +32,8 @@ export interface StreamProjection { readonly excludedIds?: readonly string[]; }; readonly baselineIds: readonly string[]; + /** Historical completed calls may echo, but cannot become a new invocation. */ + readonly baselineCallIds?: readonly string[]; readonly currentAssistantId?: string; readonly sawAssistant: boolean; readonly terminal: boolean; @@ -120,6 +122,34 @@ export function projectStream( role === 'assistant' && raw['type'] !== 'AIMessageChunk' && Array.isArray(raw['tool_calls']); + const position = incoming.indexOf(raw); + const content = textContent(raw['content']); + const executionCandidate = + (!baseline || projection.resume?.turnIds.includes(id)) && + !projection.resume?.excludedIds?.includes(id) && + (!turnEnded || + anchor >= 0 || + projection.resume?.turnIds.includes(id) || + projection.toolAssistantIds?.includes(id)) && + (anchor < 0 || position > anchor) && + (nextUser < 0 || position < nextUser); + if ( + finalizedCalls && + (executionCandidate || + (baseline && content !== previous?.content) || + (anchor >= 0 && + position > anchor && + (nextUser < 0 || position < nextUser))) + ) + for (const call of calls) + if ( + typeof call['id'] === 'string' && + projection.baselineCallIds?.includes(call['id']) + ) + state = reduceMessages(state, { + type: 'tool-conflict', + id: call['id'], + }); const callIds = calls.flatMap((call) => typeof call['id'] === 'string' ? [call['id']] : [] ); @@ -129,7 +159,6 @@ export function projectStream( (callId) => !callIds.includes(callId) ) ); - const content = textContent(raw['content']); const current = !baseline || (!!projection.resume && @@ -227,17 +256,7 @@ export function projectStream( }), }); } - const position = incoming.indexOf(raw); - if ( - (!baseline || projection.resume?.turnIds.includes(id)) && - !projection.resume?.excludedIds?.includes(id) && - (!turnEnded || - anchor >= 0 || - projection.resume?.turnIds.includes(id) || - projection.toolAssistantIds?.includes(id)) && - (anchor < 0 || position > anchor) && - (nextUser < 0 || position < nextUser) - ) + if (executionCandidate) projection = { ...projection, toolAssistantIds: [ diff --git a/libs/langgraph/src/runtime/testing/checkpoint-fixture.ts b/libs/langgraph/src/runtime/testing/checkpoint-fixture.ts new file mode 100644 index 000000000..e098326d1 --- /dev/null +++ b/libs/langgraph/src/runtime/testing/checkpoint-fixture.ts @@ -0,0 +1,87 @@ +import type { ThreadState } from '@langchain/langgraph-sdk'; +import { vi } from 'vitest'; +import type { + AgentTransport, + LangGraphSubmitOptions, + StreamEvent, +} from '../transport.types'; + +export const position = (id: string) => ({ + thread_id: 'thread', + checkpoint_id: id, + checkpoint_ns: '' as const, + checkpoint_map: {}, +}); +export function saved( + id: string, + messages: unknown[] = [], + changes: Record = {} +): ThreadState { + return { + checkpoint: position(id), + values: { messages }, + next: [], + tasks: [], + created_at: '', + parent_checkpoint: null, + metadata: { run_id: `run-${id}` }, + ...changes, + } as ThreadState; +} +export function checkpointEvent(state: ThreadState): StreamEvent { + return { + type: 'checkpoints', + sseId: `cursor-${state.checkpoint.checkpoint_id}`, + data: { + config: { + configurable: { + ...state.checkpoint, + run_id: state.metadata?.['run_id'], + }, + }, + values: state.values, + next: state.next, + tasks: state.tasks.map(({ id, name }) => ({ id, name })), + }, + }; +} +export function fixture() { + const source = saved('a', [ + { id: 'a-user', type: 'human', content: 'Source A' }, + { id: 'a-answer', type: 'ai', content: 'Answer A' }, + ]); + const states = new Map([['a', source]]); + const requests: { input: unknown; options?: LangGraphSubmitOptions }[] = []; + const transport: AgentTransport = { + getState: vi.fn( + async (_thread, checkpoint) => states.get(checkpoint.checkpoint_id)! + ), + getHistory: vi.fn(async () => [ + saved('competing', [{ id: 'b', type: 'ai', content: 'Wrong tip' }]), + ]), + getRunStatus: vi.fn(async () => 'success' as const), + joinStream: vi.fn(async function* () { + /* retained EOF */ + }), + stream: vi.fn(async function* ( + _assistant, + _thread, + input, + _signal, + options + ) { + requests.push({ input, options }); + const id = `result-${requests.length}`; + options?.onRunCreated?.({ run_id: `run-${id}`, thread_id: 'thread' }); + const messages = [ + ...(source.values as { messages: unknown[] }).messages, + ...(input as { messages: unknown[] }).messages, + { id: `assistant-${id}`, type: 'ai', content: id }, + ]; + const result = saved(id, messages); + states.set(id, result); + yield checkpointEvent(result); + }), + }; + return { source, states, requests, transport }; +} diff --git a/libs/langgraph/src/runtime/tool-persistence.ts b/libs/langgraph/src/runtime/tool-persistence.ts index 661465453..d03286994 100644 --- a/libs/langgraph/src/runtime/tool-persistence.ts +++ b/libs/langgraph/src/runtime/tool-persistence.ts @@ -6,11 +6,12 @@ type ToolBatch = ReturnType; /** Fixed-thread persistence effects. Capture a batch only when its write starts; * an earlier write must acknowledge its exact entries before the next capture. * Rejected acknowledgement blocks automatic retries, including later arrivals. */ -export function createToolPersistence( +export function createToolPersistence( buffer: ToolBuffer, write: ( messages: readonly ToolMessage[], - signal: AbortSignal + signal: AbortSignal, + context: Context | undefined ) => Promise ) { let pending = 0; @@ -27,7 +28,7 @@ export function createToolPersistence( batch.acknowledge(); if (!buffer.snapshot().messages.length) failed = false; }, - flush(signal: AbortSignal): Promise { + flush(signal: AbortSignal, context?: Context): Promise { // Admission must see queued work before any promise continuation runs. pending += 1; const operation = tail.then(async () => { @@ -39,7 +40,7 @@ export function createToolPersistence( const batch = buffer.snapshot(); if (!batch.messages.length) return; try { - await write(batch.messages, signal); + await write(batch.messages, signal, context); } catch (error) { // A rejected response may follow a committed remote write. Do not // retry it merely because another durable result arrives afterward. diff --git a/libs/langgraph/src/runtime/transport.types.ts b/libs/langgraph/src/runtime/transport.types.ts index bff57d5b2..f47463965 100644 --- a/libs/langgraph/src/runtime/transport.types.ts +++ b/libs/langgraph/src/runtime/transport.types.ts @@ -7,6 +7,8 @@ import type { StreamMode, ThreadState, } from '@langchain/langgraph-sdk'; +import type { OwnedCheckpointPosition } from '../lib/transport/checkpoint-position'; +export type { OwnedCheckpointPosition } from '../lib/transport/checkpoint-position'; /** An event emitted by a LangGraph stream. */ export interface StreamEvent { @@ -127,11 +129,16 @@ export interface AgentTransport { threadId: string, runId: string, lastEventId: string | undefined, - signal: AbortSignal + signal: AbortSignal, + options?: { streamMode?: StreamMode[] } ): AsyncIterable; /** @internal Inspect the exact owned physical run, without thread-history inference. */ - getRunStatus?(threadId: string, runId: string, signal: AbortSignal): Promise; + getRunStatus?( + threadId: string, + runId: string, + signal: AbortSignal + ): Promise; /** Optional: create a server-side queued run without joining it immediately. */ createQueuedRun?( @@ -152,6 +159,13 @@ export interface AgentTransport { /** Optional: load persisted checkpoint history for a thread. */ getHistory?(threadId: string, signal: AbortSignal): Promise; + /** Exact saved root read. Branch effects must not infer authority from latest. */ + getState?( + threadId: string, + checkpoint: OwnedCheckpointPosition, + signal: AbortSignal + ): Promise; + /** * Optional: update server-side thread state (e.g. to emit RemoveMessage * entries for regenerate rollback). Forwards to the LangGraph @@ -167,8 +181,8 @@ export interface AgentTransport { threadId: string, values: Record, signal: AbortSignal, - options?: { asNode?: string } - ): Promise; + options?: { asNode?: string; checkpoint?: OwnedCheckpointPosition } + ): Promise; } /** diff --git a/scripts/react-parity/baseline.json b/scripts/react-parity/baseline.json index 45799db47..a53fa14ca 100644 --- a/scripts/react-parity/baseline.json +++ b/scripts/react-parity/baseline.json @@ -1,23 +1,21 @@ { "schemaVersion": 1, - "baselineHead": "22d837101fdd621d878f58079d55946a1948a845", + "baselineHead": "f8de78635270a3158de156aa060dce2bdd875fd6", "sourceState": { "modified": [ - ".github/workflows/ci.yml", - "apps/website/content/docs/middleware/getting-started/introduction.mdx", + "libs/langgraph/src/lib/transport/fetch-stream.transport.ts", "libs/langgraph/src/runtime/create-session.ts", - "libs/langgraph/src/runtime/function-tools.ts", - "libs/langgraph/src/runtime/message-reducer.ts", - "libs/langgraph/src/runtime/ownership.ts", - "libs/langgraph/src/runtime/tool-invocations.ts", - "libs/middleware/README.md", - "libs/middleware/src/langgraph/client-tool-execution-store.ts", - "libs/middleware/src/langgraph/index.ts", - "libs/middleware/src/langgraph/postgres-client-tool-execution-store.ts" + "libs/langgraph/src/runtime/stream-projection.ts", + "libs/langgraph/src/runtime/tool-persistence.ts", + "libs/langgraph/src/runtime/transport.types.ts" ], "untracked": [ - "libs/langgraph/src/runtime/tool-provenance.ts", - "libs/middleware/src/langgraph/postgres-tool-execution-schema.ts" + "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" ] }, "scope": { @@ -603,6 +601,12 @@ "path": "libs/langgraph/project.json", "sha256": "4fa0e5f4b0d4b2a7fd8f8de7270dcce1f21b033c86bb031d75480e44667c29d3" }, + { + "id": "asset:libs/langgraph/src/runtime/README.md", + "kind": "asset", + "path": "libs/langgraph/src/runtime/README.md", + "sha256": "45d9aad401d27f6e53dbe2aa5f495e4665a89a69bdf4208525b7e7ca619c818b" + }, { "id": "asset:libs/langgraph/test/fixtures/streaming-reasoning-puzzle.json", "kind": "asset", @@ -8315,7 +8319,7 @@ "path": "libs/langgraph/src/runtime/transport.types.ts", "symbol": "AgentTransport", "syntaxKind": "InterfaceDeclaration", - "signature": "export interface AgentTransport {\n stream(assistantId: string, threadId: string | null, payload: unknown, signal: AbortSignal, options?: LangGraphSubmitOptions): AsyncIterable;\n joinStream?(threadId: string, runId: string, lastEventId: string | undefined, signal: AbortSignal): AsyncIterable;\n getRunStatus?(threadId: string, runId: string, signal: AbortSignal): Promise;\n createQueuedRun?(assistantId: string, threadId: string, payload: unknown, signal: AbortSignal, options?: LangGraphSubmitOptions): Promise;\n cancelRun?(threadId: string, runId: string, signal: AbortSignal): Promise;\n getHistory?(threadId: string, signal: AbortSignal): Promise;\n updateState?(threadId: string, values: Record, signal: AbortSignal, options?: {\n asNode?: string;\n }): Promise;\n}" + "signature": "export interface AgentTransport {\n stream(assistantId: string, threadId: string | null, payload: unknown, signal: AbortSignal, options?: LangGraphSubmitOptions): AsyncIterable;\n joinStream?(threadId: string, runId: string, lastEventId: string | undefined, signal: AbortSignal, options?: {\n streamMode?: StreamMode[];\n }): AsyncIterable;\n getRunStatus?(threadId: string, runId: string, signal: AbortSignal): Promise;\n createQueuedRun?(assistantId: string, threadId: string, payload: unknown, signal: AbortSignal, options?: LangGraphSubmitOptions): Promise;\n cancelRun?(threadId: string, runId: string, signal: AbortSignal): Promise;\n getHistory?(threadId: string, signal: AbortSignal): Promise;\n getState?(threadId: string, checkpoint: OwnedCheckpointPosition, signal: AbortSignal): Promise;\n updateState?(threadId: string, values: Record, signal: AbortSignal, options?: {\n asNode?: string;\n checkpoint?: OwnedCheckpointPosition;\n }): Promise;\n}" } ] }, @@ -8371,7 +8375,7 @@ "path": "libs/langgraph/src/lib/transport/fetch-stream.transport.ts", "symbol": "FetchStreamTransport", "syntaxKind": "ClassDeclaration", - "signature": "export class FetchStreamTransport implements AgentTransport {\n private client: Client;\n private onThreadId?: (id: string) => void;\n private readonly protectErrors: boolean;\n private readonly reportOperationFailure?: RuntimeOperationFailureReporter;\n readonly protectsOperationErrors: boolean;\n constructor(apiUrl: string, onThreadId?: (id: string) => void, clientOptions?: LangGraphClientOptions, reportOperationFailure?: RuntimeOperationFailureReporter) {\n this.protectErrors = clientOptions?.defaultHeaders !== undefined || reportOperationFailure !== undefined;\n this.protectsOperationErrors = this.protectErrors;\n this.reportOperationFailure = reportOperationFailure;\n this.client = this.protectErrors\n ? ɵcreateProtectedLangGraphClient(apiUrl, clientOptions, createLangGraphRuntimeFetch(reportOperationFailure))\n : createLangGraphClient(apiUrl, clientOptions);\n this.onThreadId = onThreadId;\n }\n async *stream(assistantId: string, threadId: string | null, payload: unknown, signal: AbortSignal, options?: LangGraphSubmitOptions): AsyncIterable {\n let thread = threadId;\n if (!thread) {\n try {\n const t = await this.client.threads.create();\n thread = t.thread_id;\n }\n catch (error) {\n this.rethrowOperationError(error, signal);\n }\n try {\n this.onThreadId?.(thread);\n }\n catch (error) {\n this.rethrowLocalError(error, signal);\n }\n }\n let runPayload: ReturnType;\n try {\n runPayload = buildRunPayload(payload, signal, options);\n }\n catch (error) {\n return this.rethrowLocalError(error, signal);\n }\n let run: ReturnType;\n try {\n run = this.client.runs.stream(thread, assistantId, runPayload);\n }\n catch (error) {\n this.rethrowOperationError(error, signal);\n }\n yield* this.iterateSdkRun(run, signal);\n }\n async *joinStream(threadId: string, runId: string, lastEventId: string | undefined, signal: AbortSignal): AsyncIterable {\n let run: ReturnType;\n try {\n run = this.client.runs.joinStream(threadId, runId, {\n signal,\n ...(lastEventId !== undefined ? { lastEventId } : {}),\n });\n }\n catch (error) {\n this.rethrowOperationError(error, signal);\n }\n yield* this.iterateSdkRun(run, signal);\n }\n async getRunStatus(threadId: string, runId: string, signal: AbortSignal): Promise {\n try {\n const run = await this.client.runs.get(threadId, runId, { signal });\n if (run.run_id !== runId || run.thread_id !== threadId ||\n !['pending', 'running', 'success', 'error', 'timeout', 'interrupted'].includes(run.status))\n throw new Error('Invalid LangGraph run status response.');\n return run.status;\n }\n catch (error) {\n return this.rethrowOperationError(error, signal);\n }\n }\n async createQueuedRun(assistantId: string, threadId: string, payload: unknown, signal: AbortSignal, options?: LangGraphSubmitOptions): Promise {\n let runPayload: ReturnType & {\n multitaskStrategy: 'enqueue';\n };\n try {\n runPayload = {\n ...buildRunPayload(payload, signal, options),\n multitaskStrategy: 'enqueue',\n };\n }\n catch (error) {\n return this.rethrowLocalError(error, signal);\n }\n let run: Awaited>;\n try {\n run = await this.client.runs.create(threadId, assistantId, runPayload);\n }\n catch (error) {\n return this.rethrowOperationError(error, signal);\n }\n try {\n return {\n id: run.run_id,\n threadId: run.thread_id ?? threadId,\n values: payload,\n options: { multitaskStrategy: 'enqueue', signal },\n createdAt: run.created_at ? new Date(run.created_at) : new Date(),\n };\n }\n catch (error) {\n return this.rethrowLocalError(error, signal);\n }\n }\n async cancelRun(threadId: string, runId: string, signal: AbortSignal): Promise {\n try {\n await this.client.runs.cancel(threadId, runId, false, 'interrupt', { signal });\n }\n catch (error) {\n this.rethrowOperationError(error, signal);\n }\n }\n async getHistory(threadId: string, signal: AbortSignal): Promise {\n try {\n return await this.client.threads.getHistory(threadId, { signal });\n }\n catch (error) {\n return this.rethrowOperationError(error, signal);\n }\n }\n async updateState(threadId: string, values: Record, signal: AbortSignal, options?: {\n asNode?: string;\n }): Promise {\n const body: {\n values: Record;\n signal: AbortSignal;\n asNode?: string;\n } = { values, signal };\n if (options?.asNode !== undefined) {\n body.asNode = options.asNode;\n }\n try {\n await this.client.threads.updateState(threadId, body);\n }\n catch (error) {\n this.rethrowOperationError(error, signal);\n }\n }\n private rethrowOperationError(error: unknown, signal: AbortSignal): never {\n if (!this.protectErrors)\n throw error;\n return projectLangGraphOperationFailure(error, signal, this.reportOperationFailure);\n }\n private rethrowLocalError(error: unknown, signal: AbortSignal): never {\n if (!this.protectErrors)\n throw error;\n return projectLangGraphOperationFailure(error, signal, undefined);\n }\n private async *iterateSdkRun(run: ReturnType | ReturnType, signal: AbortSignal): AsyncIterable {\n let iterator: AsyncIterator<{\n event: string;\n data: unknown;\n id?: unknown;\n }>;\n try {\n iterator = run[Symbol.asyncIterator]();\n }\n catch (error) {\n return this.rethrowOperationError(error, signal);\n }\n let completed = false;\n let failed = false;\n try {\n while (true) {\n let next: IteratorResult<{\n event: string;\n data: unknown;\n id?: unknown;\n }>;\n try {\n next = await iterator.next();\n }\n catch (error) {\n return this.rethrowOperationError(error, signal);\n }\n if (next.done) {\n completed = true;\n return;\n }\n try {\n yield { ...normalizeSdkEvent(next.value.event as StreamEvent['type'], next.value.data), sseId: typeof next.value.id === 'string' ? next.value.id : undefined };\n }\n catch (error) {\n return this.rethrowLocalError(error, signal);\n }\n }\n }\n catch (error) {\n failed = true;\n throw error;\n }\n finally {\n if (!completed) {\n try {\n await iterator.return?.();\n }\n catch (error) {\n if (!failed)\n this.rethrowOperationError(error, signal);\n }\n }\n }\n }\n}" + "signature": "export class FetchStreamTransport implements AgentTransport {\n private client: Client;\n private onThreadId?: (id: string) => void;\n private readonly protectErrors: boolean;\n private readonly reportOperationFailure?: RuntimeOperationFailureReporter;\n readonly protectsOperationErrors: boolean;\n constructor(apiUrl: string, onThreadId?: (id: string) => void, clientOptions?: LangGraphClientOptions, reportOperationFailure?: RuntimeOperationFailureReporter) {\n this.protectErrors = clientOptions?.defaultHeaders !== undefined || reportOperationFailure !== undefined;\n this.protectsOperationErrors = this.protectErrors;\n this.reportOperationFailure = reportOperationFailure;\n this.client = this.protectErrors\n ? ɵcreateProtectedLangGraphClient(apiUrl, clientOptions, createLangGraphRuntimeFetch(reportOperationFailure))\n : createLangGraphClient(apiUrl, clientOptions);\n this.onThreadId = onThreadId;\n }\n async *stream(assistantId: string, threadId: string | null, payload: unknown, signal: AbortSignal, options?: LangGraphSubmitOptions): AsyncIterable {\n let thread = threadId;\n if (!thread) {\n try {\n const t = await this.client.threads.create();\n thread = t.thread_id;\n }\n catch (error) {\n this.rethrowOperationError(error, signal);\n }\n try {\n this.onThreadId?.(thread);\n }\n catch (error) {\n this.rethrowLocalError(error, signal);\n }\n }\n let runPayload: ReturnType;\n try {\n runPayload = buildRunPayload(payload, signal, options);\n }\n catch (error) {\n return this.rethrowLocalError(error, signal);\n }\n let run: ReturnType;\n try {\n run = this.client.runs.stream(thread, assistantId, runPayload);\n }\n catch (error) {\n this.rethrowOperationError(error, signal);\n }\n yield* this.iterateSdkRun(run, signal);\n }\n async *joinStream(threadId: string, runId: string, lastEventId: string | undefined, signal: AbortSignal, options?: {\n streamMode?: StreamMode[];\n }): AsyncIterable {\n let run: ReturnType;\n try {\n run = this.client.runs.joinStream(threadId, runId, {\n signal,\n ...(options?.streamMode ? { streamMode: options.streamMode } : {}),\n ...(lastEventId !== undefined ? { lastEventId } : {}),\n });\n }\n catch (error) {\n this.rethrowOperationError(error, signal);\n }\n yield* this.iterateSdkRun(run, signal);\n }\n async getRunStatus(threadId: string, runId: string, signal: AbortSignal): Promise {\n try {\n const run = await this.client.runs.get(threadId, runId, { signal });\n if (run.run_id !== runId || run.thread_id !== threadId ||\n !['pending', 'running', 'success', 'error', 'timeout', 'interrupted'].includes(run.status))\n throw new Error('Invalid LangGraph run status response.');\n return run.status;\n }\n catch (error) {\n return this.rethrowOperationError(error, signal);\n }\n }\n async createQueuedRun(assistantId: string, threadId: string, payload: unknown, signal: AbortSignal, options?: LangGraphSubmitOptions): Promise {\n let runPayload: ReturnType & {\n multitaskStrategy: 'enqueue';\n };\n try {\n runPayload = {\n ...buildRunPayload(payload, signal, options),\n multitaskStrategy: 'enqueue',\n };\n }\n catch (error) {\n return this.rethrowLocalError(error, signal);\n }\n let run: Awaited>;\n try {\n run = await this.client.runs.create(threadId, assistantId, runPayload);\n }\n catch (error) {\n return this.rethrowOperationError(error, signal);\n }\n try {\n return {\n id: run.run_id,\n threadId: run.thread_id ?? threadId,\n values: payload,\n options: { multitaskStrategy: 'enqueue', signal },\n createdAt: run.created_at ? new Date(run.created_at) : new Date(),\n };\n }\n catch (error) {\n return this.rethrowLocalError(error, signal);\n }\n }\n async cancelRun(threadId: string, runId: string, signal: AbortSignal): Promise {\n try {\n await this.client.runs.cancel(threadId, runId, false, 'interrupt', { signal });\n }\n catch (error) {\n this.rethrowOperationError(error, signal);\n }\n }\n async getHistory(threadId: string, signal: AbortSignal): Promise {\n try {\n return await this.client.threads.getHistory(threadId, { signal });\n }\n catch (error) {\n return this.rethrowOperationError(error, signal);\n }\n }\n async getState(threadId: string, checkpoint: OwnedCheckpointPosition, signal: AbortSignal): Promise {\n try {\n return await this.client.threads.getState(threadId, checkpoint, {\n signal,\n });\n }\n catch (error) {\n return this.rethrowOperationError(error, signal);\n }\n }\n async updateState(threadId: string, values: Record, signal: AbortSignal, options?: {\n asNode?: string;\n checkpoint?: OwnedCheckpointPosition;\n }): Promise {\n const body: {\n values: Record;\n signal: AbortSignal;\n asNode?: string;\n checkpoint?: OwnedCheckpointPosition;\n } = { values, signal };\n if (options?.asNode !== undefined) {\n body.asNode = options.asNode;\n }\n if (options?.checkpoint)\n body.checkpoint = options.checkpoint;\n try {\n const result = await this.client.threads.updateState(threadId, body);\n try {\n return captureCheckpoint(result.configurable, threadId);\n }\n catch {\n return undefined;\n }\n }\n catch (error) {\n this.rethrowOperationError(error, signal);\n }\n }\n private rethrowOperationError(error: unknown, signal: AbortSignal): never {\n if (!this.protectErrors)\n throw error;\n return projectLangGraphOperationFailure(error, signal, this.reportOperationFailure);\n }\n private rethrowLocalError(error: unknown, signal: AbortSignal): never {\n if (!this.protectErrors)\n throw error;\n return projectLangGraphOperationFailure(error, signal, undefined);\n }\n private async *iterateSdkRun(run: ReturnType | ReturnType, signal: AbortSignal): AsyncIterable {\n let iterator: AsyncIterator<{\n event: string;\n data: unknown;\n id?: unknown;\n }>;\n try {\n iterator = run[Symbol.asyncIterator]();\n }\n catch (error) {\n return this.rethrowOperationError(error, signal);\n }\n let completed = false;\n let failed = false;\n try {\n while (true) {\n let next: IteratorResult<{\n event: string;\n data: unknown;\n id?: unknown;\n }>;\n try {\n next = await iterator.next();\n }\n catch (error) {\n return this.rethrowOperationError(error, signal);\n }\n if (next.done) {\n completed = true;\n return;\n }\n try {\n yield { ...normalizeSdkEvent(next.value.event as StreamEvent['type'], next.value.data), sseId: typeof next.value.id === 'string' ? next.value.id : undefined };\n }\n catch (error) {\n return this.rethrowLocalError(error, signal);\n }\n }\n }\n catch (error) {\n failed = true;\n throw error;\n }\n finally {\n if (!completed) {\n try {\n await iterator.return?.();\n }\n catch (error) {\n if (!failed)\n this.rethrowOperationError(error, signal);\n }\n }\n }\n }\n}" } ] }, @@ -12447,11 +12451,17 @@ "path": "libs/langgraph/src/lib/threads/threads-adapter.ts", "sha256": "125a0d11f9078c6b9effa0a32e27ac786f8982514beec0ca5d22344fcce6ee7c" }, + { + "id": "source:libs/langgraph/src/lib/transport/checkpoint-position.ts", + "kind": "source", + "path": "libs/langgraph/src/lib/transport/checkpoint-position.ts", + "sha256": "da46d1e7945d6277d489535c8fc4ab026ecae72df246258b4731225f0c3de9ad" + }, { "id": "source:libs/langgraph/src/lib/transport/fetch-stream.transport.ts", "kind": "source", "path": "libs/langgraph/src/lib/transport/fetch-stream.transport.ts", - "sha256": "2b899f3fb3758f8c57c500b902d16d826fbfd1090c3b160692f0b01456024981" + "sha256": "fcf23e540b9ad402403578a1c78e016392e1a41cafa663d08a31d333bd7d45d2" }, { "id": "source:libs/langgraph/src/lib/transport/mock-stream.transport.ts", @@ -12471,17 +12481,35 @@ "path": "libs/langgraph/src/public-api.ts", "sha256": "3effe8f90bcd9967bb99b9902b3bb6ffbe2769c8c5b46824d2742e2906f6536b" }, + { + "id": "source:libs/langgraph/src/runtime/checkpoint-authority.ts", + "kind": "source", + "path": "libs/langgraph/src/runtime/checkpoint-authority.ts", + "sha256": "f56912f077e4ba891f9a55419e89743d62f573c435c3984fb64ac0803271c3dc" + }, { "id": "source:libs/langgraph/src/runtime/checkpoint-history.ts", "kind": "source", "path": "libs/langgraph/src/runtime/checkpoint-history.ts", "sha256": "59618157c20ec9da6bfe0e7080876b78de5e773991e68eaf92fe81c77e7d981d" }, + { + "id": "source:libs/langgraph/src/runtime/checkpoint-state.ts", + "kind": "source", + "path": "libs/langgraph/src/runtime/checkpoint-state.ts", + "sha256": "5daf89cc20549d1707f8589eedbb30ff9e7ebcb7d99965a8181f324a928c70ef" + }, + { + "id": "source:libs/langgraph/src/runtime/checkpoint-tool-evidence.ts", + "kind": "source", + "path": "libs/langgraph/src/runtime/checkpoint-tool-evidence.ts", + "sha256": "833f052dee8017dd6bdf1968457d85e0efa7a204dbc0d75b490b5ec82e6d2db2" + }, { "id": "source:libs/langgraph/src/runtime/create-session.ts", "kind": "source", "path": "libs/langgraph/src/runtime/create-session.ts", - "sha256": "9660c52244fca8fd8712abfe2683aea6a95d8aa13f840c0fbb7b706f87d03106" + "sha256": "f4c6a1e693d205cf48a7c9ec357ddf037bdcf4ceabe315c6529c6806fbdaefb1" }, { "id": "source:libs/langgraph/src/runtime/function-tools.ts", @@ -12547,7 +12575,7 @@ "id": "source:libs/langgraph/src/runtime/stream-projection.ts", "kind": "source", "path": "libs/langgraph/src/runtime/stream-projection.ts", - "sha256": "4271e1db01f40d5f349e6b0307d01114fb5fbe5dc9c6dae6223457d73580c4e9" + "sha256": "85ac624f5921d63ba5bff724c14612279934447c5178654b6d5db2d257f99635" }, { "id": "source:libs/langgraph/src/runtime/subgraph-projection.ts", @@ -12567,6 +12595,12 @@ "path": "libs/langgraph/src/runtime/testing/binding-fixture.ts", "sha256": "3f54126919f9d1fd8b5c966ffb6ee3196ef3ea6b605ab09b560857b5a1d7b056" }, + { + "id": "source:libs/langgraph/src/runtime/testing/checkpoint-fixture.ts", + "kind": "source", + "path": "libs/langgraph/src/runtime/testing/checkpoint-fixture.ts", + "sha256": "3afa6abd58e1ec630f1b66ce7f4ff7da83023c559564f00dd4f9f6df58434fb3" + }, { "id": "source:libs/langgraph/src/runtime/testing/controlled-transport.ts", "kind": "source", @@ -12589,7 +12623,7 @@ "id": "source:libs/langgraph/src/runtime/tool-persistence.ts", "kind": "source", "path": "libs/langgraph/src/runtime/tool-persistence.ts", - "sha256": "d65667e465e91dee1e2ac2154e43ad72d788dbc2fad5fa04fca8a749d08bd0da" + "sha256": "520ce165b602918e172470a6d9d80c75f8458d3bd0c3dee68f5761f4e6e2c6ac" }, { "id": "source:libs/langgraph/src/runtime/tool-provenance.ts", @@ -12601,7 +12635,7 @@ "id": "source:libs/langgraph/src/runtime/transport.types.ts", "kind": "source", "path": "libs/langgraph/src/runtime/transport.types.ts", - "sha256": "0a147a2a93936592afb73b662737b94a40b9640a122efc092107bbbdeddd979b" + "sha256": "de35976c3107e32df45e67403db6e74af7b961c79dd77cbd39aeaf036ffc0cae" }, { "id": "source:libs/langgraph/src/runtime/values-projection.ts", diff --git a/scripts/react-parity/dispositions.json b/scripts/react-parity/dispositions.json index 80ae8ac8d..4741082a2 100644 --- a/scripts/react-parity/dispositions.json +++ b/scripts/react-parity/dispositions.json @@ -850,6 +850,13 @@ "status": "in-progress", "note": "Bounded runtime verification/configuration subset updated; broader assigned migration and release work remain open." }, + { + "id": "asset:libs/langgraph/src/runtime/README.md", + "taskIds": ["T08", "T10", "T36"], + "treatment": "infrastructure", + "status": "in-progress", + "note": "Contributor guidance for private checkpoint authority, transport obligations and consumed-task limits; no public package migration or release claim." + }, { "id": "asset:libs/langgraph/test/fixtures/streaming-reasoning-puzzle.json", "taskIds": [ @@ -11298,6 +11305,13 @@ "treatment": "shared", "status": "planned" }, + { + "id": "source:libs/langgraph/src/lib/transport/checkpoint-position.ts", + "taskIds": ["T08"], + "treatment": "shared", + "status": "in-progress", + "note": "Shared SDK routing capture owns root checkpoint identity and string map entries without importing private session or core implementation." + }, { "id": "source:libs/langgraph/src/lib/transport/fetch-stream.transport.ts", "taskIds": [ @@ -11305,7 +11319,7 @@ ], "treatment": "shared", "status": "in-progress", - "note": "SDK event normalization now protects protocol type and namespace from raw application fields. This narrow routing correction retains legacy transport ownership; broader T08 migration remains open." + "note": "SDK event normalization protects protocol type and namespace. Exact checkpoint reads, routed writes with returned positions, and checkpoint delivery on join retain legacy transport ownership; broader T08 migration remains open." }, { "id": "source:libs/langgraph/src/lib/transport/mock-stream.transport.ts", @@ -11343,6 +11357,14 @@ "reason": "Private staged LangGraph runtime subset; fixture-only and not a neutral public LangGraph root or tarball.", "note": "Owned per-command configuration, context and metadata; routing, defaults and complete adapter migration remain open." }, + { + "id": "source:libs/langgraph/src/runtime/checkpoint-authority.ts", + "taskIds": ["T04", "T09", "T10"], + "treatment": "internal", + "status": "in-progress", + "reason": "Private completed-root execution evidence; excluded from public exports.", + "note": "Checks source eligibility, physical-run checkpoint correlation and unconsumed dynamic pauses. Historical replay and public backend cutover remain open." + }, { "id": "source:libs/langgraph/src/runtime/checkpoint-history.ts", "taskIds": [ @@ -11354,6 +11376,22 @@ "reason": "Private staged LangGraph runtime subset; fixture-only and not a neutral public LangGraph root or tarball.", "note": "Compact explicit checkpoint history observation; execution from checkpoints, pagination and complete adapter migration remain open." }, + { + "id": "source:libs/langgraph/src/runtime/checkpoint-state.ts", + "taskIds": ["T04", "T09", "T10"], + "treatment": "internal", + "status": "in-progress", + "reason": "Private branch owner state; excluded from public exports.", + "note": "Retains ready, in-flight and uncertain authority for the existing session owner; no second scheduler, branch cache or global-tip recovery." + }, + { + "id": "source:libs/langgraph/src/runtime/checkpoint-tool-evidence.ts", + "taskIds": ["T04", "T09", "T10"], + "treatment": "internal", + "status": "in-progress", + "reason": "Private saved-state tool evidence assertion; excluded from public exports.", + "note": "Compares current-turn invocation and resolution evidence before granting checkpoint authority; preserves existing projection ownership and sticky conflicts." + }, { "id": "source:libs/langgraph/src/runtime/create-session.ts", "taskIds": [ @@ -11366,7 +11404,7 @@ "treatment": "internal", "status": "in-progress", "reason": "Private staged LangGraph runtime subset; fixture-only and not a neutral public LangGraph root or tarball.", - "note": "Private explicit resume continues observed pauses with owned opaque responses and null input, shares local cancellation/tool ownership, and never applies submission checkpoint correlation to resume. No automatic targeting, public backend cutover, or whole-task completion." + "note": "Private completed-checkpoint fork retains resulting authority across submit, exact load, dynamic resume, reconnect and tool effects. Consumed-task preflight is not atomic cross-client ownership. No historical pending replay, public backend cutover or whole-task completion." }, { "id": "source:libs/langgraph/src/runtime/function-tools.ts", @@ -11506,6 +11544,14 @@ "reason": "Private controlled runtime/binding test helper; excluded from public exports and production packages.", "note": "Controlled native history fixtures include two task interrupt payloads, paused delivery, equal identity, values-only sharing and empty replacement. Broader T05/T06 coverage remains open." }, + { + "id": "source:libs/langgraph/src/runtime/testing/checkpoint-fixture.ts", + "taskIds": ["T09", "T10", "T36"], + "treatment": "internal", + "status": "in-progress", + "reason": "Private controlled checkpoint test helper; excluded from public exports and production packages.", + "note": "Unit evidence for exact position, run and persistence ownership; actual protocol and server regressions provide separate integration evidence." + }, { "id": "source:libs/langgraph/src/runtime/testing/controlled-transport.ts", "taskIds": [ @@ -11548,7 +11594,7 @@ "treatment": "internal", "status": "in-progress", "reason": "Private fixed-thread persistence owner; excluded from public exports.", - "note": "Serializes terminal writes with execution-time capture, exact-entry acknowledgement and failure retention across queued and later cleanup. Explicit successful handoff releases retained failure; distributed ownership, branch-safe uncertain writes and broader T15/T16 remain open." + "note": "Serializes terminal writes with execution-time batch capture and an originating owner captured at enqueue. Returned checkpoints chain late writes; uncertain acknowledgments retain staged batches without retry. Broader T15/T16 and cross-client coordination remain open." }, { "id": "source:libs/langgraph/src/runtime/tool-provenance.ts", From d3d8e018b31ca6ab71f238df38f8429e4c153765 Mon Sep 17 00:00:00 2001 From: Brian Love Date: Thu, 24 Sep 2026 10:55:44 -0700 Subject: [PATCH 2/2] fix(examples): preserve state write transport signatures --- .../chat/angular/src/app/hero/hero-recording.transport.ts | 4 ++-- .../chat/angular/src/app/stage/stage-recording.transport.ts | 4 ++-- 2 files changed, 4 insertions(+), 4 deletions(-) diff --git a/examples/chat/angular/src/app/hero/hero-recording.transport.ts b/examples/chat/angular/src/app/hero/hero-recording.transport.ts index 4b5789048..5f44684c2 100644 --- a/examples/chat/angular/src/app/hero/hero-recording.transport.ts +++ b/examples/chat/angular/src/app/hero/hero-recording.transport.ts @@ -58,8 +58,8 @@ export class HeroRecordingTransport implements AgentTransport { getHistory(threadId: string, signal: AbortSignal): Promise { return this.inner.getHistory ? this.inner.getHistory(threadId, signal) : Promise.resolve([]); } - updateState(threadId: string, values: Record, signal: AbortSignal, options?: { asNode?: string }): Promise { - return this.inner.updateState ? this.inner.updateState(threadId, values, signal, options) : Promise.resolve(); + updateState(...args: Parameters>): ReturnType> { + return this.inner.updateState ? this.inner.updateState(...args) : Promise.resolve(); } private publish(): void { if (typeof window !== 'undefined') window.__heroRecording = this.recording(); } } diff --git a/examples/chat/angular/src/app/stage/stage-recording.transport.ts b/examples/chat/angular/src/app/stage/stage-recording.transport.ts index 9111b205f..13873b0c0 100644 --- a/examples/chat/angular/src/app/stage/stage-recording.transport.ts +++ b/examples/chat/angular/src/app/stage/stage-recording.transport.ts @@ -100,8 +100,8 @@ export class StageRecordingTransport implements AgentTransport { return states; } - updateState(threadId: string, values: Record, signal: AbortSignal, options?: { asNode?: string }): Promise { - return this.inner.updateState ? this.inner.updateState(threadId, values, signal, options) : Promise.resolve(); + updateState(...args: Parameters>): ReturnType> { + return this.inner.updateState ? this.inner.updateState(...args) : Promise.resolve(); } private publish(): void {