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..77feb41971b 100644 --- a/apps/sim/lib/mothership/request/lifecycle/run.test.ts +++ b/apps/sim/lib/mothership/request/lifecycle/run.test.ts @@ -2037,6 +2037,57 @@ 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('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( @@ -3326,6 +3377,7 @@ describe('runCopilotLifecycle', () => { ) expect(result.success).toBe(false) + expect(result.error).toBeUndefined() expect(result.errors).toEqual(['The provider is overloaded']) }) diff --git a/apps/sim/lib/mothership/request/lifecycle/run.ts b/apps/sim/lib/mothership/request/lifecycle/run.ts index b141f7d7e06..7bd939a2a49 100644 --- a/apps/sim/lib/mothership/request/lifecycle/run.ts +++ b/apps/sim/lib/mothership/request/lifecycle/run.ts @@ -90,6 +90,12 @@ 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. 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() { super(COPILOT_MODEL_CONTENT_PROJECTION_ERROR) @@ -573,6 +579,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 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 = + !turnWasAborted && + context.completionStatus === MothershipStreamV1CompletionStatus.error && + context.errors.length === 0 const result: OrchestratorResult = { success: succeeded, @@ -591,6 +605,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