Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 5 additions & 5 deletions apps/sim/lib/billing/core/usage-analytics.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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

Expand Down
3 changes: 2 additions & 1 deletion apps/sim/lib/core/redis/byte-budget.server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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.
*
Expand Down
9 changes: 7 additions & 2 deletions apps/sim/lib/mothership/constants.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
4 changes: 2 additions & 2 deletions apps/sim/lib/mothership/request/go/stream.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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,
Expand Down
52 changes: 52 additions & 0 deletions apps/sim/lib/mothership/request/lifecycle/run.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<void> => {
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<void> => {
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<void> => {
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(
Expand Down Expand Up @@ -3326,6 +3377,7 @@ describe('runCopilotLifecycle', () => {
)

expect(result.success).toBe(false)
expect(result.error).toBeUndefined()
expect(result.errors).toEqual(['The provider is overloaded'])
})

Expand Down
15 changes: 15 additions & 0 deletions apps/sim/lib/mothership/request/lifecycle/run.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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,
Expand All @@ -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,
Expand Down
2 changes: 1 addition & 1 deletion apps/sim/lib/mothership/request/session/buffer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading