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
46 changes: 45 additions & 1 deletion apps/sim/app/api/copilot/chat/stream/route.test.ts
Original file line number Diff line number Diff line change
@@ -1,11 +1,20 @@
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 { 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,
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(() => ({
Expand Down Expand Up @@ -62,6 +71,10 @@ async function readAllChunks(response: Response): Promise<string[]> {
}

describe('copilot chat stream replay route', () => {
afterEach(() => {
vi.useRealTimers()
})

beforeEach(() => {
authMockFns.mockGetSession.mockResolvedValue({
user: { id: 'user-1' },
Expand Down Expand Up @@ -171,4 +184,35 @@ 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 () => {
const exporter = new InMemorySpanExporter()
trace.setGlobalTracerProvider(
new BasicTracerProvider({ spanProcessors: [new SimpleSpanProcessor(exporter)] })
)
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}"`)
const resume = exporter
.getFinishedSpans()
.find((span) => span.name === TraceSpan.CopilotResumeRequest)
expect(resume?.attributes[TraceAttr.CopilotResumeOutcome]).toBe(
CopilotResumeOutcome.EndedWithoutTerminal
)
trace.disable()
})
})
18 changes: 9 additions & 9 deletions apps/sim/app/api/copilot/chat/stream/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Comment thread
waleedlatif1 marked this conversation as resolved.

function extractCanonicalRequestId(value: unknown): string {
return typeof value === 'string' && value.length > 0 ? value : ''
Expand Down Expand Up @@ -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', {
Expand All @@ -473,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,
Expand Down
5 changes: 5 additions & 0 deletions apps/sim/app/api/mothership/chat/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
1 change: 1 addition & 0 deletions apps/sim/app/api/mothership/execute/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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')
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<Uint8Array>({
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<string>()
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)
Expand Down
11 changes: 10 additions & 1 deletion apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down Expand Up @@ -2868,6 +2876,7 @@ export function useChat(
error: toError(err).message,
})
}
attempt = lastCursorRef.current !== cursorBeforeAttempt ? 1 : attempt + 1
}

logger.error('All reconnect attempts exhausted', {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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 = {
Expand Down Expand Up @@ -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' },
Expand Down
10 changes: 7 additions & 3 deletions apps/sim/lib/mothership/auth/application-delegation.ts
Original file line number Diff line number Diff line change
@@ -1,8 +1,12 @@
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 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

export interface CopilotExecutionContext {
requestMode?: string
Expand Down
4 changes: 2 additions & 2 deletions apps/sim/lib/mothership/auth/file-delegation.test.ts
Original file line number Diff line number Diff line change
@@ -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',
Expand Down Expand Up @@ -35,7 +35,7 @@ describe('Copilot file delegation', () => {
},
})
expect(principal.expiresAt.getTime() - principal.issuedAt.getTime()).toBe(
ORCHESTRATION_TIMEOUT_MS
COPILOT_APPLICATION_DELEGATION_TTL_MS
)
})

Expand Down
24 changes: 18 additions & 6 deletions apps/sim/lib/mothership/constants.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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. 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

/**
* Watchdog cap for a single sim-executed copilot tool. A tool that neither
Expand All @@ -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
Expand All @@ -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.
*/
Expand Down
Loading
Loading