From e35bb916af3a7f4e682beedc6c60971a3deb573e Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 30 Sep 2026 14:43:32 -0700 Subject: [PATCH 1/3] fix(mothership): explain a Chat turn the worker ends without a reason When the worker rebuilds an ended run from its log (a resume or reattach that reaches a run that already ended, for example at its deadline), it sends an error terminal with no error event. Sim then fell back to the generic "An unexpected error occurred while processing the response." The turn now says the run had already ended and can be continued by sending a message. A reason the worker reports, a replay refusal, and a Stop all still take precedence. Also: - Log the Go stream's error text under errorMessage/detail so it no longer overwrites the log line's own message (stream.ts, buffer.ts). - Rename STREAM_TIMEOUT_MS to CHAT_RUN_DEADLINE_MS and document it as the worker's default run deadline, now only the base of USAGE_SETTLE_MS. - Update the byte-budget doc: a reader behind the ring trim is re-synced from the worker log; replay_gap is only the fallback. --- apps/sim/lib/billing/core/usage-analytics.ts | 10 +-- apps/sim/lib/core/redis/byte-budget.server.ts | 3 +- apps/sim/lib/mothership/constants.ts | 9 ++- apps/sim/lib/mothership/request/go/stream.ts | 4 +- .../mothership/request/lifecycle/run.test.ts | 73 +++++++++++++++++++ .../lib/mothership/request/lifecycle/run.ts | 13 ++++ .../lib/mothership/request/session/buffer.ts | 2 +- 7 files changed, 103 insertions(+), 11 deletions(-) diff --git a/apps/sim/lib/billing/core/usage-analytics.ts b/apps/sim/lib/billing/core/usage-analytics.ts index 60d5ff1a1e9..49c6b4417aa 100644 --- a/apps/sim/lib/billing/core/usage-analytics.ts +++ b/apps/sim/lib/billing/core/usage-analytics.ts @@ -8,7 +8,7 @@ import { } from '@/lib/billing/core/reporting-period' import type { BillingEntity } from '@/lib/billing/core/usage-log' import { zonedWallClockToUtc } from '@/lib/core/utils/timezone' -import { STREAM_TIMEOUT_MS } from '@/lib/mothership/constants' +import { CHAT_RUN_DEADLINE_MS } from '@/lib/mothership/constants' /** * Pure half of organization usage analytics: window resolution, the ledger scope @@ -510,15 +510,15 @@ export function usageBucketTimestamps( * How long after a stretch of time ends before its ledger rows are final. * * Rows are stamped when inserted, but a cumulative model charge tops up its row's - * cost in place for as long as its stream runs — which {@link STREAM_TIMEOUT_MS} - * caps — plus the retry flushes that follow it. Past the cap and this margin a day or - * hour can no longer change and is treated as settled. + * cost in place for as long as its run lasts — which the worker's run deadline + * ({@link CHAT_RUN_DEADLINE_MS}) caps — plus the retry flushes that follow it. Past + * the cap and this margin a day or hour can no longer change and is treated as settled. * * Without a run deadline a Chat turn can top up its row for longer than that, so a * settled hour's cached aggregate can under-report that turn's later spend. This is * display only: invoices, threshold billing, and the usage gate read live ledger sums. */ -export const USAGE_SETTLE_MS = STREAM_TIMEOUT_MS + 2 * 60 * 60 * 1000 +export const USAGE_SETTLE_MS = CHAT_RUN_DEADLINE_MS + 2 * 60 * 60 * 1000 const HOUR_MS = 60 * 60 * 1000 diff --git a/apps/sim/lib/core/redis/byte-budget.server.ts b/apps/sim/lib/core/redis/byte-budget.server.ts index e6124bbd00d..c8554ffb400 100644 --- a/apps/sim/lib/core/redis/byte-budget.server.ts +++ b/apps/sim/lib/core/redis/byte-budget.server.ts @@ -12,7 +12,8 @@ import type { Logger } from '@sim/logger' * execution's event history is read from a cursor, so the write that would breach * the ceiling is refused and the buffer stops growing. The copilot replay ring trims * its oldest events by bytes below its ceiling instead, refunding what it drops, so a - * long run slides rather than refuses; a reader behind the trim gets a replay gap. + * long run slides rather than refuses; a reader behind the trim is re-synced from the + * worker's run log, and ends with a replay gap only when that log cannot serve it. * A live-update feed is bounded differently — see `lib/realtime/event-log.ts`, whose * readers already handle a prune by refetching, so it drops oldest-first instead. * diff --git a/apps/sim/lib/mothership/constants.ts b/apps/sim/lib/mothership/constants.ts index 7f81ab2c245..f5b6a7b95f5 100644 --- a/apps/sim/lib/mothership/constants.ts +++ b/apps/sim/lib/mothership/constants.ts @@ -45,8 +45,13 @@ 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 -/** Timeout for the client-side streaming response handler (60 min). */ -export const STREAM_TIMEOUT_MS = 3_600_000 +/** + * The worker's default deadline for one Chat run (60 min). + * + * Sim does not enforce it: stream legs have no wall clock. It is the base of + * `USAGE_SETTLE_MS`, since it bounds how long a run tops up its model charge. + */ +export const CHAT_RUN_DEADLINE_MS = 3_600_000 /** * How long a workflow tool call waits for a browser to pick it up before the diff --git a/apps/sim/lib/mothership/request/go/stream.ts b/apps/sim/lib/mothership/request/go/stream.ts index cd16ae492ad..a7a5ba38b29 100644 --- a/apps/sim/lib/mothership/request/go/stream.ts +++ b/apps/sim/lib/mothership/request/go/stream.ts @@ -413,7 +413,7 @@ export async function runStreamLoop( context.errors.push(failureMessage) logger.error('Received invalid stream event on shared path', { reason: parsedEvent.reason, - message: parsedEvent.message, + detail: parsedEvent.message, errors: parsedEvent.errors, }) throw new FatalSseEventError(failureMessage) @@ -458,7 +458,7 @@ export async function runStreamLoop( agentId: streamEvent.scope?.agentId, code: errorPayload.code, provider: errorPayload.provider, - message: errorPayload.message, + errorMessage: errorPayload.message, error: errorPayload.error, displayMessage: errorPayload.displayMessage, data: errorPayload.data, diff --git a/apps/sim/lib/mothership/request/lifecycle/run.test.ts b/apps/sim/lib/mothership/request/lifecycle/run.test.ts index d9872ecabf2..ae30c882bfe 100644 --- a/apps/sim/lib/mothership/request/lifecycle/run.test.ts +++ b/apps/sim/lib/mothership/request/lifecycle/run.test.ts @@ -3326,9 +3326,82 @@ describe('runCopilotLifecycle', () => { ) expect(result.success).toBe(false) + expect(result.error).toBeUndefined() expect(result.errors).toEqual(['The provider is overloaded']) }) + it('explains an error terminal that arrives without a reason as an already-ended run', async () => { + const executionContext: ExecutionContext = { + userId: 'user-1', + workflowId: '', + workspaceId: 'ws-1', + chatId: 'chat-1', + } + + mockRunStreamLoop.mockImplementationOnce( + async ( + _fetchUrl: string, + _fetchOptions: RequestInit, + context: StreamingContext + ): Promise => { + context.completionStatus = MothershipStreamV1CompletionStatus.error + } + ) + + const result = await runCopilotLifecycle( + { message: 'hello', messageId: 'stream-1' }, + { + userId: 'user-1', + workspaceId: 'ws-1', + chatId: 'chat-1', + executionId: 'exec-1', + runId: 'run-1', + executionContext, + } + ) + + expect(result.success).toBe(false) + expect(result.cancelled).toBe(false) + expect(result.error).toEqual(expect.stringContaining('already ended')) + }) + + it('keeps a Stop a cancellation when the error terminal carries no reason', async () => { + const executionContext: ExecutionContext = { + userId: 'user-1', + workflowId: '', + workspaceId: 'ws-1', + chatId: 'chat-1', + } + const abortController = new AbortController() + + mockRunStreamLoop.mockImplementationOnce( + async ( + _fetchUrl: string, + _fetchOptions: RequestInit, + context: StreamingContext + ): Promise => { + context.completionStatus = MothershipStreamV1CompletionStatus.error + abortController.abort() + } + ) + + const result = await runCopilotLifecycle( + { message: 'hello', messageId: 'stream-1' }, + { + userId: 'user-1', + workspaceId: 'ws-1', + chatId: 'chat-1', + executionId: 'exec-1', + runId: 'run-1', + executionContext, + abortSignal: abortController.signal, + } + ) + + expect(result.cancelled).toBe(true) + expect(result.error).toBeUndefined() + }) + it('force-fails a hung tool promise and resumes with an error result instead of wedging', async () => { vi.useFakeTimers() try { diff --git a/apps/sim/lib/mothership/request/lifecycle/run.ts b/apps/sim/lib/mothership/request/lifecycle/run.ts index b141f7d7e06..371db81866f 100644 --- a/apps/sim/lib/mothership/request/lifecycle/run.ts +++ b/apps/sim/lib/mothership/request/lifecycle/run.ts @@ -90,6 +90,10 @@ const logger = createLogger('CopilotLifecycle') const COPILOT_MODEL_CONTENT_PROJECTION_ERROR = 'Copilot model input could not be safely projected' +/** Shown when the worker ends a turn with an error terminal but gives no reason. */ +const ENDED_RUN_MESSAGE = + 'This run had already ended before it could continue. Send a message to pick up where it left off.' + class CopilotModelContentProjectionError extends Error { constructor() { super(COPILOT_MODEL_CONTENT_PROJECTION_ERROR) @@ -573,6 +577,14 @@ export async function runCopilotLifecycle( !refusal && !turnWasAborted && (backendFinishedTurn || (!context.completionStatus && context.errors.length === 0)) + // The worker sends an error terminal with no `error` event only when it rebuilds an + // ended run from its log, which drops the stored reason: a resume or reattach that + // reaches a run that already ended, for example at its deadline. Say so rather than + // leave the turn to a generic failure; a reported reason always wins. + const endedWithoutReason = + !turnWasAborted && + context.completionStatus === MothershipStreamV1CompletionStatus.error && + context.errors.length === 0 const result: OrchestratorResult = { success: succeeded, @@ -591,6 +603,7 @@ export async function runCopilotLifecycle( toolCalls: buildToolCallSummaries(context), chatId: context.chatId, requestId: context.requestId, + ...(endedWithoutReason ? { error: ENDED_RUN_MESSAGE } : {}), ...(refusal ? { error: refusal.userMessage, errorCode: refusal.code } : {}), errors: !succeeded && context.errors.length ? context.errors : undefined, usage: context.usage, diff --git a/apps/sim/lib/mothership/request/session/buffer.ts b/apps/sim/lib/mothership/request/session/buffer.ts index 9e0f1559f80..b990f543262 100644 --- a/apps/sim/lib/mothership/request/session/buffer.ts +++ b/apps/sim/lib/mothership/request/session/buffer.ts @@ -487,7 +487,7 @@ export async function readEvents( logger.warn('Skipping corrupt outbox entry', { streamId, reason: parsed.reason, - message: parsed.message, + detail: parsed.message, errors: parsed.errors, }) continue From b499e5fce5519e16c7528842a3593f4ff9caccc3 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 30 Sep 2026 14:55:06 -0700 Subject: [PATCH 2/3] fix(mothership): keep the ended-run message surface-neutral and below a replay refusal The fallback is shared by Chat, workflow execute and inbox, so it no longer tells the reader to send a message. A replay refusal now wins by guard rather than by spread order, with a test that covers the combination. --- .../mothership/request/lifecycle/run.test.ts | 22 +++++++++++++++++++ .../lib/mothership/request/lifecycle/run.ts | 17 ++++++++------ 2 files changed, 32 insertions(+), 7 deletions(-) diff --git a/apps/sim/lib/mothership/request/lifecycle/run.test.ts b/apps/sim/lib/mothership/request/lifecycle/run.test.ts index ae30c882bfe..ce6ef811d70 100644 --- a/apps/sim/lib/mothership/request/lifecycle/run.test.ts +++ b/apps/sim/lib/mothership/request/lifecycle/run.test.ts @@ -2037,6 +2037,28 @@ describe('runCopilotLifecycle', () => { } ) + it('reports a replay refusal over a reasonless error terminal', async () => { + const abortController = new AbortController() + const refusal = ownerRefusal() + mockRunStreamLoop.mockImplementationOnce( + async (_url: string, _init: RequestInit, context: StreamingContext): Promise => { + context.completionStatus = MothershipStreamV1CompletionStatus.error + abortController.abort(refusal) + } + ) + + const result = await runWithStreamAbort(abortController) + + expect(result).toEqual( + expect.objectContaining({ + success: false, + cancelled: false, + error: refusal.userMessage, + errorCode: REPLAY_BUDGET_EXHAUSTED_CODE, + }) + ) + }) + it('keeps a Stop a cancellation when a replay refusal follows it', async () => { const abortController = new AbortController() mockRunStreamLoop.mockImplementationOnce( diff --git a/apps/sim/lib/mothership/request/lifecycle/run.ts b/apps/sim/lib/mothership/request/lifecycle/run.ts index 371db81866f..28357ae10d0 100644 --- a/apps/sim/lib/mothership/request/lifecycle/run.ts +++ b/apps/sim/lib/mothership/request/lifecycle/run.ts @@ -90,9 +90,11 @@ const logger = createLogger('CopilotLifecycle') const COPILOT_MODEL_CONTENT_PROJECTION_ERROR = 'Copilot model input could not be safely projected' -/** Shown when the worker ends a turn with an error terminal but gives no reason. */ -const ENDED_RUN_MESSAGE = - 'This run had already ended before it could continue. Send a message to pick up where it left off.' +/** + * Shown when the worker ends a turn with an error terminal but gives no reason. Every surface + * (Chat, workflow execute, inbox) reports it, so it carries no surface-specific next step. + */ +const ENDED_RUN_MESSAGE = 'This run had already ended before it could continue.' class CopilotModelContentProjectionError extends Error { constructor() { @@ -577,11 +579,12 @@ export async function runCopilotLifecycle( !refusal && !turnWasAborted && (backendFinishedTurn || (!context.completionStatus && context.errors.length === 0)) - // The worker sends an error terminal with no `error` event only when it rebuilds an - // ended run from its log, which drops the stored reason: a resume or reattach that - // reaches a run that already ended, for example at its deadline. Say so rather than - // leave the turn to a generic failure; a reported reason always wins. + // The worker sends an error terminal with no `error` event only when it replays a run + // that already ended (for example at its deadline) to a resume or reattach, because + // that replay does not carry the run's stored reason. Say so rather than leave the turn + // to a generic failure; a reported reason or a replay refusal always wins. const endedWithoutReason = + !refusal && !turnWasAborted && context.completionStatus === MothershipStreamV1CompletionStatus.error && context.errors.length === 0 From 5a34b3e64ca057e6c396dba74b9cd19967853031 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 30 Sep 2026 16:09:49 -0700 Subject: [PATCH 3/3] refactor(mothership): drop the unreachable refusal guard and reuse the stream-abort test helper --- .../mothership/request/lifecycle/run.test.ts | 101 +++++------------- .../lib/mothership/request/lifecycle/run.ts | 1 - 2 files changed, 29 insertions(+), 73 deletions(-) diff --git a/apps/sim/lib/mothership/request/lifecycle/run.test.ts b/apps/sim/lib/mothership/request/lifecycle/run.test.ts index ce6ef811d70..77feb41971b 100644 --- a/apps/sim/lib/mothership/request/lifecycle/run.test.ts +++ b/apps/sim/lib/mothership/request/lifecycle/run.test.ts @@ -2059,6 +2059,35 @@ describe('runCopilotLifecycle', () => { ) }) + it('explains an error terminal that arrives without a reason as an already-ended run', async () => { + mockRunStreamLoop.mockImplementationOnce( + async (_url: string, _init: RequestInit, context: StreamingContext): Promise => { + context.completionStatus = MothershipStreamV1CompletionStatus.error + } + ) + + const result = await runWithStreamAbort(new AbortController()) + + expect(result.success).toBe(false) + expect(result.cancelled).toBe(false) + expect(result.error).toEqual(expect.stringContaining('already ended')) + }) + + it('keeps a Stop a cancellation when the error terminal carries no reason', async () => { + const abortController = new AbortController() + mockRunStreamLoop.mockImplementationOnce( + async (_url: string, _init: RequestInit, context: StreamingContext): Promise => { + context.completionStatus = MothershipStreamV1CompletionStatus.error + abortController.abort() + } + ) + + const result = await runWithStreamAbort(abortController) + + expect(result.cancelled).toBe(true) + expect(result.error).toBeUndefined() + }) + it('keeps a Stop a cancellation when a replay refusal follows it', async () => { const abortController = new AbortController() mockRunStreamLoop.mockImplementationOnce( @@ -3352,78 +3381,6 @@ describe('runCopilotLifecycle', () => { expect(result.errors).toEqual(['The provider is overloaded']) }) - it('explains an error terminal that arrives without a reason as an already-ended run', async () => { - const executionContext: ExecutionContext = { - userId: 'user-1', - workflowId: '', - workspaceId: 'ws-1', - chatId: 'chat-1', - } - - mockRunStreamLoop.mockImplementationOnce( - async ( - _fetchUrl: string, - _fetchOptions: RequestInit, - context: StreamingContext - ): Promise => { - context.completionStatus = MothershipStreamV1CompletionStatus.error - } - ) - - const result = await runCopilotLifecycle( - { message: 'hello', messageId: 'stream-1' }, - { - userId: 'user-1', - workspaceId: 'ws-1', - chatId: 'chat-1', - executionId: 'exec-1', - runId: 'run-1', - executionContext, - } - ) - - expect(result.success).toBe(false) - expect(result.cancelled).toBe(false) - expect(result.error).toEqual(expect.stringContaining('already ended')) - }) - - it('keeps a Stop a cancellation when the error terminal carries no reason', async () => { - const executionContext: ExecutionContext = { - userId: 'user-1', - workflowId: '', - workspaceId: 'ws-1', - chatId: 'chat-1', - } - const abortController = new AbortController() - - mockRunStreamLoop.mockImplementationOnce( - async ( - _fetchUrl: string, - _fetchOptions: RequestInit, - context: StreamingContext - ): Promise => { - context.completionStatus = MothershipStreamV1CompletionStatus.error - abortController.abort() - } - ) - - const result = await runCopilotLifecycle( - { message: 'hello', messageId: 'stream-1' }, - { - userId: 'user-1', - workspaceId: 'ws-1', - chatId: 'chat-1', - executionId: 'exec-1', - runId: 'run-1', - executionContext, - abortSignal: abortController.signal, - } - ) - - expect(result.cancelled).toBe(true) - expect(result.error).toBeUndefined() - }) - it('force-fails a hung tool promise and resumes with an error result instead of wedging', async () => { vi.useFakeTimers() try { diff --git a/apps/sim/lib/mothership/request/lifecycle/run.ts b/apps/sim/lib/mothership/request/lifecycle/run.ts index 28357ae10d0..7bd939a2a49 100644 --- a/apps/sim/lib/mothership/request/lifecycle/run.ts +++ b/apps/sim/lib/mothership/request/lifecycle/run.ts @@ -584,7 +584,6 @@ export async function runCopilotLifecycle( // that replay does not carry the run's stored reason. Say so rather than leave the turn // to a generic failure; a reported reason or a replay refusal always wins. const endedWithoutReason = - !refusal && !turnWasAborted && context.completionStatus === MothershipStreamV1CompletionStatus.error && context.errors.length === 0