From 7b64fcb51e06bbf96c3ca0f33a861fd0b2527436 Mon Sep 17 00:00:00 2001 From: Brian Love Date: Thu, 24 Sep 2026 14:19:12 -0700 Subject: [PATCH] feat(ag-ui): own private HTTP request lifetimes --- .github/workflows/ci.yml | 5 +- libs/ag-ui/project.json | 12 + .../src/runtime/create-http-request.spec.ts | 840 ++++++++++++++++++ libs/ag-ui/src/runtime/create-http-request.ts | 142 +++ libs/ag-ui/tsconfig.runtime-tests.json | 21 + libs/ag-ui/vite.config.mts | 1 + libs/ag-ui/vite.runtime.config.mts | 11 + scripts/ci-scope.spec.mjs | 4 + scripts/ci-workflow.spec.mjs | 3 + scripts/react-parity/baseline.json | 43 +- scripts/react-parity/dispositions.json | 22 + scripts/react-parity/package-policy.mjs | 15 +- scripts/react-parity/package-policy.spec.mjs | 26 +- scripts/react-parity/verify-boundaries.mjs | 37 +- .../react-parity/verify-boundaries.spec.mjs | 183 +++- 15 files changed, 1331 insertions(+), 34 deletions(-) create mode 100644 libs/ag-ui/src/runtime/create-http-request.spec.ts create mode 100644 libs/ag-ui/src/runtime/create-http-request.ts create mode 100644 libs/ag-ui/tsconfig.runtime-tests.json create mode 100644 libs/ag-ui/vite.runtime.config.mts diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 0e3a3e74c..36d710b14 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -137,11 +137,14 @@ jobs: node scripts/react-parity/verify-boundaries.mjs - name: Build and validate private React foundations run: npx nx run-many -t lint test type-tests build --projects=$FOUNDATIONS --parallel=2 - - name: Verify isolated LangGraph runtime and types + - name: Verify isolated backend runtimes and types run: | npx nx run langgraph:runtime-quality npx nx run langgraph:runtime-type-tests npx nx run langgraph:type-tests + npx nx run ag-ui:runtime-quality + npx nx run ag-ui:runtime-type-tests + npx nx run ag-ui:type-tests - run: npx nx test langgraph --coverage --maxWorkers=2 --reporter=default - name: Verify durable tool ownership and PostgreSQL migration run: node node_modules/tsx/dist/cli.mjs --tsconfig tsconfig.base.json scripts/react-parity/verify-tool-claims.ts diff --git a/libs/ag-ui/project.json b/libs/ag-ui/project.json index b7ee02082..3f54a2b73 100644 --- a/libs/ag-ui/project.json +++ b/libs/ag-ui/project.json @@ -60,6 +60,18 @@ "command": "node ./node_modules/typescript/bin/tsc --noEmit -p libs/ag-ui/tsconfig.type-tests.json" } }, + "runtime-quality": { + "executor": "@nx/vitest:test", + "options": { + "configFile": "libs/ag-ui/vite.runtime.config.mts" + } + }, + "runtime-type-tests": { + "executor": "nx:run-commands", + "options": { + "command": "node ./node_modules/typescript/bin/tsc --project libs/ag-ui/tsconfig.runtime-tests.json" + } + }, "prepare-install": { "executor": "nx:run-commands", "cache": true, diff --git a/libs/ag-ui/src/runtime/create-http-request.spec.ts b/libs/ag-ui/src/runtime/create-http-request.spec.ts new file mode 100644 index 000000000..4633b8b48 --- /dev/null +++ b/libs/ag-ui/src/runtime/create-http-request.spec.ts @@ -0,0 +1,840 @@ +import { ok } from 'node:assert/strict'; +import { once } from 'node:events'; +import { + createServer, + type IncomingMessage, + type ServerResponse, +} from 'node:http'; +import type { AddressInfo } from 'node:net'; +import { describe, expect, it, vi } from 'vitest'; +import { + EventType, + HttpAgent, + type BaseEvent, + type HttpAgentConfig, + type RunAgentInput, +} from '@ag-ui/client'; +import { Observable } from 'rxjs'; +import { createHttpRequest } from './create-http-request'; + +function deferred() { + let resolve!: (value: T) => void; + const promise = new Promise((yes) => { + resolve = yes; + }); + return { promise, resolve }; +} + +async function bounded(promise: Promise): Promise { + let timer: ReturnType | undefined; + try { + return await Promise.race([ + promise, + new Promise((_, reject) => { + timer = setTimeout( + () => reject(new Error('HTTP milestone timed out')), + 1500 + ); + }), + ]); + } finally { + clearTimeout(timer); + } +} + +interface Exchange { + request: IncomingMessage; + response: ServerResponse; + body: unknown; + closed: Promise; + send: (...events: unknown[]) => void; +} + +async function serve(status = 200, headers: Record = {}) { + const exchanges: Exchange[] = []; + const arrivals = new Map>>(); + const server = createServer(async (request, response) => { + const closed = deferred(); + response.on('close', () => closed.resolve()); + const chunks: Buffer[] = []; + for await (const chunk of request) chunks.push(Buffer.from(chunk)); + const exchange = { + request, + response, + closed: closed.promise, + body: JSON.parse(Buffer.concat(chunks).toString()), + send: (...events: unknown[]) => { + for (const event of events) + response.write(`data: ${JSON.stringify(event)}\n\n`); + }, + }; + response.writeHead(status, { + 'content-type': 'text/event-stream', + ...headers, + }); + response.flushHeaders(); + exchanges.push(exchange); + arrivals.get(exchanges.length - 1)?.resolve(exchange); + }); + server.listen(0, '127.0.0.1'); + await once(server, 'listening'); + return { + url: `http://127.0.0.1:${(server.address() as AddressInfo).port}/events`, + exchanges, + next: (index = 0) => { + if (exchanges[index]) return Promise.resolve(exchanges[index]); + const arrival = arrivals.get(index) ?? deferred(); + arrivals.set(index, arrival); + return bounded(arrival.promise); + }, + close: async () => { + const closed = once(server, 'close'); + server.close(); + server.closeAllConnections(); + await bounded(closed); + }, + }; +} + +function input(runId = 'run'): RunAgentInput { + return { + threadId: 'thread', + runId, + state: { count: 1 }, + messages: [{ id: 'user', role: 'user', content: 'hello' }], + tools: [], + context: [], + forwardedProps: { mode: 'test' }, + }; +} + +const started = { + type: EventType.RUN_STARTED, + threadId: 'thread', + runId: 'run', +}; +const finished = { + type: EventType.RUN_FINISHED, + threadId: 'thread', + runId: 'run', +}; + +describe('private HTTP request owner', () => { + it('pre-aborted admission makes zero fetches or POSTs', async () => { + const server = await serve(); + const fetcher = vi.fn((url: string, init: RequestInit) => fetch(url, init)); + const signal = AbortSignal.abort('already left'); + const handle = createHttpRequest({ url: server.url, fetch: fetcher }).start( + input(), + () => { + throw new Error('unexpected event'); + }, + signal + ); + try { + expect(fetcher).not.toHaveBeenCalled(); + expect(await bounded(handle.done)).toEqual({ status: 'aborted' }); + expect(server.exchanges).toHaveLength(0); + } finally { + handle.abort(); + await server.close(); + } + }); + + it('abort detaches events, aborts the source signal, settles locally and closes the open response before cleanup', async () => { + const server = await serve(); + // Test-only emergency cleanup keeps the negative control bounded without + // confusing cleanup closure with the closure asserted before finally. + let source: HttpAgent | undefined; + const rawRun = HttpAgent.prototype.run; + const run = vi + .spyOn(HttpAgent.prototype, 'run') + .mockImplementation(function (this: HttpAgent, value) { + // eslint-disable-next-line @typescript-eslint/no-this-alias -- Test-only SDK resource observation. + source = this; + return rawRun.call(this, value); + }); + let sourceSignal: AbortSignal | null | undefined; + const observed = deferred(); + const events: BaseEvent[] = []; + const handle = createHttpRequest({ + url: server.url, + fetch: (url, init) => { + sourceSignal = init.signal; + return fetch(url, init); + }, + }).start(input(), (event) => { + events.push(event); + observed.resolve(); + }); + const done = handle.done.then((outcome) => outcome); + try { + const exchange = await server.next(); + exchange.send(started); + await bounded(observed.promise); + handle.abort(); + expect(await bounded(done)).toEqual({ status: 'aborted' }); + await bounded(exchange.closed); + expect(sourceSignal?.aborted).toBe(true); + expect(exchange.response.destroyed).toBe(true); + exchange.send(finished); + handle.abort(); + expect(await done).toEqual({ status: 'aborted' }); + expect(events).toEqual([started]); + } finally { + handle.abort(); + source?.abortController.abort(); + run.mockRestore(); + await server.close(); + } + }); + + it('normalizes text and reasoning chunks without writing SDK messages or state', async () => { + const server = await serve(); + let source: HttpAgent | undefined; + const rawRun = HttpAgent.prototype.run; + const run = vi + .spyOn(HttpAgent.prototype, 'run') + .mockImplementation(function (this: HttpAgent, value) { + // eslint-disable-next-line @typescript-eslint/no-this-alias -- Test-only SDK state observation. + source = this; + const fetchResponse = this.fetch; + this.fetch = async (...args) => { + const response = await fetchResponse(...args); + expect(response).toBeInstanceOf(Response); + expect(response.url).toBe(server.url); + expect(response.status).toBe(200); + expect(response.ok).toBe(true); + expect(response.redirected).toBe(false); + expect(response.headers).toBeInstanceOf(Headers); + expect(response.headers.get('content-type')).toBe( + 'text/event-stream' + ); + expect(response.body).toBeInstanceOf(ReadableStream); + expect(response.body).toBe(response.body); + expect(response.body?.locked).toBe(false); + return response; + }; + return rawRun.call(this, value); + }); + const events: BaseEvent[] = []; + const handle = createHttpRequest({ url: server.url }).start( + input(), + (event) => events.push(event) + ); + const done = handle.done.then((outcome) => outcome); + try { + const exchange = await server.next(); + exchange.send( + started, + { type: 'STATE_SNAPSHOT', snapshot: { count: 3 } }, + { + type: 'TEXT_MESSAGE_CHUNK', + messageId: 'answer', + role: 'assistant', + delta: 'hello', + }, + { type: 'TEXT_MESSAGE_CHUNK', messageId: 'answer', delta: ' world' }, + { + type: 'REASONING_MESSAGE_CHUNK', + messageId: 'reason', + delta: 'thinking', + }, + finished + ); + exchange.response.end(); + expect(await bounded(done)).toEqual({ status: 'closed' }); + expect(events.map((event) => event.type)).toEqual([ + 'RUN_STARTED', + 'STATE_SNAPSHOT', + 'TEXT_MESSAGE_START', + 'TEXT_MESSAGE_CONTENT', + 'TEXT_MESSAGE_CONTENT', + 'TEXT_MESSAGE_END', + 'REASONING_MESSAGE_START', + 'REASONING_MESSAGE_CONTENT', + 'REASONING_MESSAGE_END', + 'RUN_FINISHED', + ]); + expect(events[2]).toMatchObject({ + messageId: 'answer', + role: 'assistant', + }); + expect(events[3]).toMatchObject({ delta: 'hello' }); + expect(events[4]).toMatchObject({ delta: ' world' }); + expect(events[7]).toMatchObject({ + messageId: 'reason', + delta: 'thinking', + }); + expect(source?.messages).toEqual([]); + expect(source?.state).toEqual({}); + } finally { + handle.abort(); + run.mockRestore(); + await server.close(); + } + }); + + for (const domainEvent of [ + undefined, + finished, + { type: 'RUN_ERROR', message: 'domain failure', code: 'domain' }, + ]) { + it(`separates physical EOF from ${ + domainEvent?.type ?? 'missing terminal event' + }`, async () => { + const server = await serve(); + const received = deferred(); + const events: BaseEvent[] = []; + let settled = false; + const handle = createHttpRequest({ url: server.url }).start( + input(), + (event) => { + events.push(event); + if (event.type === (domainEvent?.type ?? 'RUN_STARTED')) + received.resolve(); + } + ); + const done = handle.done.then((outcome) => { + settled = true; + return outcome; + }); + try { + const exchange = await server.next(); + exchange.send(started, ...(domainEvent ? [domainEvent] : [])); + await bounded(received.promise); + expect(settled).toBe(false); + exchange.response.end(); + expect(await bounded(done)).toEqual({ status: 'closed' }); + expect(events).toEqual([ + started, + ...(domainEvent ? [domainEvent] : []), + ]); + } finally { + handle.abort(); + await server.close(); + } + }); + } + + for (const failure of [ + 'decoder', + 'schema', + 'verifier', + 'http-json', + 'http-text', + ] as const) { + it(`contains ${failure} failure and aborts its captured source`, async () => { + const server = await serve( + failure.startsWith('http') ? 503 : 200, + failure === 'http-json' ? { 'content-type': 'application/json' } : {} + ); + let signal: AbortSignal | null | undefined; + const external = new AbortController(); + const remove = vi.spyOn(external.signal, 'removeEventListener'); + const handle = createHttpRequest({ + url: server.url, + fetch: (url, init) => { + signal = init.signal; + return fetch(url, init); + }, + }).start(input(), () => undefined, external.signal); + const done = handle.done.then((outcome) => outcome); + try { + const exchange = await server.next(); + if (failure === 'decoder') + exchange.response.write('data: {invalid json}\n\n'); + if (failure === 'schema') + exchange.send({ type: 'TEXT_MESSAGE_CONTENT', delta: 42 }); + if (failure === 'verifier') + exchange.send({ + type: 'TEXT_MESSAGE_CONTENT', + messageId: 'missing-start', + delta: 'invalid', + }); + if (failure === 'http-json') + exchange.response.end(JSON.stringify({ message: 'unavailable' })); + if (failure === 'http-text') exchange.response.end('unavailable'); + const outcome = await bounded(done); + expect(outcome.status).toBe('failed'); + if (outcome.status === 'failed') { + expect(outcome.error).toBeDefined(); + if (failure.startsWith('http')) + expect(outcome.error).toMatchObject({ + status: 503, + payload: + failure === 'http-json' + ? { message: 'unavailable' } + : 'unavailable', + }); + } + expect(signal?.aborted).toBe(true); + expect(remove).toHaveBeenCalledWith('abort', expect.any(Function)); + await bounded(exchange.closed); + expect(exchange.response.destroyed).toBe(true); + handle.abort(); + external.abort(); + expect(await done).toBe(outcome); + } finally { + handle.abort(); + await server.close(); + } + }); + } + + for (const abortFirst of [false, true]) { + it(`contains consumer throw${ + abortFirst ? ' after reentrant abort' : '' + } and closes the still-open HTTP response`, async () => { + const server = await serve(); + const error = { opaque: 'consumer failure' }; + let signal: AbortSignal | null | undefined; + const external = new AbortController(); + const remove = vi.spyOn(external.signal, 'removeEventListener'); + const events: BaseEvent[] = []; + const cleanupAbort = vi.fn(() => { + external.abort(); + handle.abort(); + }); + const handle = createHttpRequest({ + url: server.url, + fetch: (url, init) => { + signal = init.signal; + signal?.addEventListener('abort', cleanupAbort, { once: true }); + return fetch(url, init); + }, + }).start( + input(), + (event) => { + events.push(event); + if (abortFirst) handle.abort(); + throw error; + }, + external.signal + ); + const done = handle.done.then((outcome) => outcome); + try { + const exchange = await server.next(); + exchange.send(started, finished); + const outcome = await bounded(done); + expect(outcome).toEqual( + abortFirst ? { status: 'aborted' } : { status: 'failed', error } + ); + if (outcome.status === 'failed') expect(outcome.error).toBe(error); + expect(signal?.aborted).toBe(true); + expect(cleanupAbort).toHaveBeenCalledOnce(); + await bounded(exchange.closed); + expect(exchange.response.destroyed).toBe(true); + expect(remove).toHaveBeenCalledWith('abort', expect.any(Function)); + handle.abort(); + expect(await done).toBe(outcome); + expect(events).toEqual([started]); + } finally { + handle.abort(); + await server.close(); + } + }); + } + + it('captures the default fetch at factory construction', async () => { + const server = await serve(); + const factory = createHttpRequest({ url: server.url }); + const replaced = vi + .spyOn(globalThis, 'fetch') + .mockRejectedValue(new Error('late replacement')); + const handle = factory.start(input(), () => undefined); + const done = handle.done.then((outcome) => outcome); + try { + expect(replaced).not.toHaveBeenCalled(); + const exchange = await server.next(); + exchange.send(started, finished); + exchange.response.end(); + expect(await bounded(done)).toEqual({ status: 'closed' }); + } finally { + replaced.mockRestore(); + handle.abort(); + await server.close(); + } + }); + + it('contains an abrupt socket failure after SSE delivery without an unhandled reader cleanup rejection', async () => { + const server = await serve(); + const received = deferred(); + let signal: AbortSignal | null | undefined; + const cancelObserved = deferred(); + let readFailure: unknown; + const factory = createHttpRequest({ + url: server.url, + fetch: async (url, init) => { + signal = init.signal; + const response = await fetch(url, init); + const body = response.body; + ok(body); + const getReader = body.getReader.bind(body); + vi.spyOn(body, 'getReader').mockImplementation(() => { + const reader = getReader(); + const read = reader.read.bind(reader); + vi.spyOn(reader, 'read').mockImplementation(() => + read().catch((error: unknown) => { + readFailure = error; + throw error; + }) + ); + const cancel = reader.cancel.bind(reader); + vi.spyOn(reader, 'cancel').mockImplementation((reason) => { + cancelObserved.resolve(signal?.aborted === true); + return cancel(reason); + }); + return reader; + }); + return response; + }, + }); + const handle = factory.start(input(), () => received.resolve()); + let fresh: ReturnType | undefined; + const done = handle.done.then((outcome) => outcome); + try { + const exchange = await server.next(); + exchange.send(started); + await bounded(received.promise); + exchange.response.destroy(); + const outcome = await bounded(done); + expect(outcome.status).toBe('failed'); + if (outcome.status === 'failed') { + expect(outcome.error).toBeInstanceOf(TypeError); + expect(outcome.error).toBe(readFailure); + } + expect(await bounded(cancelObserved.promise)).toBe(true); + await bounded(exchange.closed); + handle.abort(); + expect(await done).toBe(outcome); + fresh = factory.start(input('fresh'), () => undefined); + const freshDone = fresh.done.then((value) => value); + const freshExchange = await server.next(1); + handle.abort(); + freshExchange.send( + { ...started, runId: 'fresh' }, + { ...finished, runId: 'fresh' } + ); + freshExchange.response.end(); + expect(await bounded(freshDone)).toEqual({ status: 'closed' }); + } finally { + handle.abort(); + fresh?.abort(); + await server.close(); + } + }); + + it('keeps the first failure and aborts even if subscription cleanup throws', async () => { + const failure = { opaque: 'consumer' }; + const cleanupFailure = { opaque: 'cleanup' }; + let emit!: () => void; + let source: HttpAgent | undefined; + const run = vi + .spyOn(HttpAgent.prototype, 'run') + .mockImplementation(function (this: HttpAgent) { + // eslint-disable-next-line @typescript-eslint/no-this-alias -- Test-only SDK resource observation. + source = this; + return new Observable((subscriber) => { + emit = () => subscriber.next(started); + return () => { + throw cleanupFailure; + }; + }); + }); + const handle = createHttpRequest({ url: 'http://unused.invalid' }).start( + input(), + () => { + throw failure; + } + ); + const done = handle.done.then((outcome) => outcome); + try { + emit(); + expect(await bounded(done)).toEqual({ status: 'failed', error: failure }); + expect(source?.abortController.signal.aborted).toBe(true); + handle.abort(); + } finally { + source?.abortController.abort(); + run.mockRestore(); + } + }); + + it('does not suppress cancellation begun before settlement when it rejects afterward', async () => { + // A real HTTP reader does not expose control of the underlying cancel + // promise; this local stream pins the cleanup correction's timing boundary. + const error = { opaque: 'earlier cancellation' }; + let rejectCancel!: (reason: unknown) => void; + const cancellation = new Promise((_, reject) => { + rejectCancel = reject; + }); + const body = new ReadableStream({ cancel: () => cancellation }); + const response = new Response(body, { + headers: { 'content-type': 'text/event-stream' }, + }); + const originalGetReader = body.getReader; + let cancelResult: Promise | undefined; + const rawRun = HttpAgent.prototype.run; + const run = vi + .spyOn(HttpAgent.prototype, 'run') + .mockImplementation(function (this: HttpAgent, value) { + const fetchResponse = this.fetch; + this.fetch = async (...args) => { + const ownedResponse = await fetchResponse(...args); + const reader = ownedResponse.body?.getReader(); + ok(reader); + cancelResult = reader.cancel().then( + () => 'resolved', + (error: unknown) => error + ); + reader.releaseLock(); + expect(response.body).toBe(body); + expect(body.getReader).toBe(originalGetReader); + return ownedResponse; + }; + return rawRun.call(this, value); + }); + const handle = createHttpRequest({ + url: 'http://unused.invalid', + fetch: async () => response, + }).start(input(), () => undefined); + try { + expect(await bounded(handle.done)).toEqual({ status: 'closed' }); + rejectCancel(error); + ok(cancelResult); + expect(await bounded(cancelResult)).toBe(error); + } finally { + handle.abort(); + rejectCancel(error); + run.mockRestore(); + } + }); + + it('isolates overlapping starts, external abort, stale abort and a fresh request after cancellation', async () => { + const server = await serve(); + const signals: AbortSignal[] = []; + const factory = createHttpRequest({ + url: server.url, + fetch: (url, init) => { + ok(init.signal); + signals.push(init.signal); + return fetch(url, init); + }, + }); + const external = new AbortController(); + const remove = vi.spyOn(external.signal, 'removeEventListener'); + const firstEvents: BaseEvent[] = []; + const secondEvents: BaseEvent[] = []; + const first = factory.start( + input('first'), + (event) => firstEvents.push(event), + external.signal + ); + const second = factory.start(input('second'), (event) => + secondEvents.push(event) + ); + const firstDone = first.done.then((outcome) => outcome); + const secondDone = second.done.then((outcome) => outcome); + let fresh: ReturnType | undefined; + try { + const requests = await Promise.all([server.next(0), server.next(1)]); + const firstExchange = requests.find( + (exchange) => (exchange.body as RunAgentInput).runId === 'first' + ); + const secondExchange = requests.find( + (exchange) => (exchange.body as RunAgentInput).runId === 'second' + ); + ok(firstExchange); + ok(secondExchange); + external.abort(); + expect(await bounded(firstDone)).toEqual({ status: 'aborted' }); + await bounded(firstExchange.closed); + expect(remove).toHaveBeenCalledWith('abort', expect.any(Function)); + expect(signals[0].aborted).toBe(true); + expect(signals[1].aborted).toBe(false); + fresh = factory.start(input('fresh'), () => undefined); + const freshDone = fresh.done.then((outcome) => outcome); + const freshExchange = await server.next(2); + first.abort(); + secondExchange.send( + { ...started, runId: 'second' }, + { ...finished, runId: 'second' } + ); + secondExchange.response.end(); + freshExchange.send( + { ...started, runId: 'fresh' }, + { ...finished, runId: 'fresh' } + ); + freshExchange.response.end(); + expect(await bounded(secondDone)).toEqual({ status: 'closed' }); + expect(await bounded(freshDone)).toEqual({ status: 'closed' }); + expect(firstEvents).toEqual([]); + expect(secondEvents).toHaveLength(2); + expect(new Set(signals).size).toBe(3); + expect(signals.slice(1).every((signal) => !signal.aborted)).toBe(true); + } finally { + first.abort(); + second.abort(); + fresh?.abort(); + await server.close(); + } + }); + + it('lets a callback start a replacement before aborting its old physical request', async () => { + const server = await serve(); + const factory = createHttpRequest({ url: server.url }); + let replacement: ReturnType | undefined; + let replacementDone: Promise | undefined; + const replacementEvents: BaseEvent[] = []; + const first = factory.start(input(), () => { + replacement = factory.start(input('replacement'), (event) => + replacementEvents.push(event) + ); + replacementDone = replacement.done.then((outcome) => outcome); + first.abort(); + }); + const firstDone = first.done.then((outcome) => outcome); + try { + const oldExchange = await server.next(); + oldExchange.send(started); + expect(await bounded(firstDone)).toEqual({ status: 'aborted' }); + await bounded(oldExchange.closed); + const exchange = await server.next(1); + first.abort(); + exchange.send( + { ...started, runId: 'replacement' }, + { ...finished, runId: 'replacement' } + ); + exchange.response.end(); + ok(replacementDone); + expect(await bounded(replacementDone)).toEqual({ status: 'closed' }); + expect(replacementEvents).toHaveLength(2); + } finally { + first.abort(); + replacement?.abort(); + await server.close(); + } + }); + + for (const mode of [ + 'setup', + 'synchronous-consumer', + 'synchronous-abort', + 'synchronous-error', + 'synchronous-complete', + ] as const) { + it(`settles ${mode} before subscription assignment without leaking delivery or cleanup`, async () => { + // HTTP cannot synchronously emit inside subscribe. Keep this timing seam + // local to the test while retaining the actual SDK normalization/verifier. + const error = { opaque: mode }; + const external = new AbortController(); + const remove = vi.spyOn(external.signal, 'removeEventListener'); + const teardown = vi.fn(); + let source: HttpAgent | undefined; + const run = vi + .spyOn(HttpAgent.prototype, 'run') + .mockImplementation(function (this: HttpAgent) { + // eslint-disable-next-line @typescript-eslint/no-this-alias -- Test-only SDK resource observation. + source = this; + if (mode === 'setup') throw error; + return new Observable((subscriber) => { + if (mode === 'synchronous-error') subscriber.error(error); + else { + subscriber.next(started); + subscriber.next(finished); + subscriber.complete(); + subscriber.error({ late: true }); + } + return teardown; + }); + }); + const events: BaseEvent[] = []; + try { + const handle = createHttpRequest({ + url: 'http://unused.invalid', + }).start( + input(), + (event) => { + events.push(event); + if (mode === 'synchronous-abort') external.abort(); + if (mode === 'synchronous-consumer' || mode === 'synchronous-abort') + throw error; + }, + external.signal + ); + const expected = + mode === 'synchronous-abort' + ? { status: 'aborted' } + : mode === 'synchronous-complete' + ? { status: 'closed' } + : { status: 'failed', error }; + const outcome = await bounded(handle.done); + expect(outcome).toEqual(expected); + if (outcome.status === 'failed') expect(outcome.error).toBe(error); + expect(source?.abortController.signal.aborted).toBe( + mode !== 'synchronous-complete' + ); + expect(remove).toHaveBeenCalledWith('abort', expect.any(Function)); + expect(teardown).toHaveBeenCalledTimes(mode === 'setup' ? 0 : 1); + expect(events).toHaveLength( + mode === 'setup' || mode === 'synchronous-error' + ? 0 + : mode === 'synchronous-complete' + ? 2 + : 1 + ); + handle.abort(); + expect(await handle.done).toBe(outcome); + } finally { + run.mockRestore(); + } + }); + } + + it('is inert until start, then serializes one POST from captured configuration and admitted input', async () => { + const server = await serve(); + const signals: AbortSignal[] = []; + const config: Pick = { + url: server.url, + headers: { 'x-request': 'captured' }, + fetch: (url, init) => { + ok(init.signal); + signals.push(init.signal); + return fetch(url, init); + }, + }; + const factory = createHttpRequest(config); + const start = factory.start; + expect(signals).toHaveLength(0); + expect(server.exchanges).toHaveLength(0); + config.url = 'http://127.0.0.1:1/wrong'; + ok(config.headers); + config.headers['x-request'] = 'changed'; + config.fetch = () => { + throw new Error('replacement fetch'); + }; + const admitted = input(); + const expected = structuredClone(admitted); + const events: BaseEvent[] = []; + const handle = start(admitted, (event) => events.push(event)); + const done = handle.done.then((outcome) => outcome); + admitted.messages[0].content = 'changed'; + admitted.state.count = 2; + try { + const exchange = await server.next(); + expect(exchange.request.method).toBe('POST'); + expect(exchange.request.url).toBe('/events'); + expect(exchange.request.headers['x-request']).toBe('captured'); + expect(exchange.body).toEqual(expected); + exchange.send(started, finished); + exchange.response.end(); + expect(await bounded(done)).toEqual({ status: 'closed' }); + expect(events).toEqual([started, finished]); + expect(signals).toHaveLength(1); + expect(server.exchanges).toHaveLength(1); + } finally { + handle.abort(); + await server.close(); + } + }); +}); diff --git a/libs/ag-ui/src/runtime/create-http-request.ts b/libs/ag-ui/src/runtime/create-http-request.ts new file mode 100644 index 000000000..61adad563 --- /dev/null +++ b/libs/ag-ui/src/runtime/create-http-request.ts @@ -0,0 +1,142 @@ +import { + HttpAgent, + transformChunks, + verifyEvents, + type BaseEvent, + type HttpAgentConfig, + type RunAgentInput, +} from '@ag-ui/client'; + +// Transport completion is separate from RUN_FINISHED / RUN_ERROR domain events. +export type HttpRequestOutcome = + | { status: 'closed' } + | { status: 'aborted' } + | { status: 'failed'; error: unknown }; + +export interface HttpRequestHandle { + abort(): void; + done: Promise; +} + +function withOwnedReaderCleanup( + response: Response, + isSettled: () => boolean +): Response { + const body = response.body; + if (!body) return response; + const ownedBody = new Proxy(body, { + get(target, property) { + if (property === 'getReader') { + return (...args: Parameters) => { + const reader = body.getReader(...args); + return new Proxy(reader, { + get(target, property) { + if (property === 'cancel') { + return (reason?: unknown) => { + const settledAtCancel = isSettled(); + // SDK 0.0.59 rethrows a rejected cancel in a detached promise. + // An errored reader rejects again with the original read + // error, which has already selected this request's outcome. + return reader.cancel(reason).catch((error: unknown) => { + if (!settledAtCancel) throw error; + }); + }; + } + const value = Reflect.get(target, property, target); + return typeof value === 'function' ? value.bind(target) : value; + }, + }); + }; + } + const value = Reflect.get(target, property, target); + return typeof value === 'function' ? value.bind(target) : value; + }, + }); + return new Proxy(response, { + get(target, property) { + if (property === 'body') return ownedBody; + const value = Reflect.get(target, property, target); + return typeof value === 'function' ? value.bind(target) : value; + }, + }); +} + +export function createHttpRequest( + config: Pick +) { + const captured = { + url: config.url, + headers: { ...config.headers }, + fetch: config.fetch ?? globalThis.fetch.bind(globalThis), + }; + return { + start( + input: RunAgentInput, + onEvent: (event: BaseEvent) => void, + signal?: AbortSignal + ): HttpRequestHandle { + let resolve!: (outcome: HttpRequestOutcome) => void; + const done = new Promise((yes) => { + resolve = yes; + }); + let settled = false; + let controller: AbortController | undefined; + let subscription: + | ReturnType['subscribe']> + | undefined; + const settle = (outcome: HttpRequestOutcome) => { + if (settled) return; + settled = true; + signal?.removeEventListener('abort', abort); + try { + subscription?.unsubscribe(); + } catch { + /* Cleanup cannot replace the first outcome or prevent abort. */ + } + if (outcome.status !== 'closed') controller?.abort(); + resolve(outcome); + }; + const abort = () => settle({ status: 'aborted' }); + if (signal?.aborted) { + abort(); + return { abort, done }; + } + signal?.addEventListener('abort', abort, { once: true }); + try { + // Raw run eagerly dispatches and cannot be reused after cancellation. + // It does not apply SDK events to a second messages/state authority. + const source = new HttpAgent({ + ...captured, + fetch: async (url, init) => + withOwnedReaderCleanup( + await captured.fetch(url, init), + () => settled + ), + }); + controller = source.abortController; + const events = source + .run(input) + .pipe(transformChunks(), verifyEvents()); + if (!settled) { + subscription = events.subscribe({ + next: (event) => { + if (settled) return; + try { + onEvent(event); + } catch (error) { + settle({ status: 'failed', error }); + } + }, + error: (error: unknown) => settle({ status: 'failed', error }), + complete: () => settle({ status: 'closed' }), + }); + // A synchronous callback can settle before subscribe returns. + if (settled) subscription.unsubscribe(); + } + } catch (error) { + settle({ status: 'failed', error }); + } + return { abort, done }; + }, + }; +} diff --git a/libs/ag-ui/tsconfig.runtime-tests.json b/libs/ag-ui/tsconfig.runtime-tests.json new file mode 100644 index 000000000..ee7470e3a --- /dev/null +++ b/libs/ag-ui/tsconfig.runtime-tests.json @@ -0,0 +1,21 @@ +{ + "extends": "../../tsconfig.base.json", + "compilerOptions": { + "rootDir": "../..", + "noEmit": true, + "emitDeclarationOnly": false, + "declarationMap": false, + "composite": false, + "incremental": false, + "strict": true, + "lib": ["es2022"], + "types": ["node"], + "paths": { + "@threadplane/core": ["libs/core/src/index.ts"], + "@threadplane/core/tools": ["libs/core/src/tools/index.ts"] + } + }, + "include": ["src/runtime/**/*.ts"], + "exclude": [], + "references": [] +} diff --git a/libs/ag-ui/vite.config.mts b/libs/ag-ui/vite.config.mts index ce406638a..2c7183c71 100644 --- a/libs/ag-ui/vite.config.mts +++ b/libs/ag-ui/vite.config.mts @@ -7,6 +7,7 @@ export default defineConfig({ globals: true, environment: 'jsdom', include: ['src/**/*.spec.ts'], + exclude: ['src/runtime/**'], setupFiles: ['src/test-setup.ts'], passWithNoTests: true, }, diff --git a/libs/ag-ui/vite.runtime.config.mts b/libs/ag-ui/vite.runtime.config.mts new file mode 100644 index 000000000..e5a83861b --- /dev/null +++ b/libs/ag-ui/vite.runtime.config.mts @@ -0,0 +1,11 @@ +import { defineConfig } from 'vitest/config'; + +export default defineConfig({ + root: import.meta.dirname, + test: { + environment: 'node', + reporters: ['default'], + include: ['src/runtime/**/*.spec.ts'], + passWithNoTests: false, + }, +}); diff --git a/scripts/ci-scope.spec.mjs b/scripts/ci-scope.spec.mjs index 9df343987..b0120f60d 100644 --- a/scripts/ci-scope.spec.mjs +++ b/scripts/ci-scope.spec.mjs @@ -55,6 +55,10 @@ describe('React migration baseline scope', () => { 'libs/langgraph/vite.runtime.config.mts', 'libs/langgraph/tsconfig.runtime-tests.json', 'libs/langgraph/src/runtime/testing/controlled-transport.ts', + 'libs/ag-ui/vite.runtime.config.mts', + 'libs/ag-ui/tsconfig.runtime-tests.json', + 'libs/ag-ui/src/runtime/create-http-request.ts', + 'libs/ag-ui/src/runtime/create-http-request.spec.ts', 'libs/core/src/index.ts', 'libs/angular/src/public-api.ts', 'libs/react/src/index.ts', diff --git a/scripts/ci-workflow.spec.mjs b/scripts/ci-workflow.spec.mjs index 8ced71ec5..3f7b2e99d 100644 --- a/scripts/ci-workflow.spec.mjs +++ b/scripts/ci-workflow.spec.mjs @@ -114,6 +114,9 @@ describe('CI workflow', () => { assert.match(job, /npx nx run langgraph:runtime-quality/); assert.match(job, /npx nx run langgraph:runtime-type-tests/); assert.match(job, /npx nx run langgraph:type-tests/); + assert.match(job, /npx nx run ag-ui:runtime-quality/); + assert.match(job, /npx nx run ag-ui:runtime-type-tests/); + assert.match(job, /npx nx run ag-ui:type-tests/); const install = job.indexOf('npx playwright install --with-deps chromium'); assert.ok(install >= 0, 'library browser checks need Chromium and OS dependencies'); for (const script of ['verify-packages.mjs', 'verify-angular-package.mjs']) { diff --git a/scripts/react-parity/baseline.json b/scripts/react-parity/baseline.json index e13a78709..859a9a550 100644 --- a/scripts/react-parity/baseline.json +++ b/scripts/react-parity/baseline.json @@ -1,18 +1,17 @@ { "schemaVersion": 1, - "baselineHead": "900f1e60be7edaf242af7dc7f634eac6e54b1f7f", + "baselineHead": "3a50f98a33dbc83cda3259840329a6c921e08ced", "sourceState": { "modified": [ - "libs/langgraph/src/runtime/README.md", - "libs/langgraph/src/runtime/history-projection.ts", - "libs/langgraph/src/runtime/message-reducer.ts", - "libs/langgraph/src/runtime/ownership.ts", - "libs/langgraph/src/runtime/reasoning-projection.ts", - "libs/langgraph/src/runtime/stream-projection.ts", - "libs/langgraph/src/runtime/testing/controlled-session.ts", - "libs/langgraph/src/runtime/wire-message.ts" + ".github/workflows/ci.yml", + "libs/ag-ui/project.json", + "libs/ag-ui/vite.config.mts" ], - "untracked": [] + "untracked": [ + "libs/ag-ui/src/runtime/create-http-request.ts", + "libs/ag-ui/tsconfig.runtime-tests.json", + "libs/ag-ui/vite.runtime.config.mts" + ] }, "scope": { "libraries": [ @@ -211,7 +210,7 @@ "id": "asset:libs/ag-ui/project.json", "kind": "asset", "path": "libs/ag-ui/project.json", - "sha256": "6bbdf56f7f8275a191b3bb24550f62022bead7d2d4b173ed3067e6dcfc89b2fa" + "sha256": "5a2901f7c18951ffb31a736d1a78af41760c0020eb834bffe1565bd4bf698bfc" }, { "id": "asset:libs/ag-ui/tsconfig.json", @@ -231,6 +230,12 @@ "path": "libs/ag-ui/tsconfig.lib.prod.json", "sha256": "4815c5ece87bc626f79d81495b3b60bce7bca2dc355106a9b77bea8d11855a8b" }, + { + "id": "asset:libs/ag-ui/tsconfig.runtime-tests.json", + "kind": "asset", + "path": "libs/ag-ui/tsconfig.runtime-tests.json", + "sha256": "a692d9d309837d9f37f0e18779a22d7a6f3a2fff1c96ce95a93ccbd0b81397dc" + }, { "id": "asset:libs/ag-ui/tsconfig.type-tests.json", "kind": "asset", @@ -1873,7 +1878,7 @@ "id": "config:.github/workflows/ci.yml", "kind": "config", "path": ".github/workflows/ci.yml", - "sha256": "bf216c2aadd0e521c312d70ea23065873eec1cfeb0d86f5c797137c6711229e3" + "sha256": "45d3465a27b053b180b3507f5993be73c6ee48a2f2b69f7f7ba65b676392392c" }, { "id": "config:.github/workflows/publish-middleware-npm.yml", @@ -10611,6 +10616,12 @@ "path": "libs/ag-ui/src/public-api.ts", "sha256": "3a6e03c590c4ed74eca832b04c206b8efa6544232618a8f573981e26ed8ee180" }, + { + "id": "source:libs/ag-ui/src/runtime/create-http-request.ts", + "kind": "source", + "path": "libs/ag-ui/src/runtime/create-http-request.ts", + "sha256": "c91749ff830d7e3500606d51e8c687657820ac1c5f4b6eee080ca5285d2c2373" + }, { "id": "source:libs/ag-ui/src/test-setup.ts", "kind": "source", @@ -10627,7 +10638,13 @@ "id": "source:libs/ag-ui/vite.config.mts", "kind": "source", "path": "libs/ag-ui/vite.config.mts", - "sha256": "3c93dbe6de30075f257d09a992fa84ee272078a8dea4520ad7e40afc4be72624" + "sha256": "c9b37e87270940807ac16e0ab5826ad1a2b71228e76188f304079062abe28b43" + }, + { + "id": "source:libs/ag-ui/vite.runtime.config.mts", + "kind": "source", + "path": "libs/ag-ui/vite.runtime.config.mts", + "sha256": "abbc5ed9727c0a4ff268d09cf18c22280b8093b1d16ac4d7440b27bf2c89a591" }, { "id": "source:libs/chat/.install-collector/development-install.d.ts", diff --git a/scripts/react-parity/dispositions.json b/scripts/react-parity/dispositions.json index 6be0819ec..47e8562e2 100644 --- a/scripts/react-parity/dispositions.json +++ b/scripts/react-parity/dispositions.json @@ -255,6 +255,13 @@ "treatment": "infrastructure", "status": "planned" }, + { + "id": "asset:libs/ag-ui/tsconfig.runtime-tests.json", + "taskIds": ["T36"], + "treatment": "infrastructure", + "status": "in-progress", + "note": "Isolated Node verification for the private AG-UI request owner; broader CI and release migration remain open." + }, { "id": "asset:libs/ag-ui/tsconfig.type-tests.json", "taskIds": [ @@ -8651,6 +8658,14 @@ "treatment": "shared", "status": "planned" }, + { + "id": "source:libs/ag-ui/src/runtime/create-http-request.ts", + "taskIds": ["T12", "T13"], + "treatment": "internal", + "status": "in-progress", + "reason": "Private physical HTTP request ownership only; excluded from public exports and legacy Angular reachability.", + "note": "Uses the locked SDK raw stream, normalization and verification. Domain run authority and framework-neutral session integration remain separate migration work." + }, { "id": "source:libs/ag-ui/src/test-setup.ts", "taskIds": [ @@ -8679,6 +8694,13 @@ "treatment": "infrastructure", "status": "planned" }, + { + "id": "source:libs/ag-ui/vite.runtime.config.mts", + "taskIds": ["T36"], + "treatment": "infrastructure", + "status": "in-progress", + "note": "Isolated Node verification for the private AG-UI request owner; broader CI and release migration remain open." + }, { "id": "source:libs/chat/.install-collector/development-install.d.ts", "taskIds": [ diff --git a/scripts/react-parity/package-policy.mjs b/scripts/react-parity/package-policy.mjs index b79bc932b..552a1c6b8 100644 --- a/scripts/react-parity/package-policy.mjs +++ b/scripts/react-parity/package-policy.mjs @@ -19,12 +19,15 @@ export const neutralLangGraphRoots = ['src/lib/transport/fetch-stream.transport. // Exact source-sharing exceptions during the Angular transition. Their imports // remain guarded; neither file may expose the private session through a bridge. export const sharedLangGraphRuntimeSources = ['src/runtime/transport.types.ts', 'src/runtime/operation-errors.ts']; +export const neutralRuntimeSdkEntries = { langgraph: '@langchain/langgraph-sdk', 'ag-ui': '@ag-ui/client' }; // TypeScript resolution and filesystem enumeration can use different separators. -export function langGraphRuntimeSourceKind(root, path) { - const directory = `${root.replaceAll('\\', '/').replace(/\/$/, '')}/libs/langgraph/`; +export function backendRuntimeSourceKind(root, path) { + const directory = `${root.replaceAll('\\', '/').replace(/\/$/, '')}/libs/`; const normalized = path.replaceAll('\\', '/'); - if (!normalized.startsWith(`${directory}src/runtime/`)) return undefined; - return sharedLangGraphRuntimeSources.includes(normalized.slice(directory.length)) ? 'shared' : 'private'; + if (normalized.startsWith(`${directory}ag-ui/src/runtime/`)) return 'private'; + const langGraphDirectory = `${directory}langgraph/`; + if (!normalized.startsWith(`${langGraphDirectory}src/runtime/`)) return undefined; + return sharedLangGraphRuntimeSources.includes(normalized.slice(langGraphDirectory.length)) ? 'shared' : 'private'; } const retiredProjects = ['chat', 'langgraph-core', 'ag-ui-core', 'react-render']; @@ -40,9 +43,9 @@ export function forbiddenDependency(project, specifier, { angularTransitions = [ const pkg = packageOf(specifier); const internal = pkg.startsWith('@threadplane/') ? pkg.slice('@threadplane/'.length) : undefined; if (neutralRuntime) { - if (internal) return !['langgraph', 'core'].includes(internal); + if (internal) return ![project, 'core'].includes(internal); if (specifier.startsWith('.')) return false; - return specifier !== '@langchain/langgraph-sdk'; + return specifier !== neutralRuntimeSdkEntries[project]; } const angular = pkg.startsWith('@angular/'); const react = ['react', 'react-dom', '@types/react', '@types/react-dom'].includes(pkg) || ['react', 'react-render', 'ui-react', 'workspace-react'].includes(internal); diff --git a/scripts/react-parity/package-policy.spec.mjs b/scripts/react-parity/package-policy.spec.mjs index bf9f47a8d..b5ced3507 100644 --- a/scripts/react-parity/package-policy.spec.mjs +++ b/scripts/react-parity/package-policy.spec.mjs @@ -13,8 +13,8 @@ for (const root of ['/repo', '/repo/', 'C:\\repo', 'C:/repo/']) { `${prefix}src/runtime-copy/create-session.ts`, ]) { const expected = path.endsWith('/transport.types.ts') || path.endsWith('/operation-errors.ts') ? 'shared' : path.includes('/runtime/') ? 'private' : undefined; - assert.equal(policy.langGraphRuntimeSourceKind(root, path), expected); - assert.equal(policy.langGraphRuntimeSourceKind(root, path.replaceAll('/', '\\')), expected); + assert.equal(policy.backendRuntimeSourceKind(root, path), expected); + assert.equal(policy.backendRuntimeSourceKind(root, path.replaceAll('/', '\\')), expected); } }); } @@ -28,6 +28,28 @@ test('final internal dependencies follow the approved package roles', () => { } }); +for (const project of ['langgraph', 'ag-ui']) { + test(`${project} neutral dependency policy admits only its own backend and exact SDK`, () => { + const sdk = project === 'langgraph' ? '@langchain/langgraph-sdk' : '@ag-ui/client'; + const policyOptions = { neutralRuntime: true, angularTransitions: [project] }; + for (const allowed of [sdk, `@threadplane/${project}`, '@threadplane/core', './local']) { + assert.equal(policy.forbiddenDependency(project, allowed, policyOptions), false, allowed); + } + for (const denied of ['@angular/core', '@threadplane/chat', '@threadplane/angular', 'rxjs', 'zod', 'unreviewed', `${sdk}/private`, + ...(project === 'ag-ui' ? ['@threadplane/langgraph', '@langchain/langgraph-sdk', '@ag-ui/core'] : ['@threadplane/ag-ui', '@ag-ui/client'])]) { + assert.equal(policy.forbiddenDependency(project, denied, policyOptions), true, denied); + } + }); +} + +test('AG-UI runtime has no shared source exceptions', () => { + for (const filename of ['create-http-request.ts', 'transport.types.ts', 'operation-errors.ts']) { + assert.equal(policy.backendRuntimeSourceKind('/repo', `/repo/libs/ag-ui/src/runtime/${filename}`), 'private'); + assert.equal(policy.backendRuntimeSourceKind('C:\\repo', `C:\\repo\\libs\\ag-ui\\src\\runtime\\${filename}`), 'private'); + } + assert.equal(policy.backendRuntimeSourceKind('/repo', '/repo/libs/ag-ui/src/runtime-copy/owner.ts'), undefined); +}); + test('private scaffolds and temporary Angular exceptions are separate inventories', () => { assert.deepEqual(policy.privateScaffoldProjects, ['core', 'content', 'angular', 'react']); for (const retired of ['langgraph-core', 'ag-ui-core', 'react-render']) assert.ok(!policy.scanProjects.includes(retired)); diff --git a/scripts/react-parity/verify-boundaries.mjs b/scripts/react-parity/verify-boundaries.mjs index e7bbe30de..5411a0e44 100644 --- a/scripts/react-parity/verify-boundaries.mjs +++ b/scripts/react-parity/verify-boundaries.mjs @@ -2,13 +2,14 @@ import { existsSync, readFileSync, readdirSync, realpathSync } from 'node:fs'; import { dirname, join, relative, resolve } from 'node:path'; import { fileURLToPath } from 'node:url'; import ts from 'typescript'; -import { angularTransitionProjects, assertFinalRelease, emittedEntries, forbiddenDependency, langGraphRuntimeSourceKind, manifestViolations, neutralLangGraphRoots, packageOf, privateScaffoldProjects, scanProjects, sourceEntry } from './package-policy.mjs'; +import { angularTransitionProjects, assertFinalRelease, backendRuntimeSourceKind, emittedEntries, forbiddenDependency, manifestViolations, neutralLangGraphRoots, neutralRuntimeSdkEntries, packageOf, privateScaffoldProjects, scanProjects, sourceEntry } from './package-policy.mjs'; export const foundationProjects = privateScaffoldProjects; const optional = /(?:^|\/)(?:testing|zod|math)(?:\/|$)/; const reactFeature = /^(?:chat|markdown|a2ui|debug|tools|testing|render)(?:\/|$)/; const sourceFile = /\.(?:[cm]?[jt]sx?)$/; const testFile = /(?:\.(?:spec|test|type-test)\.[cm]?[jt]sx?$|\/test-setup\.)/; +const coreEntries = new Map([['@threadplane/core', 'src/index.ts'], ['@threadplane/core/tools', 'src/tools/index.ts']]); const readJson = (path) => JSON.parse(readFileSync(path, 'utf8')); function filesIn(directory) { @@ -55,6 +56,7 @@ export function verifyBoundaries({ root = process.cwd(), mode = 'source', projec const options = ts.convertCompilerOptionsFromJson(config.compilerOptions ?? {}, root).options; options.pathsBasePath = root; options.moduleResolution = ts.ModuleResolutionKind.Bundler; + const sdkOptions = { ...options, paths: undefined, baseUrl: undefined }; const cache = new Map(); const prefix = mode === 'built' ? 'dist/libs' : 'libs'; // Resolve filesystem identity before applying exact source exceptions. Import @@ -65,7 +67,7 @@ export function verifyBoundaries({ root = process.cwd(), mode = 'source', projec return identities.get(path); }; const canonicalRoot = canonical(root); - const runtimeSourceKind = (path) => langGraphRuntimeSourceKind(canonicalRoot, canonical(path)); + const runtimeSourceKind = (path) => backendRuntimeSourceKind(canonicalRoot, canonical(path)); const manifestFor = (project) => { const path = join(root, prefix, project, 'package.json'); return existsSync(path) ? readJson(path) : undefined; @@ -113,33 +115,46 @@ export function verifyBoundaries({ root = process.cwd(), mode = 'source', projec } for (const specifier of dependencies.imports) { const target = resolveImport(specifier, path); - const targetProject = target && projectOf(target); - const normalized = targetProject ? `@threadplane/${targetProject}` : specifier; + const canonicalTarget = target && canonical(target); + const targetPaths = target ? [...new Set([target, canonicalTarget])] : []; + const targetProjects = [...new Set(targetPaths.map(projectOf).filter(Boolean))]; + const targetLocations = target ? [relative(directory, target), relative(canonical(directory), canonicalTarget)] : []; const trail = [...ancestry, relative(root, path), specifier].join(' -> '); if (legacyRuntime && target && runtimeSourceKind(target) === 'private') errors.add(`${project}: private runtime reachable from legacy source: ${trail}`); const policy = { angularTransitions, rootRuntime, neutralRuntime, browserTransition: browserTransition && browserPath(path) }; const label = neutralRuntime ? 'neutral runtime forbidden dependency' : 'forbidden dependency'; - if (forbiddenDependency(project, specifier, policy) || forbiddenDependency(project, normalized, policy)) errors.add(`${project}: ${label} ${trail}`); + if (forbiddenDependency(project, specifier, policy) || targetProjects.some((owner) => forbiddenDependency(project, `@threadplane/${owner}`, policy))) errors.add(`${project}: ${label} ${trail}`); + // Neither import aliases nor filesystem aliases may turn a private + // external file into an approved own-package, core or SDK-root edge. + if (neutralRuntime && target && + targetPaths.some((path) => path.replaceAll('\\', '/').includes('/node_modules/')) + ) { + const sdk = neutralRuntimeSdkEntries[project]; + const entry = specifier === sdk ? ts.resolveModuleName(sdk, path, sdkOptions, ts.sys).resolvedModule?.resolvedFileName : undefined; + if (!entry || canonical(entry) !== canonicalTarget) errors.add(`${project}: neutral runtime forbidden dependency ${trail}`); + } + const coreEntry = coreEntries.get(specifier); + const crossesCoreBoundary = targetProjects.includes('core') && projectOf(canonical(path)) !== 'core'; if (neutralRuntime && ( - optional.test(specifier) || (target && optional.test(relative(directory, target))) || + optional.test(specifier) || targetLocations.some((location) => optional.test(location)) || (specifier.startsWith('@threadplane/core/') && specifier !== '@threadplane/core/tools') || - (targetProject === 'core' && projectOf(path) !== 'core' && !['@threadplane/core', '@threadplane/core/tools'].includes(specifier)) + (crossesCoreBoundary && (!coreEntry || canonicalTarget !== canonical(join(root, 'libs/core', coreEntry)))) )) errors.add(`${project}: neutral runtime private/testing dependency ${trail}`); // Every core entry is dependency-free. Explicitly // review any future external dependency instead of allowing a wrapper // package to hide a framework/parser dependency behind its own imports. if (project === 'core' && (!target || target.includes('/node_modules/')) && !specifier.startsWith('.')) errors.add(`${project}: unreviewed dependency ${trail}`); if (rootRuntime && (optional.test(specifier) || ['zod', 'katex'].includes(packageOf(specifier)) || (target && optional.test(relative(directory, target))))) errors.add(`${project}: optional/testing dependency reachable from root: ${trail}`); - if (project === 'react' && rootRuntime && ((specifier.startsWith('@threadplane/react/') && reactFeature.test(specifier.slice('@threadplane/react/'.length))) || (targetProject === 'react' && reactFeature.test(relative(join(directory, 'src'), target))))) errors.add(`${project}: feature dependency reachable from root: ${trail}`); + if (project === 'react' && rootRuntime && ((specifier.startsWith('@threadplane/react/') && reactFeature.test(specifier.slice('@threadplane/react/'.length))) || (target && projectOf(target) === 'react' && reactFeature.test(relative(join(directory, 'src'), target))))) errors.add(`${project}: feature dependency reachable from root: ${trail}`); if (target && !target.includes('/node_modules/')) visit(target, rootRuntime, [...ancestry, relative(root, path)], browserTransition && browserPath(target), neutralRuntime, legacyRuntime); else if (!target && (specifier.startsWith('.') || specifier.startsWith('@threadplane/'))) errors.add(`${project}: unresolved dependency ${trail}`); } } - for (const path of allFiles) visit(path, false, [], browserPath(path), false, mode === 'source' && angularTransitions.includes(project) && runtimeSourceKind(path) === undefined); - if (mode === 'source' && project === 'langgraph') { + for (const path of allFiles) visit(path, false, [], browserPath(path), false, mode === 'source' && (project === 'angular' || angularTransitions.includes(project)) && runtimeSourceKind(path) === undefined); + if (mode === 'source' && ['langgraph', 'ag-ui'].includes(project)) { const neutralRoots = [ ...filesIn(join(directory, 'src/runtime')).filter((path) => !optional.test(relative(directory, path))), - ...neutralLangGraphRoots.map((path) => join(directory, path)).filter(existsSync), + ...(project === 'langgraph' ? neutralLangGraphRoots.map((path) => join(directory, path)).filter(existsSync) : []), ]; for (const path of neutralRoots) visit(path, false, [], false, true); } diff --git a/scripts/react-parity/verify-boundaries.spec.mjs b/scripts/react-parity/verify-boundaries.spec.mjs index 74b365a59..3fb7bdac1 100644 --- a/scripts/react-parity/verify-boundaries.spec.mjs +++ b/scripts/react-parity/verify-boundaries.spec.mjs @@ -1,5 +1,5 @@ import assert from 'node:assert/strict'; -import { existsSync, mkdtempSync, mkdirSync, writeFileSync, rmSync } from 'node:fs'; +import { existsSync, mkdtempSync, mkdirSync, writeFileSync, rmSync, symlinkSync } from 'node:fs'; import { tmpdir } from 'node:os'; import { dirname, join } from 'node:path'; import test from 'node:test'; @@ -7,6 +7,187 @@ import { verifyBoundaries } from './verify-boundaries.mjs'; const finalOptions = { angularTransitions: [], telemetryBrowserTransition: false }; +test('React root retains lexical feature restrictions through filesystem aliases', (t) => { + const root = fixture(t, { + 'libs/react/src/index.ts': "export * from './chat/linked';", + 'libs/react/src/shared.ts': 'export interface Shared {}', + }); + mkdirSync(join(root, 'libs/react/src/chat')); + symlinkSync('../shared.ts', join(root, 'libs/react/src/chat/linked.ts')); + assert.ok(verifyBoundaries({ root, projects: ['react'], ...finalOptions }).some((error) => error.includes('feature dependency reachable from root'))); +}); + +for (const project of ['langgraph', 'ag-ui']) { + const sdk = project === 'langgraph' ? '@langchain/langgraph-sdk' : '@ag-ui/client'; + const otherBackend = project === 'langgraph' ? 'ag-ui' : 'langgraph'; + for (const [target, diagnostic] of [ + [`libs/${otherBackend}/src/runtime/private.ts`, 'neutral runtime forbidden dependency'], + ['libs/core/src/private.ts', 'neutral runtime private/testing dependency'], + ['libs/core/src/testing/helper.ts', 'neutral runtime private/testing dependency'], + [`libs/${project}/src/runtime/testing/helper.ts`, 'neutral runtime private/testing dependency'], + ]) { + test(`${project} neutral runtime rejects a local symlink to ${target}`, (t) => { + const root = fixture(t, { + [`libs/${project}/src/public-api.ts`]: 'export {};', + [`libs/${project}/src/runtime/owner.ts`]: "export type { Foreign } from './linked';", + [target]: 'export interface Foreign {}', + }); + symlinkSync(join(root, target), join(root, `libs/${project}/src/runtime/linked.ts`)); + assert.ok(verifyBoundaries({ root, projects: [project] }).some((error) => error.includes(diagnostic) && error.includes('./linked'))); + }); + } + for (const entry of ['@threadplane/core', '@threadplane/core/tools']) { + test(`${project} neutral runtime rejects ${entry} mapped to a private core file`, (t) => { + const root = fixture(t, { + 'tsconfig.base.json': JSON.stringify({ compilerOptions: { paths: { [entry]: ['./libs/core/src/private.ts'] } } }), + [`libs/${project}/src/public-api.ts`]: 'export {};', + [`libs/${project}/src/runtime/owner.ts`]: `export type { Core } from '${entry}';`, + 'libs/core/src/index.ts': 'export interface Core {}', + 'libs/core/src/tools/index.ts': 'export interface Core {}', + 'libs/core/src/private.ts': 'export interface Core {}', + }); + assert.ok(verifyBoundaries({ root, projects: [project] }).some((error) => error.includes('neutral runtime private/testing dependency') && error.includes(entry))); + }); + } + test(`${project} neutral runtime permits exact core entries and their internal source traversal`, (t) => { + const root = fixture(t, { + 'tsconfig.base.json': JSON.stringify({ compilerOptions: { paths: { + '@threadplane/core': ['./libs/core/src/index.ts'], + '@threadplane/core/tools': ['./libs/core/src/tools/index.ts'], + } } }), + [`libs/${project}/src/public-api.ts`]: 'export {};', + [`libs/${project}/src/runtime/owner.ts`]: "export type { Core } from '@threadplane/core'; export type { Tool } from '@threadplane/core/tools';", + 'libs/core/src/index.ts': "export type { Core } from './private';", + 'libs/core/src/tools/index.ts': "export type { Tool } from '../private';", + 'libs/core/src/private.ts': 'export interface Core {} export interface Tool {}', + }); + assert.deepEqual(verifyBoundaries({ root, projects: [project] }), []); + }); + for (const alias of [`@threadplane/${project}`, `@threadplane/${project}/hidden`, '@threadplane/core', '@threadplane/core/tools', sdk]) { + test(`${project} neutral runtime rejects ${alias} remapped to a private external SDK file`, (t) => { + const root = fixture(t, { + 'tsconfig.base.json': JSON.stringify({ compilerOptions: { paths: { [alias]: [`./node_modules/${sdk}/private.ts`] } } }), + [`libs/${project}/src/public-api.ts`]: 'export {};', + [`libs/${project}/src/runtime/owner.ts`]: `export type { SDK } from '${alias}';`, + [`node_modules/${sdk}/package.json`]: JSON.stringify({ name: sdk, types: 'index.d.ts' }), + [`node_modules/${sdk}/index.d.ts`]: 'export interface SDK {}', + [`node_modules/${sdk}/private.ts`]: 'export interface SDK {}', + }); + assert.ok(verifyBoundaries({ root, projects: [project] }).some((error) => error.includes('neutral runtime forbidden dependency') && error.includes(alias))); + }); + } + for (const alias of [false, true]) { + for (const dependency of [`${sdk}/private`, 'rxjs/index']) { + test(`${project} neutral runtime rejects ${alias ? 'symlinked' : 'relative'} external ${dependency}`, (t) => { + const edge = alias ? './external-alias' : `../../../../node_modules/${dependency}`; + const root = fixture(t, { + [`libs/${project}/src/public-api.ts`]: 'export {};', + [`libs/${project}/src/runtime/owner.ts`]: `export * from '${edge}';`, + [`node_modules/${dependency}.ts`]: 'export interface External {}', + }); + if (alias) symlinkSync(`../../../../node_modules/${dependency}.ts`, join(root, `libs/${project}/src/runtime/external-alias.ts`)); + assert.ok(verifyBoundaries({ root, projects: [project] }).some((error) => error.includes('neutral runtime forbidden dependency') && error.includes(edge))); + }); + } + } + test(`${project} neutral runtime still permits the exact installed SDK root`, (t) => { + const root = fixture(t, { + [`libs/${project}/src/public-api.ts`]: 'export {};', + [`libs/${project}/src/runtime/owner.ts`]: `export type { SDK } from '${sdk}';`, + [`node_modules/${sdk}/package.json`]: JSON.stringify({ name: sdk, types: 'index.d.ts' }), + [`node_modules/${sdk}/index.d.ts`]: 'export interface SDK {}', + }); + assert.deepEqual(verifyBoundaries({ root, projects: [project] }), []); + }); +} + +for (const dependency of ['@angular/core', '@threadplane/chat', '@threadplane/langgraph', '@langchain/langgraph-sdk', '@ag-ui/client/private', '@ag-ui/core', 'rxjs', 'unreviewed']) { + test(`AG-UI neutral runtime rejects transitive ${dependency} despite prior legacy visits`, (t) => { + const root = fixture(t, { + 'libs/ag-ui/src/public-api.ts': "export * from './lib/shared';", + 'libs/ag-ui/src/lib/shared.ts': `export type { X } from '${dependency}';`, + 'libs/ag-ui/src/runtime/create-http-request.ts': "export * from '../lib/shared';", + }); + assert.ok(verifyBoundaries({ root, projects: ['ag-ui'] }).some((error) => error.includes('neutral runtime') && error.includes(dependency))); + }); +} + +for (const edge of [ + "export * from './runtime/create-http-request';", + "export type { Request } from './runtime/create-http-request';", + "export * from './bridge';", + "export * from '@private-request';", + "void import('./runtime/create-http-request');", + "type Request = import('./runtime/create-http-request').Request;", + "const request = require('./runtime/create-http-request');", +]) { + test(`AG-UI legacy source rejects private request reachability: ${edge}`, (t) => { + const root = fixture(t, { + 'tsconfig.base.json': JSON.stringify({ compilerOptions: { paths: { '@private-request': ['./libs/ag-ui/src/runtime/create-http-request.ts'] } } }), + 'libs/ag-ui/src/public-api.ts': edge, + ...(edge.includes('./bridge') ? { 'libs/ag-ui/src/bridge.ts': "export * from './runtime/create-http-request';" } : {}), + 'libs/ag-ui/src/runtime/create-http-request.ts': 'export interface Request {}', + }); + assert.ok(verifyBoundaries({ root, projects: ['ag-ui'] }).some((error) => error.includes('private runtime reachable from legacy source'))); + }); +} + +for (const project of ['chat', 'langgraph', 'render', 'angular']) { + test(`${project} cannot reach AG-UI private runtime through indirection`, (t) => { + const root = fixture(t, { + [`libs/${project}/src/public-api.ts`]: "export * from './bridge';", + [`libs/${project}/src/bridge.ts`]: "export * from '../../ag-ui/src/runtime/create-http-request';", + 'libs/ag-ui/src/runtime/create-http-request.ts': 'export interface Request {}', + }); + assert.ok(verifyBoundaries({ root, projects: [project] }).some((error) => error.includes('private runtime reachable from legacy source'))); + }); +} + +test('AG-UI filesystem aliases retain private runtime identity', (t) => { + const root = fixture(t, { + 'libs/ag-ui/src/public-api.ts': "export * from './linked-request';", + 'libs/ag-ui/src/runtime/create-http-request.ts': 'export interface Request {}', + }); + symlinkSync('runtime/create-http-request.ts', join(root, 'libs/ag-ui/src/linked-request.ts')); + assert.ok(verifyBoundaries({ root, projects: ['ag-ui'] }).some((error) => error.includes('private runtime reachable from legacy source'))); + if (existsSync(join(root, 'libs/ag-ui/src/RUNTIME/create-http-request.ts'))) { + writeFileSync(join(root, 'libs/ag-ui/src/public-api.ts'), "export * from './RUNTIME/create-http-request';"); + assert.ok(verifyBoundaries({ root, projects: ['ag-ui'] }).some((error) => error.includes('private runtime reachable from legacy source'))); + } +}); + +test('AG-UI public transition remains allowed while its runtime permits only public core and exact SDK', (t) => { + const root = fixture(t, { + 'tsconfig.base.json': JSON.stringify({ compilerOptions: { paths: { + '@threadplane/core': ['./libs/core/src/index.ts'], + '@backend-alias': ['./libs/langgraph/src/index.ts'], + } } }), + 'libs/ag-ui/src/public-api.ts': "export type { Signal } from '@angular/core';", + 'libs/ag-ui/src/runtime/create-http-request.ts': "export type { X } from '@threadplane/core'; export { HttpAgent } from '@ag-ui/client'; export * from './local';", + 'libs/ag-ui/src/runtime/local.ts': 'export interface Local {}', + 'libs/ag-ui/src/runtime/testing/helper.ts': 'export interface Helper {}', + 'libs/ag-ui/src/runtime/create-http-request.spec.ts': "import '@angular/core';", + 'libs/core/src/index.ts': 'export interface X {}', + 'libs/core/src/private.ts': 'export interface X {}', + 'libs/langgraph/src/index.ts': 'export interface X {}', + }); + assert.deepEqual(verifyBoundaries({ root, projects: ['ag-ui'] }), []); + for (const dependency of ['@backend-alias', '../../../langgraph/src/index', '../../../core/src/private', './testing/helper']) { + writeFileSync(join(root, 'libs/ag-ui/src/runtime/create-http-request.ts'), `export * from '${dependency}';`); + assert.ok(verifyBoundaries({ root, projects: ['ag-ui'] }).some((error) => error.includes('neutral runtime')), dependency); + } + assert.ok(verifyBoundaries({ root, projects: ['ag-ui'], finalRelease: true }).some((error) => error.includes('ag-ui: Angular transition remains enabled'))); +}); + +for (const path of ['public-api.ts', 'runtime/create-http-request.ts']) { + for (const code of ["const target = '@angular/core'; void import(target);", "const target = '@angular/core'; require(target);", 'export const broken = ;']) { + test(`AG-UI guarded ${path} rejects unsupported dependency analysis: ${code}`, (t) => { + const root = fixture(t, { [`libs/ag-ui/src/${path}`]: code }); + assert.ok(verifyBoundaries({ root, projects: ['ag-ui'] }).some((error) => error.includes('cannot analyze guarded dependencies'))); + }); + } +} + for (const edge of [ "export * from './runtime/create-session';", "export type { Session } from './runtime/create-session';",