diff --git a/packages/cli-v3/src/entryPoints/dev-run-worker.ts b/packages/cli-v3/src/entryPoints/dev-run-worker.ts index 613790e30..4619b990b 100644 --- a/packages/cli-v3/src/entryPoints/dev-run-worker.ts +++ b/packages/cli-v3/src/entryPoints/dev-run-worker.ts @@ -391,6 +391,7 @@ const zodIpc = new ZodIpcConnection({ } resetExecutionEnvironment(); + standardInputStreamManager.setRunId(execution.run.id); standardTraceContextManager.traceContext = traceContext; standardRunTimelineMetricsManager.registerMetricsFromExecution(metrics, isWarmStart); @@ -645,9 +646,6 @@ const zodIpc = new ZodIpcConnection({ RESOLVE_WAITPOINT: async ({ waitpoint }) => { _sharedWorkerRuntime?.resolveWaitpoints([waitpoint]); }, - INPUT_STREAM_CREATED: async ({ runId }) => { - standardInputStreamManager.connectTail(runId); - }, }, }); diff --git a/packages/cli-v3/src/entryPoints/managed-run-worker.ts b/packages/cli-v3/src/entryPoints/managed-run-worker.ts index cb80e2139..ab7b8d080 100644 --- a/packages/cli-v3/src/entryPoints/managed-run-worker.ts +++ b/packages/cli-v3/src/entryPoints/managed-run-worker.ts @@ -375,6 +375,7 @@ const zodIpc = new ZodIpcConnection({ } resetExecutionEnvironment(); + standardInputStreamManager.setRunId(execution.run.id); standardTraceContextManager.traceContext = traceContext; @@ -638,9 +639,6 @@ const zodIpc = new ZodIpcConnection({ RESOLVE_WAITPOINT: async ({ waitpoint }) => { _sharedWorkerRuntime?.resolveWaitpoints([waitpoint]); }, - INPUT_STREAM_CREATED: async ({ runId }) => { - standardInputStreamManager.connectTail(runId); - }, }, }); diff --git a/packages/core/src/v3/inputStreams/index.ts b/packages/core/src/v3/inputStreams/index.ts index d2e904e6d..e6bc01aab 100644 --- a/packages/core/src/v3/inputStreams/index.ts +++ b/packages/core/src/v3/inputStreams/index.ts @@ -28,6 +28,10 @@ export class InputStreamsAPI implements InputStreamManager { return getGlobal(API_NAME) ?? NOOP_MANAGER; } + public setRunId(runId: string): void { + this.#getManager().setRunId(runId); + } + public on( streamId: string, handler: (data: unknown) => void | Promise diff --git a/packages/core/src/v3/inputStreams/manager.ts b/packages/core/src/v3/inputStreams/manager.ts index 964ca474f..1916a4cc3 100644 --- a/packages/core/src/v3/inputStreams/manager.ts +++ b/packages/core/src/v3/inputStreams/manager.ts @@ -26,6 +26,7 @@ export class StandardInputStreamManager implements InputStreamManager { private buffer = new Map(); private tailAbortController: AbortController | null = null; private tailPromise: Promise | null = null; + private currentRunId: string | null = null; constructor( private apiClient: ApiClient, @@ -33,6 +34,10 @@ export class StandardInputStreamManager implements InputStreamManager { private debug: boolean = false ) {} + setRunId(runId: string): void { + this.currentRunId = runId; + } + on(streamId: string, handler: InputStreamHandler): { off: () => void } { let handlerSet = this.handlers.get(streamId); if (!handlerSet) { @@ -41,6 +46,9 @@ export class StandardInputStreamManager implements InputStreamManager { } handlerSet.add(handler); + // Lazily connect the tail on first listener registration + this.#ensureTailConnected(); + // Flush any buffered data for this stream const buffered = this.buffer.get(streamId); if (buffered && buffered.length > 0) { @@ -61,6 +69,9 @@ export class StandardInputStreamManager implements InputStreamManager { } once(streamId: string, options?: InputStreamOnceOptions): Promise { + // Lazily connect the tail on first listener registration + this.#ensureTailConnected(); + // Check buffer first const buffered = this.buffer.get(streamId); if (buffered && buffered.length > 0) { @@ -140,6 +151,7 @@ export class StandardInputStreamManager implements InputStreamManager { reset(): void { this.disconnect(); + this.currentRunId = null; this.handlers.clear(); // Reject all pending once waiters @@ -155,6 +167,12 @@ export class StandardInputStreamManager implements InputStreamManager { this.buffer.clear(); } + #ensureTailConnected(): void { + if (!this.tailAbortController && this.currentRunId) { + this.connectTail(this.currentRunId); + } + } + async #runTail(runId: string, signal: AbortSignal): Promise { try { const stream = await this.apiClient.fetchStream( diff --git a/packages/core/src/v3/inputStreams/noopManager.ts b/packages/core/src/v3/inputStreams/noopManager.ts index 24a40efcf..1511506a4 100644 --- a/packages/core/src/v3/inputStreams/noopManager.ts +++ b/packages/core/src/v3/inputStreams/noopManager.ts @@ -2,6 +2,8 @@ import { InputStreamManager } from "./types.js"; import { InputStreamOnceOptions } from "../realtimeStreams/types.js"; export class NoopInputStreamManager implements InputStreamManager { + setRunId(_runId: string): void {} + on(_streamId: string, _handler: (data: unknown) => void | Promise): { off: () => void } { return { off: () => {} }; } diff --git a/packages/core/src/v3/inputStreams/types.ts b/packages/core/src/v3/inputStreams/types.ts index acdedbbf5..df88245af 100644 --- a/packages/core/src/v3/inputStreams/types.ts +++ b/packages/core/src/v3/inputStreams/types.ts @@ -1,6 +1,12 @@ import { InputStreamOnceOptions } from "../realtimeStreams/types.js"; export interface InputStreamManager { + /** + * Set the current run ID. The tail connection will be established lazily + * when `on()` or `once()` is first called. + */ + setRunId(runId: string): void; + /** * Register a handler that fires every time data arrives on the given input stream. */ diff --git a/packages/core/src/v3/schemas/messages.ts b/packages/core/src/v3/schemas/messages.ts index 04b8b929e..b58babba7 100644 --- a/packages/core/src/v3/schemas/messages.ts +++ b/packages/core/src/v3/schemas/messages.ts @@ -229,12 +229,6 @@ export const WorkerToExecutorMessageCatalog = { waitpoint: CompletedWaitpoint, }), }, - INPUT_STREAM_CREATED: { - message: z.object({ - version: z.literal("v1").default("v1"), - runId: z.string(), - }), - }, }; export const ProviderToPlatformMessages = {