fix(orchestration): scope the worker-list anchor cursor to the requested Run

This commit is contained in:
Jinwoo-H
2026-09-04 04:17:18 -04:00
parent a68ef063a0
commit 1047dee3d4
2 changed files with 89 additions and 8 deletions
@@ -52,10 +52,26 @@ export function listWorkerTerminalReleaseBacklog(
export const WORKER_LIST_CURSOR_EXPIRED_MESSAGE =
'The worker inventory changed destructively while paging. Restart without --cursor.'
function resolveDispatchRowId(this: OrchestrationDb, after: WorkerTerminalOrderingKey): number {
/** The anchor must still belong to the filtered set. An anchor from another Run resolved to a
* rowid past this Run's rows, so the page read as a finished, empty inventory. */
function resolveAnchorRowId(
this: OrchestrationDb,
after: WorkerTerminalOrderingKey,
runId: string | undefined
): number {
const conditions = ['id = ?']
const values: (string | number)[] = [after.dispatchId]
if (runId) {
conditions.push('run_id = ?')
values.push(runId)
}
if (after.databaseId !== undefined) {
conditions.push('rowid = ?')
values.push(after.databaseId)
}
const anchor = this.db
.prepare('SELECT rowid AS rowid FROM dispatch_contexts WHERE id = ?')
.get(after.dispatchId) as { rowid: number } | undefined
.prepare(`SELECT rowid AS rowid FROM dispatch_contexts WHERE ${conditions.join(' AND ')}`)
.get(...values) as { rowid: number } | undefined
if (!anchor) {
throw new OrchestrationError('worker_list_cursor_expired', WORKER_LIST_CURSOR_EXPIRED_MESSAGE)
}
@@ -116,12 +132,10 @@ export function listWorkerTerminalResources(
}
if (params.after) {
// Order and fence must share one key, or a row created between pages moves across the cut.
// A pre-v3 cursor carries no rowid, so it has to be resolved from its anchor row; when a
// reset deleted that row `rowid > NULL` matched nothing and the page read as a finished,
// empty inventory instead of an expired cursor.
const anchorRowId = params.after.databaseId ?? resolveDispatchRowId.call(this, params.after)
// A pre-v3 cursor is resolved from its anchor row; when a reset deleted that row
// `rowid > NULL` matched nothing and the page read as a finished, empty inventory.
where.push('d.rowid > ?')
values.push(anchorRowId)
values.push(resolveAnchorRowId.call(this, params.after, params.runId))
}
let detailWhere = where
let detailValues = values
@@ -499,6 +499,73 @@ describe('orchestration worker-list pagination', () => {
expect(active.workers.map((worker) => worker.dispatchId)).toEqual(['dispatch-unsupervised'])
expect(active.page.total).toBe(1)
})
describe('a legacy cursor anchored outside the requested Run', () => {
function twoRuns(): { runA: string; runB: string } {
db = new OrchestrationDb(':memory:')
const runA = db.createRun({
objective: 'A',
coordinatorHandle: 'term-a',
coordinatorPaneKey: 'tab-a:aaaaaaaa-aaaa-4aaa-8aaa-aaaaaaaaaaaa'
})
const runB = db.createRun({
objective: 'B',
coordinatorHandle: 'term-b',
coordinatorPaneKey: 'tab-b:bbbbbbbb-bbbb-4bbb-8bbb-bbbbbbbbbbbb'
})
insertDispatch(db, runA.id, 'a-1')
insertDispatch(db, runA.id, 'a-2')
insertDispatch(db, runB.id, 'b-1')
insertDispatch(db, runB.id, 'b-2')
return { runA: runA.id, runB: runB.id }
}
// Both shapes used to resolve to a rowid past Run A's rows and report a finished, empty page.
it('expires a v2 cursor that must be resolved from a foreign anchor', async () => {
const { runA } = twoRuns()
const runtime = new OrcaRuntimeService()
runtime.setOrchestrationDb(db!)
const foreign = encodeWorkerListCursor({
version: 2,
snapshot: { databaseId: 4 },
after: { createdAt: '2026-08-27 00:00:00', dispatchId: 'b-1' }
})
await expect(
callWorkerList(runtime, { run: runA, limit: 10, cursor: foreign })
).rejects.toThrow(/changed destructively/u)
})
it('expires a v2 cursor that carries a foreign rowid', async () => {
const { runA } = twoRuns()
const runtime = new OrcaRuntimeService()
runtime.setOrchestrationDb(db!)
const foreign = encodeWorkerListCursor({
version: 2,
snapshot: { databaseId: 4 },
after: { createdAt: '2026-08-27 00:00:00', dispatchId: 'b-1', databaseId: 3 }
})
await expect(
callWorkerList(runtime, { run: runA, limit: 10, cursor: foreign })
).rejects.toThrow(/changed destructively/u)
})
it('still pages the requested Run from its own anchor', async () => {
const { runA } = twoRuns()
const runtime = new OrcaRuntimeService()
runtime.setOrchestrationDb(db!)
const own = encodeWorkerListCursor({
version: 2,
snapshot: { databaseId: 4 },
after: { createdAt: '2026-08-27 00:00:00', dispatchId: 'a-1', databaseId: 1 }
})
const page = await callWorkerList(runtime, { run: runA, limit: 10, cursor: own })
expect(page.workers.map((worker) => worker.dispatchId)).toEqual(['a-2'])
})
})
})
async function callWorkerList(