Skip to content

Commit bc9369c

Browse files
committed
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.
1 parent d78d860 commit bc9369c

2 files changed

Lines changed: 213 additions & 54 deletions

File tree

‎apps/sim/lib/mothership/async-runs/orphaned-runs.integration.ts‎

Lines changed: 100 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -29,13 +29,24 @@ const { redisUrl, inheritedEnv, worker } = await vi.hoisted(async () => {
2929
})
3030

3131
import { db } from '@sim/db'
32-
import { copilotChats, copilotRuns, permissions, user, workspace } from '@sim/db/schema'
32+
import {
33+
copilotChats,
34+
copilotRequestStops,
35+
copilotRuns,
36+
permissions,
37+
user,
38+
workspace,
39+
} from '@sim/db/schema'
40+
import { sleep } from '@sim/utils/helpers'
3341
import { generateId } from '@sim/utils/id'
42+
import { randomInt } from '@sim/utils/random'
3443
import { eq, inArray, sql } from 'drizzle-orm'
3544
import { closeRedisConnection, getRedisClient } from '@/lib/core/config/redis'
3645
import {
46+
ORPHANED_RUN_ERROR,
3747
settleStoppedRunWithoutController,
3848
sweepOrphanedRuns,
49+
UNLEASED_RUN_ERROR,
3950
} from '@/lib/mothership/async-runs/orphaned-runs'
4051
import { updateRunStatus } from '@/lib/mothership/async-runs/repository'
4152
import { abortRun } from '@/lib/mothership/request/application/controls'
@@ -98,13 +109,16 @@ describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => {
98109
* A run as the chat POST admits it: the chat marker names its stream and the run
99110
* records the lock value its first controller held. `idleMinutes` backdates its
100111
* last durable write; `controllerToken: null` is a run with no lease protocol.
112+
* `stopped` records the user's Stop intent; `superseded` only closes tool admission,
113+
* as a newer turn's workbench does to older runs.
101114
*/
102115
async function admittedRun(
103116
options: {
104117
idleMinutes?: number
105118
status?: 'active' | 'paused_waiting_for_tool'
106119
controllerToken?: string | null
107120
stopped?: boolean
121+
superseded?: boolean
108122
} = {}
109123
) {
110124
const chatId = generateId()
@@ -137,8 +151,10 @@ describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => {
137151
: { source: 'headless_lifecycle' },
138152
startedAt: idle,
139153
updatedAt: idle,
140-
...(options.stopped ? { toolAdmissionClosedAt: idle } : {}),
154+
...(options.stopped || options.superseded ? { toolAdmissionClosedAt: idle } : {}),
141155
})
156+
if (options.stopped)
157+
await db.insert(copilotRequestStops).values({ userId, workspaceId, streamId })
142158
return { chatId, streamId, runId, controllerToken }
143159
}
144160

@@ -174,15 +190,35 @@ describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => {
174190
expect((await stored(orphan.runId)).status).toBe('cancelled')
175191
})
176192

193+
it('settles a run a newer turn superseded, without a Stop, as an error', async () => {
194+
const orphan = await admittedRun({ idleMinutes: 90, superseded: true })
195+
196+
const { settledRunIds } = await sweepOrphanedRuns()
197+
198+
expect(settledRunIds).toContain(orphan.runId)
199+
const run = await stored(orphan.runId)
200+
expect(run.status).toBe('error')
201+
expect(run.error).toBeTruthy()
202+
})
203+
177204
it('settles a run without a lease only after the unleased ceiling', async () => {
178205
const recent = await admittedRun({ idleMinutes: 90, controllerToken: null })
179206
const abandoned = await admittedRun({ idleMinutes: 25 * 60, controllerToken: null })
180-
207+
const [before] = await db
208+
.select({ updatedAt: copilotRuns.updatedAt })
209+
.from(copilotRuns)
210+
.where(eq(copilotRuns.id, abandoned.runId))
181211
const { settledRunIds } = await sweepOrphanedRuns()
182212

183213
expect(settledRunIds).toContain(abandoned.runId)
184214
expect(settledRunIds).not.toContain(recent.runId)
185-
expect((await stored(abandoned.runId)).status).toBe('error')
215+
const settled = await stored(abandoned.runId)
216+
expect(settled.status).toBe('error')
217+
/** Its retention clock keeps running from its last real write, and it reads as never finalized. */
218+
expect(settled.updatedAt).toEqual(before.updatedAt)
219+
expect(settled.completedAt).toEqual(before.updatedAt)
220+
expect(settled.error).toBe(UNLEASED_RUN_ERROR)
221+
expect(settled.error).not.toBe(ORPHANED_RUN_ERROR)
186222
expect((await stored(recent.runId)).status).toBe('active')
187223
})
188224

@@ -283,4 +319,64 @@ describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => {
283319
await redis().del(chatStreamLockKey(owned.chatId))
284320
}
285321
})
322+
323+
it('never deadlocks a sweep against recovering controllers claiming the same runs', async () => {
324+
for (let attempt = 0; attempt < 30; attempt++) {
325+
const orphans = await Promise.all(
326+
Array.from({ length: 20 }, () => admittedRun({ idleMinutes: 90 }))
327+
)
328+
329+
const [sweep, ...claims] = await Promise.all([
330+
sweepOrphanedRuns(),
331+
/** Staggered so claims land while the sweep's settling transaction holds its locks. */
332+
...orphans.map((orphan) =>
333+
sleep(randomInt(0, 40)).then(() =>
334+
claimRunController({
335+
runId: orphan.runId,
336+
chatId: orphan.chatId,
337+
previousToken: orphan.controllerToken!,
338+
token: `${orphan.streamId}\n${generateId()}`,
339+
})
340+
)
341+
),
342+
])
343+
344+
orphans.forEach((orphan, index) => {
345+
expect(claims[index] !== sweep.settledRunIds.includes(orphan.runId)).toBe(true)
346+
})
347+
}
348+
})
349+
350+
it('settles a stopped run exactly once when Stop races its own controller finalizing', async () => {
351+
for (let attempt = 0; attempt < 50; attempt++) {
352+
const orphan = await admittedRun()
353+
354+
const [finalized, stopped] = await Promise.all([
355+
updateRunStatus(orphan.runId, 'complete', {}, orphan.controllerToken!),
356+
settleStoppedRunWithoutController(orphan.runId),
357+
])
358+
359+
expect(Boolean(finalized) !== stopped).toBe(true)
360+
expect((await stored(orphan.runId)).status).toBe(stopped ? 'cancelled' : 'complete')
361+
}
362+
})
363+
364+
it('settles a stopped run exactly once when Stop races a recovering controller claiming it', async () => {
365+
for (let attempt = 0; attempt < 50; attempt++) {
366+
const orphan = await admittedRun()
367+
368+
const [claimed, stopped] = await Promise.all([
369+
claimRunController({
370+
runId: orphan.runId,
371+
chatId: orphan.chatId,
372+
previousToken: orphan.controllerToken!,
373+
token: `${orphan.streamId}\n${generateId()}`,
374+
}),
375+
settleStoppedRunWithoutController(orphan.runId),
376+
])
377+
378+
expect(claimed !== stopped).toBe(true)
379+
expect((await stored(orphan.runId)).status).toBe(stopped ? 'cancelled' : 'active')
380+
}
381+
})
286382
})

0 commit comments

Comments
 (0)