diff --git a/apps/sim/app/api/copilot/chat/stream/route.test.ts b/apps/sim/app/api/copilot/chat/stream/route.test.ts index d7de092d2c1..9a020efeda3 100644 --- a/apps/sim/app/api/copilot/chat/stream/route.test.ts +++ b/apps/sim/app/api/copilot/chat/stream/route.test.ts @@ -1,11 +1,20 @@ +import { trace } from '@opentelemetry/api' +import { + BasicTracerProvider, + InMemorySpanExporter, + SimpleSpanProcessor, +} from '@opentelemetry/sdk-trace-base' import { authMockFns } from '@sim/testing' import { NextRequest } from 'next/server' -import { beforeEach, describe, expect, it, vi } from 'vitest' +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' import { OrchestrationError } from '@/lib/core/orchestration/types' import { MothershipStreamV1CompletionStatus, MothershipStreamV1EventType, } from '@/lib/mothership/generated/mothership-stream-v1' +import { CopilotResumeOutcome } from '@/lib/mothership/generated/trace-attribute-values-v1' +import { TraceAttr } from '@/lib/mothership/generated/trace-attributes-v1' +import { TraceSpan } from '@/lib/mothership/generated/trace-spans-v1' const { getLatestRunForStream, readEvents, readFilePreviewSessions, checkForReplayGap } = vi.hoisted(() => ({ @@ -62,6 +71,10 @@ async function readAllChunks(response: Response): Promise { } describe('copilot chat stream replay route', () => { + afterEach(() => { + vi.useRealTimers() + }) + beforeEach(() => { authMockFns.mockGetSession.mockResolvedValue({ user: { id: 'user-1' }, @@ -171,4 +184,35 @@ describe('copilot chat stream replay route', () => { expect(body).toContain('"code":"resume_run_unavailable"') expect(body).toContain(`"type":"${MothershipStreamV1EventType.complete}"`) }) + + it('ends a still-running replay at its cap without a terminal so the client re-attaches', async () => { + const exporter = new InMemorySpanExporter() + trace.setGlobalTracerProvider( + new BasicTracerProvider({ spanProcessors: [new SimpleSpanProcessor(exporter)] }) + ) + vi.useFakeTimers() + getLatestRunForStream.mockResolvedValue({ + status: 'active', + executionId: 'exec-1', + id: 'run-1', + }) + + const response = await GET( + new NextRequest('http://localhost:3000/api/copilot/chat/stream?streamId=stream-1&after=7') + ) + const body = readAllChunks(response) + await vi.advanceTimersByTimeAsync(61 * 60 * 1000) + const text = (await body).join('') + + expect(text).toContain(': keepalive') + expect(text).not.toContain(`"type":"${MothershipStreamV1EventType.error}"`) + expect(text).not.toContain(`"type":"${MothershipStreamV1EventType.complete}"`) + const resume = exporter + .getFinishedSpans() + .find((span) => span.name === TraceSpan.CopilotResumeRequest) + expect(resume?.attributes[TraceAttr.CopilotResumeOutcome]).toBe( + CopilotResumeOutcome.EndedWithoutTerminal + ) + trace.disable() + }) }) diff --git a/apps/sim/app/api/copilot/chat/stream/route.ts b/apps/sim/app/api/copilot/chat/stream/route.ts index c84cad4d102..b1023a791b9 100644 --- a/apps/sim/app/api/copilot/chat/stream/route.ts +++ b/apps/sim/app/api/copilot/chat/stream/route.ts @@ -43,7 +43,12 @@ const logger = createLogger('CopilotChatStreamAPI') const POLL_INTERVAL_MS = 250 const POLL_INTERVAL_MAX_MS = 2_000 const REPLAY_KEEPALIVE_INTERVAL_MS = 15_000 -const MAX_STREAM_MS = 60 * 60 * 1000 +/** + * One replay response stays open at most this long, inside the route's `maxDuration`. + * A run still going at the cap is not over: the response ends without a terminal + * event and the client re-attaches from its cursor. + */ +const MAX_STREAM_MS = 60 * 60 * 1000 - 60_000 function extractCanonicalRequestId(value: unknown): string { return typeof value === 'string' && value.length > 0 ? value : '' @@ -451,13 +456,6 @@ async function handleResumeRequestBody({ await sleep(pollDelayMs) } - if (!controllerClosed && Date.now() - startTime >= MAX_STREAM_MS) { - emitTerminalIfMissing(MothershipStreamV1CompletionStatus.error, { - message: 'The stream recovery timed out before completion.', - code: 'resume_timeout', - reason: 'timeout', - }) - } } catch (error) { if (!controllerClosed && !request.signal.aborted) { logger.warn('Stream replay failed', { @@ -473,11 +471,13 @@ async function handleResumeRequestBody({ markSpanForError(rootSpan, error) } finally { request.signal.removeEventListener('abort', abortListener) + // Read before closing: closing the controller here is this route ending, not the client. + const clientDisconnected = controllerClosed closeController() rootSpan.setAttributes({ [TraceAttr.CopilotResumeOutcome]: sawTerminalEvent ? CopilotResumeOutcome.TerminalDelivered - : controllerClosed + : clientDisconnected ? CopilotResumeOutcome.ClientDisconnected : CopilotResumeOutcome.EndedWithoutTerminal, [TraceAttr.CopilotResumeEventCount]: totalEventsFlushed, diff --git a/apps/sim/app/api/mothership/chat/route.ts b/apps/sim/app/api/mothership/chat/route.ts index 9251aa74f12..62316bdd8f8 100644 --- a/apps/sim/app/api/mothership/chat/route.ts +++ b/apps/sim/app/api/mothership/chat/route.ts @@ -10,6 +10,11 @@ import { handleUnifiedChatPost } from '@/lib/mothership/chat/post' import { validateShimEnvelope } from '@/lib/mothership/request/http' import { GET as copilotChatGet } from '@/app/api/copilot/chat/queries' +/** + * Caps this response on serverless hosts only; the Node server bounds nothing with it. + * The run does not depend on this response: one that ends without a terminal is + * re-attached through the replay stream. + */ export const maxDuration = 3600 // Unified chat route surface. diff --git a/apps/sim/app/api/mothership/execute/route.ts b/apps/sim/app/api/mothership/execute/route.ts index 3e397eb70ac..3b460d1e910 100644 --- a/apps/sim/app/api/mothership/execute/route.ts +++ b/apps/sim/app/api/mothership/execute/route.ts @@ -48,6 +48,7 @@ import { import type { ChatContext } from '@/stores/panel' import { hasToolId } from '@/tools/tool-ids' +/** Caps this response on serverless hosts only; the Node server bounds nothing with it. */ export const maxDuration = 3600 const logger = createLogger('MothershipExecuteAPI') diff --git a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx index 95fe0c04e6c..cebc1911396 100644 --- a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx +++ b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx @@ -1089,6 +1089,74 @@ describe('useChat remount send recovery', () => { } }) + it('keeps re-attaching a long turn whose tails deliver events between separate network failures', async () => { + vi.useFakeTimers({ toFake: ['setTimeout', 'clearTimeout'] }) + try { + let tails = 0 + const history: MothershipChatHistory = { + id: 'chat-long-turn', + mode: 'agent', + title: 'Long turn', + messages: [], + activeStreamId: null, + resources: [], + } + mockRequestJson.mockImplementation(() => + Promise.resolve({ + chat: { ...history, activeStreamId: state.postBodies[0]?.userMessageId ?? null }, + }) + ) + state.postBehavior = 'accept' + vi.stubGlobal('fetch', async (input: RequestInfo | URL, init?: RequestInit) => { + const url = String(input) + if (!url.includes('/api/mothership/chat/stream')) return fetchStub(input, init) + if (url.includes('batch=true')) { + return Response.json({ success: true, events: [], status: 'streaming' }) + } + tails++ + const streamId = state.postBodies[0]?.userMessageId ?? '' + const event: MothershipStreamV1EventEnvelope = { + v: 1, + seq: tails, + ts: new Date().toISOString(), + type: 'text', + stream: { streamId, cursor: String(tails) }, + payload: { channel: 'assistant', text: `part ${tails} ` }, + } + return new Response( + new ReadableStream({ + start(controller) { + controller.enqueue(new TextEncoder().encode(`data: ${JSON.stringify(event)}\n\n`)) + }, + pull(controller) { + controller.error(new TypeError('network error')) + }, + }), + { headers: { 'Content-Type': 'text/event-stream' } } + ) + }) + const { getResult } = renderUseChatInChat(history.id, history) + await act(async () => { + void getResult().sendMessage('Keep going for hours') + }) + const errors = new Set() + let seconds = 0 + for (; seconds < 600 && tails < 15; seconds++) { + await act(async () => vi.advanceTimersByTimeAsync(1_000)) + const error = getResult().error + if (error) errors.add(error) + } + + expect(tails).toBeGreaterThanOrEqual(15) + /* Each failure after a tail that delivered events retries at the base delay. */ + expect(seconds).toBeLessThan(60) + expect([...errors]).toEqual([]) + expect(getResult().isSending).toBe(true) + } finally { + vi.useRealTimers() + } + }) + it('sends a queued correction after stopping with more than 10 MiB of tool input', async () => { state.postBehavior = 'tool' state.toolInputPadding = 'x'.repeat(11 * 1024 * 1024) diff --git a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts index 92da00740fe..69707e2bf95 100644 --- a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts +++ b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts @@ -2765,8 +2765,16 @@ export function useChat( abortControllerRef.current?.signal.aborted === true || shouldContinue?.() === false - for (let attempt = 0; attempt <= MAX_RECONNECT_ATTEMPTS; attempt++) { + /** + * An attempt whose tail delivered new events re-attached successfully, so + * the failure after it starts a fresh budget at the base delay. Only + * failures without progress count toward exhaustion, which keeps separate + * network drops hours apart in a long turn from adding up. + */ + let attempt = 0 + while (attempt <= MAX_RECONNECT_ATTEMPTS) { if (isStaleReconnect()) return true + const cursorBeforeAttempt = lastCursorRef.current if (attempt > 0) { const delayMs = Math.min( @@ -2868,6 +2876,7 @@ export function useChat( error: toError(err).message, }) } + attempt = lastCursorRef.current !== cursorBeforeAttempt ? 1 : attempt + 1 } logger.error('All reconnect attempts exhausted', { diff --git a/apps/sim/lib/mothership/application/execute-workflow-use-case.test.ts b/apps/sim/lib/mothership/application/execute-workflow-use-case.test.ts index b219e6524da..4a7e13dc775 100644 --- a/apps/sim/lib/mothership/application/execute-workflow-use-case.test.ts +++ b/apps/sim/lib/mothership/application/execute-workflow-use-case.test.ts @@ -12,7 +12,7 @@ import { executeCopilotWorkflowUseCase, messageForCopilotWorkflowError, } from '@/lib/mothership/application/execute-workflow-use-case' -import { ORCHESTRATION_TIMEOUT_MS } from '@/lib/mothership/constants' +import { COPILOT_APPLICATION_DELEGATION_TTL_MS } from '@/lib/mothership/auth/application-delegation' import { workflowOperations } from '@/lib/workflows/application/operations' const trustedContext = { @@ -54,7 +54,7 @@ describe('Copilot Workflow application adapter', () => { delegationId: 'copilot-tool:tool-call-1', audience: 'sim:workflows', issuedAt: new Date('2026-01-01T00:00:00Z'), - expiresAt: new Date(Date.now() + ORCHESTRATION_TIMEOUT_MS), + expiresAt: new Date(Date.now() + COPILOT_APPLICATION_DELEGATION_TTL_MS), resourceScope: { chatId: 'chat-1', executionId: 'execution-1' }, }, input: { workflowId: 'workflow-1', assertedWorkspaceId: 'workspace-1' }, diff --git a/apps/sim/lib/mothership/auth/application-delegation.ts b/apps/sim/lib/mothership/auth/application-delegation.ts index 6b00108f18a..62993efc220 100644 --- a/apps/sim/lib/mothership/auth/application-delegation.ts +++ b/apps/sim/lib/mothership/auth/application-delegation.ts @@ -1,8 +1,12 @@ import type { DelegatedPrincipal, OrganizationDelegatedPrincipal } from '@sim/auth/principal' -import { ORCHESTRATION_TIMEOUT_MS } from '@/lib/mothership/constants' +import { TOOL_WATCHDOG_LONG_RUNNING_MS } from '@/lib/mothership/constants' -/** Keeps delegated authority valid for the full bounded Copilot orchestration lifetime. */ -export const COPILOT_APPLICATION_DELEGATION_TTL_MS = ORCHESTRATION_TIMEOUT_MS +/** + * Delegated authority is minted per operation and must outlive the longest single + * tool call it authorizes, never the whole run, so it reuses the long-running tool + * watchdog's cap rather than a lifetime of its own. + */ +export const COPILOT_APPLICATION_DELEGATION_TTL_MS = TOOL_WATCHDOG_LONG_RUNNING_MS export interface CopilotExecutionContext { requestMode?: string diff --git a/apps/sim/lib/mothership/auth/file-delegation.test.ts b/apps/sim/lib/mothership/auth/file-delegation.test.ts index 0cbcdc942be..9810c493aee 100644 --- a/apps/sim/lib/mothership/auth/file-delegation.test.ts +++ b/apps/sim/lib/mothership/auth/file-delegation.test.ts @@ -1,12 +1,12 @@ import { describe, expect, it } from 'vitest' import { OrchestrationError } from '@/lib/core/orchestration/types' +import { COPILOT_APPLICATION_DELEGATION_TTL_MS } from '@/lib/mothership/auth/application-delegation' import { createCopilotChatFilePrincipal, createCopilotWorkspaceContextFilePrincipal, messageForCopilotFileError, resolveCopilotFilePrincipal, } from '@/lib/mothership/auth/file-delegation' -import { ORCHESTRATION_TIMEOUT_MS } from '@/lib/mothership/constants' const trustedContext = { userId: 'user-1', @@ -35,7 +35,7 @@ describe('Copilot file delegation', () => { }, }) expect(principal.expiresAt.getTime() - principal.issuedAt.getTime()).toBe( - ORCHESTRATION_TIMEOUT_MS + COPILOT_APPLICATION_DELEGATION_TTL_MS ) }) diff --git a/apps/sim/lib/mothership/constants.ts b/apps/sim/lib/mothership/constants.ts index d61f1bc0d54..2fda1e8ef9a 100644 --- a/apps/sim/lib/mothership/constants.ts +++ b/apps/sim/lib/mothership/constants.ts @@ -10,8 +10,15 @@ export const SIM_AGENT_API_URL = ? rawAgentUrl : SIM_AGENT_API_URL_DEFAULT -/** Default timeout for the copilot orchestration stream loop (60 min). */ -export const ORCHESTRATION_TIMEOUT_MS = 3_600_000 +/** + * How long a worker SSE leg may stay silent before Sim treats the connection as + * lost and re-attaches. The worker writes a keepalive comment whenever a leg has + * been quiet for 10 s (checked every 15 s), independent of model or tool progress, + * so a healthy leg is never silent for more than about 25 s. It stays well under + * the idle timeouts of network intermediaries, which can drop a silent connection + * without closing it. + */ +export const WORKER_STREAM_IDLE_TIMEOUT_MS = 120_000 /** * Watchdog cap for a single sim-executed copilot tool. A tool that neither @@ -25,10 +32,15 @@ export const TOOL_WATCHDOG_DEFAULT_MS = 60_000 * Watchdog cap for tool classes with legitimately long runtimes (workflow * executions, media/image generation, sandboxed code, deep research). Those * tools carry their own inner budgets (plan execution timeouts, sandbox - * timeouts), so this cap only backstops a true hang and sits above all of - * them — matching ORCHESTRATION_TIMEOUT_MS so it never undercuts a legal run. + * timeouts), so this cap only backstops a true hang and sits above all of them. */ -export const TOOL_WATCHDOG_LONG_RUNNING_MS = ORCHESTRATION_TIMEOUT_MS +export const TOOL_WATCHDOG_LONG_RUNNING_MS = 60 * 60 * 1000 + +/** How long a tool call held for the user's approval waits for an answer. */ +export const PERMISSION_WAIT_TIMEOUT_MS = 60 * 60 * 1000 + +/** How long a client-executed tool (browser or desktop app) may take to report its result. */ +export const CLIENT_TOOL_RESULT_TIMEOUT_MS = 60 * 60 * 1000 /** Extra slack the resume gate allows past the slowest pending tool's watchdog. */ export const TOOL_WATCHDOG_RESUME_GRACE_MS = 30_000 @@ -43,7 +55,7 @@ export const STREAM_TIMEOUT_MS = 3_600_000 * Workflow tools are client-routed, but the only thing that starts one is the * mounted chat view — a call frame that arrives while the user is on a * different chat is never dispatched by anyone, and the turn used to park for - * the full STREAM_TIMEOUT_MS. The real pickup path (stream frame -> execute + * the full CLIENT_TOOL_RESULT_TIMEOUT_MS. The real pickup path (stream frame -> execute * POST -> claim) lands in ~1-3s, so 30s is an order of magnitude of headroom * and cannot steal work from a live tab. */ diff --git a/apps/sim/lib/mothership/request/go/stream.test.ts b/apps/sim/lib/mothership/request/go/stream.test.ts index 51b07fb5ef7..7568ccd7626 100644 --- a/apps/sim/lib/mothership/request/go/stream.test.ts +++ b/apps/sim/lib/mothership/request/go/stream.test.ts @@ -6,7 +6,7 @@ import { workspaceFilesListMock, workspaceFilesListMockFns, } from '@sim/testing/mocks/workspace-files-list.mock' -import { beforeEach, describe, expect, it, vi } from 'vitest' +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' import { MothershipStreamV1CompletionStatus, MothershipStreamV1EventType, @@ -86,6 +86,8 @@ import { runStreamLoop, STREAM_ENDED_WITHOUT_TERMINAL_MESSAGE, StreamEndedWithoutTerminalError, + WorkerStreamInterruptedError, + WorkerUnreachableError, } from '@/lib/mothership/request/go/stream' import { createProviderToolCallIdentity, @@ -1023,4 +1025,174 @@ describe('copilot go stream helpers', () => { ) ).toBe(true) }) + + describe('worker stream liveness without a caller deadline', () => { + /** Well past the idle timeout, and under common intermediary idle cuts. */ + const INTERMEDIARY_IDLE_MS = 300_000 + const encoder = new TextEncoder() + const frame = (event: unknown) => encoder.encode(`data: ${JSON.stringify(event)}\n\n`) + const firstText = createEvent({ + streamId: 'long-stream', + cursor: '1', + seq: 1, + requestId: 'req-long', + type: MothershipStreamV1EventType.text, + payload: { channel: 'assistant', text: 'working' }, + }) + + function settle(promise: Promise) { + const state: { done: boolean; error?: unknown } = { done: false } + promise.then( + () => { + state.done = true + }, + (error: unknown) => { + state.done = true + state.error = error + } + ) + return state + } + + afterEach(() => { + vi.useRealTimers() + }) + + it('keeps a leg open past an hour while the worker sends keepalives', async () => { + vi.useFakeTimers() + const complete = createEvent({ + streamId: 'long-stream', + cursor: '2', + seq: 2, + requestId: 'req-long', + type: MothershipStreamV1EventType.complete, + payload: { status: MothershipStreamV1CompletionStatus.complete }, + }) + vi.mocked(fetch).mockResolvedValueOnce( + new Response( + new ReadableStream({ + start(controller) { + controller.enqueue(frame(firstText)) + const keepalive = setInterval(() => { + try { + controller.enqueue(encoder.encode(': keepalive\n\n')) + } catch { + clearInterval(keepalive) + } + }, 25_000) + setTimeout( + () => { + clearInterval(keepalive) + controller.enqueue(frame(complete)) + controller.close() + }, + 2 * 60 * 60 * 1000 + ) + }, + }), + { status: 200, headers: { 'Content-Type': 'text/event-stream' } } + ) + ) + const context = createStreamingContext() + const state = settle( + runStreamLoop( + 'https://example.com/mothership/stream', + {}, + context, + turnScopedExecContext(), + { + flushAfterEvent: false, + } + ) + ) + + await vi.advanceTimersByTimeAsync(2 * 60 * 60 * 1000 + 1_000) + + expect(state).toEqual({ done: true }) + expect(context.errors).toEqual([]) + expect(context.streamComplete).toBe(true) + }) + + it('fails a silent leg as a retryable interruption before an intermediary drops it', async () => { + vi.useFakeTimers() + vi.mocked(fetch).mockResolvedValueOnce( + new Response( + new ReadableStream({ + start(controller) { + controller.enqueue(frame(firstText)) + }, + }), + { status: 200, headers: { 'Content-Type': 'text/event-stream' } } + ) + ) + const state = settle( + runStreamLoop( + 'https://example.com/mothership/stream', + {}, + createStreamingContext(), + turnScopedExecContext(), + { flushAfterEvent: false } + ) + ) + + await vi.advanceTimersByTimeAsync(60_000) + expect(state.done).toBe(false) + await vi.advanceTimersByTimeAsync(INTERMEDIARY_IDLE_MS - 60_000) + + expect(state.done).toBe(true) + expect(state.error).toBeInstanceOf(WorkerStreamInterruptedError) + }) + + it('fails a worker whose error body stalls instead of waiting on it forever', async () => { + vi.useFakeTimers() + vi.mocked(fetch).mockResolvedValueOnce( + new Response(new ReadableStream(), { + status: 503, + headers: { 'Content-Type': 'application/json' }, + }) + ) + const state = settle( + runStreamLoop( + 'https://example.com/mothership/stream', + {}, + createStreamingContext(), + turnScopedExecContext(), + {} + ) + ) + + await vi.advanceTimersByTimeAsync(INTERMEDIARY_IDLE_MS) + + expect(state.done).toBe(true) + expect(state.error).toMatchObject({ name: 'CopilotBackendError', status: 503 }) + }) + + it('fails a worker that never answers as unreachable before an intermediary drops it', async () => { + vi.useFakeTimers() + vi.mocked(fetch).mockImplementationOnce( + (_url, options) => + new Promise((_resolve, reject) => { + options?.signal?.addEventListener('abort', () => reject(options.signal?.reason), { + once: true, + }) + }) + ) + const context = createStreamingContext() + const state = settle( + runStreamLoop( + 'https://example.com/mothership/stream', + {}, + context, + turnScopedExecContext(), + {} + ) + ) + + await vi.advanceTimersByTimeAsync(INTERMEDIARY_IDLE_MS) + + expect(state.done).toBe(true) + expect(state.error).toBeInstanceOf(WorkerUnreachableError) + expect(context.wasAborted).not.toBe(true) + }) + }) }) diff --git a/apps/sim/lib/mothership/request/go/stream.ts b/apps/sim/lib/mothership/request/go/stream.ts index b55b623458d..336423cbdd7 100644 --- a/apps/sim/lib/mothership/request/go/stream.ts +++ b/apps/sim/lib/mothership/request/go/stream.ts @@ -2,7 +2,7 @@ import { type Context, SpanStatusCode } from '@opentelemetry/api' import { createLogger } from '@sim/logger' import { getErrorMessage } from '@sim/utils/errors' import { toRecordOrNull } from '@sim/utils/object' -import { ORCHESTRATION_TIMEOUT_MS } from '@/lib/mothership/constants' +import { WORKER_STREAM_IDLE_TIMEOUT_MS } from '@/lib/mothership/constants' import { MothershipStreamV1EventType } from '@/lib/mothership/generated/mothership-stream-v1' import { CopilotSseCloseReason } from '@/lib/mothership/generated/trace-attribute-values-v1' import { TraceAttr } from '@/lib/mothership/generated/trace-attributes-v1' @@ -155,6 +155,11 @@ export class StreamEndedWithoutTerminalError extends Error { } } +/** No bytes, keepalives included, arrived from the worker within the idle timeout. */ +function workerStreamIdleError(): Error { + return new Error(`No bytes from the worker in ${WORKER_STREAM_IDLE_TIMEOUT_MS / 1000} s`) +} + /** * Options for the shared stream processing loop. */ @@ -177,6 +182,11 @@ export interface StreamLoopOptions extends OrchestratorOptions { * Handles: fetch -> parse -> normalize -> dedupe -> subagent routing -> handler dispatch. * Callers provide the fetch URL/options and can intercept events via onBeforeDispatch. * Feature-specific normalization runs through dedicated adapters before the raw event is forwarded. + * + * A leg has no wall clock unless the caller sets `timeout`. Its liveness is the + * worker's own traffic: while Sim waits for response headers or the next bytes, + * {@link WORKER_STREAM_IDLE_TIMEOUT_MS} of silence fails the leg as unreachable + * or interrupted, which the caller's retry window re-attaches. */ export async function runStreamLoop( fetchUrl: string, @@ -185,9 +195,21 @@ export async function runStreamLoop( execContext: ExecutionContext, options: StreamLoopOptions ): Promise { - const { timeout = ORCHESTRATION_TIMEOUT_MS, abortSignal } = options - const timeoutSignal = AbortSignal.timeout(Math.ceil(timeout)) - const requestSignal = abortSignal ? AbortSignal.any([abortSignal, timeoutSignal]) : timeoutSignal + const { timeout, abortSignal } = options + const idle = new AbortController() + let idleTimer: ReturnType | undefined + const armIdleTimeout = (onIdle?: () => void) => { + clearTimeout(idleTimer) + idleTimer = setTimeout(() => { + idle.abort(workerStreamIdleError()) + onIdle?.() + }, WORKER_STREAM_IDLE_TIMEOUT_MS) + } + const requestSignal = AbortSignal.any([ + idle.signal, + ...(abortSignal ? [abortSignal] : []), + ...(timeout === undefined ? [] : [AbortSignal.timeout(Math.ceil(timeout))]), + ]) const filePreviewAdapterState = createFilePreviewAdapterState() const attemptedInlineImages = new Set() @@ -200,6 +222,7 @@ export async function runStreamLoop( }) const fetchStart = performance.now() let response: Response + armIdleTimeout() try { response = await fetchGo(fetchUrl, { ...fetchOptions, @@ -218,8 +241,11 @@ export async function runStreamLoop( headersMs: Math.round(performance.now() - fetchStart), } context.trace.endSpan(fetchSpan, abortSignal?.aborted ? 'cancelled' : 'error') + if (idle.signal.aborted) throw new WorkerUnreachableError(idle.signal.reason) if (requestSignal.aborted) throw error throw new WorkerUnreachableError(error) + } finally { + clearTimeout(idleTimer) } const headersElapsedMs = Math.round(performance.now() - fetchStart) fetchSpan.attributes = { @@ -230,7 +256,12 @@ export async function runStreamLoop( if (!response.ok) { context.trace.endSpan(fetchSpan, 'error') - const errorText = await response.text().catch(() => '') + // An error body is bounded by the same silence as the leg; a stalled one reads as empty. + armIdleTimeout() + const errorText = await new Promise((resolve) => { + idle.signal.addEventListener('abort', () => resolve(''), { once: true }) + response.text().then(resolve, () => resolve('')) + }).finally(() => clearTimeout(idleTimer)) if (response.status === 402) { throw new BillingLimitError(execContext.userId) @@ -306,11 +337,22 @@ export async function runStreamLoop( const reader: ReadableStreamDefaultReader = { async read() { let result: ReadableStreamReadResult + armIdleTimeout(() => rawReader.cancel(idle.signal.reason).catch(() => {})) try { result = await rawReader.read() } catch (error) { + if (idle.signal.aborted) { + endedOn = CopilotSseCloseReason.Timeout + throw new WorkerStreamInterruptedError(idle.signal.reason) + } if (requestSignal.aborted) throw error throw new WorkerStreamInterruptedError(error) + } finally { + clearTimeout(idleTimer) + } + if (idle.signal.aborted) { + endedOn = CopilotSseCloseReason.Timeout + throw new WorkerStreamInterruptedError(idle.signal.reason) } if (!result.done && result.value) { const now = performance.now() @@ -329,12 +371,15 @@ export async function runStreamLoop( }, } - const timeoutId = setTimeout(() => { - context.errors.push('Request timed out') - context.streamComplete = true - endedOn = CopilotSseCloseReason.Timeout - reader.cancel().catch(() => {}) - }, timeout) + const timeoutId = + timeout === undefined + ? undefined + : setTimeout(() => { + context.errors.push('Request timed out') + context.streamComplete = true + endedOn = CopilotSseCloseReason.Timeout + reader.cancel().catch(() => {}) + }, timeout) try { await processSSEStream(reader, abortSignal, async (raw) => { @@ -553,6 +598,7 @@ export async function runStreamLoop( flushSubagentThinkingBlock(context) flushThinkingBlock(context) clearTimeout(timeoutId) + clearTimeout(idleTimer) // Legacy TraceCollector span (consumed by the in-memory trace // collector, kept for backwards compatibility with existing diff --git a/apps/sim/lib/mothership/request/handlers/tool.ts b/apps/sim/lib/mothership/request/handlers/tool.ts index cb65ad989ef..00d828f2753 100644 --- a/apps/sim/lib/mothership/request/handlers/tool.ts +++ b/apps/sim/lib/mothership/request/handlers/tool.ts @@ -9,8 +9,8 @@ import type { } from '@/lib/mothership/async-runs/lifecycle' import { upsertAsyncToolCall } from '@/lib/mothership/async-runs/repository' import { + CLIENT_TOOL_RESULT_TIMEOUT_MS, COPILOT_WORKFLOW_TOOL_CLIENT_GRACE_MS, - STREAM_TIMEOUT_MS, } from '@/lib/mothership/constants' import { MothershipStreamV1AsyncToolRecordStatus, @@ -883,7 +883,7 @@ async function dispatchToolExecution( */ function waitForClientExecution(): Promise { toolCall.status = 'executing' - const timeoutMs = options.timeout || STREAM_TIMEOUT_MS + const timeoutMs = options.timeout || CLIENT_TOOL_RESULT_TIMEOUT_MS return withCopilotSpan( TraceSpan.CopilotToolWaitForClientResult, { diff --git a/apps/sim/lib/mothership/request/lifecycle/run.ts b/apps/sim/lib/mothership/request/lifecycle/run.ts index 9a9ac2cd667..0f2ca8e919b 100644 --- a/apps/sim/lib/mothership/request/lifecycle/run.ts +++ b/apps/sim/lib/mothership/request/lifecycle/run.ts @@ -247,7 +247,7 @@ export interface CopilotLifecycleOptions extends OrchestratorOptions { * * Beyond the flag, gating is limited to interactive mothership chats: that is * the only surface with a UI that can answer a prompt, so enabling it anywhere - * else would hang the turn until the orchestration timeout with nothing to click. + * else would hang the turn until the permission wait expires with nothing to click. */ async function resolveToolPermissions( options: CopilotLifecycleOptions diff --git a/apps/sim/lib/mothership/request/lifecycle/stream-retry.test.ts b/apps/sim/lib/mothership/request/lifecycle/stream-retry.test.ts index 7ef1e679e49..558566dbe31 100644 --- a/apps/sim/lib/mothership/request/lifecycle/stream-retry.test.ts +++ b/apps/sim/lib/mothership/request/lifecycle/stream-retry.test.ts @@ -10,7 +10,7 @@ import { StreamRetryWindow } from '@/lib/mothership/request/lifecycle/stream-ret afterEach(() => vi.useRealTimers()) describe('stream recovery budget', () => { - it('stops an ended-without-terminal stream after three retries despite a long task budget', () => { + it('stops an ended-without-terminal stream after three retries on a leg with no deadline', () => { vi.useFakeTimers() const error = new StreamEndedWithoutTerminalError('/api/mothership') const retry = new StreamRetryWindow() @@ -21,7 +21,6 @@ describe('stream recovery budget', () => { } expect(retry.nextDelay(error)).toBeNull() expect(retry.attempt).toBe(3) - expect(retry.remainingMs()).toBeGreaterThan(3_500_000) }) it.each([ @@ -121,7 +120,6 @@ describe('stream recovery budget', () => { } expect(Date.now() - firstFailure).toBeGreaterThan(110_000) expect(Date.now() - firstFailure).toBeLessThanOrEqual(120_000) - expect(retry.remainingMs()).toBeGreaterThan(2_800_000) }) it('never extends the original execution deadline', () => { @@ -158,4 +156,64 @@ describe('stream recovery budget', () => { } expect(retry.nextDelay(error)).toBeNull() }) + + it('has no leg deadline unless the caller sets one', () => { + vi.useFakeTimers() + const retry = new StreamRetryWindow() + vi.advanceTimersByTime(3 * 60 * 60 * 1000) + expect(retry.remainingMs()).toBeUndefined() + expect(retry.nextDelay(new StreamEndedWithoutTerminalError('/api/mothership'))).not.toBeNull() + }) + + it('gives an interruption hours into a healthy leg a fresh reachable budget', () => { + vi.useFakeTimers() + const error = new WorkerStreamInterruptedError(new Error('socket closed')) + const retry = new StreamRetryWindow() + for (let index = 0; index < 3; index++) { + const delay = retry.nextDelay(error) + expect(delay).not.toBeNull() + vi.advanceTimersByTime(delay ?? 0) + } + for (let minute = 0; minute < 3 * 60; minute++) { + retry.recovered() + vi.advanceTimersByTime(60_000) + } + for (let index = 0; index < 3; index++) { + const delay = retry.nextDelay(error) + expect(delay).not.toBeNull() + vi.advanceTimersByTime(delay ?? 0) + } + expect(retry.nextDelay(error)).toBeNull() + }) + + it('keeps the reachable budget spent when the leg fails again soon after re-attaching', () => { + vi.useFakeTimers() + const error = new WorkerStreamInterruptedError(new Error('socket closed')) + const retry = new StreamRetryWindow() + for (let index = 0; index < 3; index++) { + const delay = retry.nextDelay(error) + expect(delay).not.toBeNull() + vi.advanceTimersByTime(delay ?? 0) + for (let event = 0; event < 5; event++) { + retry.recovered() + vi.advanceTimersByTime(1_000) + } + } + retry.recovered() + expect(retry.nextDelay(error)).toBeNull() + }) + + it('does not refill the reachable budget for a leg that delivered one event and then went quiet', () => { + vi.useFakeTimers() + const error = new WorkerStreamInterruptedError(new Error('socket closed')) + const retry = new StreamRetryWindow() + for (let index = 0; index < 3; index++) { + const delay = retry.nextDelay(error) + expect(delay).not.toBeNull() + vi.advanceTimersByTime(delay ?? 0) + } + retry.recovered() + vi.advanceTimersByTime(30 * 60_000) + expect(retry.nextDelay(error)).toBeNull() + }) }) diff --git a/apps/sim/lib/mothership/request/lifecycle/stream-retry.ts b/apps/sim/lib/mothership/request/lifecycle/stream-retry.ts index 08da02415cc..2650ac6f195 100644 --- a/apps/sim/lib/mothership/request/lifecycle/stream-retry.ts +++ b/apps/sim/lib/mothership/request/lifecycle/stream-retry.ts @@ -1,5 +1,4 @@ import { backoffWithJitter } from '@sim/utils/retry' -import { ORCHESTRATION_TIMEOUT_MS } from '@/lib/mothership/constants' import { StreamContinuityError } from '@/lib/mothership/request/go/parser' import { CopilotBackendError, @@ -20,22 +19,34 @@ const GATEWAY_STATUSES: ReadonlySet = new Set([502, 503, 504]) * reattach, never a second run. */ const WORKER_REPLACEMENT_WINDOW_MS = 120_000 +/** + * A leg whose delivered events span this long since it last re-attached has proven + * healthy, so a later interruption, possibly hours on, gets the reachable budget + * afresh. The span runs from its first event to its latest, so a leg that delivered + * one event and then only kept alive has made no progress, and a leg that fails again + * sooner keeps spending the same three retries: a deterministic failure stays bounded. + */ +const HEALTHY_STREAM_REPLENISH_MS = 5 * 60_000 /** - * Recovery is bounded independently of the healthy run's execution budget, by - * two budgets that never share state: an unreachable worker gets a two-minute + * Recovery is bounded independently of the healthy leg's lifetime, by two + * budgets that never share state: an unreachable worker gets a two-minute * window from the moment it stopped answering, and any failure of a worker that - * did answer gets the original three retries within 30 s. + * did answer gets three retries within 30 s, replenished only after + * {@link HEALTHY_STREAM_REPLENISH_MS} of healthy streaming. A leg has no deadline + * unless the caller sets one. */ export class StreamRetryWindow { - private readonly deadline: number + private readonly deadline?: number private firstFailureAt?: number + private streamingSince?: number + private lastEventAt?: number private firstUnreachableAt?: number private unreachableAttempt = 0 private attempt = 0 - constructor(timeoutMs = ORCHESTRATION_TIMEOUT_MS) { - this.deadline = Date.now() + timeoutMs + constructor(timeoutMs?: number) { + this.deadline = timeoutMs === undefined ? undefined : Date.now() + timeoutMs } /** Retries taken across both budgets, for logs and spans. */ @@ -43,21 +54,26 @@ export class StreamRetryWindow { return this.attempt + this.unreachableAttempt } - remainingMs(): number { + /** Time left before the caller's deadline, or `undefined` when the leg has none. */ + remainingMs(): number | undefined { + if (this.deadline === undefined) return undefined const remaining = this.deadline - Date.now() if (remaining <= 0) throw new Error('The connection to the assistant could not be restored in time.') return remaining } - /** The worker answered, so a later loss of it starts a fresh unreachable window. */ + /** The worker delivered an event, so a later loss of it starts a fresh unreachable window. */ recovered(): void { - this.firstUnreachableAt = undefined - this.unreachableAttempt = 0 + this.resetUnreachable() + this.lastEventAt = Date.now() + this.streamingSince ??= this.lastEventAt } nextDelay(error: unknown, signal?: AbortSignal): number | null { if (signal?.aborted || !isRetryableStreamError(error)) return null + this.replenishAfterHealthyStreaming() + this.streamingSince = undefined if (isWorkerUnreachable(error)) { this.firstUnreachableAt ??= Date.now() const delay = backoff(this.unreachableAttempt) @@ -66,7 +82,7 @@ export class StreamRetryWindow { return delay } // Any other retryable failure is an answer from the worker. - this.recovered() + this.resetUnreachable() this.firstFailureAt ??= Date.now() if (this.attempt >= MAX_STREAM_RETRIES) return null const delay = backoff(this.attempt) @@ -75,8 +91,26 @@ export class StreamRetryWindow { return delay } + private resetUnreachable(): void { + this.firstUnreachableAt = undefined + this.unreachableAttempt = 0 + } + + private replenishAfterHealthyStreaming(): void { + if ( + this.streamingSince !== undefined && + this.lastEventAt !== undefined && + this.lastEventAt - this.streamingSince >= HEALTHY_STREAM_REPLENISH_MS + ) { + this.attempt = 0 + this.firstFailureAt = undefined + } + } + private fits(delay: number, recoveryDeadline: number): boolean { - return Date.now() + delay < Math.min(this.deadline, recoveryDeadline) + return ( + Date.now() + delay < Math.min(this.deadline ?? Number.POSITIVE_INFINITY, recoveryDeadline) + ) } } diff --git a/apps/sim/lib/mothership/request/session/abort.ts b/apps/sim/lib/mothership/request/session/abort.ts index 84f97632a53..7ff29fa5000 100644 --- a/apps/sim/lib/mothership/request/session/abort.ts +++ b/apps/sim/lib/mothership/request/session/abort.ts @@ -8,7 +8,7 @@ import { TraceAttr } from '@/lib/mothership/generated/trace-attributes-v1' import { TraceSpan } from '@/lib/mothership/generated/trace-spans-v1' import { withCopilotSpan } from '@/lib/mothership/request/otel' import { AbortReason } from './abort-reason' -import { clearAbortMarker, hasAbortMarker, writeAbortMarker } from './buffer' +import { clearAbortMarker, hasAbortMarker, refreshBufferTtl, writeAbortMarker } from './buffer' import { type ChatStreamLease, chatStreamLockKey, @@ -395,7 +395,9 @@ export function startAbortPoller( streamId, ...(requestId ? { requestId } : {}), }) + return } + await refreshBufferTtl(streamId) } catch (error) { logger.warn('Failed to extend chat stream lock TTL', { chatId, diff --git a/apps/sim/lib/mothership/request/session/buffer-ttl.integration.ts b/apps/sim/lib/mothership/request/session/buffer-ttl.integration.ts new file mode 100644 index 00000000000..9fc83b57c48 --- /dev/null +++ b/apps/sim/lib/mothership/request/session/buffer-ttl.integration.ts @@ -0,0 +1,109 @@ +/** + * The replay buffer of a live run against real Redis: it must outlive a park longer + * than its idle TTL while its controller still holds the chat lock, and a buffer + * whose numbering restarted must not pass a reconnect cursor off as in range. + */ +import { afterAll, describe, expect, it, vi } from 'vitest' + +const { redisUrl } = await vi.hoisted(async () => { + const { readTestRedisUrl } = await import('@sim/db/testing/test-infrastructure') + const url = readTestRedisUrl() + process.env.REDIS_URL = url + /** The park below outlasts it more than twice over in real time. */ + process.env.COPILOT_STREAM_TTL_SECONDS = '5' + return { redisUrl: url } +}) + +import { sleep } from '@sim/utils/helpers' +import { generateId } from '@sim/utils/id' +import { closeRedisConnection, getRedisClient } from '@/lib/core/config/redis' +import { getRedisBudgetKeys } from '@/lib/core/redis/byte-budget.server' +import { MothershipStreamV1EventType } from '@/lib/mothership/generated/mothership-stream-v1' +import { + acquirePendingChatStream, + releasePendingChatStream, + startAbortPoller, +} from '@/lib/mothership/request/session/abort' +import { + allocateCursor, + appendEvents, + getLatestSeq, + readEvents, + refreshBufferTtl, + scheduleBufferCleanup, +} from '@/lib/mothership/request/session/buffer' +import { createEvent } from '@/lib/mothership/request/session/event' +import { checkForReplayGap } from '@/lib/mothership/request/session/recovery' + +async function appendText(streamId: string, text: string): Promise { + const { seq, cursor } = await allocateCursor(streamId) + await appendEvents([ + createEvent({ + streamId, + cursor, + seq, + requestId: 'req-ttl', + type: MothershipStreamV1EventType.text, + payload: { channel: 'assistant', text }, + }), + ]) + return seq +} + +describe.runIf(Boolean(redisUrl))('replay buffer lifetime', () => { + afterAll(async () => { + await closeRedisConnection() + }) + + it('keeps a live run’s buffer and its byte counter through a park longer than their TTLs', async () => { + const chatId = generateId() + const streamId = generateId() + expect(await acquirePendingChatStream(chatId, streamId, 0)).toBe(true) + await appendText(streamId, 'before the park') + const [ownerBudgetKey] = getRedisBudgetKeys({ kind: 'copilot_stream', id: streamId }) + const chargedBytes = await getRedisClient()!.get(ownerBudgetKey) + /** The counter's own TTL is an hour; shortening it stands in for a park that long. */ + await getRedisClient()!.expire(ownerBudgetKey, 5) + + vi.useFakeTimers({ toFake: ['Date'] }) + const poller = startAbortPoller(streamId, new AbortController(), { chatId, pollMs: 50 }) + try { + for (let tick = 0; tick < 24; tick++) { + vi.setSystemTime(Date.now() + 21_000) + await sleep(500) + } + } finally { + clearInterval(poller) + vi.useRealTimers() + await releasePendingChatStream(chatId, streamId) + } + + expect(await getLatestSeq(streamId)).toBe(1) + expect((await readEvents(streamId, '0')).map((event) => event.seq)).toEqual([1]) + expect(chargedBytes).not.toBeNull() + expect(await getRedisClient()!.get(ownerBudgetKey)).toBe(chargedBytes) + expect(await appendText(streamId, 'after the park')).toBe(2) + }) + + it('never re-extends a finished stream’s buffer after its cleanup was scheduled', async () => { + const streamId = generateId() + await appendText(streamId, 'done') + await scheduleBufferCleanup(streamId) + + await refreshBufferTtl(streamId) + + const redis = getRedisClient()! + expect(await redis.ttl(`mothership_stream:${streamId}:events`)).toBeGreaterThan(250) + expect(await redis.ttl(`mothership_stream:${streamId}:seq`)).toBeGreaterThan(250) + }) + + it('reports a gap to a cursor ahead of a buffer whose numbering restarted', async () => { + const streamId = generateId() + for (let index = 0; index < 5; index++) await appendText(streamId, `part ${index}`) + const redis = getRedisClient()! + await redis.del(`mothership_stream:${streamId}:events`, `mothership_stream:${streamId}:seq`) + await appendText(streamId, 'after expiry') + + expect(await checkForReplayGap(streamId, '5')).not.toBeNull() + }) +}) diff --git a/apps/sim/lib/mothership/request/session/buffer.ts b/apps/sim/lib/mothership/request/session/buffer.ts index e1760adf11d..68fd24cd447 100644 --- a/apps/sim/lib/mothership/request/session/buffer.ts +++ b/apps/sim/lib/mothership/request/session/buffer.ts @@ -42,6 +42,11 @@ function getAbortKey(streamId: string) { return `${STREAM_OUTBOX_PREFIX}${streamId}:abort` } +/** Marks a stream whose cleanup is scheduled, so a late heartbeat cannot revive it. */ +function getClosedKey(streamId: string) { + return `${STREAM_OUTBOX_PREFIX}${streamId}:closed` +} + export type StreamConfig = { ttlSeconds: number eventLimit: number @@ -127,18 +132,59 @@ export async function clearBuffer(streamId: string, operation = 'clear_outbox'): getEventsKey(streamId), getSeqKey(streamId), getAbortKey(streamId), + getClosedKey(streamId), ownerBudgetKey ) }) } +/** + * KEYS: [events, seq, ownerBudget, closed] + * ARGV: [ttlSeconds, budgetTtlSeconds] + */ +const REFRESH_BUFFER_TTL_SCRIPT = ` +if redis.call('EXISTS', KEYS[4]) == 1 then return 0 end +redis.call('EXPIRE', KEYS[1], ARGV[1]) +redis.call('EXPIRE', KEYS[2], ARGV[1]) +redis.call('EXPIRE', KEYS[3], ARGV[2]) +return 1 +` + +/** + * Slides a live stream's replay TTLs, and its byte counter's, without an append. They + * otherwise move only when an event lands, so a run parked on a long tool call or + * approval would lose its replay history and restart its numbering while it still + * runs, or keep its history after the counter that accounts for it expired. A stream + * whose cleanup is already scheduled is left to expire. + */ +export async function refreshBufferTtl(streamId: string): Promise { + const { ttlSeconds } = getStreamConfig() + const [ownerBudgetKey] = getRedisBudgetKeys({ kind: 'copilot_stream', id: streamId }) + const budgetTtlSeconds = Math.max(getRedisBudgetLimits('copilot_stream').ttlSeconds, ttlSeconds) + await withRedisRetry({ operation: 'refresh_outbox_ttl', streamId }, async (redis) => { + await redis.eval( + REFRESH_BUFFER_TTL_SCRIPT, + 4, + getEventsKey(streamId), + getSeqKey(streamId), + ownerBudgetKey, + getClosedKey(streamId), + ttlSeconds, + budgetTtlSeconds + ) + }) +} + export async function scheduleBufferCleanup( streamId: string, ttlSeconds = DEFAULT_COMPLETED_TTL_SECONDS ): Promise { try { await withRedisRetry({ operation: 'schedule_outbox_cleanup', streamId }, async (redis) => { + // The marker goes first: a refresh that lands before it is overridden by the + // expirations below, and one that lands after it sees the marker and does nothing. const pipeline = redis.pipeline() + pipeline.set(getClosedKey(streamId), '1', 'EX', ttlSeconds) pipeline.expire(getEventsKey(streamId), ttlSeconds) pipeline.expire(getSeqKey(streamId), ttlSeconds) pipeline.expire(getAbortKey(streamId), ttlSeconds) diff --git a/apps/sim/lib/mothership/request/session/recovery.ts b/apps/sim/lib/mothership/request/session/recovery.ts index c4a74a143f0..bb0153c8a14 100644 --- a/apps/sim/lib/mothership/request/session/recovery.ts +++ b/apps/sim/lib/mothership/request/session/recovery.ts @@ -45,14 +45,16 @@ export async function checkForReplayGap( [TraceAttr.CopilotRecoveryLatestSeq]: latestSeq ?? -1, }) + /* Trimmed below the ring, or ahead of a buffer whose numbering restarted after it + expired: either way the events after the cursor are not the ones it names. */ if ( latestSeq !== null && latestSeq > 0 && oldestSeq !== null && - requestedAfterSeq < oldestSeq - 1 + (requestedAfterSeq < oldestSeq - 1 || requestedAfterSeq > latestSeq) ) { const resolvedRequestId = await resolveReplayGapRequestId(streamId, latestSeq, requestId) - logger.warn('Replay gap detected: requested cursor is below oldest available event', { + logger.warn('Replay gap detected: requested cursor is outside the retained events', { streamId, requestedAfterSeq, oldestAvailableSeq: oldestSeq, diff --git a/apps/sim/lib/mothership/request/tools/executor.ts b/apps/sim/lib/mothership/request/tools/executor.ts index c75d179d423..f63136c4fe7 100644 --- a/apps/sim/lib/mothership/request/tools/executor.ts +++ b/apps/sim/lib/mothership/request/tools/executor.ts @@ -14,7 +14,11 @@ import { upsertAsyncToolCall, } from '@/lib/mothership/async-runs/repository' import { withToolServiceMeter } from '@/lib/mothership/billing/service-meter' -import { TOOL_WATCHDOG_DEFAULT_MS, TOOL_WATCHDOG_LONG_RUNNING_MS } from '@/lib/mothership/constants' +import { + PERMISSION_WAIT_TIMEOUT_MS, + TOOL_WATCHDOG_DEFAULT_MS, + TOOL_WATCHDOG_LONG_RUNNING_MS, +} from '@/lib/mothership/constants' import { MothershipStreamV1AsyncToolRecordStatus, MothershipStreamV1EventType, @@ -250,7 +254,7 @@ export function toolWatchdogTimeoutMs(toolName: string | undefined): number { /** * How long the resume gate may wait on one pending tool call. Permission - * prompts receive the long-running budget. Browser calls share the renderer's + * prompts wait as long as the permission wait itself. Browser calls share the renderer's * budget so authorization and native queueing cannot outlive the resume gate. */ export function pendingToolWaitBudgetMs( @@ -258,7 +262,7 @@ export function pendingToolWaitBudgetMs( | (Pick & Partial>) | undefined ): number { - if (toolCall?.status === 'awaiting_approval') return TOOL_WATCHDOG_LONG_RUNNING_MS + if (toolCall?.status === 'awaiting_approval') return PERMISSION_WAIT_TIMEOUT_MS const executableName = toolCall?.execName ?? toolCall?.name if (executableName && isCurrentBrowserToolName(executableName)) { return browserToolRendererTimeoutMs(executableName, toolCall?.params) diff --git a/apps/sim/lib/mothership/request/tools/permission.ts b/apps/sim/lib/mothership/request/tools/permission.ts index 7eb360f67cb..850a1f5da87 100644 --- a/apps/sim/lib/mothership/request/tools/permission.ts +++ b/apps/sim/lib/mothership/request/tools/permission.ts @@ -2,7 +2,7 @@ import { createLogger } from '@sim/logger' import { TERMINAL_TOOL_NAME } from '@sim/terminal-protocol' import { getErrorMessage } from '@sim/utils/errors' import type { AsyncCompletionSignal } from '@/lib/mothership/async-runs/lifecycle' -import { ORCHESTRATION_TIMEOUT_MS } from '@/lib/mothership/constants' +import { PERMISSION_WAIT_TIMEOUT_MS } from '@/lib/mothership/constants' import { MothershipStreamV1EventType, type MothershipStreamV1ToolExecutor, @@ -43,19 +43,13 @@ function terminalOperationNeedsApproval(args: Record | undefine return args?.operation === 'run' } -/** - * A human can take as long as they like to answer, so the wait is bounded only - * by the overall orchestration budget rather than a per-tool watchdog. - */ -const PERMISSION_WAIT_TIMEOUT_MS = ORCHESTRATION_TIMEOUT_MS - export const TOOL_AWAITING_APPROVAL_STATUS = MothershipStreamV1ToolStatus.awaiting_approval /** * Whether this call must be held for an explicit user decision. * * Headless one-shot executions are never gated: nobody is there to answer, and - * blocking them would hang the run until the orchestration timeout. + * blocking them would hang the run until the permission wait expires. * * This is the dispatch lane's answer, and it needs a streaming context. A lane * that has none — the in-band route — asks `toolRequiresApprovalLane` instead,