Skip to content

Commit e35bb91

Browse files
committed
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.
1 parent be50bc6 commit e35bb91

7 files changed

Lines changed: 103 additions & 11 deletions

File tree

‎apps/sim/lib/billing/core/usage-analytics.ts‎

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -8,7 +8,7 @@ import {
88
} from '@/lib/billing/core/reporting-period'
99
import type { BillingEntity } from '@/lib/billing/core/usage-log'
1010
import { zonedWallClockToUtc } from '@/lib/core/utils/timezone'
11-
import { STREAM_TIMEOUT_MS } from '@/lib/mothership/constants'
11+
import { CHAT_RUN_DEADLINE_MS } from '@/lib/mothership/constants'
1212

1313
/**
1414
* Pure half of organization usage analytics: window resolution, the ledger scope
@@ -510,15 +510,15 @@ export function usageBucketTimestamps(
510510
* How long after a stretch of time ends before its ledger rows are final.
511511
*
512512
* Rows are stamped when inserted, but a cumulative model charge tops up its row's
513-
* cost in place for as long as its stream runs — which {@link STREAM_TIMEOUT_MS}
514-
* caps — plus the retry flushes that follow it. Past the cap and this margin a day or
515-
* hour can no longer change and is treated as settled.
513+
* cost in place for as long as its run lasts — which the worker's run deadline
514+
* ({@link CHAT_RUN_DEADLINE_MS}) caps — plus the retry flushes that follow it. Past
515+
* the cap and this margin a day or hour can no longer change and is treated as settled.
516516
*
517517
* Without a run deadline a Chat turn can top up its row for longer than that, so a
518518
* settled hour's cached aggregate can under-report that turn's later spend. This is
519519
* display only: invoices, threshold billing, and the usage gate read live ledger sums.
520520
*/
521-
export const USAGE_SETTLE_MS = STREAM_TIMEOUT_MS + 2 * 60 * 60 * 1000
521+
export const USAGE_SETTLE_MS = CHAT_RUN_DEADLINE_MS + 2 * 60 * 60 * 1000
522522

523523
const HOUR_MS = 60 * 60 * 1000
524524

‎apps/sim/lib/core/redis/byte-budget.server.ts‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -12,7 +12,8 @@ import type { Logger } from '@sim/logger'
1212
* execution's event history is read from a cursor, so the write that would breach
1313
* the ceiling is refused and the buffer stops growing. The copilot replay ring trims
1414
* its oldest events by bytes below its ceiling instead, refunding what it drops, so a
15-
* long run slides rather than refuses; a reader behind the trim gets a replay gap.
15+
* long run slides rather than refuses; a reader behind the trim is re-synced from the
16+
* worker's run log, and ends with a replay gap only when that log cannot serve it.
1617
* A live-update feed is bounded differently — see `lib/realtime/event-log.ts`, whose
1718
* readers already handle a prune by refetching, so it drops oldest-first instead.
1819
*

‎apps/sim/lib/mothership/constants.ts‎

Lines changed: 7 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -45,8 +45,13 @@ export const CLIENT_TOOL_RESULT_TIMEOUT_MS = 60 * 60 * 1000
4545
/** Extra slack the resume gate allows past the slowest pending tool's watchdog. */
4646
export const TOOL_WATCHDOG_RESUME_GRACE_MS = 30_000
4747

48-
/** Timeout for the client-side streaming response handler (60 min). */
49-
export const STREAM_TIMEOUT_MS = 3_600_000
48+
/**
49+
* The worker's default deadline for one Chat run (60 min).
50+
*
51+
* Sim does not enforce it: stream legs have no wall clock. It is the base of
52+
* `USAGE_SETTLE_MS`, since it bounds how long a run tops up its model charge.
53+
*/
54+
export const CHAT_RUN_DEADLINE_MS = 3_600_000
5055

5156
/**
5257
* How long a workflow tool call waits for a browser to pick it up before the

‎apps/sim/lib/mothership/request/go/stream.ts‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -413,7 +413,7 @@ export async function runStreamLoop(
413413
context.errors.push(failureMessage)
414414
logger.error('Received invalid stream event on shared path', {
415415
reason: parsedEvent.reason,
416-
message: parsedEvent.message,
416+
detail: parsedEvent.message,
417417
errors: parsedEvent.errors,
418418
})
419419
throw new FatalSseEventError(failureMessage)
@@ -458,7 +458,7 @@ export async function runStreamLoop(
458458
agentId: streamEvent.scope?.agentId,
459459
code: errorPayload.code,
460460
provider: errorPayload.provider,
461-
message: errorPayload.message,
461+
errorMessage: errorPayload.message,
462462
error: errorPayload.error,
463463
displayMessage: errorPayload.displayMessage,
464464
data: errorPayload.data,

‎apps/sim/lib/mothership/request/lifecycle/run.test.ts‎

Lines changed: 73 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3326,9 +3326,82 @@ describe('runCopilotLifecycle', () => {
33263326
)
33273327

33283328
expect(result.success).toBe(false)
3329+
expect(result.error).toBeUndefined()
33293330
expect(result.errors).toEqual(['The provider is overloaded'])
33303331
})
33313332

3333+
it('explains an error terminal that arrives without a reason as an already-ended run', async () => {
3334+
const executionContext: ExecutionContext = {
3335+
userId: 'user-1',
3336+
workflowId: '',
3337+
workspaceId: 'ws-1',
3338+
chatId: 'chat-1',
3339+
}
3340+
3341+
mockRunStreamLoop.mockImplementationOnce(
3342+
async (
3343+
_fetchUrl: string,
3344+
_fetchOptions: RequestInit,
3345+
context: StreamingContext
3346+
): Promise<void> => {
3347+
context.completionStatus = MothershipStreamV1CompletionStatus.error
3348+
}
3349+
)
3350+
3351+
const result = await runCopilotLifecycle(
3352+
{ message: 'hello', messageId: 'stream-1' },
3353+
{
3354+
userId: 'user-1',
3355+
workspaceId: 'ws-1',
3356+
chatId: 'chat-1',
3357+
executionId: 'exec-1',
3358+
runId: 'run-1',
3359+
executionContext,
3360+
}
3361+
)
3362+
3363+
expect(result.success).toBe(false)
3364+
expect(result.cancelled).toBe(false)
3365+
expect(result.error).toEqual(expect.stringContaining('already ended'))
3366+
})
3367+
3368+
it('keeps a Stop a cancellation when the error terminal carries no reason', async () => {
3369+
const executionContext: ExecutionContext = {
3370+
userId: 'user-1',
3371+
workflowId: '',
3372+
workspaceId: 'ws-1',
3373+
chatId: 'chat-1',
3374+
}
3375+
const abortController = new AbortController()
3376+
3377+
mockRunStreamLoop.mockImplementationOnce(
3378+
async (
3379+
_fetchUrl: string,
3380+
_fetchOptions: RequestInit,
3381+
context: StreamingContext
3382+
): Promise<void> => {
3383+
context.completionStatus = MothershipStreamV1CompletionStatus.error
3384+
abortController.abort()
3385+
}
3386+
)
3387+
3388+
const result = await runCopilotLifecycle(
3389+
{ message: 'hello', messageId: 'stream-1' },
3390+
{
3391+
userId: 'user-1',
3392+
workspaceId: 'ws-1',
3393+
chatId: 'chat-1',
3394+
executionId: 'exec-1',
3395+
runId: 'run-1',
3396+
executionContext,
3397+
abortSignal: abortController.signal,
3398+
}
3399+
)
3400+
3401+
expect(result.cancelled).toBe(true)
3402+
expect(result.error).toBeUndefined()
3403+
})
3404+
33323405
it('force-fails a hung tool promise and resumes with an error result instead of wedging', async () => {
33333406
vi.useFakeTimers()
33343407
try {

‎apps/sim/lib/mothership/request/lifecycle/run.ts‎

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -90,6 +90,10 @@ const logger = createLogger('CopilotLifecycle')
9090

9191
const COPILOT_MODEL_CONTENT_PROJECTION_ERROR = 'Copilot model input could not be safely projected'
9292

93+
/** Shown when the worker ends a turn with an error terminal but gives no reason. */
94+
const ENDED_RUN_MESSAGE =
95+
'This run had already ended before it could continue. Send a message to pick up where it left off.'
96+
9397
class CopilotModelContentProjectionError extends Error {
9498
constructor() {
9599
super(COPILOT_MODEL_CONTENT_PROJECTION_ERROR)
@@ -573,6 +577,14 @@ export async function runCopilotLifecycle(
573577
!refusal &&
574578
!turnWasAborted &&
575579
(backendFinishedTurn || (!context.completionStatus && context.errors.length === 0))
580+
// The worker sends an error terminal with no `error` event only when it rebuilds an
581+
// ended run from its log, which drops the stored reason: a resume or reattach that
582+
// reaches a run that already ended, for example at its deadline. Say so rather than
583+
// leave the turn to a generic failure; a reported reason always wins.
584+
const endedWithoutReason =
585+
!turnWasAborted &&
586+
context.completionStatus === MothershipStreamV1CompletionStatus.error &&
587+
context.errors.length === 0
576588

577589
const result: OrchestratorResult = {
578590
success: succeeded,
@@ -591,6 +603,7 @@ export async function runCopilotLifecycle(
591603
toolCalls: buildToolCallSummaries(context),
592604
chatId: context.chatId,
593605
requestId: context.requestId,
606+
...(endedWithoutReason ? { error: ENDED_RUN_MESSAGE } : {}),
594607
...(refusal ? { error: refusal.userMessage, errorCode: refusal.code } : {}),
595608
errors: !succeeded && context.errors.length ? context.errors : undefined,
596609
usage: context.usage,

‎apps/sim/lib/mothership/request/session/buffer.ts‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -487,7 +487,7 @@ export async function readEvents(
487487
logger.warn('Skipping corrupt outbox entry', {
488488
streamId,
489489
reason: parsed.reason,
490-
message: parsed.message,
490+
detail: parsed.message,
491491
errors: parsed.errors,
492492
})
493493
continue

0 commit comments

Comments
 (0)