Skip to content

Commit e46504b

Browse files
committed
Remove default Mship execution cutoff and expose Slack agent prerequisites
1 parent e66af75 commit e46504b

13 files changed

Lines changed: 115 additions & 68 deletions

File tree

‎apps/sim/app/api/copilot/chat/stream/route.test.ts‎

Lines changed: 35 additions & 25 deletions
Original file line numberDiff line numberDiff line change
@@ -124,33 +124,43 @@ describe('copilot chat stream replay route', () => {
124124
})
125125
})
126126

127-
it('stops replay polling when run becomes cancelled', async () => {
128-
getLatestRunForStream
129-
.mockResolvedValueOnce({
130-
status: 'active',
131-
executionId: 'exec-1',
132-
id: 'run-1',
127+
it.each([0, 2 * 60 * 60_000])(
128+
'delivers cancellation after %i ms of replay',
129+
async (elapsedMs) => {
130+
const now = Date.now()
131+
const clock = vi.spyOn(Date, 'now').mockReturnValue(now)
132+
readEvents.mockImplementationOnce(async () => {
133+
clock.mockReturnValue(now + elapsedMs)
134+
return []
133135
})
134-
.mockResolvedValueOnce({
135-
status: 'cancelled',
136-
executionId: 'exec-1',
137-
id: 'run-1',
138-
})
139-
140-
const response = await GET(
141-
new NextRequest('http://localhost:3000/api/copilot/chat/stream?streamId=stream-1&after=0')
142-
)
136+
getLatestRunForStream
137+
.mockResolvedValueOnce({
138+
status: 'active',
139+
executionId: 'exec-1',
140+
id: 'run-1',
141+
})
142+
.mockResolvedValueOnce({
143+
status: 'cancelled',
144+
executionId: 'exec-1',
145+
id: 'run-1',
146+
})
147+
148+
const response = await GET(
149+
new NextRequest('http://localhost:3000/api/copilot/chat/stream?streamId=stream-1&after=0')
150+
)
143151

144-
const chunks = await readAllChunks(response)
145-
expect(chunks[0]).toBe(': accepted\n\n')
146-
expect(chunks.join('')).toContain(
147-
JSON.stringify({
148-
status: MothershipStreamV1CompletionStatus.cancelled,
149-
reason: 'terminal_status',
150-
})
151-
)
152-
expect(getLatestRunForStream).toHaveBeenCalledTimes(2)
153-
})
152+
const chunks = await readAllChunks(response)
153+
expect(chunks[0]).toBe(': accepted\n\n')
154+
expect(chunks.join('')).toContain(
155+
JSON.stringify({
156+
status: MothershipStreamV1CompletionStatus.cancelled,
157+
reason: 'terminal_status',
158+
})
159+
)
160+
expect(getLatestRunForStream).toHaveBeenCalledTimes(2)
161+
clock.mockRestore()
162+
}
163+
)
154164

155165
it('emits structured terminal replay error when run metadata disappears', async () => {
156166
getLatestRunForStream

‎apps/sim/app/api/copilot/chat/stream/route.ts‎

Lines changed: 2 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -37,13 +37,10 @@ import {
3737
} from '@/lib/mothership/request/session'
3838
import { toReplayEnvelope, toStreamBatchEvent } from '@/lib/mothership/request/session/types'
3939

40-
export const maxDuration = 3600
41-
4240
const logger = createLogger('CopilotChatStreamAPI')
4341
const POLL_INTERVAL_MS = 250
4442
const POLL_INTERVAL_MAX_MS = 2_000
4543
const REPLAY_KEEPALIVE_INTERVAL_MS = 15_000
46-
const MAX_STREAM_MS = 60 * 60 * 1000
4744

4845
function extractCanonicalRequestId(value: unknown): string {
4946
return typeof value === 'string' && value.length > 0 ? value : ''
@@ -396,7 +393,7 @@ async function handleResumeRequestBody({
396393
await flushEvents()
397394

398395
let pollDelayMs = POLL_INTERVAL_MS
399-
while (!controllerClosed && Date.now() - startTime < MAX_STREAM_MS) {
396+
while (!controllerClosed) {
400397
pollIterations += 1
401398
const currentRun = await readRun().catch((err) => {
402399
logger.warn('Failed to poll latest run for stream', {
@@ -419,7 +416,7 @@ async function handleResumeRequestBody({
419416
const flushed = await flushEvents()
420417
/* Adaptive tail: 4 Hz only while events are actually flowing; a quiet stream
421418
decays toward the cap so an attached client doesn't hammer Postgres + Redis
422-
at 4 Hz for up to an hour. Any flushed event snaps back to full rate. */
419+
at 4 Hz during long runs. Any flushed event snaps back to full rate. */
423420
pollDelayMs =
424421
flushed > 0 ? POLL_INTERVAL_MS : Math.min(pollDelayMs * 2, POLL_INTERVAL_MAX_MS)
425422

@@ -451,13 +448,6 @@ async function handleResumeRequestBody({
451448

452449
await sleep(pollDelayMs)
453450
}
454-
if (!controllerClosed && Date.now() - startTime >= MAX_STREAM_MS) {
455-
emitTerminalIfMissing(MothershipStreamV1CompletionStatus.error, {
456-
message: 'The stream recovery timed out before completion.',
457-
code: 'resume_timeout',
458-
reason: 'timeout',
459-
})
460-
}
461451
} catch (error) {
462452
if (!controllerClosed && !request.signal.aborted) {
463453
logger.warn('Stream replay failed', {

‎apps/sim/app/api/mothership/chat/route.ts‎

Lines changed: 0 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -10,8 +10,6 @@ import { handleUnifiedChatPost } from '@/lib/mothership/chat/post'
1010
import { validateShimEnvelope } from '@/lib/mothership/request/http'
1111
import { GET as copilotChatGet } from '@/app/api/copilot/chat/queries'
1212

13-
export const maxDuration = 3600
14-
1513
// Unified chat route surface.
1614
export const GET = withRouteHandler((request: NextRequest) => {
1715
const validation = mothershipChatGetQuerySchema.safeParse(

‎apps/sim/app/api/mothership/chat/stream/route.ts‎

Lines changed: 1 addition & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,9 +1,7 @@
11
import type { NextRequest } from 'next/server'
22
import { mothershipChatStreamQuerySchema } from '@/lib/api/contracts/mothership-chats'
33
import { validationErrorResponse } from '@/lib/api/server'
4-
import { GET as copilotStreamGet, maxDuration } from '@/app/api/copilot/chat/stream/route'
5-
6-
export { maxDuration }
4+
import { GET as copilotStreamGet } from '@/app/api/copilot/chat/stream/route'
75

86
export function GET(request: NextRequest) {
97
const validation = mothershipChatStreamQuerySchema.safeParse(

‎apps/sim/app/api/mothership/execute/route.ts‎

Lines changed: 0 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -48,8 +48,6 @@ import {
4848
import type { ChatContext } from '@/stores/panel'
4949
import { hasToolId } from '@/tools/tool-ids'
5050

51-
export const maxDuration = 3600
52-
5351
const logger = createLogger('MothershipExecuteAPI')
5452
const MOTHERSHIP_EXECUTE_STREAM_HEADER = 'x-mothership-execute-stream'
5553
const MOTHERSHIP_EXECUTE_STREAM_VALUE = 'ndjson'

‎apps/sim/blocks/blocks/slack.ts‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,7 @@ import {
2020
} from '@/blocks/utils'
2121
import type { SlackResponse } from '@/tools/slack/types'
2222
import { getTrigger } from '@/triggers'
23+
import { SLACK_AGENT_SCOPES } from '@/triggers/slack/capabilities'
2324

2425
/**
2526
* Canonical basic/advanced pair for the channel target, shared by the card
@@ -3594,6 +3595,11 @@ export const SlackV2Block: BlockConfig<SlackResponse> = {
35943595
description: 'Manage Slack messages, channels, users, files, Lists, canvases, and Agent Sessions',
35953596
longDescription:
35963597
'Build Slack workflows with messages, conversations, files, reactions, pins, bookmarks, user groups, profiles, Lists, canvases, and Agent Sessions. Operations that need additional app scopes use custom Slack bots. Lists require lists:read/lists:write and a paid Slack plan. Native Sim connections retain their existing permissions. Page through list outputs explicitly.',
3598+
bestPractices: `${SlackBlock.bestPractices}
3599+
Native agent-session streaming uses a custom Slack bot and a supported trigger event. Enable streamResponse and select the intended outputs in streamOutputs; Sim Chat (mothership) streams its content output. The trigger owns the streamed reply, so another send needs a separate purpose.
3600+
The custom-bot manifest's baseline agent scopes are ${SLACK_AGENT_SCOPES.join(', ')}. Selected capabilities add their required scopes and event subscriptions. Verify the installed app's grants, resource access and event configuration; a saved credential or edited manifest alone does not prove access.
3601+
streamIncludeToolCalls displays tool activity; it does not authorize tool execution. Configure Sim Chat's selected operations and their credentials separately. Choose streamIncludeThinking deliberately for the destination audience.
3602+
Block Kit messages can contain buttons, selects and forms, but interactions require the app's interactivity request URL and the corresponding interaction trigger. Treat message events, button/select callbacks and modal submissions as separate paths. A control's action_id/value correlates work; runtime actor and state checks authorize it.`,
35973603
hideFromToolbar: false,
35983604
sunset: undefined,
35993605
canvasPresentation: {

‎apps/sim/lib/mothership/auth/application-delegation.ts‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
11
import type { DelegatedPrincipal, OrganizationDelegatedPrincipal } from '@sim/auth/principal'
22
import { ORCHESTRATION_TIMEOUT_MS } from '@/lib/mothership/constants'
33

4-
/** Keeps delegated authority valid for the full bounded Copilot orchestration lifetime. */
4+
/** Per-operation authority expires independently of the assistant run lifetime. */
55
export const COPILOT_APPLICATION_DELEGATION_TTL_MS = ORCHESTRATION_TIMEOUT_MS
66

77
export interface CopilotExecutionContext {

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

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -10,7 +10,7 @@ export const SIM_AGENT_API_URL =
1010
? rawAgentUrl
1111
: SIM_AGENT_API_URL_DEFAULT
1212

13-
/** Default timeout for the copilot orchestration stream loop (60 min). */
13+
/** Bounded per-operation credential and permission lifetime; not a total run limit. */
1414
export const ORCHESTRATION_TIMEOUT_MS = 3_600_000
1515

1616
/**
@@ -26,14 +26,14 @@ export const TOOL_WATCHDOG_DEFAULT_MS = 60_000
2626
* executions, media/image generation, sandboxed code, deep research). Those
2727
* tools carry their own inner budgets (plan execution timeouts, sandbox
2828
* timeouts), so this cap only backstops a true hang and sits above all of
29-
* them — matching ORCHESTRATION_TIMEOUT_MS so it never undercuts a legal run.
29+
* them. This limits one operation, not the entire assistant run.
3030
*/
3131
export const TOOL_WATCHDOG_LONG_RUNNING_MS = ORCHESTRATION_TIMEOUT_MS
3232

3333
/** Extra slack the resume gate allows past the slowest pending tool's watchdog. */
3434
export const TOOL_WATCHDOG_RESUME_GRACE_MS = 30_000
3535

36-
/** Timeout for the client-side streaming response handler (60 min). */
36+
/** Maximum wait for one client-executed workflow tool (60 min). */
3737
export const STREAM_TIMEOUT_MS = 3_600_000
3838

3939
/**

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

Lines changed: 36 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -179,6 +179,7 @@ describe('copilot go stream helpers', () => {
179179
})
180180

181181
afterEach(() => {
182+
vi.useRealTimers()
182183
vi.unstubAllGlobals()
183184
})
184185

@@ -423,6 +424,41 @@ describe('copilot go stream helpers', () => {
423424
}
424425
)
425426

427+
it('keeps a healthy stream open beyond an hour when no deadline was requested', async () => {
428+
vi.useFakeTimers()
429+
let streamController: ReadableStreamDefaultController<Uint8Array> | undefined
430+
vi.mocked(fetch).mockResolvedValueOnce(
431+
new Response(
432+
new ReadableStream<Uint8Array>({
433+
start(controller) {
434+
streamController = controller
435+
},
436+
}),
437+
{ headers: { 'content-type': 'text/event-stream' } }
438+
)
439+
)
440+
const context = createStreamingContext()
441+
const pending = runStreamLoop(
442+
'https://example.com/mothership/stream',
443+
{},
444+
context,
445+
turnScopedExecContext(),
446+
{ flushAfterEvent: false }
447+
)
448+
await vi.advanceTimersByTimeAsync(2 * 60 * 60_000)
449+
expect(context.streamComplete).toBe(false)
450+
expect(context.errors).toEqual([])
451+
streamController?.enqueue(
452+
new TextEncoder().encode(
453+
'data: {"v":1,"type":"complete","seq":1,"ts":"","stream":{"streamId":"s"},"payload":{"status":"complete"}}\n\n'
454+
)
455+
)
456+
streamController?.close()
457+
await pending
458+
expect(context.errors).toEqual([])
459+
vi.useRealTimers()
460+
})
461+
426462
it('bounds response-header waits without classifying the deadline as user Stop', async () => {
427463
vi.mocked(fetch).mockImplementationOnce(
428464
(_url, options) =>

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

Lines changed: 16 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,6 @@
11
import { type Context, SpanStatusCode } from '@opentelemetry/api'
22
import { createLogger } from '@sim/logger'
33
import { getErrorMessage } from '@sim/utils/errors'
4-
import { ORCHESTRATION_TIMEOUT_MS } from '@/lib/mothership/constants'
54
import { MothershipStreamV1EventType } from '@/lib/mothership/generated/mothership-stream-v1'
65
import { CopilotSseCloseReason } from '@/lib/mothership/generated/trace-attribute-values-v1'
76
import { TraceAttr } from '@/lib/mothership/generated/trace-attributes-v1'
@@ -125,9 +124,13 @@ export async function runStreamLoop(
125124
execContext: ExecutionContext,
126125
options: StreamLoopOptions
127126
): Promise<void> {
128-
const { timeout = ORCHESTRATION_TIMEOUT_MS, abortSignal } = options
129-
const timeoutSignal = AbortSignal.timeout(Math.ceil(timeout))
130-
const requestSignal = abortSignal ? AbortSignal.any([abortSignal, timeoutSignal]) : timeoutSignal
127+
const { timeout, abortSignal } = options
128+
const timeoutSignal = timeout === undefined ? undefined : AbortSignal.timeout(Math.ceil(timeout))
129+
const requestSignal = timeoutSignal
130+
? abortSignal
131+
? AbortSignal.any([abortSignal, timeoutSignal])
132+
: timeoutSignal
133+
: abortSignal
131134
const filePreviewAdapterState = createFilePreviewAdapterState()
132135
const attemptedInlineImages = new Set<string>()
133136

@@ -262,12 +265,15 @@ export async function runStreamLoop(
262265
},
263266
}
264267

265-
const timeoutId = setTimeout(() => {
266-
context.errors.push('Request timed out')
267-
context.streamComplete = true
268-
endedOn = CopilotSseCloseReason.Timeout
269-
reader.cancel().catch(() => {})
270-
}, timeout)
268+
const timeoutId =
269+
timeout === undefined
270+
? undefined
271+
: setTimeout(() => {
272+
context.errors.push('Request timed out')
273+
context.streamComplete = true
274+
endedOn = CopilotSseCloseReason.Timeout
275+
reader.cancel().catch(() => {})
276+
}, timeout)
271277

272278
try {
273279
await processSSEStream(reader, abortSignal, async (raw) => {

0 commit comments

Comments
 (0)