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/knowledge/application/slack-search/assistant.integration.ts b/apps/sim/lib/knowledge/application/slack-search/assistant.integration.ts new file mode 100644 index 00000000000..0923e5b5bbb --- /dev/null +++ b/apps/sim/lib/knowledge/application/slack-search/assistant.integration.ts @@ -0,0 +1,287 @@ +/** + * 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), + outcome: vi.fn(async () => undefined), + finalize: vi.fn(async () => ({ appendedAssistant: true })), +})) +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: hoisted.outcome, +})) +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: hoisted.finalize, +})) +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) + 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, + 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') + }) + + 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 5e3b1039c2d..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' @@ -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] : [] @@ -378,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/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 new file mode 100644 index 00000000000..dae25a4dfa6 --- /dev/null +++ b/apps/sim/lib/mothership/async-runs/orphaned-runs.integration.ts @@ -0,0 +1,562 @@ +/** + * 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, beforeEach, 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, + 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 { + LEGACY_RUN_ERROR, + ORPHANED_RUN_ERROR, + settleStoppedRunWithoutController, + 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 { + acquirePendingChatStream, + getLocalChatStreamLease, + releasePendingChatStream, +} from '@/lib/mothership/request/session/abort' +import { + assertChatStreamLease, + 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 +} + +/** 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)) { + 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. + * `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: { + idleMinutes?: number + status?: 'active' | 'paused_waiting_for_tool' + controllerToken?: string | null + stopped?: boolean + superseded?: boolean + /** Admitted by code predating the current tool-execution protocol. */ + legacy?: boolean + id?: string + } = {} + ) { + const chatId = generateId() + const streamId = generateId() + const runId = options.id ?? 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: options.legacy ? 0 : 2, + status: options.status ?? 'active', + requestContext: controllerToken + ? { requestId: generateId(), controllerToken, recovery: { kind: 'interactive_stream' } } + : { source: 'headless_lifecycle' }, + startedAt: idle, + updatedAt: idle, + ...(options.stopped || options.superseded ? { toolAdmissionClosedAt: idle } : {}), + }) + if (options.stopped) + await db.insert(copilotRequestStops).values({ userId, workspaceId, streamId }) + 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 + .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 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 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) + .where(eq(copilotRuns.id, abandoned.runId)) + const { settledRunIds } = await sweepOrphanedRuns() + + expect(settledRunIds).toContain(abandoned.runId) + expect(settledRunIds).not.toContain(recent.runId) + 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(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) + /** 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('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 { + 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)) + } + }) + + 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() + await stop(orphan) + + 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() + await stop(orphan) + + 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') + } + }) + + 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 new file mode 100644 index 00000000000..f656084f646 --- /dev/null +++ b/apps/sim/lib/mothership/async-runs/orphaned-runs.ts @@ -0,0 +1,414 @@ +import { db } from '@sim/db' +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, + eq, + gt, + inArray, + 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 { + 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') + +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 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 = 60 * 60 * 1000 + +/** + * 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 LEGACY_RUN_GRACE_MS = 24 * 60 * 60 * 1000 + +export const ORPHANED_RUN_ERROR = 'This response was interrupted before it finished.' +export const LEGACY_RUN_ERROR = 'Run was never finalized (pre-lease run).' + +const SWEEP_BATCH_SIZE = 500 +/** 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 + +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 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 + 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 + 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, +} + +/** 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`, + 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. + * + * 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 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( + tx: DbTransaction, + runs: UnownedRun[], + guard: SQL | undefined +): 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 }) + .from(copilotChats) + .where(inArray(copilotChats.id, chatIds)) + .orderBy(asc(copilotChats.id)) + .for('update') + + const settled: UnownedRun[] = [] + const apply = async (batch: UnownedRun[], owner: SQL, values: object) => { + if (batch.length === 0) return + const rows = await tx + .update(copilotRuns) + .set({ + ...values, + toolAdmissionClosedAt: sql`coalesce(${copilotRuns.toolAdmissionClosedAt}, 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( + runs.filter((run) => run.controllerToken === null), + isNull(controllerToken), + { ...terminalValues('legacy'), completedAt: sql`${copilotRuns.updatedAt}` } + ) + for (const run of runs) { + if (run.controllerToken === null) continue + await apply([run], eq(controllerToken, run.controllerToken), { + ...terminalValues('orphaned'), + completedAt: sql`now()`, + updatedAt: sql`now()`, + }) + } + + 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)) } +} + +/** 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, { + 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), + }) + } + } +} + +/** + * 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) + 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 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 legacy + } +} + +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. 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[] = [] + const start = await readSweepCursor() + let cursor = start + let wrapped = start === undefined + let examined = 0 + + 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) + .where( + and( + inArray(copilotRuns.status, UNFINISHED_RUN_STATUSES), + orphanIdle, + cursor ? gt(copilotRuns.id, cursor) : undefined, + wrapped && start ? lte(copilotRuns.id, start) : undefined + ) + ) + .orderBy(asc(copilotRuns.id)) + .limit(limit) + 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 + } + await sleep(SWEEP_BATCH_PAUSE_MS) + } + + await writeSweepCursor(cursor) + 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. 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 + .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, released } = await db.transaction((tx) => settleRuns(tx, [run], stopRequested)) + announceReleased(released) + return settled.length > 0 +} 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 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) => {