Skip to content

Commit fd4fc18

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

3 files changed

Lines changed: 46 additions & 35 deletions

File tree

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

Lines changed: 19 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -48,7 +48,7 @@ import {
4848
settleStoppedRunWithoutController,
4949
sweepOrphanedRuns,
5050
} from '@/lib/mothership/async-runs/orphaned-runs'
51-
import { updateRunStatus } from '@/lib/mothership/async-runs/repository'
51+
import { requestRunStop, updateRunStatus } from '@/lib/mothership/async-runs/repository'
5252
import { abortRun } from '@/lib/mothership/request/application/controls'
5353
import { claimRunController } from '@/lib/mothership/request/lifecycle/controller-ownership'
5454
import { chatStreamLockKey } from '@/lib/mothership/request/session/controller-lease'
@@ -160,6 +160,11 @@ describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => {
160160
return { chatId, streamId, runId, controllerToken }
161161
}
162162

163+
/** Records the user's Stop the way the abort use case does before it settles anything. */
164+
async function stop(run: { streamId: string; chatId: string }) {
165+
await requestRunStop({ userId, workspaceId, streamId: run.streamId, chatId: run.chatId })
166+
}
167+
163168
async function stored(runId: string) {
164169
const [run] = await db.select().from(copilotRuns).where(eq(copilotRuns.id, runId))
165170
const [chat] = await db
@@ -322,8 +327,19 @@ describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => {
322327
expect(run.marker).toBeNull()
323328
})
324329

330+
it('never cancels a run nobody stopped', async () => {
331+
const orphan = await admittedRun({ superseded: true })
332+
333+
expect(await settleStoppedRunWithoutController(orphan.runId)).toBe(false)
334+
335+
const run = await stored(orphan.runId)
336+
expect(run.status).toBe('active')
337+
expect(run.marker).toBe(orphan.streamId)
338+
})
339+
325340
it('leaves a stopped run to the controller of its stream that holds the chat lock', async () => {
326341
const owned = await admittedRun()
342+
await stop(owned)
327343
await redis().set(chatStreamLockKey(owned.chatId), owned.controllerToken!, 'EX', 60)
328344

329345
try {
@@ -366,6 +382,7 @@ describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => {
366382
it('settles a stopped run exactly once when Stop races its own controller finalizing', async () => {
367383
for (let attempt = 0; attempt < 50; attempt++) {
368384
const orphan = await admittedRun()
385+
await stop(orphan)
369386

370387
const [finalized, stopped] = await Promise.all([
371388
updateRunStatus(orphan.runId, 'complete', {}, orphan.controllerToken!),
@@ -380,6 +397,7 @@ describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => {
380397
it('settles a stopped run exactly once when Stop races a recovering controller claiming it', async () => {
381398
for (let attempt = 0; attempt < 50; attempt++) {
382399
const orphan = await admittedRun()
400+
await stop(orphan)
383401

384402
const [claimed, stopped] = await Promise.all([
385403
claimRunController({

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

Lines changed: 27 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -23,10 +23,11 @@ import {
2323
type SQL,
2424
sql,
2525
} from 'drizzle-orm'
26+
import type { DbTransaction } from '@/lib/db/types'
2627
import { SIM_TOOL_EXECUTION_VERSION } from '@/lib/mothership/async-runs/lifecycle'
2728
import { publishChatStatusChanged } from '@/lib/mothership/chat-status'
29+
import { getChatStreamLockOwners } from '@/lib/mothership/request/session/abort'
2830
import { findStreamsWithReplay } from '@/lib/mothership/request/session/buffer'
29-
import { findStreamsHoldingChatLock } from '@/lib/mothership/request/session/controller-lease'
3031

3132
const logger = createLogger('OrphanedCopilotRuns')
3233

@@ -110,13 +111,8 @@ const unownedRunColumns = {
110111
controllerToken,
111112
}
112113

113-
type Transaction = Parameters<Parameters<typeof db.transaction>[0]>[0]
114-
115-
/** A stopped run ends cancelled; the sweep ends it cancelled only if its user pressed Stop. */
116-
function terminalValues(reason: 'stopped' | 'orphaned' | 'legacy') {
117-
if (reason === 'stopped') {
118-
return { status: sql`'cancelled'::copilot_run_status`, error: sql`NULL::text` }
119-
}
114+
/** A run ends cancelled only if its user pressed Stop, whoever settles it. */
115+
function terminalValues(reason: 'orphaned' | 'legacy') {
120116
const error = reason === 'orphaned' ? ORPHANED_RUN_ERROR : LEGACY_RUN_ERROR
121117
return {
122118
status: sql`(CASE WHEN ${stopRequested} THEN 'cancelled' ELSE 'error' END)::copilot_run_status`,
@@ -135,9 +131,8 @@ function terminalValues(reason: 'stopped' | 'orphaned' | 'legacy') {
135131
* touching the chat's ordering timestamp.
136132
*/
137133
async function settleRuns(
138-
tx: Transaction,
134+
tx: DbTransaction,
139135
runs: UnownedRun[],
140-
reason: 'stopped' | 'orphaned',
141136
guard: SQL | undefined
142137
): Promise<UnownedRun[]> {
143138
if (runs.length === 0) return []
@@ -182,7 +177,7 @@ async function settleRuns(
182177
for (const run of runs) {
183178
if (run.controllerToken === null) continue
184179
await apply([run], eq(controllerToken, run.controllerToken), {
185-
...terminalValues(reason),
180+
...terminalValues('orphaned'),
186181
completedAt: sql`now()`,
187182
updatedAt: sql`now()`,
188183
})
@@ -217,6 +212,23 @@ function announceSettled(runs: UnownedRun[]): void {
217212
}
218213
}
219214

215+
/**
216+
* The streams among these whose own controller holds its chat lock, under any token: a
217+
* recovering controller locks the chat before it claims the run. Throws unless the
218+
* locks were read, since otherwise no stream is provably unowned.
219+
*/
220+
async function findStreamsHoldingChatLock(
221+
runs: Array<{ chatId: string; streamId: string }>
222+
): Promise<Set<string>> {
223+
const { status, ownersByChatId } = await getChatStreamLockOwners([
224+
...new Set(runs.map((run) => run.chatId)),
225+
])
226+
if (status !== 'verified') throw new Error('Chat stream locks are unreadable')
227+
return new Set(
228+
runs.filter((run) => ownersByChatId.get(run.chatId) === run.streamId).map((run) => run.streamId)
229+
)
230+
}
231+
220232
/** The candidates no controller owns; leased runs are skipped when ownership is unreadable. */
221233
async function withoutOwners(candidates: UnownedRun[]): Promise<UnownedRun[]> {
222234
const leased = candidates.filter((run) => run.controllerToken !== null)
@@ -268,7 +280,7 @@ export async function sweepOrphanedRuns(): Promise<{ settledRunIds: string[] }>
268280

269281
try {
270282
const unowned = await withoutOwners(candidates)
271-
const settled = await db.transaction((tx) => settleRuns(tx, unowned, 'orphaned', orphanIdle))
283+
const settled = await db.transaction((tx) => settleRuns(tx, unowned, orphanIdle))
272284
announceSettled(settled.filter((run) => run.controllerToken !== null))
273285
settledRunIds.push(...settled.map((run) => run.id))
274286
} catch (error) {
@@ -289,7 +301,8 @@ export async function sweepOrphanedRuns(): Promise<{ settledRunIds: string[] }>
289301

290302
/**
291303
* Settles a stopped run as cancelled when no controller of its stream holds the chat
292-
* lock. A live controller observes the Stop and settles its own run.
304+
* lock. A live controller observes the Stop and settles its own run. The update itself
305+
* requires the recorded Stop, so this can never settle a run nobody stopped.
293306
*/
294307
export async function settleStoppedRunWithoutController(runId: string): Promise<boolean> {
295308
const [run] = await db
@@ -299,7 +312,7 @@ export async function settleStoppedRunWithoutController(runId: string): Promise<
299312
.limit(1)
300313
if (!run?.controllerToken) return false
301314
if ((await findStreamsHoldingChatLock([run])).has(run.streamId)) return false
302-
const settled = await db.transaction((tx) => settleRuns(tx, [run], 'stopped', undefined))
315+
const settled = await db.transaction((tx) => settleRuns(tx, [run], stopRequested))
303316
announceSettled(settled)
304317
return settled.length > 0
305318
}

‎apps/sim/lib/mothership/request/session/controller-lease.ts‎

Lines changed: 0 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -28,26 +28,6 @@ export async function assertChatStreamLease(lease: ChatStreamLease): Promise<voi
2828
}
2929
}
3030

31-
/**
32-
* The streams among these whose own controller holds its chat lock, under any token: a
33-
* recovering controller locks the chat before it claims the run. Throws when the locks
34-
* cannot be read, since then no stream is provably unowned.
35-
*/
36-
export async function findStreamsHoldingChatLock(
37-
streams: Array<{ chatId: string; streamId: string }>
38-
): Promise<Set<string>> {
39-
const redis = getRedisClient()
40-
if (!redis) throw new Error('Chat stream locks are unreadable without Redis')
41-
if (streams.length === 0) return new Set()
42-
const values = await redis.mget(streams.map(({ chatId }) => chatStreamLockKey(chatId)))
43-
const held = new Set<string>()
44-
streams.forEach(({ streamId }, index) => {
45-
const value = values[index]
46-
if (value && streamIdFromLock(value) === streamId) held.add(streamId)
47-
})
48-
return held
49-
}
50-
5131
/** Whether this lease still holds its chat lock; an unreadable lock counts as lost. */
5232
export async function holdsChatStreamLease(lease: ChatStreamLease): Promise<boolean> {
5333
try {

0 commit comments

Comments
 (0)