Skip to content

Commit 34768f4

Browse files
committed
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.
1 parent fd4fc18 commit 34768f4

2 files changed

Lines changed: 274 additions & 40 deletions

File tree

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

Lines changed: 140 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -49,9 +49,18 @@ import {
4949
sweepOrphanedRuns,
5050
} from '@/lib/mothership/async-runs/orphaned-runs'
5151
import { requestRunStop, updateRunStatus } from '@/lib/mothership/async-runs/repository'
52+
import { chatPubSub } from '@/lib/mothership/chat-status'
5253
import { abortRun } from '@/lib/mothership/request/application/controls'
5354
import { claimRunController } from '@/lib/mothership/request/lifecycle/controller-ownership'
54-
import { chatStreamLockKey } from '@/lib/mothership/request/session/controller-lease'
55+
import {
56+
acquirePendingChatStream,
57+
getLocalChatStreamLease,
58+
releasePendingChatStream,
59+
} from '@/lib/mothership/request/session/abort'
60+
import {
61+
assertChatStreamLease,
62+
chatStreamLockKey,
63+
} from '@/lib/mothership/request/session/controller-lease'
5564

5665
function redis() {
5766
const client = getRedisClient()
@@ -60,6 +69,7 @@ function redis() {
6069
}
6170

6271
afterAll(async () => {
72+
chatPubSub?.dispose()
6373
await closeRedisConnection()
6474
await new Promise<void>((resolve) => worker.server.close(() => resolve()))
6575
for (const [key, value] of Object.entries(inheritedEnv)) {
@@ -121,11 +131,12 @@ describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => {
121131
superseded?: boolean
122132
/** Admitted by code predating the current tool-execution protocol. */
123133
legacy?: boolean
134+
id?: string
124135
} = {}
125136
) {
126137
const chatId = generateId()
127138
const streamId = generateId()
128-
const runId = generateId()
139+
const runId = options.id ?? generateId()
129140
chatIds.push(chatId)
130141
const controllerToken =
131142
options.controllerToken === undefined
@@ -413,4 +424,131 @@ describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => {
413424
expect((await stored(orphan.runId)).status).toBe(stopped ? 'cancelled' : 'active')
414425
}
415426
})
427+
428+
it('never takes a run from a reconnect that locked its chat while the sweep was settling', async () => {
429+
for (let attempt = 0; attempt < 20; attempt++) {
430+
const orphans = await Promise.all(
431+
Array.from({ length: 20 }, () => admittedRun({ idleMinutes: 90 }))
432+
)
433+
434+
/**
435+
* Each reconnect locks the chat, proves its lease, then claims the run, as recovery
436+
* does. A run still unfinished once its reconnect holds the lock belongs to it.
437+
*/
438+
const reconnect = async (orphan: (typeof orphans)[number]) => {
439+
await sleep(randomInt(0, 40))
440+
if (!(await acquirePendingChatStream(orphan.chatId, orphan.streamId, 0))) {
441+
return { owned: false, claimed: false }
442+
}
443+
const lease = getLocalChatStreamLease(orphan.chatId, orphan.streamId)!
444+
try {
445+
await assertChatStreamLease(lease)
446+
const owned = (await stored(orphan.runId)).status === 'active'
447+
await sleep(randomInt(0, 10))
448+
const claimed = await claimRunController({
449+
runId: orphan.runId,
450+
chatId: orphan.chatId,
451+
previousToken: orphan.controllerToken!,
452+
token: lease.value,
453+
})
454+
return { owned, claimed }
455+
} finally {
456+
await releasePendingChatStream(orphan.chatId, orphan.streamId, lease)
457+
}
458+
}
459+
const [sweep, ...reconnects] = await Promise.all([
460+
sweepOrphanedRuns(),
461+
...orphans.map(reconnect),
462+
])
463+
464+
orphans.forEach((orphan, index) => {
465+
const swept = sweep.settledRunIds.includes(orphan.runId)
466+
if (reconnects[index].owned) expect(reconnects[index].claimed).toBe(true)
467+
expect(reconnects[index].claimed !== swept).toBe(true)
468+
})
469+
}
470+
})
471+
472+
it('announces every settled run whose chat it released, legacy runs included', async () => {
473+
const legacy = await admittedRun({
474+
idleMinutes: 25 * 60,
475+
controllerToken: null,
476+
legacy: true,
477+
})
478+
const announced: string[] = []
479+
const unsubscribe = chatPubSub!.onStatusChanged((event) => {
480+
if (event.type === 'completed') announced.push(event.chatId)
481+
})
482+
483+
try {
484+
const { settledRunIds } = await sweepOrphanedRuns()
485+
expect(settledRunIds).toContain(legacy.runId)
486+
for (let wait = 0; wait < 50 && !announced.includes(legacy.chatId); wait++) await sleep(20)
487+
expect(announced).toContain(legacy.chatId)
488+
} finally {
489+
unsubscribe()
490+
}
491+
})
492+
493+
it('reaches an orphan behind more unsettleable runs than one sweep examines', async () => {
494+
/** Runs whose replay is still live, all sorting before the orphan. */
495+
const blockers = Array.from({ length: 10_500 }, (_, index) => ({
496+
runId: `00000000-0000-4000-8000-${index.toString(16).padStart(12, '0')}`,
497+
chatId: generateId(),
498+
streamId: generateId(),
499+
}))
500+
const orphan = await admittedRun({
501+
idleMinutes: 90,
502+
id: 'ffffffff-ffff-4fff-bfff-ffffffffffff',
503+
})
504+
const blockerChatIds = blockers.map((blocker) => blocker.chatId)
505+
try {
506+
for (let start = 0; start < blockers.length; start += 1000) {
507+
const page = blockers.slice(start, start + 1000)
508+
await db.insert(copilotChats).values(
509+
page.map((blocker) => ({
510+
id: blocker.chatId,
511+
userId,
512+
workspaceId,
513+
type: 'mothership' as const,
514+
}))
515+
)
516+
await db.insert(copilotRuns).values(
517+
page.map((blocker) => ({
518+
id: blocker.runId,
519+
executionId: generateId(),
520+
chatId: blocker.chatId,
521+
userId,
522+
workspaceId,
523+
streamId: blocker.streamId,
524+
toolExecutionVersion: 2,
525+
status: 'active' as const,
526+
requestContext: { controllerToken: `${blocker.streamId}\n${generateId()}` },
527+
startedAt: sql`now() - interval '2 hours'`,
528+
updatedAt: sql`now() - interval '2 hours'`,
529+
}))
530+
)
531+
const pipeline = redis().pipeline()
532+
for (const blocker of page) {
533+
pipeline.set(`mothership_stream:${blocker.streamId}:seq`, '1', 'EX', 600)
534+
}
535+
await pipeline.exec()
536+
}
537+
538+
const first = await sweepOrphanedRuns()
539+
const second = first.settledRunIds.includes(orphan.runId) ? first : await sweepOrphanedRuns()
540+
541+
expect(second.settledRunIds).toContain(orphan.runId)
542+
expect((await stored(orphan.runId)).status).toBe('error')
543+
} finally {
544+
for (let start = 0; start < blockerChatIds.length; start += 1000) {
545+
await db
546+
.delete(copilotChats)
547+
.where(inArray(copilotChats.id, blockerChatIds.slice(start, start + 1000)))
548+
}
549+
const pipeline = redis().pipeline()
550+
for (const blocker of blockers) pipeline.del(`mothership_stream:${blocker.streamId}:seq`)
551+
await pipeline.exec()
552+
}
553+
}, 120_000)
416554
})

0 commit comments

Comments
 (0)