Skip to content

Commit cebe8a3

Browse files
authored
fix(mothership): settle chat runs no controller owns (#8457)
* 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. * 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. * 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. * 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. * 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. * 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. * 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.
1 parent e09970a commit cebe8a3

9 files changed

Lines changed: 1338 additions & 4 deletions

File tree

‎apps/sim/app/api/cron/cleanup-stale-executions/route.ts‎

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,7 @@ import {
3232
STALE_SWEEPABLE_EXECUTION_STATUSES,
3333
type StaleSweepableExecutionStatus,
3434
} from '@/lib/logs/types'
35+
import { sweepOrphanedRuns } from '@/lib/mothership/async-runs/orphaned-runs'
3536
import { cancelStaleDispatches } from '@/lib/table/dispatcher'
3637
import { deleteFile } from '@/lib/uploads/core/storage-service'
3738
import {
@@ -738,6 +739,20 @@ export const GET = withRouteHandler(async (request: NextRequest) => {
738739
})
739740
}
740741

742+
/**
743+
* Settle Chat runs no controller will finish: their process died, their
744+
* controller was superseded without a successor, or Stop found none. Without
745+
* this they stay unfinished forever and keep their chat marked as busy.
746+
*/
747+
let orphanedRunsSettled = 0
748+
try {
749+
orphanedRunsSettled = (await sweepOrphanedRuns()).settledRunIds.length
750+
} catch (error) {
751+
logger.error('Failed to settle orphaned Chat runs:', {
752+
error: toError(error).message,
753+
})
754+
}
755+
741756
return NextResponse.json({
742757
success: true,
743758
executions: {
@@ -768,6 +783,9 @@ export const GET = withRouteHandler(async (request: NextRequest) => {
768783
pruned: deploymentOperationsPruned,
769784
retentionDays: DEPLOYMENT_OPERATION_RETENTION_DAYS,
770785
},
786+
chatRuns: {
787+
orphanedSettled: orphanedRunsSettled,
788+
},
771789
})
772790
} catch (error) {
773791
logger.error('Error in stale execution cleanup job:', error)
Lines changed: 287 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,287 @@
1+
/**
2+
* The Slack Assistant's run record against real PostgreSQL: the run row it admits and
3+
* the terminal status it records are production code. Slack delivery, the worker
4+
* lifecycle, identity, and chat locking are stood in, since only the record's outcome
5+
* is under test.
6+
*/
7+
import { billingAttributionMock } from '@sim/testing/mocks/billing-attribution.mock'
8+
import {
9+
mothershipChatPayloadMock,
10+
mothershipChatPayloadMockFns,
11+
} from '@sim/testing/mocks/mothership-chat-payload.mock'
12+
import { mothershipEnvironmentContextMock } from '@sim/testing/mocks/mothership-environment-context.mock'
13+
import { organizationAuthorizationMock } from '@sim/testing/mocks/organization-authorization.mock'
14+
import { afterAll, beforeAll, beforeEach, describe, expect, it, vi } from 'vitest'
15+
16+
const hoisted = vi.hoisted(() => ({
17+
chat: { id: '', model: 'default' },
18+
turnId: '',
19+
userId: '',
20+
organizationId: '',
21+
lifecycle: vi.fn(),
22+
/** The turn's own controller, which a Stop aborts through its registered stream. */
23+
controller: new AbortController(),
24+
stopped: vi.fn(async () => false),
25+
outcome: vi.fn(async () => undefined),
26+
finalize: vi.fn(async () => ({ appendedAssistant: true })),
27+
}))
28+
vi.mock('@/lib/knowledge/application/slack-search/authorization', () => ({
29+
authorizeSlackSearchInstallation: async () => ({
30+
installation: { id: 'i1', organizationId: hoisted.organizationId, teamId: 'T1' },
31+
secret: { botToken: 'token' },
32+
}),
33+
}))
34+
vi.mock('@/lib/internal/slack/search-client', () => ({
35+
getSlackSearchSender: async () => ({ email: 'member@example.com' }),
36+
}))
37+
vi.mock('@/lib/knowledge/application/slack-search/identity', () => ({
38+
resolveSlackSearchMember: async () => hoisted.userId,
39+
SlackSearchIdentityError: class extends Error {},
40+
}))
41+
vi.mock('@/lib/knowledge/application/slack-search/chat', () => ({
42+
resolveSlackSearchChat: async () => hoisted.chat,
43+
persistSlackSearchQuestion: async () => undefined,
44+
slackSearchChatOperation: { id: 'organization.chats.slack' },
45+
}))
46+
vi.mock('@/lib/core/application/organization-authorization', () => organizationAuthorizationMock)
47+
vi.mock('@/lib/knowledge/application/operations', () => ({
48+
knowledgeOperations: { search: { organizationOperation: { id: 'knowledge.search' } } },
49+
}))
50+
vi.mock('@/lib/knowledge/application/slack-search/repository', () => ({
51+
recordSlackSearchOutcome: hoisted.outcome,
52+
}))
53+
vi.mock('@/lib/knowledge/application/slack-search/turns', () => ({
54+
requireSlackSearchTurnLease: async () => undefined,
55+
wasSlackSearchTurnStopped: hoisted.stopped,
56+
}))
57+
vi.mock('@/lib/knowledge/application/slack-search/onboarding', () => ({
58+
sendSlackSearchOnboarding: vi.fn(),
59+
}))
60+
vi.mock('@/lib/knowledge/application/slack-search/title', () => ({
61+
generateSlackSearchChatTitle: async () => undefined,
62+
}))
63+
vi.mock('@/lib/billing/core/billing-attribution', () => billingAttributionMock)
64+
vi.mock('@/lib/mothership/application/load-search-integrations', () => ({
65+
loadCopilotSearchIntegrations: async () => '{"connections":[],"available":[]}',
66+
}))
67+
vi.mock('@/lib/mothership/chat/payload', () => mothershipChatPayloadMock)
68+
vi.mock('@/lib/mothership/chat/terminal-state', () => ({
69+
finalizeAssistantTurn: hoisted.finalize,
70+
}))
71+
vi.mock('@/lib/mothership/environment-context', () => mothershipEnvironmentContextMock)
72+
vi.mock('@/lib/mothership/request/lifecycle/headless', () => ({
73+
runHeadlessCopilotLifecycle: hoisted.lifecycle,
74+
}))
75+
vi.mock('@/lib/mothership/request/session/abort', () => ({
76+
acquirePendingChatStream: async () => true,
77+
cleanupAbortMarker: async () => undefined,
78+
getChatStreamLockOwners: async () => ({
79+
status: 'verified',
80+
ownersByChatId: new Map([[hoisted.chat.id, hoisted.turnId]]),
81+
}),
82+
registerActiveStream: vi.fn(),
83+
releasePendingChatStream: async () => undefined,
84+
startAbortPoller: () => 0,
85+
unregisterActiveStream: vi.fn(),
86+
}))
87+
vi.mock('@/executor/utils/resolved-secret-content-projection', () => ({
88+
projectResolvedSecretDiagnosticContent: (value: unknown) => ({ safe: true, value }),
89+
}))
90+
vi.mock('@/lib/slack-search/connections', () => ({ deliverSlackSearchConnections: vi.fn() }))
91+
vi.mock('@/lib/slack-search/assistant-stream', () => ({
92+
SlackSearchAssistantStream: class {
93+
start = async () => undefined
94+
finish = async () => undefined
95+
finishWithError = async () => undefined
96+
terminateAfterFailure = async () => undefined
97+
onEvent = async () => undefined
98+
assertHealthy = () => undefined
99+
},
100+
}))
101+
102+
import { db } from '@sim/db'
103+
import { copilotChats, copilotRuns, organization, user } from '@sim/db/schema'
104+
import { generateId } from '@sim/utils/id'
105+
import { eq } from 'drizzle-orm'
106+
import { runSlackSearchAssistant } from '@/lib/knowledge/application/slack-search/assistant'
107+
import { AbortReason } from '@/lib/mothership/request/session/abort-reason'
108+
109+
const principal = {
110+
kind: 'slack_installation',
111+
credentialId: 'c1',
112+
credentialVersion: 'v1',
113+
appId: 'A1',
114+
teamId: 'T1',
115+
eventId: 'Ev1',
116+
receivedAt: new Date(),
117+
} as const
118+
119+
function job() {
120+
return {
121+
installationId: 'i1',
122+
revision: 'r1',
123+
credentialId: 'c1',
124+
credentialVersion: 'v1',
125+
receivedAt: Date.now(),
126+
message: {
127+
appId: 'A1',
128+
teamId: 'T1',
129+
eventId: 'Ev1',
130+
channelId: 'D1',
131+
userId: 'U1',
132+
messageTs: '1800000000.000001',
133+
query: 'release notes',
134+
queryTooLong: false,
135+
},
136+
}
137+
}
138+
139+
/** Runs one Slack turn in a fresh private chat and returns the run it recorded. */
140+
async function slackTurn() {
141+
const chatId = generateId()
142+
const turnId = generateId()
143+
hoisted.chat = { id: chatId, model: 'default' }
144+
hoisted.turnId = turnId
145+
await db.insert(copilotChats).values({
146+
id: chatId,
147+
userId: hoisted.userId,
148+
organizationId: hoisted.organizationId,
149+
type: 'mothership',
150+
})
151+
const controller = new AbortController()
152+
hoisted.controller = controller
153+
const outcome = await runSlackSearchAssistant(principal, {
154+
job: job(),
155+
turnId,
156+
leaseId: generateId(),
157+
controller,
158+
}).then(
159+
() => undefined,
160+
(error: unknown) => error
161+
)
162+
const [run] = await db.select().from(copilotRuns).where(eq(copilotRuns.chatId, chatId))
163+
return { run, outcome }
164+
}
165+
166+
describe('Slack Assistant run record', () => {
167+
beforeAll(async () => {
168+
hoisted.userId = generateId()
169+
hoisted.organizationId = generateId()
170+
const now = new Date()
171+
await db.insert(user).values({
172+
id: hoisted.userId,
173+
name: 'Slack Assistant fixture',
174+
email: `${hoisted.userId}@slack-assistant.test`,
175+
emailVerified: true,
176+
createdAt: now,
177+
updatedAt: now,
178+
})
179+
await db.insert(organization).values({
180+
id: hoisted.organizationId,
181+
name: 'Slack Assistant fixture',
182+
slug: `slack-assistant-${hoisted.organizationId}`,
183+
})
184+
mothershipChatPayloadMockFns.mockBuildCopilotRequestPayload.mockResolvedValue({
185+
mode: 'assistant',
186+
})
187+
})
188+
189+
afterAll(async () => {
190+
await db.delete(copilotChats).where(eq(copilotChats.userId, hoisted.userId))
191+
await db.delete(organization).where(eq(organization.id, hoisted.organizationId))
192+
await db.delete(user).where(eq(user.id, hoisted.userId))
193+
})
194+
195+
beforeEach(() => {
196+
hoisted.stopped.mockResolvedValue(false)
197+
hoisted.outcome.mockResolvedValue(undefined)
198+
hoisted.finalize.mockResolvedValue({ appendedAssistant: true })
199+
})
200+
201+
const answered = {
202+
success: true,
203+
content: 'Answer',
204+
contentBlocks: [],
205+
toolCalls: [],
206+
}
207+
208+
it('records a completed turn as complete', async () => {
209+
hoisted.lifecycle.mockResolvedValueOnce({
210+
success: true,
211+
content: 'Answer',
212+
contentBlocks: [],
213+
toolCalls: [],
214+
})
215+
216+
const { run, outcome } = await slackTurn()
217+
218+
expect(outcome).toBeUndefined()
219+
expect(run.status).toBe('complete')
220+
expect(run.completedAt).not.toBeNull()
221+
})
222+
223+
it('records a turn its user stopped as cancelled', async () => {
224+
hoisted.lifecycle.mockImplementationOnce(async () => {
225+
/** A Slack Stop marks the turn stopped, then aborts its registered stream. */
226+
hoisted.stopped.mockResolvedValue(true)
227+
hoisted.controller.abort(AbortReason.UserStop)
228+
return { success: false, cancelled: true, content: '', contentBlocks: [], toolCalls: [] }
229+
})
230+
231+
const { run } = await slackTurn()
232+
233+
expect(run.status).toBe('cancelled')
234+
expect(run.completedAt).not.toBeNull()
235+
})
236+
237+
it('records a failed turn as an error', async () => {
238+
hoisted.lifecycle.mockResolvedValueOnce({
239+
success: false,
240+
error: 'worker failed',
241+
content: '',
242+
contentBlocks: [],
243+
toolCalls: [],
244+
})
245+
246+
const { run, outcome } = await slackTurn()
247+
248+
expect(outcome).toBeInstanceOf(Error)
249+
expect(run.status).toBe('error')
250+
})
251+
252+
it('records a failed turn as an error even when its Stop cannot be looked up', async () => {
253+
hoisted.stopped.mockRejectedValue(new Error('database unavailable'))
254+
hoisted.lifecycle.mockResolvedValueOnce({
255+
success: false,
256+
error: 'worker failed',
257+
content: '',
258+
contentBlocks: [],
259+
toolCalls: [],
260+
})
261+
262+
const { run, outcome } = await slackTurn()
263+
264+
expect(outcome).toBeInstanceOf(Error)
265+
expect(run.status).toBe('error')
266+
})
267+
268+
it('records an answered turn as an error when its response is not saved', async () => {
269+
hoisted.lifecycle.mockResolvedValueOnce(answered)
270+
hoisted.finalize.mockResolvedValue({ appendedAssistant: false })
271+
272+
const { run, outcome } = await slackTurn()
273+
274+
expect(outcome).toBeInstanceOf(Error)
275+
expect(run.status).toBe('error')
276+
})
277+
278+
it('records an answered turn as an error when its outcome is not saved', async () => {
279+
hoisted.lifecycle.mockResolvedValueOnce(answered)
280+
hoisted.outcome.mockRejectedValueOnce(new Error('outcome write failed'))
281+
282+
const { run, outcome } = await slackTurn()
283+
284+
expect(outcome).toBeInstanceOf(Error)
285+
expect(run.status).toBe('error')
286+
})
287+
})

‎apps/sim/lib/knowledge/application/slack-search/assistant.ts‎

Lines changed: 27 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -3,7 +3,7 @@ import type {
33
SlackInstallationPrincipal,
44
} from '@sim/auth/principal'
55
import { createLogger } from '@sim/logger'
6-
import { toError } from '@sim/utils/errors'
6+
import { getErrorMessage, toError } from '@sim/utils/errors'
77
import { generateId } from '@sim/utils/id'
88
import { isRecordLike } from '@sim/utils/object'
99
import { resolveOrganizationBillingAttribution } from '@/lib/billing/core/billing-attribution'
@@ -49,6 +49,7 @@ import {
4949
startAbortPoller,
5050
unregisterActiveStream,
5151
} from '@/lib/mothership/request/session/abort'
52+
import { isExplicitStopReason } from '@/lib/mothership/request/session/abort-reason'
5253
import type { OrchestratorResult } from '@/lib/mothership/request/types'
5354
import { organizationRoutes } from '@/lib/navigation/paths'
5455
import { SlackSearchAssistantStream } from '@/lib/slack-search/assistant-stream'
@@ -303,7 +304,6 @@ export async function runSlackSearchAssistant(
303304
]
304305
: []),
305306
recordSlackSearchOutcome(installation, 'assistant_or_delivery_failed'),
306-
...(runId ? [updateRunStatus(runId, 'error')] : []),
307307
])
308308
const errors = outcomes.flatMap((outcome) =>
309309
outcome.status === 'rejected' ? [outcome.reason] : []
@@ -378,6 +378,31 @@ export async function runSlackSearchAssistant(
378378
? new AggregateError([failure, error], 'Slack turn and history persistence failed')
379379
: toError(error)
380380
} finally {
381+
if (runId) {
382+
/**
383+
* This turn admitted its own run, so it records the one terminal status no other
384+
* path will, after its outcome and response were saved: any failure, including
385+
* a failed save, ends it as an error unless its user stopped it.
386+
*/
387+
let cancelled = failed && isExplicitStopReason(controller.signal.reason)
388+
if (failed && !cancelled) {
389+
try {
390+
cancelled = await wasSlackSearchTurnStopped(turnId, leaseId)
391+
} catch (error) {
392+
logger.warn('Slack turn Stop could not be read; recording its run as an error', {
393+
turnId,
394+
error: getErrorMessage(error),
395+
})
396+
}
397+
}
398+
try {
399+
await updateRunStatus(runId, !failure ? 'complete' : cancelled ? 'cancelled' : 'error')
400+
} catch (error) {
401+
failure = failure
402+
? new AggregateError([failure, error], 'Slack turn run status could not be recorded')
403+
: toError(error)
404+
}
405+
}
381406
unregisterActiveStream(messageId)
382407
await releasePendingChatStream(chat.id, messageId)
383408
await cleanupAbortMarker(messageId)

‎apps/sim/lib/mothership/async-runs/lifecycle.ts‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,9 @@ import {
44
MothershipStreamV1ToolOutcome,
55
} from '@/lib/mothership/generated/mothership-stream-v1'
66

7+
/** Recorded on every run the current code admits; older values mark runs from earlier protocols. */
8+
export const SIM_TOOL_EXECUTION_VERSION = 2
9+
710
export const ASYNC_TOOL_STATUS = MothershipStreamV1AsyncToolRecordStatus
811

912
export const EXECUTABLE_TOOL_PERMISSION_DECISIONS = [

0 commit comments

Comments
 (0)