From d78d860690263b5d9a2921db5b519e52f4706806 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 30 Sep 2026 01:18:52 -0700 Subject: [PATCH 1/7] fix(mothership): settle Chat runs no controller will finish A run whose process died before finalize, whose controller was superseded with no successor, or that was stopped while no controller existed stayed unfinished forever: its chat marker kept pointing at it, so the chat read as busy and a reconnect polled a run nothing would ever end. - Stop now settles the run as cancelled once no controller of its stream holds the chat lock, including after it force-releases a controller that did not exit in time. - The stale-execution cron settles leased runs whose stream holds no chat lock, has no replay buffer left, and has been idle past the orchestration budget (so a reconnect has nothing left to resume), and runs without a lease once idle for 24 hours. A run Stop already closed settles as cancelled, any other as error; its chat marker is released. - Each settle is one conditional update on the run row that requires it to be unfinished, idle, and still naming the controller that was observed, so a finalizing controller or a successor's claim wins or loses against it atomically and the run settles exactly once. --- .../cron/cleanup-stale-executions/route.ts | 18 ++ .../async-runs/orphaned-runs.integration.ts | 286 ++++++++++++++++++ .../mothership/async-runs/orphaned-runs.ts | 234 ++++++++++++++ .../request/application/controls.ts | 10 +- .../lib/mothership/request/session/buffer.ts | 17 ++ .../request/session/controller-lease.ts | 20 ++ 6 files changed, 584 insertions(+), 1 deletion(-) create mode 100644 apps/sim/lib/mothership/async-runs/orphaned-runs.integration.ts create mode 100644 apps/sim/lib/mothership/async-runs/orphaned-runs.ts diff --git a/apps/sim/app/api/cron/cleanup-stale-executions/route.ts b/apps/sim/app/api/cron/cleanup-stale-executions/route.ts index c920c86826d..0eba10da8cc 100644 --- a/apps/sim/app/api/cron/cleanup-stale-executions/route.ts +++ b/apps/sim/app/api/cron/cleanup-stale-executions/route.ts @@ -32,6 +32,7 @@ import { STALE_SWEEPABLE_EXECUTION_STATUSES, type StaleSweepableExecutionStatus, } from '@/lib/logs/types' +import { sweepOrphanedRuns } from '@/lib/mothership/async-runs/orphaned-runs' import { cancelStaleDispatches } from '@/lib/table/dispatcher' import { deleteFile } from '@/lib/uploads/core/storage-service' import { @@ -738,6 +739,20 @@ export const GET = withRouteHandler(async (request: NextRequest) => { }) } + /** + * Settle Chat runs no controller will finish: their process died, their + * controller was superseded without a successor, or Stop found none. Without + * this they stay unfinished forever and keep their chat marked as busy. + */ + let orphanedRunsSettled = 0 + try { + orphanedRunsSettled = (await sweepOrphanedRuns()).settledRunIds.length + } catch (error) { + logger.error('Failed to settle orphaned Chat runs:', { + error: toError(error).message, + }) + } + return NextResponse.json({ success: true, executions: { @@ -768,6 +783,9 @@ export const GET = withRouteHandler(async (request: NextRequest) => { pruned: deploymentOperationsPruned, retentionDays: DEPLOYMENT_OPERATION_RETENTION_DAYS, }, + chatRuns: { + orphanedSettled: orphanedRunsSettled, + }, }) } catch (error) { logger.error('Error in stale execution cleanup job:', error) diff --git a/apps/sim/lib/mothership/async-runs/orphaned-runs.integration.ts b/apps/sim/lib/mothership/async-runs/orphaned-runs.integration.ts new file mode 100644 index 00000000000..9cc4db35b25 --- /dev/null +++ b/apps/sim/lib/mothership/async-runs/orphaned-runs.integration.ts @@ -0,0 +1,286 @@ +/** + * Settlement of Chat runs that no controller owns, against real PostgreSQL and Redis: + * the chat stream lock, the replay buffer keys, the run and chat rows, and the Stop + * use case are production code. A local HTTP server stands in for the worker's abort + * endpoint. + */ +import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest' + +const { redisUrl, inheritedEnv, worker } = await vi.hoisted(async () => { + const { readTestRedisUrl } = await import('@sim/db/testing/test-infrastructure') + const { createServer } = await import('node:http') + const server = createServer(async (request, response) => { + for await (const _chunk of request) { + } + response.writeHead(200, { 'content-type': 'application/json' }) + response.end(JSON.stringify({ settled: true })) + }) + await new Promise((resolve) => server.listen(0, '127.0.0.1', resolve)) + const { port } = server.address() as { port: number } + const url = readTestRedisUrl() + const inheritedEnv = { + REDIS_URL: process.env.REDIS_URL, + SIM_AGENT_API_URL: process.env.SIM_AGENT_API_URL, + } + /** The real Redis module and worker URL resolution read these at import. */ + process.env.REDIS_URL = url + process.env.SIM_AGENT_API_URL = `http://127.0.0.1:${port}` + return { redisUrl: url, inheritedEnv, worker: { server } } +}) + +import { db } from '@sim/db' +import { copilotChats, copilotRuns, permissions, user, workspace } from '@sim/db/schema' +import { generateId } from '@sim/utils/id' +import { eq, inArray, sql } from 'drizzle-orm' +import { closeRedisConnection, getRedisClient } from '@/lib/core/config/redis' +import { + settleStoppedRunWithoutController, + sweepOrphanedRuns, +} from '@/lib/mothership/async-runs/orphaned-runs' +import { updateRunStatus } from '@/lib/mothership/async-runs/repository' +import { abortRun } from '@/lib/mothership/request/application/controls' +import { claimRunController } from '@/lib/mothership/request/lifecycle/controller-ownership' +import { chatStreamLockKey } from '@/lib/mothership/request/session/controller-lease' + +function redis() { + const client = getRedisClient() + if (!client) throw new Error('The integration suite requires TEST_REDIS_URL') + return client +} + +afterAll(async () => { + await closeRedisConnection() + await new Promise((resolve) => worker.server.close(() => resolve())) + for (const [key, value] of Object.entries(inheritedEnv)) { + if (value === undefined) delete process.env[key] + else process.env[key] = value + } +}) + +describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => { + const userId = generateId() + const workspaceId = generateId() + const chatIds: string[] = [] + + beforeAll(async () => { + const now = new Date() + await db.insert(user).values({ + id: userId, + name: 'Orphaned run fixture', + email: `${userId}@orphaned-runs.test`, + emailVerified: true, + createdAt: now, + updatedAt: now, + }) + await db.insert(workspace).values({ + id: workspaceId, + name: 'Orphaned run fixture', + ownerId: userId, + billedAccountUserId: userId, + }) + await db.insert(permissions).values({ + id: generateId(), + userId, + entityType: 'workspace', + entityId: workspaceId, + permissionType: 'admin', + }) + }) + + afterAll(async () => { + if (chatIds.length) await db.delete(copilotChats).where(inArray(copilotChats.id, chatIds)) + await db.delete(permissions).where(eq(permissions.userId, userId)) + await db.delete(workspace).where(eq(workspace.id, workspaceId)) + await db.delete(user).where(eq(user.id, userId)) + }) + + /** + * A run as the chat POST admits it: the chat marker names its stream and the run + * records the lock value its first controller held. `idleMinutes` backdates its + * last durable write; `controllerToken: null` is a run with no lease protocol. + */ + async function admittedRun( + options: { + idleMinutes?: number + status?: 'active' | 'paused_waiting_for_tool' + controllerToken?: string | null + stopped?: boolean + } = {} + ) { + const chatId = generateId() + const streamId = generateId() + const runId = generateId() + chatIds.push(chatId) + const controllerToken = + options.controllerToken === undefined + ? `${streamId}\n${generateId()}` + : options.controllerToken + const idle = sql`now() - make_interval(mins => ${options.idleMinutes ?? 0})` + await db.insert(copilotChats).values({ + id: chatId, + userId, + workspaceId, + type: 'mothership', + conversationId: streamId, + }) + await db.insert(copilotRuns).values({ + id: runId, + executionId: generateId(), + chatId, + userId, + workspaceId, + streamId, + toolExecutionVersion: 2, + status: options.status ?? 'active', + requestContext: controllerToken + ? { requestId: generateId(), controllerToken, recovery: { kind: 'interactive_stream' } } + : { source: 'headless_lifecycle' }, + startedAt: idle, + updatedAt: idle, + ...(options.stopped ? { toolAdmissionClosedAt: idle } : {}), + }) + return { chatId, streamId, runId, controllerToken } + } + + async function stored(runId: string) { + const [run] = await db.select().from(copilotRuns).where(eq(copilotRuns.id, runId)) + const [chat] = await db + .select({ conversationId: copilotChats.conversationId }) + .from(copilotChats) + .where(eq(copilotChats.id, run.chatId)) + return { ...run, marker: chat?.conversationId ?? null } + } + + it('settles a run whose controller died once its recovery window has passed', async () => { + const orphan = await admittedRun({ idleMinutes: 90, status: 'paused_waiting_for_tool' }) + + const { settledRunIds } = await sweepOrphanedRuns() + + expect(settledRunIds).toContain(orphan.runId) + const run = await stored(orphan.runId) + expect(run.status).toBe('error') + expect(run.error).toBeTruthy() + expect(run.completedAt).not.toBeNull() + expect(run.toolAdmissionClosedAt).not.toBeNull() + expect(run.marker).toBeNull() + }) + + it('settles a run stopped while no controller owned it as cancelled', async () => { + const orphan = await admittedRun({ idleMinutes: 90, stopped: true }) + + const { settledRunIds } = await sweepOrphanedRuns() + + expect(settledRunIds).toContain(orphan.runId) + expect((await stored(orphan.runId)).status).toBe('cancelled') + }) + + it('settles a run without a lease only after the unleased ceiling', async () => { + const recent = await admittedRun({ idleMinutes: 90, controllerToken: null }) + const abandoned = await admittedRun({ idleMinutes: 25 * 60, controllerToken: null }) + + const { settledRunIds } = await sweepOrphanedRuns() + + expect(settledRunIds).toContain(abandoned.runId) + expect(settledRunIds).not.toContain(recent.runId) + expect((await stored(abandoned.runId)).status).toBe('error') + expect((await stored(recent.runId)).status).toBe('active') + }) + + it('never settles a run whose stream holds its chat lock, or that is still recoverable', async () => { + const leased = await admittedRun({ idleMinutes: 90 }) + await redis().set(chatStreamLockKey(leased.chatId), leased.controllerToken!, 'EX', 60) + /** A recovering controller holds the lock under a new token before it claims the run. */ + const recovering = await admittedRun({ idleMinutes: 90 }) + await redis().set( + chatStreamLockKey(recovering.chatId), + `${recovering.streamId}\n${generateId()}`, + 'EX', + 60 + ) + const replayable = await admittedRun({ idleMinutes: 90 }) + await redis().set(`mothership_stream:${replayable.streamId}:seq`, '4', 'EX', 60) + const fresh = await admittedRun({ idleMinutes: 5 }) + + try { + const { settledRunIds } = await sweepOrphanedRuns() + + for (const run of [leased, recovering, replayable, fresh]) { + expect(settledRunIds).not.toContain(run.runId) + const current = await stored(run.runId) + expect(current.status).toBe('active') + expect(current.marker).toBe(run.streamId) + } + } finally { + await redis().del( + chatStreamLockKey(leased.chatId), + chatStreamLockKey(recovering.chatId), + `mothership_stream:${replayable.streamId}:seq` + ) + } + }) + + it('settles a run exactly once when a sweep races its own controller finalizing', async () => { + for (let attempt = 0; attempt < 20; attempt++) { + const orphan = await admittedRun({ idleMinutes: 90 }) + + const [finalized, sweep] = await Promise.all([ + updateRunStatus(orphan.runId, 'complete', {}, orphan.controllerToken!), + sweepOrphanedRuns(), + ]) + + const swept = sweep.settledRunIds.includes(orphan.runId) + expect(Boolean(finalized) !== swept).toBe(true) + expect((await stored(orphan.runId)).status).toBe(swept ? 'error' : 'complete') + } + }) + + it('settles a run exactly once when a sweep races a recovering controller claiming it', async () => { + for (let attempt = 0; attempt < 20; attempt++) { + const orphan = await admittedRun({ idleMinutes: 90 }) + + const [claimed, sweep] = await Promise.all([ + claimRunController({ + runId: orphan.runId, + chatId: orphan.chatId, + previousToken: orphan.controllerToken!, + token: `${orphan.streamId}\n${generateId()}`, + }), + sweepOrphanedRuns(), + ]) + + const swept = sweep.settledRunIds.includes(orphan.runId) + expect(claimed !== swept).toBe(true) + expect((await stored(orphan.runId)).status).toBe(swept ? 'error' : 'active') + } + }) + + it('settles a stopped run as cancelled when no controller owns it', async () => { + const orphan = await admittedRun({ status: 'paused_waiting_for_tool' }) + + const result = await abortRun.execute({ + principal: { kind: 'session', userId, sessionId: generateId() }, + input: { streamId: orphan.streamId, chatId: orphan.chatId, workspaceId }, + }) + + expect(result).toMatchObject({ aborted: true }) + const run = await stored(orphan.runId) + expect(run.status).toBe('cancelled') + expect(run.completedAt).not.toBeNull() + expect(run.toolAdmissionClosedAt).not.toBeNull() + expect(run.marker).toBeNull() + }) + + it('leaves a stopped run to the controller of its stream that holds the chat lock', async () => { + const owned = await admittedRun() + await redis().set(chatStreamLockKey(owned.chatId), owned.controllerToken!, 'EX', 60) + + try { + expect(await settleStoppedRunWithoutController(owned.runId)).toBe(false) + const run = await stored(owned.runId) + expect(run.status).toBe('active') + expect(run.marker).toBe(owned.streamId) + } finally { + await redis().del(chatStreamLockKey(owned.chatId)) + } + }) +}) diff --git a/apps/sim/lib/mothership/async-runs/orphaned-runs.ts b/apps/sim/lib/mothership/async-runs/orphaned-runs.ts new file mode 100644 index 00000000000..5268b6c8c5d --- /dev/null +++ b/apps/sim/lib/mothership/async-runs/orphaned-runs.ts @@ -0,0 +1,234 @@ +import { db } from '@sim/db' +import { type CopilotRunStatus, copilotChats, copilotRuns } from '@sim/db/schema' +import { createLogger } from '@sim/logger' +import { getErrorMessage } from '@sim/utils/errors' +import { + and, + asc, + eq, + gt, + inArray, + isNotNull, + isNull, + notInArray, + or, + type SQL, + sql, +} from 'drizzle-orm' +import { publishChatStatusChanged } from '@/lib/mothership/chat-status' +import { ORCHESTRATION_TIMEOUT_MS } from '@/lib/mothership/constants' +import { findStreamsWithReplay } from '@/lib/mothership/request/session/buffer' +import { findStreamsHoldingChatLock } from '@/lib/mothership/request/session/controller-lease' + +const logger = createLogger('OrphanedCopilotRuns') + +const TERMINAL_RUN_STATUSES: CopilotRunStatus[] = ['complete', 'error', 'cancelled'] +/** Listed positively so the sweep's scan stays on the status index once few runs are open. */ +const UNFINISHED_RUN_STATUSES: CopilotRunStatus[] = [ + 'active', + 'paused_waiting_for_tool', + 'resuming', +] + +/** + * How long a leased run must sit without a controller or a durable write before it is + * settled. No worker leg outlives the orchestration budget, so past it a reconnect has + * nothing left to resume. + */ +export const ORPHANED_RUN_GRACE_MS = ORCHESTRATION_TIMEOUT_MS + +/** + * Runs admitted without a chat lease (headless turns and rows from before the lease + * protocol) have no liveness signal, so only an age far past any process lifetime + * proves them dead. + */ +export const UNLEASED_RUN_GRACE_MS = 24 * 60 * 60 * 1000 + +export const ORPHANED_RUN_ERROR = 'This response was interrupted before it finished.' + +const SWEEP_BATCH_SIZE = 500 +const SWEEP_MAX_ROWS_PER_RUN = 10_000 + +const controllerToken = sql`${copilotRuns.requestContext}->>'controllerToken'` + +function idleFor(ms: number): SQL { + return sql`${copilotRuns.updatedAt} < now() - make_interval(secs => ${ms / 1000})` +} + +const leasedRunIdle = and(isNotNull(controllerToken), idleFor(ORPHANED_RUN_GRACE_MS)) +const unleasedRunIdle = and(isNull(controllerToken), idleFor(UNLEASED_RUN_GRACE_MS)) + +interface UnownedRun { + id: string + chatId: string + streamId: string + userId: string + workspaceId: string | null + organizationId: string | null + controllerToken: string | null +} + +const unownedRunColumns = { + id: copilotRuns.id, + chatId: copilotRuns.chatId, + streamId: copilotRuns.streamId, + userId: copilotRuns.userId, + workspaceId: copilotRuns.workspaceId, + organizationId: copilotRuns.organizationId, + controllerToken, +} + +type Transaction = Parameters[0]>[0] + +/** + * Settles each run only while it is still unfinished and still names the controller + * the caller observed, so a finalizing controller or a successor's claim, both of + * which write the same row, wins or loses atomically against it. A run that Stop + * already closed settles as cancelled. The chat marker is released without touching + * the chat's ordering timestamp. + */ +async function settleRuns( + tx: Transaction, + runs: UnownedRun[], + guard: SQL | undefined, + outcome: { status: 'error' | 'cancelled'; error?: string } +): Promise { + const settled: UnownedRun[] = [] + const apply = async (owner: SQL | undefined, batch: UnownedRun[]) => { + if (batch.length === 0) return + const rows = await tx + .update(copilotRuns) + .set({ + status: sql`(CASE WHEN ${copilotRuns.toolAdmissionClosedAt} IS NOT NULL THEN 'cancelled' ELSE ${outcome.status} END)::copilot_run_status`, + error: sql`CASE WHEN ${copilotRuns.toolAdmissionClosedAt} IS NOT NULL THEN NULL ELSE ${outcome.error ?? null}::text END`, + completedAt: sql`now()`, + toolAdmissionClosedAt: sql`coalesce(${copilotRuns.toolAdmissionClosedAt}, now())`, + updatedAt: sql`now()`, + }) + .where( + and( + inArray( + copilotRuns.id, + batch.map((run) => run.id) + ), + notInArray(copilotRuns.status, TERMINAL_RUN_STATUSES), + owner, + guard + ) + ) + .returning({ id: copilotRuns.id }) + const won = new Set(rows.map((row) => row.id)) + settled.push(...batch.filter((run) => won.has(run.id))) + } + + await apply( + isNull(controllerToken), + runs.filter((run) => run.controllerToken === null) + ) + for (const run of runs) { + if (run.controllerToken !== null) await apply(eq(controllerToken, run.controllerToken), [run]) + } + + for (const run of settled) { + await tx + .update(copilotChats) + .set({ conversationId: null }) + .where(and(eq(copilotChats.id, run.chatId), eq(copilotChats.conversationId, run.streamId))) + } + return settled +} + +function announceSettled(runs: UnownedRun[]): void { + for (const run of runs) { + try { + publishChatStatusChanged(run, { + chatId: run.chatId, + type: 'completed', + streamId: run.streamId, + }) + } catch (error) { + logger.warn('Settled run status could not be announced', { + runId: run.id, + error: getErrorMessage(error), + }) + } + } +} + +/** + * Settles runs that no controller will ever finish: a leased run whose stream holds no + * chat lock and has no replay buffer left, idle past the recovery window, and a run + * without a lease idle past any process lifetime. Unreadable locks skip leased runs. + */ +export async function sweepOrphanedRuns(): Promise<{ settledRunIds: string[] }> { + const settledRunIds: string[] = [] + const idle = or(leasedRunIdle, unleasedRunIdle) + let cursor: string | undefined + let considered = 0 + + while (considered < SWEEP_MAX_ROWS_PER_RUN) { + const candidates: UnownedRun[] = await db + .select(unownedRunColumns) + .from(copilotRuns) + .where( + and( + inArray(copilotRuns.status, UNFINISHED_RUN_STATUSES), + idle, + cursor ? gt(copilotRuns.id, cursor) : undefined + ) + ) + .orderBy(asc(copilotRuns.id)) + .limit(Math.min(SWEEP_BATCH_SIZE, SWEEP_MAX_ROWS_PER_RUN - considered)) + if (candidates.length === 0) break + considered += candidates.length + cursor = candidates[candidates.length - 1].id + + const leased = candidates.filter((run) => run.controllerToken !== null) + let unowned = candidates.filter((run) => run.controllerToken === null) + if (leased.length > 0) { + try { + const [locked, replayable] = await Promise.all([ + findStreamsHoldingChatLock(leased), + findStreamsWithReplay(leased.map((run) => run.streamId)), + ]) + unowned = unowned.concat( + leased.filter((run) => !locked.has(run.streamId) && !replayable.has(run.streamId)) + ) + } catch (error) { + logger.warn('Chat stream ownership is unreadable; leaving leased runs for a later sweep', { + error: getErrorMessage(error), + }) + } + } + + const settled = await db.transaction((tx) => + settleRuns(tx, unowned, idle, { status: 'error', error: ORPHANED_RUN_ERROR }) + ) + announceSettled(settled.filter((run) => run.controllerToken !== null)) + settledRunIds.push(...settled.map((run) => run.id)) + } + + if (settledRunIds.length > 0) { + logger.info('Settled runs no controller owned', { count: settledRunIds.length }) + } + return { settledRunIds } +} + +/** + * Settles a stopped run as cancelled when no controller of its stream holds the chat + * lock. A live controller observes the Stop and settles its own run. + */ +export async function settleStoppedRunWithoutController(runId: string): Promise { + const [run] = await db + .select(unownedRunColumns) + .from(copilotRuns) + .where(and(eq(copilotRuns.id, runId), notInArray(copilotRuns.status, TERMINAL_RUN_STATUSES))) + .limit(1) + if (!run?.controllerToken) return false + if ((await findStreamsHoldingChatLock([run])).has(run.streamId)) return false + const settled = await db.transaction((tx) => + settleRuns(tx, [run], undefined, { status: 'cancelled' }) + ) + announceSettled(settled) + return settled.length > 0 +} diff --git a/apps/sim/lib/mothership/request/application/controls.ts b/apps/sim/lib/mothership/request/application/controls.ts index b1f3ad07131..e911028ed88 100644 --- a/apps/sim/lib/mothership/request/application/controls.ts +++ b/apps/sim/lib/mothership/request/application/controls.ts @@ -8,6 +8,7 @@ import { defineOrganizationOperation } from '@/lib/core/application/organization import { OrchestrationError } from '@/lib/core/orchestration/types' import { markExecutionCancelled } from '@/lib/execution/cancellation' import { abortManualExecution } from '@/lib/execution/manual-cancellation' +import { settleStoppedRunWithoutController } from '@/lib/mothership/async-runs/orphaned-runs' import { areStreamToolExecutionsSettled, getLatestRunForStream, @@ -208,8 +209,15 @@ export const abortRun = defineAuthorizedChatUseCase({ if (!settled) { await releasePendingChatStream(chatId, streamId) logger.warn('Stopped stream did not settle; released its chat lock', { chatId, streamId }) - return { aborted: true, settled: false, forceReleased: true } } + /** A run with no controller left, or none to begin with, has nothing else to settle it. */ + await settleStoppedRunWithoutController(run.id).catch((error) => { + logger.warn('Stopped run without a controller could not be settled', { + streamId, + error: getErrorMessage(error), + }) + }) + if (!settled) return { aborted: true, settled: false, forceReleased: true } const toolsSettled = await areStreamToolExecutionsSettled(streamId, userId).catch((error) => { logger.warn('Stopped stream tool settlement could not be verified', { streamId, diff --git a/apps/sim/lib/mothership/request/session/buffer.ts b/apps/sim/lib/mothership/request/session/buffer.ts index e1760adf11d..0cbd02ad216 100644 --- a/apps/sim/lib/mothership/request/session/buffer.ts +++ b/apps/sim/lib/mothership/request/session/buffer.ts @@ -427,6 +427,23 @@ export async function getLatestSeq(streamId: string): Promise { }) } +/** The streams among these whose replay buffer has not yet expired. */ +export async function findStreamsWithReplay(streamIds: string[]): Promise> { + const redis = getRedisClient() + if (!redis) throw new Error('Redis is required for mothership stream durability') + if (streamIds.length === 0) return new Set() + const pipeline = redis.pipeline() + for (const streamId of streamIds) pipeline.exists(getSeqKey(streamId)) + const replies = (await pipeline.exec()) ?? [] + const withReplay = new Set() + streamIds.forEach((streamId, index) => { + const [error, count] = replies[index] ?? [new Error('Redis returned no reply')] + if (error) throw error + if (count === 1) withReplay.add(streamId) + }) + return withReplay +} + export async function writeAbortMarker(streamId: string): Promise { const ttlSeconds = getStreamConfig().ttlSeconds await withRedisRetry({ operation: 'write_abort_marker', streamId }, async (redis) => { diff --git a/apps/sim/lib/mothership/request/session/controller-lease.ts b/apps/sim/lib/mothership/request/session/controller-lease.ts index b7d89647a33..2519854a3ce 100644 --- a/apps/sim/lib/mothership/request/session/controller-lease.ts +++ b/apps/sim/lib/mothership/request/session/controller-lease.ts @@ -28,6 +28,26 @@ export async function assertChatStreamLease(lease: ChatStreamLease): Promise +): Promise> { + const redis = getRedisClient() + if (!redis) throw new Error('Chat stream locks are unreadable without Redis') + if (streams.length === 0) return new Set() + const values = await redis.mget(streams.map(({ chatId }) => chatStreamLockKey(chatId))) + const held = new Set() + streams.forEach(({ streamId }, index) => { + const value = values[index] + if (value && streamIdFromLock(value) === streamId) held.add(streamId) + }) + return held +} + /** Whether this lease still holds its chat lock; an unreadable lock counts as lost. */ export async function holdsChatStreamLease(lease: ChatStreamLease): Promise { try { From bc9369c634375ce2ec3fbcf05f40f06187218d5b Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 30 Sep 2026 01:32:45 -0700 Subject: [PATCH 2/7] fix(mothership): lock chats before runs when settling orphans, and label them precisely - Settling now locks the affected chat rows first, in id order, as a controller's claim does. Locking the run and then the chat deadlocked against a concurrent reconnect claim. - Each sweep batch fails on its own, the chat markers of a batch clear in one statement, and a sweep settles at most 5k rows with a short pause between full batches. - A run settles as cancelled only when its user pressed Stop. A newer turn also closes tool admission on older runs, and those now settle as errors. - Runs without a controller lease keep their last write as their completion and retention time and read "never finalized (no controller lease)". - The leased-run grace no longer derives from the orchestration deadline. Liveness comes only from the heartbeat-renewed chat lock; the grace and the replay TTL only bound how long a reconnect can resume a dead run. --- .../async-runs/orphaned-runs.integration.ts | 104 ++++++++++- .../mothership/async-runs/orphaned-runs.ts | 163 ++++++++++++------ 2 files changed, 213 insertions(+), 54 deletions(-) diff --git a/apps/sim/lib/mothership/async-runs/orphaned-runs.integration.ts b/apps/sim/lib/mothership/async-runs/orphaned-runs.integration.ts index 9cc4db35b25..88a7eb6b401 100644 --- a/apps/sim/lib/mothership/async-runs/orphaned-runs.integration.ts +++ b/apps/sim/lib/mothership/async-runs/orphaned-runs.integration.ts @@ -29,13 +29,24 @@ const { redisUrl, inheritedEnv, worker } = await vi.hoisted(async () => { }) import { db } from '@sim/db' -import { copilotChats, copilotRuns, permissions, user, workspace } from '@sim/db/schema' +import { + copilotChats, + copilotRequestStops, + copilotRuns, + permissions, + user, + workspace, +} from '@sim/db/schema' +import { sleep } from '@sim/utils/helpers' import { generateId } from '@sim/utils/id' +import { randomInt } from '@sim/utils/random' import { eq, inArray, sql } from 'drizzle-orm' import { closeRedisConnection, getRedisClient } from '@/lib/core/config/redis' import { + ORPHANED_RUN_ERROR, settleStoppedRunWithoutController, sweepOrphanedRuns, + UNLEASED_RUN_ERROR, } from '@/lib/mothership/async-runs/orphaned-runs' import { updateRunStatus } from '@/lib/mothership/async-runs/repository' import { abortRun } from '@/lib/mothership/request/application/controls' @@ -98,6 +109,8 @@ describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => { * A run as the chat POST admits it: the chat marker names its stream and the run * records the lock value its first controller held. `idleMinutes` backdates its * last durable write; `controllerToken: null` is a run with no lease protocol. + * `stopped` records the user's Stop intent; `superseded` only closes tool admission, + * as a newer turn's workbench does to older runs. */ async function admittedRun( options: { @@ -105,6 +118,7 @@ describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => { status?: 'active' | 'paused_waiting_for_tool' controllerToken?: string | null stopped?: boolean + superseded?: boolean } = {} ) { const chatId = generateId() @@ -137,8 +151,10 @@ describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => { : { source: 'headless_lifecycle' }, startedAt: idle, updatedAt: idle, - ...(options.stopped ? { toolAdmissionClosedAt: idle } : {}), + ...(options.stopped || options.superseded ? { toolAdmissionClosedAt: idle } : {}), }) + if (options.stopped) + await db.insert(copilotRequestStops).values({ userId, workspaceId, streamId }) return { chatId, streamId, runId, controllerToken } } @@ -174,15 +190,35 @@ describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => { expect((await stored(orphan.runId)).status).toBe('cancelled') }) + it('settles a run a newer turn superseded, without a Stop, as an error', async () => { + const orphan = await admittedRun({ idleMinutes: 90, superseded: true }) + + const { settledRunIds } = await sweepOrphanedRuns() + + expect(settledRunIds).toContain(orphan.runId) + const run = await stored(orphan.runId) + expect(run.status).toBe('error') + expect(run.error).toBeTruthy() + }) + it('settles a run without a lease only after the unleased ceiling', async () => { const recent = await admittedRun({ idleMinutes: 90, controllerToken: null }) const abandoned = await admittedRun({ idleMinutes: 25 * 60, controllerToken: null }) - + const [before] = await db + .select({ updatedAt: copilotRuns.updatedAt }) + .from(copilotRuns) + .where(eq(copilotRuns.id, abandoned.runId)) const { settledRunIds } = await sweepOrphanedRuns() expect(settledRunIds).toContain(abandoned.runId) expect(settledRunIds).not.toContain(recent.runId) - expect((await stored(abandoned.runId)).status).toBe('error') + const settled = await stored(abandoned.runId) + expect(settled.status).toBe('error') + /** Its retention clock keeps running from its last real write, and it reads as never finalized. */ + expect(settled.updatedAt).toEqual(before.updatedAt) + expect(settled.completedAt).toEqual(before.updatedAt) + expect(settled.error).toBe(UNLEASED_RUN_ERROR) + expect(settled.error).not.toBe(ORPHANED_RUN_ERROR) expect((await stored(recent.runId)).status).toBe('active') }) @@ -283,4 +319,64 @@ describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => { await redis().del(chatStreamLockKey(owned.chatId)) } }) + + it('never deadlocks a sweep against recovering controllers claiming the same runs', async () => { + for (let attempt = 0; attempt < 30; attempt++) { + const orphans = await Promise.all( + Array.from({ length: 20 }, () => admittedRun({ idleMinutes: 90 })) + ) + + const [sweep, ...claims] = await Promise.all([ + sweepOrphanedRuns(), + /** Staggered so claims land while the sweep's settling transaction holds its locks. */ + ...orphans.map((orphan) => + sleep(randomInt(0, 40)).then(() => + claimRunController({ + runId: orphan.runId, + chatId: orphan.chatId, + previousToken: orphan.controllerToken!, + token: `${orphan.streamId}\n${generateId()}`, + }) + ) + ), + ]) + + orphans.forEach((orphan, index) => { + expect(claims[index] !== sweep.settledRunIds.includes(orphan.runId)).toBe(true) + }) + } + }) + + it('settles a stopped run exactly once when Stop races its own controller finalizing', async () => { + for (let attempt = 0; attempt < 50; attempt++) { + const orphan = await admittedRun() + + const [finalized, stopped] = await Promise.all([ + updateRunStatus(orphan.runId, 'complete', {}, orphan.controllerToken!), + settleStoppedRunWithoutController(orphan.runId), + ]) + + expect(Boolean(finalized) !== stopped).toBe(true) + expect((await stored(orphan.runId)).status).toBe(stopped ? 'cancelled' : 'complete') + } + }) + + it('settles a stopped run exactly once when Stop races a recovering controller claiming it', async () => { + for (let attempt = 0; attempt < 50; attempt++) { + const orphan = await admittedRun() + + const [claimed, stopped] = await Promise.all([ + claimRunController({ + runId: orphan.runId, + chatId: orphan.chatId, + previousToken: orphan.controllerToken!, + token: `${orphan.streamId}\n${generateId()}`, + }), + settleStoppedRunWithoutController(orphan.runId), + ]) + + expect(claimed !== stopped).toBe(true) + expect((await stored(orphan.runId)).status).toBe(stopped ? 'cancelled' : 'active') + } + }) }) diff --git a/apps/sim/lib/mothership/async-runs/orphaned-runs.ts b/apps/sim/lib/mothership/async-runs/orphaned-runs.ts index 5268b6c8c5d..a64006392ec 100644 --- a/apps/sim/lib/mothership/async-runs/orphaned-runs.ts +++ b/apps/sim/lib/mothership/async-runs/orphaned-runs.ts @@ -1,7 +1,14 @@ import { db } from '@sim/db' -import { type CopilotRunStatus, copilotChats, copilotRuns } from '@sim/db/schema' +import { + type CopilotRunStatus, + copilotChats, + copilotOrganizationRequestStops, + copilotRequestStops, + copilotRuns, +} from '@sim/db/schema' import { createLogger } from '@sim/logger' import { getErrorMessage } from '@sim/utils/errors' +import { sleep } from '@sim/utils/helpers' import { and, asc, @@ -16,7 +23,6 @@ import { sql, } from 'drizzle-orm' import { publishChatStatusChanged } from '@/lib/mothership/chat-status' -import { ORCHESTRATION_TIMEOUT_MS } from '@/lib/mothership/constants' import { findStreamsWithReplay } from '@/lib/mothership/request/session/buffer' import { findStreamsHoldingChatLock } from '@/lib/mothership/request/session/controller-lease' @@ -31,11 +37,17 @@ const UNFINISHED_RUN_STATUSES: CopilotRunStatus[] = [ ] /** - * How long a leased run must sit without a controller or a durable write before it is - * settled. No worker leg outlives the orchestration budget, so past it a reconnect has - * nothing left to resume. + * How long a leased run must go without a status write before the sweep may settle it. + * + * This is a recovery window, not a liveness test, and is independent of any run + * deadline. Liveness comes only from the chat lock: a live controller renews it by + * heartbeat for as long as it runs, however long that is, and a run whose stream holds + * the lock is never settled. For a run with no lock holder, this window and the replay + * buffer's `:seq` key (whose TTL each write renews, `COPILOT_STREAM_TTL_SECONDS`, + * one hour by default) leave a reconnect time to resume it; the sweep waits for both. + * A TTL configured below this window shortens only that resume window, never safety. */ -export const ORPHANED_RUN_GRACE_MS = ORCHESTRATION_TIMEOUT_MS +export const ORPHANED_RUN_GRACE_MS = 60 * 60 * 1000 /** * Runs admitted without a chat lease (headless turns and rows from before the lease @@ -45,9 +57,12 @@ export const ORPHANED_RUN_GRACE_MS = ORCHESTRATION_TIMEOUT_MS export const UNLEASED_RUN_GRACE_MS = 24 * 60 * 60 * 1000 export const ORPHANED_RUN_ERROR = 'This response was interrupted before it finished.' +export const UNLEASED_RUN_ERROR = 'Run was never finalized (no controller lease).' const SWEEP_BATCH_SIZE = 500 -const SWEEP_MAX_ROWS_PER_RUN = 10_000 +const SWEEP_MAX_ROWS_PER_RUN = 5_000 +/** Spaces full batches so a backlog drains without a sustained burst of synchronous commits. */ +const SWEEP_BATCH_PAUSE_MS = 200 const controllerToken = sql`${copilotRuns.requestContext}->>'controllerToken'` @@ -57,6 +72,15 @@ function idleFor(ms: number): SQL { const leasedRunIdle = and(isNotNull(controllerToken), idleFor(ORPHANED_RUN_GRACE_MS)) const unleasedRunIdle = and(isNull(controllerToken), idleFor(UNLEASED_RUN_GRACE_MS)) +const orphanIdle = or(leasedRunIdle, unleasedRunIdle) + +/** The user pressed Stop on this stream; a newer turn also closes tool admission, without one. */ +const stopRequested = sql`(EXISTS (SELECT 1 FROM ${copilotRequestStops} s + WHERE s.user_id = ${copilotRuns.userId} AND s.workspace_id = ${copilotRuns.workspaceId} + AND s.stream_id = ${copilotRuns.streamId}) + OR EXISTS (SELECT 1 FROM ${copilotOrganizationRequestStops} s + WHERE s.user_id = ${copilotRuns.userId} AND s.organization_id = ${copilotRuns.organizationId} + AND s.stream_id = ${copilotRuns.streamId}))` interface UnownedRun { id: string @@ -80,30 +104,51 @@ const unownedRunColumns = { type Transaction = Parameters[0]>[0] +/** A stopped run ends cancelled; the sweep ends it cancelled only if its user pressed Stop. */ +function terminalValues(reason: 'stopped' | 'orphaned' | 'unleased') { + if (reason === 'stopped') { + return { status: sql`'cancelled'::copilot_run_status`, error: sql`NULL::text` } + } + const error = reason === 'orphaned' ? ORPHANED_RUN_ERROR : UNLEASED_RUN_ERROR + return { + status: sql`(CASE WHEN ${stopRequested} THEN 'cancelled' ELSE 'error' END)::copilot_run_status`, + error: sql`CASE WHEN ${stopRequested} THEN NULL ELSE ${error}::text END`, + } +} + /** * Settles each run only while it is still unfinished and still names the controller * the caller observed, so a finalizing controller or a successor's claim, both of - * which write the same row, wins or loses atomically against it. A run that Stop - * already closed settles as cancelled. The chat marker is released without touching - * the chat's ordering timestamp. + * which write the same row, wins or loses atomically against it. + * + * Chat rows are locked first, in id order, as a controller's claim does, so the two + * never wait on each other in opposite orders. A run without a lease keeps its last + * write as its completion and retention time. The chat marker is released without + * touching the chat's ordering timestamp. */ async function settleRuns( tx: Transaction, runs: UnownedRun[], - guard: SQL | undefined, - outcome: { status: 'error' | 'cancelled'; error?: string } + reason: 'stopped' | 'orphaned', + guard: SQL | undefined ): Promise { + if (runs.length === 0) return [] + const chatIds = [...new Set(runs.map((run) => run.chatId))] + await tx + .select({ id: copilotChats.id }) + .from(copilotChats) + .where(inArray(copilotChats.id, chatIds)) + .orderBy(asc(copilotChats.id)) + .for('update') + const settled: UnownedRun[] = [] - const apply = async (owner: SQL | undefined, batch: UnownedRun[]) => { + const apply = async (batch: UnownedRun[], owner: SQL, values: object) => { if (batch.length === 0) return const rows = await tx .update(copilotRuns) .set({ - status: sql`(CASE WHEN ${copilotRuns.toolAdmissionClosedAt} IS NOT NULL THEN 'cancelled' ELSE ${outcome.status} END)::copilot_run_status`, - error: sql`CASE WHEN ${copilotRuns.toolAdmissionClosedAt} IS NOT NULL THEN NULL ELSE ${outcome.error ?? null}::text END`, - completedAt: sql`now()`, + ...values, toolAdmissionClosedAt: sql`coalesce(${copilotRuns.toolAdmissionClosedAt}, now())`, - updatedAt: sql`now()`, }) .where( and( @@ -122,18 +167,27 @@ async function settleRuns( } await apply( + runs.filter((run) => run.controllerToken === null), isNull(controllerToken), - runs.filter((run) => run.controllerToken === null) + { ...terminalValues('unleased'), completedAt: sql`${copilotRuns.updatedAt}` } ) for (const run of runs) { - if (run.controllerToken !== null) await apply(eq(controllerToken, run.controllerToken), [run]) + if (run.controllerToken === null) continue + await apply([run], eq(controllerToken, run.controllerToken), { + ...terminalValues(reason), + completedAt: sql`now()`, + updatedAt: sql`now()`, + }) } - for (const run of settled) { + if (settled.length > 0) { + const markers = settled.map((run) => sql`(${run.chatId}::uuid, ${run.streamId}::text)`) await tx .update(copilotChats) .set({ conversationId: null }) - .where(and(eq(copilotChats.id, run.chatId), eq(copilotChats.conversationId, run.streamId))) + .where( + sql`(${copilotChats.id}, ${copilotChats.conversationId}) IN (${sql.join(markers, sql`, `)})` + ) } return settled } @@ -155,57 +209,68 @@ function announceSettled(runs: UnownedRun[]): void { } } +/** The candidates no controller owns; leased runs are skipped when ownership is unreadable. */ +async function withoutOwners(candidates: UnownedRun[]): Promise { + const leased = candidates.filter((run) => run.controllerToken !== null) + const unleased = candidates.filter((run) => run.controllerToken === null) + if (leased.length === 0) return unleased + try { + const [locked, replayable] = await Promise.all([ + findStreamsHoldingChatLock(leased), + findStreamsWithReplay(leased.map((run) => run.streamId)), + ]) + return unleased.concat( + leased.filter((run) => !locked.has(run.streamId) && !replayable.has(run.streamId)) + ) + } catch (error) { + logger.warn('Chat stream ownership is unreadable; leaving leased runs for a later sweep', { + error: getErrorMessage(error), + }) + return unleased + } +} + /** * Settles runs that no controller will ever finish: a leased run whose stream holds no * chat lock and has no replay buffer left, idle past the recovery window, and a run - * without a lease idle past any process lifetime. Unreadable locks skip leased runs. + * without a lease idle past any process lifetime. A failed batch is logged and skipped. */ export async function sweepOrphanedRuns(): Promise<{ settledRunIds: string[] }> { const settledRunIds: string[] = [] - const idle = or(leasedRunIdle, unleasedRunIdle) let cursor: string | undefined let considered = 0 while (considered < SWEEP_MAX_ROWS_PER_RUN) { + const limit = Math.min(SWEEP_BATCH_SIZE, SWEEP_MAX_ROWS_PER_RUN - considered) const candidates: UnownedRun[] = await db .select(unownedRunColumns) .from(copilotRuns) .where( and( inArray(copilotRuns.status, UNFINISHED_RUN_STATUSES), - idle, + orphanIdle, cursor ? gt(copilotRuns.id, cursor) : undefined ) ) .orderBy(asc(copilotRuns.id)) - .limit(Math.min(SWEEP_BATCH_SIZE, SWEEP_MAX_ROWS_PER_RUN - considered)) + .limit(limit) if (candidates.length === 0) break considered += candidates.length cursor = candidates[candidates.length - 1].id - const leased = candidates.filter((run) => run.controllerToken !== null) - let unowned = candidates.filter((run) => run.controllerToken === null) - if (leased.length > 0) { - try { - const [locked, replayable] = await Promise.all([ - findStreamsHoldingChatLock(leased), - findStreamsWithReplay(leased.map((run) => run.streamId)), - ]) - unowned = unowned.concat( - leased.filter((run) => !locked.has(run.streamId) && !replayable.has(run.streamId)) - ) - } catch (error) { - logger.warn('Chat stream ownership is unreadable; leaving leased runs for a later sweep', { - error: getErrorMessage(error), - }) - } + try { + const unowned = await withoutOwners(candidates) + const settled = await db.transaction((tx) => settleRuns(tx, unowned, 'orphaned', orphanIdle)) + announceSettled(settled.filter((run) => run.controllerToken !== null)) + settledRunIds.push(...settled.map((run) => run.id)) + } catch (error) { + logger.warn('A batch of orphaned runs could not be settled; a later sweep retries it', { + count: candidates.length, + error: getErrorMessage(error), + }) } - - const settled = await db.transaction((tx) => - settleRuns(tx, unowned, idle, { status: 'error', error: ORPHANED_RUN_ERROR }) - ) - announceSettled(settled.filter((run) => run.controllerToken !== null)) - settledRunIds.push(...settled.map((run) => run.id)) + if (candidates.length < limit) break + await sleep(SWEEP_BATCH_PAUSE_MS) } if (settledRunIds.length > 0) { @@ -226,9 +291,7 @@ export async function settleStoppedRunWithoutController(runId: string): Promise< .limit(1) if (!run?.controllerToken) return false if ((await findStreamsHoldingChatLock([run])).has(run.streamId)) return false - const settled = await db.transaction((tx) => - settleRuns(tx, [run], undefined, { status: 'cancelled' }) - ) + const settled = await db.transaction((tx) => settleRuns(tx, [run], 'stopped', undefined)) announceSettled(settled) return settled.length > 0 } From f1d878a428cd721da5468e0a64010d21172fd627 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 30 Sep 2026 01:40:34 -0700 Subject: [PATCH 3/7] fix(mothership): never sweep a current headless run A headless turn has no chat lease and no heartbeat, so its age says nothing about whether it is still running once runs have no deadline. The sweep's lease-less rule now applies only to runs admitted before the current tool-execution protocol: every run the current code admits records the current version, so after a deploy no such row can be live. A current headless run is left to its own lifecycle, which always settles it. The protocol version moves beside the other async-run constants so the sweep can read it without importing the repository. --- .../lib/mothership/async-runs/lifecycle.ts | 3 ++ .../async-runs/orphaned-runs.integration.ts | 28 +++++++++--- .../mothership/async-runs/orphaned-runs.ts | 44 +++++++++++-------- .../lib/mothership/async-runs/repository.ts | 2 +- 4 files changed, 52 insertions(+), 25 deletions(-) diff --git a/apps/sim/lib/mothership/async-runs/lifecycle.ts b/apps/sim/lib/mothership/async-runs/lifecycle.ts index 060c5dba3c7..73e1d4f0556 100644 --- a/apps/sim/lib/mothership/async-runs/lifecycle.ts +++ b/apps/sim/lib/mothership/async-runs/lifecycle.ts @@ -4,6 +4,9 @@ import { MothershipStreamV1ToolOutcome, } from '@/lib/mothership/generated/mothership-stream-v1' +/** Recorded on every run the current code admits; older values mark runs from earlier protocols. */ +export const SIM_TOOL_EXECUTION_VERSION = 2 + export const ASYNC_TOOL_STATUS = MothershipStreamV1AsyncToolRecordStatus export const EXECUTABLE_TOOL_PERMISSION_DECISIONS = [ diff --git a/apps/sim/lib/mothership/async-runs/orphaned-runs.integration.ts b/apps/sim/lib/mothership/async-runs/orphaned-runs.integration.ts index 88a7eb6b401..e7b080f1316 100644 --- a/apps/sim/lib/mothership/async-runs/orphaned-runs.integration.ts +++ b/apps/sim/lib/mothership/async-runs/orphaned-runs.integration.ts @@ -43,10 +43,10 @@ import { randomInt } from '@sim/utils/random' import { eq, inArray, sql } from 'drizzle-orm' import { closeRedisConnection, getRedisClient } from '@/lib/core/config/redis' import { + LEGACY_RUN_ERROR, ORPHANED_RUN_ERROR, settleStoppedRunWithoutController, sweepOrphanedRuns, - UNLEASED_RUN_ERROR, } from '@/lib/mothership/async-runs/orphaned-runs' import { updateRunStatus } from '@/lib/mothership/async-runs/repository' import { abortRun } from '@/lib/mothership/request/application/controls' @@ -119,6 +119,8 @@ describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => { controllerToken?: string | null stopped?: boolean superseded?: boolean + /** Admitted by code predating the current tool-execution protocol. */ + legacy?: boolean } = {} ) { const chatId = generateId() @@ -144,7 +146,7 @@ describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => { userId, workspaceId, streamId, - toolExecutionVersion: 2, + toolExecutionVersion: options.legacy ? 0 : 2, status: options.status ?? 'active', requestContext: controllerToken ? { requestId: generateId(), controllerToken, recovery: { kind: 'interactive_stream' } } @@ -201,9 +203,13 @@ describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => { expect(run.error).toBeTruthy() }) - it('settles a run without a lease only after the unleased ceiling', async () => { - const recent = await admittedRun({ idleMinutes: 90, controllerToken: null }) - const abandoned = await admittedRun({ idleMinutes: 25 * 60, controllerToken: null }) + it('settles a legacy run without a lease only after the legacy ceiling', async () => { + const recent = await admittedRun({ idleMinutes: 90, controllerToken: null, legacy: true }) + const abandoned = await admittedRun({ + idleMinutes: 25 * 60, + controllerToken: null, + legacy: true, + }) const [before] = await db .select({ updatedAt: copilotRuns.updatedAt }) .from(copilotRuns) @@ -217,11 +223,21 @@ describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => { /** Its retention clock keeps running from its last real write, and it reads as never finalized. */ expect(settled.updatedAt).toEqual(before.updatedAt) expect(settled.completedAt).toEqual(before.updatedAt) - expect(settled.error).toBe(UNLEASED_RUN_ERROR) + expect(settled.error).toBe(LEGACY_RUN_ERROR) expect(settled.error).not.toBe(ORPHANED_RUN_ERROR) expect((await stored(recent.runId)).status).toBe('active') }) + it('never settles a current headless run, however long it has run', async () => { + /** A headless turn has no lease or heartbeat; only its own lifecycle can end it. */ + const headless = await admittedRun({ idleMinutes: 25 * 60, controllerToken: null }) + + const { settledRunIds } = await sweepOrphanedRuns() + + expect(settledRunIds).not.toContain(headless.runId) + expect((await stored(headless.runId)).status).toBe('active') + }) + it('never settles a run whose stream holds its chat lock, or that is still recoverable', async () => { const leased = await admittedRun({ idleMinutes: 90 }) await redis().set(chatStreamLockKey(leased.chatId), leased.controllerToken!, 'EX', 60) diff --git a/apps/sim/lib/mothership/async-runs/orphaned-runs.ts b/apps/sim/lib/mothership/async-runs/orphaned-runs.ts index a64006392ec..ec46d1e9c93 100644 --- a/apps/sim/lib/mothership/async-runs/orphaned-runs.ts +++ b/apps/sim/lib/mothership/async-runs/orphaned-runs.ts @@ -17,11 +17,13 @@ import { inArray, isNotNull, isNull, + lt, notInArray, or, type SQL, sql, } from 'drizzle-orm' +import { SIM_TOOL_EXECUTION_VERSION } from '@/lib/mothership/async-runs/lifecycle' import { publishChatStatusChanged } from '@/lib/mothership/chat-status' import { findStreamsWithReplay } from '@/lib/mothership/request/session/buffer' import { findStreamsHoldingChatLock } from '@/lib/mothership/request/session/controller-lease' @@ -50,14 +52,16 @@ const UNFINISHED_RUN_STATUSES: CopilotRunStatus[] = [ export const ORPHANED_RUN_GRACE_MS = 60 * 60 * 1000 /** - * Runs admitted without a chat lease (headless turns and rows from before the lease - * protocol) have no liveness signal, so only an age far past any process lifetime - * proves them dead. + * Runs admitted by code predating the current tool-execution protocol. Every run the + * current code admits records the current version, so once a deploy has replaced the + * processes that admitted these, none can be live; the age only leaves room for a + * rollout. A current run without a lease (a headless turn) is never swept: it has no + * liveness signal and its own lifecycle always settles it. */ -export const UNLEASED_RUN_GRACE_MS = 24 * 60 * 60 * 1000 +export const LEGACY_RUN_GRACE_MS = 24 * 60 * 60 * 1000 export const ORPHANED_RUN_ERROR = 'This response was interrupted before it finished.' -export const UNLEASED_RUN_ERROR = 'Run was never finalized (no controller lease).' +export const LEGACY_RUN_ERROR = 'Run was never finalized (pre-lease run).' const SWEEP_BATCH_SIZE = 500 const SWEEP_MAX_ROWS_PER_RUN = 5_000 @@ -71,8 +75,12 @@ function idleFor(ms: number): SQL { } const leasedRunIdle = and(isNotNull(controllerToken), idleFor(ORPHANED_RUN_GRACE_MS)) -const unleasedRunIdle = and(isNull(controllerToken), idleFor(UNLEASED_RUN_GRACE_MS)) -const orphanIdle = or(leasedRunIdle, unleasedRunIdle) +const legacyRunIdle = and( + isNull(controllerToken), + lt(copilotRuns.toolExecutionVersion, SIM_TOOL_EXECUTION_VERSION), + idleFor(LEGACY_RUN_GRACE_MS) +) +const orphanIdle = or(leasedRunIdle, legacyRunIdle) /** The user pressed Stop on this stream; a newer turn also closes tool admission, without one. */ const stopRequested = sql`(EXISTS (SELECT 1 FROM ${copilotRequestStops} s @@ -105,11 +113,11 @@ const unownedRunColumns = { type Transaction = Parameters[0]>[0] /** A stopped run ends cancelled; the sweep ends it cancelled only if its user pressed Stop. */ -function terminalValues(reason: 'stopped' | 'orphaned' | 'unleased') { +function terminalValues(reason: 'stopped' | 'orphaned' | 'legacy') { if (reason === 'stopped') { return { status: sql`'cancelled'::copilot_run_status`, error: sql`NULL::text` } } - const error = reason === 'orphaned' ? ORPHANED_RUN_ERROR : UNLEASED_RUN_ERROR + const error = reason === 'orphaned' ? ORPHANED_RUN_ERROR : LEGACY_RUN_ERROR return { status: sql`(CASE WHEN ${stopRequested} THEN 'cancelled' ELSE 'error' END)::copilot_run_status`, error: sql`CASE WHEN ${stopRequested} THEN NULL ELSE ${error}::text END`, @@ -122,8 +130,8 @@ function terminalValues(reason: 'stopped' | 'orphaned' | 'unleased') { * which write the same row, wins or loses atomically against it. * * Chat rows are locked first, in id order, as a controller's claim does, so the two - * never wait on each other in opposite orders. A run without a lease keeps its last - * write as its completion and retention time. The chat marker is released without + * never wait on each other in opposite orders. A legacy run keeps its last write as its + * completion and retention time. The chat marker is released without * touching the chat's ordering timestamp. */ async function settleRuns( @@ -169,7 +177,7 @@ async function settleRuns( await apply( runs.filter((run) => run.controllerToken === null), isNull(controllerToken), - { ...terminalValues('unleased'), completedAt: sql`${copilotRuns.updatedAt}` } + { ...terminalValues('legacy'), completedAt: sql`${copilotRuns.updatedAt}` } ) for (const run of runs) { if (run.controllerToken === null) continue @@ -212,28 +220,28 @@ function announceSettled(runs: UnownedRun[]): void { /** The candidates no controller owns; leased runs are skipped when ownership is unreadable. */ async function withoutOwners(candidates: UnownedRun[]): Promise { const leased = candidates.filter((run) => run.controllerToken !== null) - const unleased = candidates.filter((run) => run.controllerToken === null) - if (leased.length === 0) return unleased + const legacy = candidates.filter((run) => run.controllerToken === null) + if (leased.length === 0) return legacy try { const [locked, replayable] = await Promise.all([ findStreamsHoldingChatLock(leased), findStreamsWithReplay(leased.map((run) => run.streamId)), ]) - return unleased.concat( + return legacy.concat( leased.filter((run) => !locked.has(run.streamId) && !replayable.has(run.streamId)) ) } catch (error) { logger.warn('Chat stream ownership is unreadable; leaving leased runs for a later sweep', { error: getErrorMessage(error), }) - return unleased + return legacy } } /** * Settles runs that no controller will ever finish: a leased run whose stream holds no - * chat lock and has no replay buffer left, idle past the recovery window, and a run - * without a lease idle past any process lifetime. A failed batch is logged and skipped. + * chat lock and has no replay buffer left, idle past the recovery window, and a legacy + * run from before the current protocol. A failed batch is logged and skipped. */ export async function sweepOrphanedRuns(): Promise<{ settledRunIds: string[] }> { const settledRunIds: string[] = [] diff --git a/apps/sim/lib/mothership/async-runs/repository.ts b/apps/sim/lib/mothership/async-runs/repository.ts index 28cc32397f9..c3d0ca09417 100644 --- a/apps/sim/lib/mothership/async-runs/repository.ts +++ b/apps/sim/lib/mothership/async-runs/repository.ts @@ -42,6 +42,7 @@ import { type AsyncTerminalStatus, DESKTOP_TOOL_CLAIM_OWNER, EXECUTABLE_TOOL_PERMISSION_DECISIONS, + SIM_TOOL_EXECUTION_VERSION, } from '@/lib/mothership/async-runs/lifecycle' import { TraceAttr } from '@/lib/mothership/generated/trace-attributes-v1' import { TraceSpan } from '@/lib/mothership/generated/trace-spans-v1' @@ -55,7 +56,6 @@ import { chatSandboxSessionKey } from '@/lib/mothership/tools/sandbox-session-ke const logger = createLogger('CopilotAsyncRunsRepo') const WORKFLOW_EXECUTION_CLAIM_PREFIX = 'workflow:' -const SIM_TOOL_EXECUTION_VERSION = 2 const TERMINAL_RUN_STATUSES: CopilotRunStatus[] = ['complete', 'error', 'cancelled'] // Resolve the tracer lazily per-call to avoid capturing the NoOp tracer // before NodeSDK installs the global TracerProvider (Next.js 16/Turbopack From fd4fc18187074bfe6ba3e939753a665a3c5cf780 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 30 Sep 2026 01:54:10 -0700 Subject: [PATCH 4/7] refactor(mothership): require the recorded Stop inside the stopped-run settle - Settling a stopped run now passes the Stop-row check as the update's own guard, so it cannot cancel a run nobody stopped; the separate stopped branch is gone and every settle derives cancelled from the Stop row. - Chat lock ownership is read through getChatStreamLockOwners and trusted only when verified, instead of a second Redis read of the same keys. - Settle transactions use the shared DbTransaction type. --- .../async-runs/orphaned-runs.integration.ts | 20 ++++++++- .../mothership/async-runs/orphaned-runs.ts | 41 ++++++++++++------- .../request/session/controller-lease.ts | 20 --------- 3 files changed, 46 insertions(+), 35 deletions(-) diff --git a/apps/sim/lib/mothership/async-runs/orphaned-runs.integration.ts b/apps/sim/lib/mothership/async-runs/orphaned-runs.integration.ts index e7b080f1316..1e325e7ccb5 100644 --- a/apps/sim/lib/mothership/async-runs/orphaned-runs.integration.ts +++ b/apps/sim/lib/mothership/async-runs/orphaned-runs.integration.ts @@ -48,7 +48,7 @@ import { settleStoppedRunWithoutController, sweepOrphanedRuns, } from '@/lib/mothership/async-runs/orphaned-runs' -import { updateRunStatus } from '@/lib/mothership/async-runs/repository' +import { requestRunStop, updateRunStatus } from '@/lib/mothership/async-runs/repository' import { abortRun } from '@/lib/mothership/request/application/controls' import { claimRunController } from '@/lib/mothership/request/lifecycle/controller-ownership' import { chatStreamLockKey } from '@/lib/mothership/request/session/controller-lease' @@ -160,6 +160,11 @@ describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => { return { chatId, streamId, runId, controllerToken } } + /** Records the user's Stop the way the abort use case does before it settles anything. */ + async function stop(run: { streamId: string; chatId: string }) { + await requestRunStop({ userId, workspaceId, streamId: run.streamId, chatId: run.chatId }) + } + async function stored(runId: string) { const [run] = await db.select().from(copilotRuns).where(eq(copilotRuns.id, runId)) const [chat] = await db @@ -322,8 +327,19 @@ describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => { expect(run.marker).toBeNull() }) + it('never cancels a run nobody stopped', async () => { + const orphan = await admittedRun({ superseded: true }) + + expect(await settleStoppedRunWithoutController(orphan.runId)).toBe(false) + + const run = await stored(orphan.runId) + expect(run.status).toBe('active') + expect(run.marker).toBe(orphan.streamId) + }) + it('leaves a stopped run to the controller of its stream that holds the chat lock', async () => { const owned = await admittedRun() + await stop(owned) await redis().set(chatStreamLockKey(owned.chatId), owned.controllerToken!, 'EX', 60) try { @@ -366,6 +382,7 @@ describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => { it('settles a stopped run exactly once when Stop races its own controller finalizing', async () => { for (let attempt = 0; attempt < 50; attempt++) { const orphan = await admittedRun() + await stop(orphan) const [finalized, stopped] = await Promise.all([ updateRunStatus(orphan.runId, 'complete', {}, orphan.controllerToken!), @@ -380,6 +397,7 @@ describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => { it('settles a stopped run exactly once when Stop races a recovering controller claiming it', async () => { for (let attempt = 0; attempt < 50; attempt++) { const orphan = await admittedRun() + await stop(orphan) const [claimed, stopped] = await Promise.all([ claimRunController({ diff --git a/apps/sim/lib/mothership/async-runs/orphaned-runs.ts b/apps/sim/lib/mothership/async-runs/orphaned-runs.ts index ec46d1e9c93..b410c317504 100644 --- a/apps/sim/lib/mothership/async-runs/orphaned-runs.ts +++ b/apps/sim/lib/mothership/async-runs/orphaned-runs.ts @@ -23,10 +23,11 @@ import { type SQL, sql, } from 'drizzle-orm' +import type { DbTransaction } from '@/lib/db/types' import { SIM_TOOL_EXECUTION_VERSION } from '@/lib/mothership/async-runs/lifecycle' import { publishChatStatusChanged } from '@/lib/mothership/chat-status' +import { getChatStreamLockOwners } from '@/lib/mothership/request/session/abort' import { findStreamsWithReplay } from '@/lib/mothership/request/session/buffer' -import { findStreamsHoldingChatLock } from '@/lib/mothership/request/session/controller-lease' const logger = createLogger('OrphanedCopilotRuns') @@ -110,13 +111,8 @@ const unownedRunColumns = { controllerToken, } -type Transaction = Parameters[0]>[0] - -/** A stopped run ends cancelled; the sweep ends it cancelled only if its user pressed Stop. */ -function terminalValues(reason: 'stopped' | 'orphaned' | 'legacy') { - if (reason === 'stopped') { - return { status: sql`'cancelled'::copilot_run_status`, error: sql`NULL::text` } - } +/** A run ends cancelled only if its user pressed Stop, whoever settles it. */ +function terminalValues(reason: 'orphaned' | 'legacy') { const error = reason === 'orphaned' ? ORPHANED_RUN_ERROR : LEGACY_RUN_ERROR return { status: sql`(CASE WHEN ${stopRequested} THEN 'cancelled' ELSE 'error' END)::copilot_run_status`, @@ -135,9 +131,8 @@ function terminalValues(reason: 'stopped' | 'orphaned' | 'legacy') { * touching the chat's ordering timestamp. */ async function settleRuns( - tx: Transaction, + tx: DbTransaction, runs: UnownedRun[], - reason: 'stopped' | 'orphaned', guard: SQL | undefined ): Promise { if (runs.length === 0) return [] @@ -182,7 +177,7 @@ async function settleRuns( for (const run of runs) { if (run.controllerToken === null) continue await apply([run], eq(controllerToken, run.controllerToken), { - ...terminalValues(reason), + ...terminalValues('orphaned'), completedAt: sql`now()`, updatedAt: sql`now()`, }) @@ -217,6 +212,23 @@ function announceSettled(runs: UnownedRun[]): void { } } +/** + * The streams among these whose own controller holds its chat lock, under any token: a + * recovering controller locks the chat before it claims the run. Throws unless the + * locks were read, since otherwise no stream is provably unowned. + */ +async function findStreamsHoldingChatLock( + runs: Array<{ chatId: string; streamId: string }> +): Promise> { + const { status, ownersByChatId } = await getChatStreamLockOwners([ + ...new Set(runs.map((run) => run.chatId)), + ]) + if (status !== 'verified') throw new Error('Chat stream locks are unreadable') + return new Set( + runs.filter((run) => ownersByChatId.get(run.chatId) === run.streamId).map((run) => run.streamId) + ) +} + /** The candidates no controller owns; leased runs are skipped when ownership is unreadable. */ async function withoutOwners(candidates: UnownedRun[]): Promise { const leased = candidates.filter((run) => run.controllerToken !== null) @@ -268,7 +280,7 @@ export async function sweepOrphanedRuns(): Promise<{ settledRunIds: string[] }> try { const unowned = await withoutOwners(candidates) - const settled = await db.transaction((tx) => settleRuns(tx, unowned, 'orphaned', orphanIdle)) + const settled = await db.transaction((tx) => settleRuns(tx, unowned, orphanIdle)) announceSettled(settled.filter((run) => run.controllerToken !== null)) settledRunIds.push(...settled.map((run) => run.id)) } catch (error) { @@ -289,7 +301,8 @@ export async function sweepOrphanedRuns(): Promise<{ settledRunIds: string[] }> /** * Settles a stopped run as cancelled when no controller of its stream holds the chat - * lock. A live controller observes the Stop and settles its own run. + * lock. A live controller observes the Stop and settles its own run. The update itself + * requires the recorded Stop, so this can never settle a run nobody stopped. */ export async function settleStoppedRunWithoutController(runId: string): Promise { const [run] = await db @@ -299,7 +312,7 @@ export async function settleStoppedRunWithoutController(runId: string): Promise< .limit(1) if (!run?.controllerToken) return false if ((await findStreamsHoldingChatLock([run])).has(run.streamId)) return false - const settled = await db.transaction((tx) => settleRuns(tx, [run], 'stopped', undefined)) + const settled = await db.transaction((tx) => settleRuns(tx, [run], stopRequested)) announceSettled(settled) return settled.length > 0 } diff --git a/apps/sim/lib/mothership/request/session/controller-lease.ts b/apps/sim/lib/mothership/request/session/controller-lease.ts index 2519854a3ce..b7d89647a33 100644 --- a/apps/sim/lib/mothership/request/session/controller-lease.ts +++ b/apps/sim/lib/mothership/request/session/controller-lease.ts @@ -28,26 +28,6 @@ export async function assertChatStreamLease(lease: ChatStreamLease): Promise -): Promise> { - const redis = getRedisClient() - if (!redis) throw new Error('Chat stream locks are unreadable without Redis') - if (streams.length === 0) return new Set() - const values = await redis.mget(streams.map(({ chatId }) => chatStreamLockKey(chatId))) - const held = new Set() - streams.forEach(({ streamId }, index) => { - const value = values[index] - if (value && streamIdFromLock(value) === streamId) held.add(streamId) - }) - return held -} - /** Whether this lease still holds its chat lock; an unreadable lock counts as lost. */ export async function holdsChatStreamLease(lease: ChatStreamLease): Promise { try { From 34768f4a32c42a29d0fac653cbee0598f8598c73 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 30 Sep 2026 02:12:04 -0700 Subject: [PATCH 5/7] fix(mothership): fence orphan settlement on the chat lock and resume sweeps where they stopped - The sweep takes each unowned leased run's chat lock under the run's own stream before settling it and releases it after the commit, so a reconnect can no longer lock the chat between the ownership check and the settle and then lose its claim; a reconnect that meets the fence retries. - A sweep examines at most 10k candidates and settles at most about 5k, resuming from a cursor saved in Redis and wrapping to the first run, so runs that cannot be settled yet never starve the ones after them. - Every settled run whose chat marker was released is announced, legacy runs included, so an open client stops showing the chat as busy. --- .../async-runs/orphaned-runs.integration.ts | 142 ++++++++++++++- .../mothership/async-runs/orphaned-runs.ts | 172 ++++++++++++++---- 2 files changed, 274 insertions(+), 40 deletions(-) diff --git a/apps/sim/lib/mothership/async-runs/orphaned-runs.integration.ts b/apps/sim/lib/mothership/async-runs/orphaned-runs.integration.ts index 1e325e7ccb5..97b916e3826 100644 --- a/apps/sim/lib/mothership/async-runs/orphaned-runs.integration.ts +++ b/apps/sim/lib/mothership/async-runs/orphaned-runs.integration.ts @@ -49,9 +49,18 @@ import { sweepOrphanedRuns, } from '@/lib/mothership/async-runs/orphaned-runs' import { requestRunStop, updateRunStatus } from '@/lib/mothership/async-runs/repository' +import { chatPubSub } from '@/lib/mothership/chat-status' import { abortRun } from '@/lib/mothership/request/application/controls' import { claimRunController } from '@/lib/mothership/request/lifecycle/controller-ownership' -import { chatStreamLockKey } from '@/lib/mothership/request/session/controller-lease' +import { + acquirePendingChatStream, + getLocalChatStreamLease, + releasePendingChatStream, +} from '@/lib/mothership/request/session/abort' +import { + assertChatStreamLease, + chatStreamLockKey, +} from '@/lib/mothership/request/session/controller-lease' function redis() { const client = getRedisClient() @@ -60,6 +69,7 @@ function redis() { } afterAll(async () => { + chatPubSub?.dispose() await closeRedisConnection() await new Promise((resolve) => worker.server.close(() => resolve())) for (const [key, value] of Object.entries(inheritedEnv)) { @@ -121,11 +131,12 @@ describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => { superseded?: boolean /** Admitted by code predating the current tool-execution protocol. */ legacy?: boolean + id?: string } = {} ) { const chatId = generateId() const streamId = generateId() - const runId = generateId() + const runId = options.id ?? generateId() chatIds.push(chatId) const controllerToken = options.controllerToken === undefined @@ -413,4 +424,131 @@ describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => { expect((await stored(orphan.runId)).status).toBe(stopped ? 'cancelled' : 'active') } }) + + it('never takes a run from a reconnect that locked its chat while the sweep was settling', async () => { + for (let attempt = 0; attempt < 20; attempt++) { + const orphans = await Promise.all( + Array.from({ length: 20 }, () => admittedRun({ idleMinutes: 90 })) + ) + + /** + * Each reconnect locks the chat, proves its lease, then claims the run, as recovery + * does. A run still unfinished once its reconnect holds the lock belongs to it. + */ + const reconnect = async (orphan: (typeof orphans)[number]) => { + await sleep(randomInt(0, 40)) + if (!(await acquirePendingChatStream(orphan.chatId, orphan.streamId, 0))) { + return { owned: false, claimed: false } + } + const lease = getLocalChatStreamLease(orphan.chatId, orphan.streamId)! + try { + await assertChatStreamLease(lease) + const owned = (await stored(orphan.runId)).status === 'active' + await sleep(randomInt(0, 10)) + const claimed = await claimRunController({ + runId: orphan.runId, + chatId: orphan.chatId, + previousToken: orphan.controllerToken!, + token: lease.value, + }) + return { owned, claimed } + } finally { + await releasePendingChatStream(orphan.chatId, orphan.streamId, lease) + } + } + const [sweep, ...reconnects] = await Promise.all([ + sweepOrphanedRuns(), + ...orphans.map(reconnect), + ]) + + orphans.forEach((orphan, index) => { + const swept = sweep.settledRunIds.includes(orphan.runId) + if (reconnects[index].owned) expect(reconnects[index].claimed).toBe(true) + expect(reconnects[index].claimed !== swept).toBe(true) + }) + } + }) + + it('announces every settled run whose chat it released, legacy runs included', async () => { + const legacy = await admittedRun({ + idleMinutes: 25 * 60, + controllerToken: null, + legacy: true, + }) + const announced: string[] = [] + const unsubscribe = chatPubSub!.onStatusChanged((event) => { + if (event.type === 'completed') announced.push(event.chatId) + }) + + try { + const { settledRunIds } = await sweepOrphanedRuns() + expect(settledRunIds).toContain(legacy.runId) + for (let wait = 0; wait < 50 && !announced.includes(legacy.chatId); wait++) await sleep(20) + expect(announced).toContain(legacy.chatId) + } finally { + unsubscribe() + } + }) + + it('reaches an orphan behind more unsettleable runs than one sweep examines', async () => { + /** Runs whose replay is still live, all sorting before the orphan. */ + const blockers = Array.from({ length: 10_500 }, (_, index) => ({ + runId: `00000000-0000-4000-8000-${index.toString(16).padStart(12, '0')}`, + chatId: generateId(), + streamId: generateId(), + })) + const orphan = await admittedRun({ + idleMinutes: 90, + id: 'ffffffff-ffff-4fff-bfff-ffffffffffff', + }) + const blockerChatIds = blockers.map((blocker) => blocker.chatId) + try { + for (let start = 0; start < blockers.length; start += 1000) { + const page = blockers.slice(start, start + 1000) + await db.insert(copilotChats).values( + page.map((blocker) => ({ + id: blocker.chatId, + userId, + workspaceId, + type: 'mothership' as const, + })) + ) + await db.insert(copilotRuns).values( + page.map((blocker) => ({ + id: blocker.runId, + executionId: generateId(), + chatId: blocker.chatId, + userId, + workspaceId, + streamId: blocker.streamId, + toolExecutionVersion: 2, + status: 'active' as const, + requestContext: { controllerToken: `${blocker.streamId}\n${generateId()}` }, + startedAt: sql`now() - interval '2 hours'`, + updatedAt: sql`now() - interval '2 hours'`, + })) + ) + const pipeline = redis().pipeline() + for (const blocker of page) { + pipeline.set(`mothership_stream:${blocker.streamId}:seq`, '1', 'EX', 600) + } + await pipeline.exec() + } + + const first = await sweepOrphanedRuns() + const second = first.settledRunIds.includes(orphan.runId) ? first : await sweepOrphanedRuns() + + expect(second.settledRunIds).toContain(orphan.runId) + expect((await stored(orphan.runId)).status).toBe('error') + } finally { + for (let start = 0; start < blockerChatIds.length; start += 1000) { + await db + .delete(copilotChats) + .where(inArray(copilotChats.id, blockerChatIds.slice(start, start + 1000))) + } + const pipeline = redis().pipeline() + for (const blocker of blockers) pipeline.del(`mothership_stream:${blocker.streamId}:seq`) + await pipeline.exec() + } + }, 120_000) }) diff --git a/apps/sim/lib/mothership/async-runs/orphaned-runs.ts b/apps/sim/lib/mothership/async-runs/orphaned-runs.ts index b410c317504..f656084f646 100644 --- a/apps/sim/lib/mothership/async-runs/orphaned-runs.ts +++ b/apps/sim/lib/mothership/async-runs/orphaned-runs.ts @@ -18,16 +18,24 @@ import { isNotNull, isNull, lt, + lte, notInArray, or, type SQL, sql, } from 'drizzle-orm' +import { getRedisClient } from '@/lib/core/config/redis' import type { DbTransaction } from '@/lib/db/types' import { SIM_TOOL_EXECUTION_VERSION } from '@/lib/mothership/async-runs/lifecycle' import { publishChatStatusChanged } from '@/lib/mothership/chat-status' -import { getChatStreamLockOwners } from '@/lib/mothership/request/session/abort' +import { + acquirePendingChatStream, + getChatStreamLockOwners, + getLocalChatStreamLease, + releasePendingChatStream, +} from '@/lib/mothership/request/session/abort' import { findStreamsWithReplay } from '@/lib/mothership/request/session/buffer' +import type { ChatStreamLease } from '@/lib/mothership/request/session/controller-lease' const logger = createLogger('OrphanedCopilotRuns') @@ -65,7 +73,16 @@ export const ORPHANED_RUN_ERROR = 'This response was interrupted before it finis export const LEGACY_RUN_ERROR = 'Run was never finalized (pre-lease run).' const SWEEP_BATCH_SIZE = 500 -const SWEEP_MAX_ROWS_PER_RUN = 5_000 +/** Settles at most about this many runs per sweep, to bound its synchronous commits. */ +const SWEEP_MAX_SETTLED_PER_RUN = 5_000 +/** Examines at most this many candidates per sweep; the next sweep resumes after them. */ +const SWEEP_MAX_EXAMINED_PER_RUN = 10_000 +/** + * Where the last sweep stopped, so runs that cannot be settled yet (their chat is locked + * or their replay is live) never starve the runs after them. It wraps to the start. + */ +const SWEEP_CURSOR_KEY = 'copilot:orphaned-runs:sweep-cursor' +const SWEEP_CURSOR_TTL_SECONDS = 7 * 24 * 60 * 60 /** Spaces full batches so a backlog drains without a sustained burst of synchronous commits. */ const SWEEP_BATCH_PAUSE_MS = 200 @@ -134,8 +151,8 @@ async function settleRuns( tx: DbTransaction, runs: UnownedRun[], guard: SQL | undefined -): Promise { - if (runs.length === 0) return [] +): Promise<{ settled: UnownedRun[]; released: UnownedRun[] }> { + if (runs.length === 0) return { settled: [], released: [] } const chatIds = [...new Set(runs.map((run) => run.chatId))] await tx .select({ id: copilotChats.id }) @@ -183,19 +200,21 @@ async function settleRuns( }) } - if (settled.length > 0) { - const markers = settled.map((run) => sql`(${run.chatId}::uuid, ${run.streamId}::text)`) - await tx - .update(copilotChats) - .set({ conversationId: null }) - .where( - sql`(${copilotChats.id}, ${copilotChats.conversationId}) IN (${sql.join(markers, sql`, `)})` - ) - } - return settled + if (settled.length === 0) return { settled, released: [] } + const markers = settled.map((run) => sql`(${run.chatId}::uuid, ${run.streamId}::text)`) + const cleared = await tx + .update(copilotChats) + .set({ conversationId: null }) + .where( + sql`(${copilotChats.id}, ${copilotChats.conversationId}) IN (${sql.join(markers, sql`, `)})` + ) + .returning({ id: copilotChats.id }) + const releasedChats = new Set(cleared.map((chat) => chat.id)) + return { settled, released: settled.filter((run) => releasedChats.has(run.chatId)) } } -function announceSettled(runs: UnownedRun[]): void { +/** Tells open clients a chat is no longer busy, for every chat whose marker was released. */ +function announceReleased(runs: UnownedRun[]): void { for (const run of runs) { try { publishChatStatusChanged(run, { @@ -250,18 +269,90 @@ async function withoutOwners(candidates: UnownedRun[]): Promise { } } +interface ChatLockFence { + run: UnownedRun + lease: ChatStreamLease +} + +/** + * Takes each unowned leased run's chat lock under the run's own stream, as a reconnect + * would, so no controller can take over between the ownership check and the settle. + * A reconnect that meets the fence retries; runs whose lock is taken are skipped. + */ +async function fenceChatLocks(runs: UnownedRun[]): Promise { + const fenced = await Promise.all( + runs.map(async (run) => { + if (!(await acquirePendingChatStream(run.chatId, run.streamId, 0))) return null + const lease = getLocalChatStreamLease(run.chatId, run.streamId) + return lease ? { run, lease } : null + }) + ) + return fenced.filter((fence): fence is ChatLockFence => fence !== null) +} + +async function releaseChatLocks(fences: ChatLockFence[]): Promise { + await Promise.all( + fences.map(({ run, lease }) => releasePendingChatStream(run.chatId, run.streamId, lease)) + ) +} + +async function readSweepCursor(): Promise { + try { + return (await getRedisClient()?.get(SWEEP_CURSOR_KEY)) ?? undefined + } catch (error) { + logger.warn('Orphaned-run sweep cursor is unreadable; starting from the first run', { + error: getErrorMessage(error), + }) + return undefined + } +} + +async function writeSweepCursor(cursor: string | undefined): Promise { + try { + const redis = getRedisClient() + if (!redis) return + if (cursor) await redis.set(SWEEP_CURSOR_KEY, cursor, 'EX', SWEEP_CURSOR_TTL_SECONDS) + else await redis.del(SWEEP_CURSOR_KEY) + } catch (error) { + logger.warn('Orphaned-run sweep cursor could not be saved', { error: getErrorMessage(error) }) + } +} + +/** Settles one examined batch, fencing leased runs on their chat locks while it commits. */ +async function settleBatch(candidates: UnownedRun[]): Promise { + const unowned = await withoutOwners(candidates) + const fences = await fenceChatLocks(unowned.filter((run) => run.controllerToken !== null)) + try { + const eligible = unowned + .filter((run) => run.controllerToken === null) + .concat(fences.map(({ run }) => run)) + const { settled, released } = await db.transaction((tx) => settleRuns(tx, eligible, orphanIdle)) + announceReleased(released) + return settled + } finally { + await releaseChatLocks(fences) + } +} + /** * Settles runs that no controller will ever finish: a leased run whose stream holds no * chat lock and has no replay buffer left, idle past the recovery window, and a legacy - * run from before the current protocol. A failed batch is logged and skipped. + * run from before the current protocol. Each sweep resumes where the last one stopped + * and wraps to the first run, so no run is starved by the ones before it. A failed + * batch is logged and skipped. */ export async function sweepOrphanedRuns(): Promise<{ settledRunIds: string[] }> { const settledRunIds: string[] = [] - let cursor: string | undefined - let considered = 0 + const start = await readSweepCursor() + let cursor = start + let wrapped = start === undefined + let examined = 0 - while (considered < SWEEP_MAX_ROWS_PER_RUN) { - const limit = Math.min(SWEEP_BATCH_SIZE, SWEEP_MAX_ROWS_PER_RUN - considered) + while ( + examined < SWEEP_MAX_EXAMINED_PER_RUN && + settledRunIds.length < SWEEP_MAX_SETTLED_PER_RUN + ) { + const limit = Math.min(SWEEP_BATCH_SIZE, SWEEP_MAX_EXAMINED_PER_RUN - examined) const candidates: UnownedRun[] = await db .select(unownedRunColumns) .from(copilotRuns) @@ -269,30 +360,35 @@ export async function sweepOrphanedRuns(): Promise<{ settledRunIds: string[] }> and( inArray(copilotRuns.status, UNFINISHED_RUN_STATUSES), orphanIdle, - cursor ? gt(copilotRuns.id, cursor) : undefined + cursor ? gt(copilotRuns.id, cursor) : undefined, + wrapped && start ? lte(copilotRuns.id, start) : undefined ) ) .orderBy(asc(copilotRuns.id)) .limit(limit) - if (candidates.length === 0) break - considered += candidates.length - cursor = candidates[candidates.length - 1].id - - try { - const unowned = await withoutOwners(candidates) - const settled = await db.transaction((tx) => settleRuns(tx, unowned, orphanIdle)) - announceSettled(settled.filter((run) => run.controllerToken !== null)) - settledRunIds.push(...settled.map((run) => run.id)) - } catch (error) { - logger.warn('A batch of orphaned runs could not be settled; a later sweep retries it', { - count: candidates.length, - error: getErrorMessage(error), - }) + examined += candidates.length + if (candidates.length > 0) { + cursor = candidates[candidates.length - 1].id + try { + settledRunIds.push(...(await settleBatch(candidates)).map((run) => run.id)) + } catch (error) { + logger.warn('A batch of orphaned runs could not be settled; a later sweep retries it', { + count: candidates.length, + error: getErrorMessage(error), + }) + } + } + if (candidates.length < limit) { + /** The end of the table: wrap once to cover the runs before the starting point. */ + cursor = undefined + if (wrapped) break + wrapped = true + continue } - if (candidates.length < limit) break await sleep(SWEEP_BATCH_PAUSE_MS) } + await writeSweepCursor(cursor) if (settledRunIds.length > 0) { logger.info('Settled runs no controller owned', { count: settledRunIds.length }) } @@ -312,7 +408,7 @@ export async function settleStoppedRunWithoutController(runId: string): Promise< .limit(1) if (!run?.controllerToken) return false if ((await findStreamsHoldingChatLock([run])).has(run.streamId)) return false - const settled = await db.transaction((tx) => settleRuns(tx, [run], stopRequested)) - announceSettled(settled) + const { settled, released } = await db.transaction((tx) => settleRuns(tx, [run], stopRequested)) + announceReleased(released) return settled.length > 0 } From 7c023cf9a274b79a236f502146b5ef1be7e110f9 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 30 Sep 2026 02:18:02 -0700 Subject: [PATCH 6/7] fix(knowledge): record a terminal status for every Slack Assistant run The Slack Assistant admits its own run row and only ever marked it as an error, so every completed or stopped Slack turn stayed active. It now records the terminal status once, after the turn ends, through the shared run-status update: complete on success, cancelled when its user stopped it (in Slack or in Sim), and error otherwise. Every other run-creating path already settles its run: interactive turns through their controller's finalize, and headless turns that admit their own run in the lifecycle's own finally. --- .../slack-search/assistant.integration.ts | 240 ++++++++++++++++++ .../application/slack-search/assistant.ts | 16 +- 2 files changed, 255 insertions(+), 1 deletion(-) create mode 100644 apps/sim/lib/knowledge/application/slack-search/assistant.integration.ts diff --git a/apps/sim/lib/knowledge/application/slack-search/assistant.integration.ts b/apps/sim/lib/knowledge/application/slack-search/assistant.integration.ts new file mode 100644 index 00000000000..f455315a0a6 --- /dev/null +++ b/apps/sim/lib/knowledge/application/slack-search/assistant.integration.ts @@ -0,0 +1,240 @@ +/** + * The Slack Assistant's run record against real PostgreSQL: the run row it admits and + * the terminal status it records are production code. Slack delivery, the worker + * lifecycle, identity, and chat locking are stood in, since only the record's outcome + * is under test. + */ +import { billingAttributionMock } from '@sim/testing/mocks/billing-attribution.mock' +import { + mothershipChatPayloadMock, + mothershipChatPayloadMockFns, +} from '@sim/testing/mocks/mothership-chat-payload.mock' +import { mothershipEnvironmentContextMock } from '@sim/testing/mocks/mothership-environment-context.mock' +import { organizationAuthorizationMock } from '@sim/testing/mocks/organization-authorization.mock' +import { afterAll, beforeAll, beforeEach, describe, expect, it, vi } from 'vitest' + +const hoisted = vi.hoisted(() => ({ + chat: { id: '', model: 'default' }, + turnId: '', + userId: '', + organizationId: '', + lifecycle: vi.fn(), + /** The turn's own controller, which a Stop aborts through its registered stream. */ + controller: new AbortController(), + stopped: vi.fn(async () => false), +})) +vi.mock('@/lib/knowledge/application/slack-search/authorization', () => ({ + authorizeSlackSearchInstallation: async () => ({ + installation: { id: 'i1', organizationId: hoisted.organizationId, teamId: 'T1' }, + secret: { botToken: 'token' }, + }), +})) +vi.mock('@/lib/internal/slack/search-client', () => ({ + getSlackSearchSender: async () => ({ email: 'member@example.com' }), +})) +vi.mock('@/lib/knowledge/application/slack-search/identity', () => ({ + resolveSlackSearchMember: async () => hoisted.userId, + SlackSearchIdentityError: class extends Error {}, +})) +vi.mock('@/lib/knowledge/application/slack-search/chat', () => ({ + resolveSlackSearchChat: async () => hoisted.chat, + persistSlackSearchQuestion: async () => undefined, + slackSearchChatOperation: { id: 'organization.chats.slack' }, +})) +vi.mock('@/lib/core/application/organization-authorization', () => organizationAuthorizationMock) +vi.mock('@/lib/knowledge/application/operations', () => ({ + knowledgeOperations: { search: { organizationOperation: { id: 'knowledge.search' } } }, +})) +vi.mock('@/lib/knowledge/application/slack-search/repository', () => ({ + recordSlackSearchOutcome: async () => undefined, +})) +vi.mock('@/lib/knowledge/application/slack-search/turns', () => ({ + requireSlackSearchTurnLease: async () => undefined, + wasSlackSearchTurnStopped: hoisted.stopped, +})) +vi.mock('@/lib/knowledge/application/slack-search/onboarding', () => ({ + sendSlackSearchOnboarding: vi.fn(), +})) +vi.mock('@/lib/knowledge/application/slack-search/title', () => ({ + generateSlackSearchChatTitle: async () => undefined, +})) +vi.mock('@/lib/billing/core/billing-attribution', () => billingAttributionMock) +vi.mock('@/lib/mothership/application/load-search-integrations', () => ({ + loadCopilotSearchIntegrations: async () => '{"connections":[],"available":[]}', +})) +vi.mock('@/lib/mothership/chat/payload', () => mothershipChatPayloadMock) +vi.mock('@/lib/mothership/chat/terminal-state', () => ({ + finalizeAssistantTurn: async () => ({ appendedAssistant: true }), +})) +vi.mock('@/lib/mothership/environment-context', () => mothershipEnvironmentContextMock) +vi.mock('@/lib/mothership/request/lifecycle/headless', () => ({ + runHeadlessCopilotLifecycle: hoisted.lifecycle, +})) +vi.mock('@/lib/mothership/request/session/abort', () => ({ + acquirePendingChatStream: async () => true, + cleanupAbortMarker: async () => undefined, + getChatStreamLockOwners: async () => ({ + status: 'verified', + ownersByChatId: new Map([[hoisted.chat.id, hoisted.turnId]]), + }), + registerActiveStream: vi.fn(), + releasePendingChatStream: async () => undefined, + startAbortPoller: () => 0, + unregisterActiveStream: vi.fn(), +})) +vi.mock('@/executor/utils/resolved-secret-content-projection', () => ({ + projectResolvedSecretDiagnosticContent: (value: unknown) => ({ safe: true, value }), +})) +vi.mock('@/lib/slack-search/connections', () => ({ deliverSlackSearchConnections: vi.fn() })) +vi.mock('@/lib/slack-search/assistant-stream', () => ({ + SlackSearchAssistantStream: class { + start = async () => undefined + finish = async () => undefined + finishWithError = async () => undefined + terminateAfterFailure = async () => undefined + onEvent = async () => undefined + assertHealthy = () => undefined + }, +})) + +import { db } from '@sim/db' +import { copilotChats, copilotRuns, organization, user } from '@sim/db/schema' +import { generateId } from '@sim/utils/id' +import { eq } from 'drizzle-orm' +import { runSlackSearchAssistant } from '@/lib/knowledge/application/slack-search/assistant' +import { AbortReason } from '@/lib/mothership/request/session/abort-reason' + +const principal = { + kind: 'slack_installation', + credentialId: 'c1', + credentialVersion: 'v1', + appId: 'A1', + teamId: 'T1', + eventId: 'Ev1', + receivedAt: new Date(), +} as const + +function job() { + return { + installationId: 'i1', + revision: 'r1', + credentialId: 'c1', + credentialVersion: 'v1', + receivedAt: Date.now(), + message: { + appId: 'A1', + teamId: 'T1', + eventId: 'Ev1', + channelId: 'D1', + userId: 'U1', + messageTs: '1800000000.000001', + query: 'release notes', + queryTooLong: false, + }, + } +} + +/** Runs one Slack turn in a fresh private chat and returns the run it recorded. */ +async function slackTurn() { + const chatId = generateId() + const turnId = generateId() + hoisted.chat = { id: chatId, model: 'default' } + hoisted.turnId = turnId + await db.insert(copilotChats).values({ + id: chatId, + userId: hoisted.userId, + organizationId: hoisted.organizationId, + type: 'mothership', + }) + const controller = new AbortController() + hoisted.controller = controller + const outcome = await runSlackSearchAssistant(principal, { + job: job(), + turnId, + leaseId: generateId(), + controller, + }).then( + () => undefined, + (error: unknown) => error + ) + const [run] = await db.select().from(copilotRuns).where(eq(copilotRuns.chatId, chatId)) + return { run, outcome } +} + +describe('Slack Assistant run record', () => { + beforeAll(async () => { + hoisted.userId = generateId() + hoisted.organizationId = generateId() + const now = new Date() + await db.insert(user).values({ + id: hoisted.userId, + name: 'Slack Assistant fixture', + email: `${hoisted.userId}@slack-assistant.test`, + emailVerified: true, + createdAt: now, + updatedAt: now, + }) + await db.insert(organization).values({ + id: hoisted.organizationId, + name: 'Slack Assistant fixture', + slug: `slack-assistant-${hoisted.organizationId}`, + }) + mothershipChatPayloadMockFns.mockBuildCopilotRequestPayload.mockResolvedValue({ + mode: 'assistant', + }) + }) + + afterAll(async () => { + await db.delete(copilotChats).where(eq(copilotChats.userId, hoisted.userId)) + await db.delete(organization).where(eq(organization.id, hoisted.organizationId)) + await db.delete(user).where(eq(user.id, hoisted.userId)) + }) + + beforeEach(() => { + hoisted.stopped.mockResolvedValue(false) + }) + + it('records a completed turn as complete', async () => { + hoisted.lifecycle.mockResolvedValueOnce({ + success: true, + content: 'Answer', + contentBlocks: [], + toolCalls: [], + }) + + const { run, outcome } = await slackTurn() + + expect(outcome).toBeUndefined() + expect(run.status).toBe('complete') + expect(run.completedAt).not.toBeNull() + }) + + it('records a turn its user stopped as cancelled', async () => { + hoisted.lifecycle.mockImplementationOnce(async () => { + /** A Slack Stop marks the turn stopped, then aborts its registered stream. */ + hoisted.stopped.mockResolvedValue(true) + hoisted.controller.abort(AbortReason.UserStop) + return { success: false, cancelled: true, content: '', contentBlocks: [], toolCalls: [] } + }) + + const { run } = await slackTurn() + + expect(run.status).toBe('cancelled') + expect(run.completedAt).not.toBeNull() + }) + + it('records a failed turn as an error', async () => { + hoisted.lifecycle.mockResolvedValueOnce({ + success: false, + error: 'worker failed', + content: '', + contentBlocks: [], + toolCalls: [], + }) + + const { run, outcome } = await slackTurn() + + expect(outcome).toBeInstanceOf(Error) + expect(run.status).toBe('error') + }) +}) diff --git a/apps/sim/lib/knowledge/application/slack-search/assistant.ts b/apps/sim/lib/knowledge/application/slack-search/assistant.ts index 5e3b1039c2d..2381846278d 100644 --- a/apps/sim/lib/knowledge/application/slack-search/assistant.ts +++ b/apps/sim/lib/knowledge/application/slack-search/assistant.ts @@ -49,6 +49,7 @@ import { startAbortPoller, unregisterActiveStream, } from '@/lib/mothership/request/session/abort' +import { isExplicitStopReason } from '@/lib/mothership/request/session/abort-reason' import type { OrchestratorResult } from '@/lib/mothership/request/types' import { organizationRoutes } from '@/lib/navigation/paths' import { SlackSearchAssistantStream } from '@/lib/slack-search/assistant-stream' @@ -303,7 +304,6 @@ export async function runSlackSearchAssistant( ] : []), recordSlackSearchOutcome(installation, 'assistant_or_delivery_failed'), - ...(runId ? [updateRunStatus(runId, 'error')] : []), ]) const errors = outcomes.flatMap((outcome) => outcome.status === 'rejected' ? [outcome.reason] : [] @@ -317,6 +317,20 @@ export async function runSlackSearchAssistant( await titleTask clearInterval(accessPoller) clearInterval(abortPoller) + try { + if (runId) { + /** This turn admitted its own run, so it records the terminal status no other path will. */ + const cancelled = + failed && + (isExplicitStopReason(controller.signal.reason) || + (await wasSlackSearchTurnStopped(turnId, leaseId))) + await updateRunStatus(runId, failed ? (cancelled ? 'cancelled' : 'error') : 'complete') + } + } catch (error) { + failure = failure + ? new AggregateError([failure, error], 'Slack turn run status could not be recorded') + : toError(error) + } try { if (questionPersisted) { const stopped = failed && (await wasSlackSearchTurnStopped(turnId, leaseId)) From 3ff5a3f8d5d70c747d7d5bc21ca0932ae09e6eda Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 30 Sep 2026 02:43:03 -0700 Subject: [PATCH 7/7] fix(knowledge): record the Slack run's status after its turn is saved - The Slack Assistant now writes its run's one terminal status after its outcome and response are persisted, from the final outcome, so a failed save ends the run as an error instead of complete. - A Stop lookup that fails no longer skips that write: the turn is treated as not stopped, logged, and settled as an error. - The orphaned-run suite deletes the sweep cursor before each test and in teardown, so no later suite starts from its leftover position. --- .../slack-search/assistant.integration.ts | 51 ++++++++++++++++++- .../application/slack-search/assistant.ts | 41 +++++++++------ .../async-runs/orphaned-runs.integration.ts | 10 +++- 3 files changed, 84 insertions(+), 18 deletions(-) diff --git a/apps/sim/lib/knowledge/application/slack-search/assistant.integration.ts b/apps/sim/lib/knowledge/application/slack-search/assistant.integration.ts index f455315a0a6..0923e5b5bbb 100644 --- a/apps/sim/lib/knowledge/application/slack-search/assistant.integration.ts +++ b/apps/sim/lib/knowledge/application/slack-search/assistant.integration.ts @@ -22,6 +22,8 @@ const hoisted = vi.hoisted(() => ({ /** The turn's own controller, which a Stop aborts through its registered stream. */ controller: new AbortController(), stopped: vi.fn(async () => false), + outcome: vi.fn(async () => undefined), + finalize: vi.fn(async () => ({ appendedAssistant: true })), })) vi.mock('@/lib/knowledge/application/slack-search/authorization', () => ({ authorizeSlackSearchInstallation: async () => ({ @@ -46,7 +48,7 @@ vi.mock('@/lib/knowledge/application/operations', () => ({ knowledgeOperations: { search: { organizationOperation: { id: 'knowledge.search' } } }, })) vi.mock('@/lib/knowledge/application/slack-search/repository', () => ({ - recordSlackSearchOutcome: async () => undefined, + recordSlackSearchOutcome: hoisted.outcome, })) vi.mock('@/lib/knowledge/application/slack-search/turns', () => ({ requireSlackSearchTurnLease: async () => undefined, @@ -64,7 +66,7 @@ vi.mock('@/lib/mothership/application/load-search-integrations', () => ({ })) vi.mock('@/lib/mothership/chat/payload', () => mothershipChatPayloadMock) vi.mock('@/lib/mothership/chat/terminal-state', () => ({ - finalizeAssistantTurn: async () => ({ appendedAssistant: true }), + finalizeAssistantTurn: hoisted.finalize, })) vi.mock('@/lib/mothership/environment-context', () => mothershipEnvironmentContextMock) vi.mock('@/lib/mothership/request/lifecycle/headless', () => ({ @@ -192,8 +194,17 @@ describe('Slack Assistant run record', () => { beforeEach(() => { hoisted.stopped.mockResolvedValue(false) + hoisted.outcome.mockResolvedValue(undefined) + hoisted.finalize.mockResolvedValue({ appendedAssistant: true }) }) + const answered = { + success: true, + content: 'Answer', + contentBlocks: [], + toolCalls: [], + } + it('records a completed turn as complete', async () => { hoisted.lifecycle.mockResolvedValueOnce({ success: true, @@ -237,4 +248,40 @@ describe('Slack Assistant run record', () => { expect(outcome).toBeInstanceOf(Error) expect(run.status).toBe('error') }) + + it('records a failed turn as an error even when its Stop cannot be looked up', async () => { + hoisted.stopped.mockRejectedValue(new Error('database unavailable')) + hoisted.lifecycle.mockResolvedValueOnce({ + success: false, + error: 'worker failed', + content: '', + contentBlocks: [], + toolCalls: [], + }) + + const { run, outcome } = await slackTurn() + + expect(outcome).toBeInstanceOf(Error) + expect(run.status).toBe('error') + }) + + it('records an answered turn as an error when its response is not saved', async () => { + hoisted.lifecycle.mockResolvedValueOnce(answered) + hoisted.finalize.mockResolvedValue({ appendedAssistant: false }) + + const { run, outcome } = await slackTurn() + + expect(outcome).toBeInstanceOf(Error) + expect(run.status).toBe('error') + }) + + it('records an answered turn as an error when its outcome is not saved', async () => { + hoisted.lifecycle.mockResolvedValueOnce(answered) + hoisted.outcome.mockRejectedValueOnce(new Error('outcome write failed')) + + const { run, outcome } = await slackTurn() + + expect(outcome).toBeInstanceOf(Error) + expect(run.status).toBe('error') + }) }) diff --git a/apps/sim/lib/knowledge/application/slack-search/assistant.ts b/apps/sim/lib/knowledge/application/slack-search/assistant.ts index 2381846278d..0906f3ea506 100644 --- a/apps/sim/lib/knowledge/application/slack-search/assistant.ts +++ b/apps/sim/lib/knowledge/application/slack-search/assistant.ts @@ -3,7 +3,7 @@ import type { SlackInstallationPrincipal, } from '@sim/auth/principal' import { createLogger } from '@sim/logger' -import { toError } from '@sim/utils/errors' +import { getErrorMessage, toError } from '@sim/utils/errors' import { generateId } from '@sim/utils/id' import { isRecordLike } from '@sim/utils/object' import { resolveOrganizationBillingAttribution } from '@/lib/billing/core/billing-attribution' @@ -317,20 +317,6 @@ export async function runSlackSearchAssistant( await titleTask clearInterval(accessPoller) clearInterval(abortPoller) - try { - if (runId) { - /** This turn admitted its own run, so it records the terminal status no other path will. */ - const cancelled = - failed && - (isExplicitStopReason(controller.signal.reason) || - (await wasSlackSearchTurnStopped(turnId, leaseId))) - await updateRunStatus(runId, failed ? (cancelled ? 'cancelled' : 'error') : 'complete') - } - } catch (error) { - failure = failure - ? new AggregateError([failure, error], 'Slack turn run status could not be recorded') - : toError(error) - } try { if (questionPersisted) { const stopped = failed && (await wasSlackSearchTurnStopped(turnId, leaseId)) @@ -392,6 +378,31 @@ export async function runSlackSearchAssistant( ? new AggregateError([failure, error], 'Slack turn and history persistence failed') : toError(error) } finally { + if (runId) { + /** + * This turn admitted its own run, so it records the one terminal status no other + * path will, after its outcome and response were saved: any failure, including + * a failed save, ends it as an error unless its user stopped it. + */ + let cancelled = failed && isExplicitStopReason(controller.signal.reason) + if (failed && !cancelled) { + try { + cancelled = await wasSlackSearchTurnStopped(turnId, leaseId) + } catch (error) { + logger.warn('Slack turn Stop could not be read; recording its run as an error', { + turnId, + error: getErrorMessage(error), + }) + } + } + try { + await updateRunStatus(runId, !failure ? 'complete' : cancelled ? 'cancelled' : 'error') + } catch (error) { + failure = failure + ? new AggregateError([failure, error], 'Slack turn run status could not be recorded') + : toError(error) + } + } unregisterActiveStream(messageId) await releasePendingChatStream(chat.id, messageId) await cleanupAbortMarker(messageId) diff --git a/apps/sim/lib/mothership/async-runs/orphaned-runs.integration.ts b/apps/sim/lib/mothership/async-runs/orphaned-runs.integration.ts index 97b916e3826..dae25a4dfa6 100644 --- a/apps/sim/lib/mothership/async-runs/orphaned-runs.integration.ts +++ b/apps/sim/lib/mothership/async-runs/orphaned-runs.integration.ts @@ -4,7 +4,7 @@ * use case are production code. A local HTTP server stands in for the worker's abort * endpoint. */ -import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest' +import { afterAll, beforeAll, beforeEach, describe, expect, it, vi } from 'vitest' const { redisUrl, inheritedEnv, worker } = await vi.hoisted(async () => { const { readTestRedisUrl } = await import('@sim/db/testing/test-infrastructure') @@ -68,8 +68,16 @@ function redis() { return client } +/** The sweep's resume point lives in shared Redis; each test and the next suite start fresh. */ +async function resetSweepCursor() { + await getRedisClient()?.del('copilot:orphaned-runs:sweep-cursor') +} + +beforeEach(resetSweepCursor) + afterAll(async () => { chatPubSub?.dispose() + await resetSweepCursor() await closeRedisConnection() await new Promise((resolve) => worker.server.close(() => resolve())) for (const [key, value] of Object.entries(inheritedEnv)) {