From f376c316401b2138e7a67fcf8ce050d7a329cf10 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 30 Sep 2026 01:47:59 -0700 Subject: [PATCH 1/5] fix(mothership): keep Chat stream legs alive past the hour without a wall clock - End the replay GET at its cap without a terminal event, so the client re-attaches from its cursor instead of ending a live turn with resume_timeout. - Replace the per-leg worker SSE wall clock with an idle timeout (WORKER_STREAM_IDLE_TIMEOUT_MS, 120 s) that fails a silent leg as a retryable interruption; a caller-set timeout still applies. - Let StreamRetryWindow run without a deadline by default and replenish its reachable budget only after five minutes of healthy streaming, never per event. - Split the tool watchdog, permission wait, client tool wait, and delegation TTL onto their own constants and remove ORCHESTRATION_TIMEOUT_MS. - Refresh the replay buffer TTLs from the chat-lock heartbeat, and report a replay gap for a cursor ahead of a buffer whose numbering restarted. - Document why maxDuration stays on the chat POST and execute routes. --- .../app/api/copilot/chat/stream/route.test.ts | 26 ++- apps/sim/app/api/copilot/chat/stream/route.ts | 14 +- apps/sim/app/api/mothership/chat/route.ts | 5 + apps/sim/app/api/mothership/execute/route.ts | 1 + .../execute-workflow-use-case.test.ts | 4 +- .../mothership/auth/application-delegation.ts | 9 +- .../mothership/auth/file-delegation.test.ts | 4 +- apps/sim/lib/mothership/constants.ts | 24 ++- .../lib/mothership/request/go/stream.test.ts | 149 +++++++++++++++++- apps/sim/lib/mothership/request/go/stream.ts | 61 +++++-- .../lib/mothership/request/handlers/tool.ts | 4 +- .../lib/mothership/request/lifecycle/run.ts | 2 +- .../request/lifecycle/stream-retry.test.ts | 50 +++++- .../request/lifecycle/stream-retry.ts | 56 +++++-- .../lib/mothership/request/session/abort.ts | 4 +- .../request/session/buffer-ttl.integration.ts | 88 +++++++++++ .../lib/mothership/request/session/buffer.ts | 15 ++ .../mothership/request/session/recovery.ts | 6 +- .../lib/mothership/request/tools/executor.ts | 10 +- .../mothership/request/tools/permission.ts | 10 +- 20 files changed, 476 insertions(+), 66 deletions(-) create mode 100644 apps/sim/lib/mothership/request/session/buffer-ttl.integration.ts 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..93190e4faaa 100644 --- a/apps/sim/app/api/copilot/chat/stream/route.test.ts +++ b/apps/sim/app/api/copilot/chat/stream/route.test.ts @@ -1,6 +1,6 @@ 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, @@ -62,6 +62,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 +175,24 @@ 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 () => { + 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}"`) + }) }) diff --git a/apps/sim/app/api/copilot/chat/stream/route.ts b/apps/sim/app/api/copilot/chat/stream/route.ts index c84cad4d102..26c33b17d31 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', { 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/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..9b0b1b0e7fa 100644 --- a/apps/sim/lib/mothership/auth/application-delegation.ts +++ b/apps/sim/lib/mothership/auth/application-delegation.ts @@ -1,8 +1,11 @@ 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, so it must outlive the longest + * single tool call it authorizes, never the whole run. + */ +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..40d7dc54595 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. This must stay under + * the NAT's fixed 350 s idle timeout, which drops 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..bf246a35a8c 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,149 @@ describe('copilot go stream helpers', () => { ) ).toBe(true) }) + + describe('worker stream liveness without a caller deadline', () => { + const NAT_IDLE_MS = 350_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 the NAT 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(NAT_IDLE_MS - 60_000) + + expect(state.done).toBe(true) + expect(state.error).toBeInstanceOf(WorkerStreamInterruptedError) + }) + + it('fails a worker that never answers as unreachable before the NAT 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(NAT_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..8a88e03b1aa 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 = { @@ -306,11 +332,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 +366,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 +593,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..3187d4eb14f 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,50 @@ 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() + }) }) diff --git a/apps/sim/lib/mothership/request/lifecycle/stream-retry.ts b/apps/sim/lib/mothership/request/lifecycle/stream-retry.ts index 08da02415cc..109c498fcdf 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,32 @@ const GATEWAY_STATUSES: ReadonlySet = new Set([502, 503, 504]) * reattach, never a second run. */ const WORKER_REPLACEMENT_WINDOW_MS = 120_000 +/** + * A leg that has delivered events for this long since it last re-attached has + * proven healthy, so a later interruption, possibly hours on, gets the reachable + * budget afresh. Events alone never replenish it: a leg that fails again sooner + * keeps spending the same three retries, so 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 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 +52,25 @@ 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.streamingSince ??= Date.now() } 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 +79,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 +88,25 @@ export class StreamRetryWindow { return delay } + private resetUnreachable(): void { + this.firstUnreachableAt = undefined + this.unreachableAttempt = 0 + } + + private replenishAfterHealthyStreaming(): void { + if ( + this.streamingSince !== undefined && + Date.now() - 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..239b71cb8f4 --- /dev/null +++ b/apps/sim/lib/mothership/request/session/buffer-ttl.integration.ts @@ -0,0 +1,88 @@ +/** + * 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 + /** Short enough that the park below outlasts it several times over. */ + process.env.COPILOT_STREAM_TTL_SECONDS = '2' + return { redisUrl: url } +}) + +import { sleep } from '@sim/utils/helpers' +import { generateId } from '@sim/utils/id' +import { closeRedisConnection, getRedisClient } from '@/lib/core/config/redis' +import { MothershipStreamV1EventType } from '@/lib/mothership/generated/mothership-stream-v1' +import { + acquirePendingChatStream, + releasePendingChatStream, + startAbortPoller, +} from '@/lib/mothership/request/session/abort' +import { + allocateCursor, + appendEvents, + getLatestSeq, + readEvents, +} 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 through a park longer than its TTL', async () => { + const chatId = generateId() + const streamId = generateId() + expect(await acquirePendingChatStream(chatId, streamId, 0)).toBe(true) + await appendText(streamId, 'before the park') + + vi.useFakeTimers({ toFake: ['Date'] }) + const poller = startAbortPoller(streamId, new AbortController(), { chatId, pollMs: 50 }) + try { + for (let tick = 0; tick < 10; 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(await appendText(streamId, 'after the park')).toBe(2) + }) + + 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..8f99a36d040 100644 --- a/apps/sim/lib/mothership/request/session/buffer.ts +++ b/apps/sim/lib/mothership/request/session/buffer.ts @@ -132,6 +132,21 @@ export async function clearBuffer(streamId: string, operation = 'clear_outbox'): }) } +/** + * Slides a live stream's replay TTLs without an append. The TTLs 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 is still running. + */ +export async function refreshBufferTtl(streamId: string): Promise { + const { ttlSeconds } = getStreamConfig() + await withRedisRetry({ operation: 'refresh_outbox_ttl', streamId }, async (redis) => { + const pipeline = redis.pipeline() + pipeline.expire(getEventsKey(streamId), ttlSeconds) + pipeline.expire(getSeqKey(streamId), ttlSeconds) + await pipeline.exec() + }) +} + export async function scheduleBufferCleanup( streamId: string, ttlSeconds = DEFAULT_COMPLETED_TTL_SECONDS 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, From be26984cc05f6f21b1c7abefeb4b3aab71b673eb Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 30 Sep 2026 01:52:59 -0700 Subject: [PATCH 2/5] fix(mothership): renew the Chat reconnect budget after a tail delivers events A reconnect attempt whose re-attached tail delivered new events now restarts the retry budget at the base delay, so separate network drops hours apart in a long turn no longer add up to the ten-attempt exhaustion. --- .../home/hooks/use-chat.dom.test.tsx | 68 +++++++++++++++++++ .../[workspaceId]/home/hooks/use-chat.ts | 11 ++- 2 files changed, 78 insertions(+), 1 deletion(-) 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', { From 25a7d2d568a0423260ea0c000895fb0c87f048b6 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 30 Sep 2026 03:31:52 -0700 Subject: [PATCH 3/5] fix(mothership): refresh the stream byte counter with its buffer, and never revive a closed buffer - The chat-lock heartbeat now slides the owner byte counter's TTL along with the replay buffer's. Otherwise, after a park longer than an hour, the counter expired while the ring survived: the ring could grow to about twice its target, and its refunds drained the user counter. - Scheduling a finished stream's cleanup now marks it closed, and the refresh leaves a closed stream alone. A heartbeat still in flight at teardown can no longer re-extend a buffer whose cleanup was already scheduled. - Describe the worker idle timeout in terms of network intermediaries, and state that the delegation TTL reuses the long-running tool watchdog's cap. --- .../mothership/auth/application-delegation.ts | 5 ++- apps/sim/lib/mothership/constants.ts | 6 +-- .../lib/mothership/request/go/stream.test.ts | 11 ++--- .../request/session/buffer-ttl.integration.ts | 23 +++++++++- .../lib/mothership/request/session/buffer.ts | 43 ++++++++++++++++--- 5 files changed, 70 insertions(+), 18 deletions(-) diff --git a/apps/sim/lib/mothership/auth/application-delegation.ts b/apps/sim/lib/mothership/auth/application-delegation.ts index 9b0b1b0e7fa..62993efc220 100644 --- a/apps/sim/lib/mothership/auth/application-delegation.ts +++ b/apps/sim/lib/mothership/auth/application-delegation.ts @@ -2,8 +2,9 @@ import type { DelegatedPrincipal, OrganizationDelegatedPrincipal } from '@sim/au import { TOOL_WATCHDOG_LONG_RUNNING_MS } from '@/lib/mothership/constants' /** - * Delegated authority is minted per operation, so it must outlive the longest - * single tool call it authorizes, never the whole run. + * 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 diff --git a/apps/sim/lib/mothership/constants.ts b/apps/sim/lib/mothership/constants.ts index 40d7dc54595..2fda1e8ef9a 100644 --- a/apps/sim/lib/mothership/constants.ts +++ b/apps/sim/lib/mothership/constants.ts @@ -14,9 +14,9 @@ export const SIM_AGENT_API_URL = * 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. This must stay under - * the NAT's fixed 350 s idle timeout, which drops a silent connection without - * closing it. + * 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 diff --git a/apps/sim/lib/mothership/request/go/stream.test.ts b/apps/sim/lib/mothership/request/go/stream.test.ts index bf246a35a8c..3fc11c17721 100644 --- a/apps/sim/lib/mothership/request/go/stream.test.ts +++ b/apps/sim/lib/mothership/request/go/stream.test.ts @@ -1027,7 +1027,8 @@ describe('copilot go stream helpers', () => { }) describe('worker stream liveness without a caller deadline', () => { - const NAT_IDLE_MS = 350_000 + /** 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({ @@ -1112,7 +1113,7 @@ describe('copilot go stream helpers', () => { expect(context.streamComplete).toBe(true) }) - it('fails a silent leg as a retryable interruption before the NAT drops it', async () => { + it('fails a silent leg as a retryable interruption before an intermediary drops it', async () => { vi.useFakeTimers() vi.mocked(fetch).mockResolvedValueOnce( new Response( @@ -1136,13 +1137,13 @@ describe('copilot go stream helpers', () => { await vi.advanceTimersByTimeAsync(60_000) expect(state.done).toBe(false) - await vi.advanceTimersByTimeAsync(NAT_IDLE_MS - 60_000) + await vi.advanceTimersByTimeAsync(INTERMEDIARY_IDLE_MS - 60_000) expect(state.done).toBe(true) expect(state.error).toBeInstanceOf(WorkerStreamInterruptedError) }) - it('fails a worker that never answers as unreachable before the NAT drops it', async () => { + it('fails a worker that never answers as unreachable before an intermediary drops it', async () => { vi.useFakeTimers() vi.mocked(fetch).mockImplementationOnce( (_url, options) => @@ -1163,7 +1164,7 @@ describe('copilot go stream helpers', () => { ) ) - await vi.advanceTimersByTimeAsync(NAT_IDLE_MS) + await vi.advanceTimersByTimeAsync(INTERMEDIARY_IDLE_MS) expect(state.done).toBe(true) expect(state.error).toBeInstanceOf(WorkerUnreachableError) diff --git a/apps/sim/lib/mothership/request/session/buffer-ttl.integration.ts b/apps/sim/lib/mothership/request/session/buffer-ttl.integration.ts index 239b71cb8f4..e05d7ad9995 100644 --- a/apps/sim/lib/mothership/request/session/buffer-ttl.integration.ts +++ b/apps/sim/lib/mothership/request/session/buffer-ttl.integration.ts @@ -17,6 +17,7 @@ const { redisUrl } = await vi.hoisted(async () => { 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, @@ -28,6 +29,8 @@ import { 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' @@ -52,11 +55,15 @@ describe.runIf(Boolean(redisUrl))('replay buffer lifetime', () => { await closeRedisConnection() }) - it('keeps a live run’s buffer through a park longer than its TTL', async () => { + 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, 2) vi.useFakeTimers({ toFake: ['Date'] }) const poller = startAbortPoller(streamId, new AbortController(), { chatId, pollMs: 50 }) @@ -73,9 +80,23 @@ describe.runIf(Boolean(redisUrl))('replay buffer lifetime', () => { 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}`) diff --git a/apps/sim/lib/mothership/request/session/buffer.ts b/apps/sim/lib/mothership/request/session/buffer.ts index 8f99a36d040..0fe1f14ebc3 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,23 +132,46 @@ export async function clearBuffer(streamId: string, operation = 'clear_outbox'): getEventsKey(streamId), getSeqKey(streamId), getAbortKey(streamId), + getClosedKey(streamId), ownerBudgetKey ) }) } /** - * Slides a live stream's replay TTLs without an append. The TTLs 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 is still running. + * 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) => { - const pipeline = redis.pipeline() - pipeline.expire(getEventsKey(streamId), ttlSeconds) - pipeline.expire(getSeqKey(streamId), ttlSeconds) - await pipeline.exec() + await redis.eval( + REFRESH_BUFFER_TTL_SCRIPT, + 4, + getEventsKey(streamId), + getSeqKey(streamId), + ownerBudgetKey, + getClosedKey(streamId), + ttlSeconds, + budgetTtlSeconds + ) }) } @@ -157,6 +185,7 @@ export async function scheduleBufferCleanup( pipeline.expire(getEventsKey(streamId), ttlSeconds) pipeline.expire(getSeqKey(streamId), ttlSeconds) pipeline.expire(getAbortKey(streamId), ttlSeconds) + pipeline.set(getClosedKey(streamId), '1', 'EX', ttlSeconds) await pipeline.exec() }) } catch (error) { From b020bbb02fcbaafed0bf2979707240d3afdb8121 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 30 Sep 2026 07:41:58 -0700 Subject: [PATCH 4/5] fix(mothership): bound a stalled worker error body and refill retries only on real progress - A non-OK worker response's body is read under the same idle bound as the leg, so a stalled error body can no longer block the turn forever. - The reachable retry budget refills only when a leg's delivered events span five minutes, from its first event to its latest. A leg that delivered one event and then only kept alive has made no progress and no longer refills it. --- .../lib/mothership/request/go/stream.test.ts | 24 +++++++++++++++++++ apps/sim/lib/mothership/request/go/stream.ts | 7 +++++- .../request/lifecycle/stream-retry.test.ts | 14 +++++++++++ .../request/lifecycle/stream-retry.ts | 16 ++++++++----- 4 files changed, 54 insertions(+), 7 deletions(-) diff --git a/apps/sim/lib/mothership/request/go/stream.test.ts b/apps/sim/lib/mothership/request/go/stream.test.ts index 3fc11c17721..7568ccd7626 100644 --- a/apps/sim/lib/mothership/request/go/stream.test.ts +++ b/apps/sim/lib/mothership/request/go/stream.test.ts @@ -1143,6 +1143,30 @@ describe('copilot go stream helpers', () => { 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( diff --git a/apps/sim/lib/mothership/request/go/stream.ts b/apps/sim/lib/mothership/request/go/stream.ts index 8a88e03b1aa..336423cbdd7 100644 --- a/apps/sim/lib/mothership/request/go/stream.ts +++ b/apps/sim/lib/mothership/request/go/stream.ts @@ -256,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) 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 3187d4eb14f..558566dbe31 100644 --- a/apps/sim/lib/mothership/request/lifecycle/stream-retry.test.ts +++ b/apps/sim/lib/mothership/request/lifecycle/stream-retry.test.ts @@ -202,4 +202,18 @@ describe('stream recovery budget', () => { 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 109c498fcdf..2650ac6f195 100644 --- a/apps/sim/lib/mothership/request/lifecycle/stream-retry.ts +++ b/apps/sim/lib/mothership/request/lifecycle/stream-retry.ts @@ -20,10 +20,11 @@ const GATEWAY_STATUSES: ReadonlySet = new Set([502, 503, 504]) */ const WORKER_REPLACEMENT_WINDOW_MS = 120_000 /** - * A leg that has delivered events for this long since it last re-attached has - * proven healthy, so a later interruption, possibly hours on, gets the reachable - * budget afresh. Events alone never replenish it: a leg that fails again sooner - * keeps spending the same three retries, so a deterministic failure stays bounded. + * 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 @@ -39,6 +40,7 @@ export class StreamRetryWindow { private readonly deadline?: number private firstFailureAt?: number private streamingSince?: number + private lastEventAt?: number private firstUnreachableAt?: number private unreachableAttempt = 0 private attempt = 0 @@ -64,7 +66,8 @@ export class StreamRetryWindow { /** The worker delivered an event, so a later loss of it starts a fresh unreachable window. */ recovered(): void { this.resetUnreachable() - this.streamingSince ??= Date.now() + this.lastEventAt = Date.now() + this.streamingSince ??= this.lastEventAt } nextDelay(error: unknown, signal?: AbortSignal): number | null { @@ -96,7 +99,8 @@ export class StreamRetryWindow { private replenishAfterHealthyStreaming(): void { if ( this.streamingSince !== undefined && - Date.now() - this.streamingSince >= HEALTHY_STREAM_REPLENISH_MS + this.lastEventAt !== undefined && + this.lastEventAt - this.streamingSince >= HEALTHY_STREAM_REPLENISH_MS ) { this.attempt = 0 this.firstFailureAt = undefined From 5beb5ddcb687d7ef9d1044f0f47c266c29989230 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 30 Sep 2026 07:44:33 -0700 Subject: [PATCH 5/5] fix(mothership): mark a closing buffer before expiring it, and report a capped replay as ended - scheduleBufferCleanup sets the closed marker before shortening the TTLs, so a heartbeat refresh that lands mid-pipeline is either overridden or sees the marker. - A reconnect that ends at its cap is reported as ended without a terminal, not as a client disconnect: the outcome is read before the route closes its own stream. - The buffer TTL integration test keeps a 5 s TTL against a 12 s park, so a slow runner cannot expire the buffer between heartbeats. --- .../app/api/copilot/chat/stream/route.test.ts | 20 +++++++++++++++++++ apps/sim/app/api/copilot/chat/stream/route.ts | 4 +++- .../request/session/buffer-ttl.integration.ts | 8 ++++---- .../lib/mothership/request/session/buffer.ts | 4 +++- 4 files changed, 30 insertions(+), 6 deletions(-) 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 93190e4faaa..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,3 +1,9 @@ +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 { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' @@ -6,6 +12,9 @@ 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(() => ({ @@ -177,6 +186,10 @@ describe('copilot chat stream replay route', () => { }) 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', @@ -194,5 +207,12 @@ describe('copilot chat stream replay route', () => { 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 26c33b17d31..b1023a791b9 100644 --- a/apps/sim/app/api/copilot/chat/stream/route.ts +++ b/apps/sim/app/api/copilot/chat/stream/route.ts @@ -471,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/lib/mothership/request/session/buffer-ttl.integration.ts b/apps/sim/lib/mothership/request/session/buffer-ttl.integration.ts index e05d7ad9995..9fc83b57c48 100644 --- a/apps/sim/lib/mothership/request/session/buffer-ttl.integration.ts +++ b/apps/sim/lib/mothership/request/session/buffer-ttl.integration.ts @@ -9,8 +9,8 @@ const { redisUrl } = await vi.hoisted(async () => { const { readTestRedisUrl } = await import('@sim/db/testing/test-infrastructure') const url = readTestRedisUrl() process.env.REDIS_URL = url - /** Short enough that the park below outlasts it several times over. */ - process.env.COPILOT_STREAM_TTL_SECONDS = '2' + /** The park below outlasts it more than twice over in real time. */ + process.env.COPILOT_STREAM_TTL_SECONDS = '5' return { redisUrl: url } }) @@ -63,12 +63,12 @@ describe.runIf(Boolean(redisUrl))('replay buffer lifetime', () => { 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, 2) + await getRedisClient()!.expire(ownerBudgetKey, 5) vi.useFakeTimers({ toFake: ['Date'] }) const poller = startAbortPoller(streamId, new AbortController(), { chatId, pollMs: 50 }) try { - for (let tick = 0; tick < 10; tick++) { + for (let tick = 0; tick < 24; tick++) { vi.setSystemTime(Date.now() + 21_000) await sleep(500) } diff --git a/apps/sim/lib/mothership/request/session/buffer.ts b/apps/sim/lib/mothership/request/session/buffer.ts index 0fe1f14ebc3..68fd24cd447 100644 --- a/apps/sim/lib/mothership/request/session/buffer.ts +++ b/apps/sim/lib/mothership/request/session/buffer.ts @@ -181,11 +181,13 @@ export async function scheduleBufferCleanup( ): 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) - pipeline.set(getClosedKey(streamId), '1', 'EX', ttlSeconds) await pipeline.exec() }) } catch (error) {