mirror of
https://github.com/mvschwarz/openrig.git
synced 2026-10-02 08:35:15 +08:00
fix(daemon): keep a blocked row's park timer through a seat handover (#242)
* fix(daemon): keep a blocked row's park timer through a seat handover At a seat swap, the occupant invalidator stops every watchdog job the retiring generation registered. A park timer (rig queue block --wake-after, with or without --wake-max) is registered by the parking seat, so a handover stopped it too. Nothing re-armed it, and the row stayed blocked with no timer, which is the strand reported in korallis/agent-stack#61. The rig's parked-owner consumer eventually noticed the dead wake and sent one generic wake once the new occupant was idle. The recorded wake time, message and backoff cadence were still lost. A blocked row's current queue-generated timer wakes the seat that owns the row, and the successor inherits that row, so the swap now keeps it. The queue lists those timers: each is still the row's current armed wake, is still active, and still targets the row's destination. The watchdog drop leaves them running. Every other job the retiring generation registered still stops: operator-attached watchdogs, superseded or stale timers, timers of rows rerouted to another seat, and occupant-only jobs. Re-parking still retires the kept timer, so no duplicate is created. A failed handover never reaches the invalidator, as before. * fix(daemon): retire a park timer when its row is rerouted routeToFallback moved a blocked row to a new owner but left the row's park timer active and aimed at the old owner, so it could fire to the old owner. The pre-delivery check only refuses timers whose rows are all terminal. That was already reachable on main (park, reroute, fire), and keeping park timers through a handover also let it happen after a swap (review-r2). The reroute now retires the row's current park-generated timer through the existing retireParkGeneratedTimer path, reason park_rerouted, the same one re-park and exit use. The new owner has no dedicated wake from a reroute, as before; the retired timer leaves the row's wake dead, so the parked-owner backstop covers it. * fix(daemon): leave a repeating queue wait running when its row is rerouted The previous commit retired every park timer on reroute. A repeating queue wait already delivers to the row's current owner (evaluateQueueWait sends to view.owner), so it woke the new owner correctly, and retiring it took that wake away. Only a periodic park timer is aimed at the old owner, so only that is retired now. Tests now fire through the engine with resolveQueueWait wired, as startup.ts wires it. After a periodic reroute the row shows no live park timer, and the new owner's parked diagnosis sees it as unhealthy, so the backstop covers the new owner. A repeating wait survives the reroute, with or without swaps, and delivers only to the new owner. --------- Co-authored-by: OpenRig contributors <noreply@openrig.dev>
This commit is contained in:
co-authored by
OpenRig contributors
parent
cc2c37dc2f
commit
afde814f5b
@@ -30,10 +30,11 @@ export interface OccupantInvalidatorDeps {
|
||||
/** (e/Class-B) durable watchdog_jobs store — armed jobs registered by the retiring generation are
|
||||
* stopped at swap (a stale wake into the successor's context is the ghost). Optional: absent ⇒ the
|
||||
* watchdog branch is skipped (never a name-scoped fallback). */
|
||||
watchdog?: { dropArmedByRegisteringGeneration(generationUuid: string): number };
|
||||
watchdog?: { dropArmedByRegisteringGeneration(generationUuid: string, keepJobIds?: readonly string[]): number };
|
||||
/** (e/Class-B) durable queue_items store — in-progress items claimed by the retiring generation are
|
||||
* RELEASED to pending (never dropped: the role work is durable, the successor re-claims). Optional. */
|
||||
queue?: { releaseClaimsByGeneration(generationUuid: string): number };
|
||||
* RELEASED to pending (never dropped: the role work is durable, the successor re-claims). Its current
|
||||
* park timers (blocked rows' wakes) are role work too, so the watchdog drop keeps them. Optional. */
|
||||
queue?: { releaseClaimsByGeneration(generationUuid: string): number; currentParkTimerIds?(): string[] };
|
||||
log?: (msg: string) => void;
|
||||
}
|
||||
|
||||
@@ -65,8 +66,10 @@ export class DefaultOccupantInvalidator implements OccupantInvalidator {
|
||||
// atom-B present → Class-B gen-scoped invalidation.
|
||||
// Watchdog (3b): stop every ARMED job registered by the retiring generation — a stale wake firing
|
||||
// into the successor's context is the specimen; the successor re-arms its own. Gen-scoped so the
|
||||
// successor's OWN armed jobs (same name, live gen) are untouched.
|
||||
const stopped = this.deps.watchdog?.dropArmedByRegisteringGeneration(retiringGeneration) ?? 0;
|
||||
// successor's OWN armed jobs (same name, live gen) are untouched. A blocked row's current park timer
|
||||
// is kept: it wakes the seat that owns the row, and nothing re-arms it if it stops.
|
||||
const keepJobIds = this.deps.queue?.currentParkTimerIds?.() ?? [];
|
||||
const stopped = this.deps.watchdog?.dropArmedByRegisteringGeneration(retiringGeneration, keepJobIds) ?? 0;
|
||||
if (stopped > 0) {
|
||||
log(
|
||||
`[occupant-invalidator] Class-B: stopped ${stopped} armed watchdog job(s) registered by retired ` +
|
||||
|
||||
@@ -25,7 +25,7 @@ import {
|
||||
type ParkWakeStatus,
|
||||
} from "./queue-wake-repository.js";
|
||||
import { WatchdogJobsRepository } from "./watchdog-jobs-repository.js";
|
||||
import { armQueueWait, backOffQueueWait, refreshQueueWaits, evaluateQueueWait, retargetQueueWait } from "./queue-wait-backoff.js";
|
||||
import { armQueueWait, backOffQueueWait, refreshQueueWaits, evaluateQueueWait, retargetQueueWait, isQueueWait } from "./queue-wait-backoff.js";
|
||||
|
||||
export const QUEUE_STATES = [
|
||||
"pending",
|
||||
@@ -2849,6 +2849,12 @@ export class QueueRepository {
|
||||
return this.wakeRepo.getStatus(qitemId);
|
||||
}
|
||||
|
||||
/** Park timers a seat swap keeps: each is still its blocked row's current wake and still targets the
|
||||
* row's owner (see QueueWakeRepository.currentParkTimerIds). */
|
||||
currentParkTimerIds(): string[] {
|
||||
return this.wakeRepo.currentParkTimerIds();
|
||||
}
|
||||
|
||||
/** Refuse a legacy park-generated timer only when every row bound to it is
|
||||
* terminal. Current exits retire these timers transactionally; this is the
|
||||
* delivery-seam backstop for residue persisted by an older daemon. A timer
|
||||
@@ -3385,6 +3391,12 @@ export class QueueRepository {
|
||||
WHERE qitem_id = ?`
|
||||
)
|
||||
.run(fallbackDestination, ts, newChain, `fallback: ${reason}`, qitemId);
|
||||
// A park timer targets the owner that parked the row. Once the row belongs to someone else it must not keep
|
||||
// waking the old owner, so it ends here like any other exit from the park. A repeating wait is left running:
|
||||
// it resolves its recipient from the row's current owner at delivery (evaluateQueueWait).
|
||||
const armed = this.wakeRepo.getStatus(qitemId);
|
||||
const armedJob = armed?.kind === "timer" ? (this.watchdogJobsRepo ?? new WatchdogJobsRepository(this.db)).getById(armed.ref) : null;
|
||||
if (!armedJob || !isQueueWait(armedJob.specYaml)) this.retireParkGeneratedTimer(qitemId, "park_rerouted");
|
||||
|
||||
this.transitionLog.append({
|
||||
qitemId,
|
||||
|
||||
@@ -174,6 +174,25 @@ export class QueueWakeRepository {
|
||||
});
|
||||
}
|
||||
|
||||
/** Active park-generated timers that are still their blocked row's current wake and still target the
|
||||
* row's owner. These wake the seat that owns the row, not a particular occupant, so a seat swap keeps
|
||||
* them; any other job the retiring occupant registered still stops. */
|
||||
currentParkTimerIds(): string[] {
|
||||
if (!this.available) return [];
|
||||
return this.db.prepare(
|
||||
`SELECT w.wake_ref
|
||||
FROM queue_transition_wakes w
|
||||
JOIN queue_items q ON q.qitem_id = w.qitem_id
|
||||
JOIN watchdog_jobs j ON j.job_id = w.wake_ref
|
||||
WHERE w.phase = 'armed' AND w.wake_kind = 'timer'
|
||||
AND q.state = 'blocked' AND j.state = 'active' AND j.target_session = q.destination_session
|
||||
AND w.transition_id = (
|
||||
SELECT MAX(a.transition_id) FROM queue_transition_wakes a
|
||||
WHERE a.qitem_id = w.qitem_id AND a.phase = 'armed'
|
||||
)`,
|
||||
).all().map((row) => (row as { wake_ref: string }).wake_ref);
|
||||
}
|
||||
|
||||
private isLive(kind: ParkWakeKind, ref: string): boolean {
|
||||
if (kind === "blocker") {
|
||||
const row = this.db.prepare("SELECT state FROM queue_items WHERE qitem_id = ?").get(ref) as
|
||||
|
||||
@@ -593,16 +593,18 @@ export class WatchdogJobsRepository {
|
||||
* successor re-arms its own. Gen-scoped, NEVER name-scoped (the successor shares the seat name, so a
|
||||
* name-scoped stop would kill the successor's own jobs). A NULL/empty generation never matches
|
||||
* (UNKNOWN ≠ retired — note-2). Returns the count stopped. Auditable (terminal_reason), not a hard
|
||||
* delete — the row stays for forensics. Pre-063 dbs no-op.
|
||||
* delete — the row stays for forensics. Pre-063 dbs no-op. `keepJobIds` are left running: the queue's
|
||||
* current park timers wake the seat that owns a still-blocked row, which the successor inherits.
|
||||
*/
|
||||
dropArmedByRegisteringGeneration(generationUuid: string): number {
|
||||
dropArmedByRegisteringGeneration(generationUuid: string, keepJobIds: readonly string[] = []): number {
|
||||
if (!this.hasGenColumn || !generationUuid) return 0;
|
||||
const keep = keepJobIds.length > 0 ? ` AND job_id NOT IN (${keepJobIds.map(() => "?").join(", ")})` : "";
|
||||
const res = this.db
|
||||
.prepare(
|
||||
`UPDATE watchdog_jobs SET state = 'stopped', terminal_reason = 'registering generation retired (seat handover)'
|
||||
WHERE state = 'active' AND registered_by_generation_uuid = ?`,
|
||||
WHERE state = 'active' AND registered_by_generation_uuid = ?${keep}`,
|
||||
)
|
||||
.run(generationUuid);
|
||||
.run(generationUuid, ...keepJobIds);
|
||||
return res.changes;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,209 @@
|
||||
// A seat swap stops the watchdog jobs the retiring occupant registered, so a stale wake doesn't reach the successor.
|
||||
// A blocked row's current park timer is different: it wakes the seat that still owns the row, and the successor
|
||||
// inherits that row. It survives the swap; every other job the retiring generation registered still stops.
|
||||
// Real queue, watchdog and session repositories on a full test DB, wired as startup.ts wires them.
|
||||
import { describe, it, expect, beforeEach, afterEach } from "vitest";
|
||||
import type Database from "better-sqlite3";
|
||||
import { createFullTestDb } from "./helpers/test-app.js";
|
||||
import { RigRepository } from "../src/domain/rig-repository.js";
|
||||
import { SessionRegistry } from "../src/domain/session-registry.js";
|
||||
import { EventBus } from "../src/domain/event-bus.js";
|
||||
import { QueueRepository } from "../src/domain/queue-repository.js";
|
||||
import { WatchdogJobsRepository } from "../src/domain/watchdog-jobs-repository.js";
|
||||
import { DefaultOccupantInvalidator } from "../src/domain/occupant-invalidator.js";
|
||||
import { diagnoseSeatParked } from "../src/domain/parked-query.js";
|
||||
import { WatchdogPolicyEngine } from "../src/domain/watchdog-policy-engine.js";
|
||||
import { WatchdogHistoryLog } from "../src/domain/watchdog-history-log.js";
|
||||
import { watchdogHistorySchema } from "../src/db/migrations/032_watchdog_history.js";
|
||||
|
||||
const SEAT = "dev-impl@seat-rig";
|
||||
const OTHER = "dev-qa@seat-rig";
|
||||
|
||||
describe("#handover-wake: a blocked row's park timer across an occupant swap", () => {
|
||||
let db: Database.Database;
|
||||
let queue: QueueRepository;
|
||||
let jobs: WatchdogJobsRepository;
|
||||
let invalidator: DefaultOccupantInvalidator;
|
||||
let gen: (session: string) => string | null;
|
||||
let seatNodeId: string;
|
||||
let otherNodeId: string;
|
||||
let sessionRegistry: SessionRegistry;
|
||||
let bus: EventBus;
|
||||
|
||||
beforeEach(() => {
|
||||
db = createFullTestDb();
|
||||
db.exec(watchdogHistorySchema.sql); // the engine audits each evaluation; the full test DB does not include this table
|
||||
const rigRepo = new RigRepository(db);
|
||||
sessionRegistry = new SessionRegistry(db);
|
||||
const rig = rigRepo.createRig("seat-rig");
|
||||
for (const [logicalId, session] of [["dev.impl", SEAT], ["dev.qa", OTHER]] as const) {
|
||||
const node = rigRepo.addNode(rig.id, logicalId, { runtime: "codex" });
|
||||
sessionRegistry.updateStatus(sessionRegistry.registerSession(node.id, session).id, "running");
|
||||
if (session === SEAT) seatNodeId = node.id; else otherNodeId = node.id;
|
||||
}
|
||||
gen = (session) => sessionRegistry.currentOccupantGenerationForSession(session);
|
||||
bus = new EventBus(db);
|
||||
queue = new QueueRepository(db, bus, { validateRig: () => true, resolveOccupantGeneration: gen });
|
||||
jobs = new WatchdogJobsRepository(db, undefined, gen);
|
||||
queue.attachWatchdogJobsRepository(jobs);
|
||||
invalidator = new DefaultOccupantInvalidator({
|
||||
enforcer: { invalidateOccupant() {} }, contextUsage: { invalidateOccupantSidecar() {} }, watchdog: jobs, queue,
|
||||
});
|
||||
});
|
||||
afterEach(() => db.close());
|
||||
|
||||
async function park(opts: { destination?: string; actor?: string; repeating?: boolean; watchdogId?: string } = {}) {
|
||||
const row = await queue.create({ sourceSession: "orch-lead@seat-rig", destinationSession: opts.destination ?? SEAT, body: "work" } as never);
|
||||
queue.update({
|
||||
qitemId: row.qitemId, actorSession: opts.actor ?? SEAT, state: "blocked", blockedOn: "external:ci",
|
||||
transitionNote: "continuation: check CI",
|
||||
...(opts.watchdogId ? { wakeWatchdogId: opts.watchdogId } : { wakeAfterSeconds: 1800, ...(opts.repeating ? { wakeMaxSeconds: 7200 } : {}) }),
|
||||
} as never);
|
||||
return { qitemId: row.qitemId, wakeRef: (queue.getParkWakeStatus(row.qitemId) as { ref: string }).ref };
|
||||
}
|
||||
const swap = () => invalidator.invalidateRetiringOccupant({ retiringSessionName: SEAT, successorSessionName: SEAT, retiringGeneration: gen(SEAT)! });
|
||||
const jobState = (jobId: string) => (db.prepare("SELECT state FROM watchdog_jobs WHERE job_id = ?").get(jobId) as { state: string }).state;
|
||||
const activeJobsFor = (session: string) => (db.prepare("SELECT COUNT(*) AS n FROM watchdog_jobs WHERE state = 'active' AND target_session = ?").get(session) as { n: number }).n;
|
||||
const wakeLive = (qitemId: string) => (queue.getParkWakeStatus(qitemId) as { live: boolean }).live;
|
||||
const parkedWhenIdle = (session = SEAT) => diagnoseSeatParked({
|
||||
getSeatState: () => ({ activity: "idle-at-prompt", needsInput: { count: 0, reason: null }, decidedBy: "test" }) as never,
|
||||
listOpenObligations: (destinationSession, limit) => ({
|
||||
rows: queue.list({ destinationSession, state: ["pending", "in-progress", "blocked"], limit })
|
||||
.map((r) => ({ qitemId: r.qitemId, state: r.state as "pending" | "in-progress" | "blocked", summary: r.summary ?? null })),
|
||||
limit,
|
||||
}),
|
||||
getParkWake: (qitemId) => queue.getParkWakeStatus(qitemId),
|
||||
}, { seatNodeId: session === SEAT ? seatNodeId : otherNodeId, sessionName: session });
|
||||
|
||||
it.each([false, true])("the still-owned blocked row keeps its one live timer (repeating=%s)", async (repeating) => {
|
||||
const { qitemId, wakeRef } = await park({ repeating });
|
||||
|
||||
swap();
|
||||
|
||||
expect(jobState(wakeRef)).toBe("active");
|
||||
expect(wakeLive(qitemId)).toBe(true);
|
||||
expect(activeJobsFor(SEAT)).toBe(1);
|
||||
expect(parkedWhenIdle().obligations.unhealthyHeldCount).toBe(0);
|
||||
});
|
||||
|
||||
it("a successor re-park still supersedes the kept timer: one live timer, not two", async () => {
|
||||
const { qitemId, wakeRef } = await park();
|
||||
swap();
|
||||
|
||||
queue.update({ qitemId, actorSession: SEAT, state: "blocked", blockedOn: "external:ci", transitionNote: "continuation: re-parked", wakeAfterSeconds: 600 } as never);
|
||||
|
||||
expect(jobState(wakeRef)).not.toBe("active");
|
||||
expect(activeJobsFor(SEAT)).toBe(1);
|
||||
});
|
||||
|
||||
it("another seat's blocked row parked by the retiring occupant keeps its live timer", async () => {
|
||||
const { qitemId, wakeRef } = await park({ destination: OTHER, actor: SEAT });
|
||||
|
||||
swap();
|
||||
|
||||
expect(jobState(wakeRef)).toBe("active");
|
||||
expect(wakeLive(qitemId)).toBe(true);
|
||||
});
|
||||
|
||||
it("control: an occupant-only job the retiring generation registered still stops", () => {
|
||||
const own = jobs.register({
|
||||
policy: "periodic-reminder", targetSession: SEAT, intervalSeconds: 600, registeredBySession: SEAT,
|
||||
specYaml: ["policy: periodic-reminder", "target:", ` session: ${JSON.stringify(SEAT)}`, 'message: "own reminder"', ""].join("\n"),
|
||||
});
|
||||
|
||||
swap();
|
||||
|
||||
expect(jobState(own.jobId)).toBe("stopped");
|
||||
});
|
||||
|
||||
it("control: an operator watchdog attached to the blocked row still stops (custom job, not queue-generated)", async () => {
|
||||
const attached = jobs.register({
|
||||
policy: "periodic-reminder", targetSession: SEAT, intervalSeconds: 600, registeredBySession: SEAT,
|
||||
specYaml: ["policy: periodic-reminder", "target:", ` session: ${JSON.stringify(SEAT)}`, 'message: "operator watchdog"', ""].join("\n"),
|
||||
});
|
||||
await park({ watchdogId: attached.jobId });
|
||||
|
||||
swap();
|
||||
|
||||
expect(jobState(attached.jobId)).toBe("stopped");
|
||||
});
|
||||
|
||||
it("control: a row rerouted to another seat does not keep the old seat's timer", async () => {
|
||||
const { qitemId, wakeRef } = await park();
|
||||
queue.routeToFallback(qitemId, OTHER, "test transfer");
|
||||
|
||||
swap();
|
||||
|
||||
expect(jobState(wakeRef)).not.toBe("active");
|
||||
});
|
||||
|
||||
// review-r2 (PR #242): a rerouted row's park timer must not keep waking the OLD owner. The scheduler evaluates
|
||||
// active jobs; this evaluates every active job, due or not, through the real engine and pre-delivery check.
|
||||
async function fireAllActive() {
|
||||
const deliveries: Array<{ targetSession: string; message: string }> = [];
|
||||
const engine = new WatchdogPolicyEngine({
|
||||
jobsRepo: jobs, historyLog: new WatchdogHistoryLog(db), eventBus: bus,
|
||||
deliver: async (request) => { deliveries.push(request); return { status: "ok" }; },
|
||||
resolvePreDeliveryTerminalReason: ({ jobId }: { jobId: string }) => queue.resolveWatchdogPreDeliveryTerminalReason(jobId),
|
||||
resolveQueueWait: (input: { jobId: string }) => queue.evaluateWaitReminder(input),
|
||||
onWakeAttempt: ({ jobId, deliveryStatus }) => queue.recordWatchdogWakeAttempt(jobId, deliveryStatus),
|
||||
});
|
||||
for (const job of jobs.listActive()) await engine.evaluate(job);
|
||||
return deliveries;
|
||||
}
|
||||
|
||||
it("park, reroute, fire: nothing reaches the old owner (no handover)", async () => {
|
||||
const { qitemId, wakeRef } = await park();
|
||||
|
||||
queue.routeToFallback(qitemId, OTHER, "test reroute");
|
||||
|
||||
expect(jobState(wakeRef)).not.toBe("active");
|
||||
expect(wakeLive(qitemId)).toBe(false);
|
||||
expect(parkedWhenIdle(OTHER).obligations.unhealthyHeldCount).toBe(1);
|
||||
expect((await fireAllActive()).filter((d) => d.targetSession === SEAT)).toEqual([]);
|
||||
});
|
||||
|
||||
it("park, swap (timer kept), reroute, swap, fire: nothing reaches the old owner", async () => {
|
||||
const { qitemId, wakeRef } = await park();
|
||||
swap();
|
||||
expect(jobState(wakeRef)).toBe("active");
|
||||
queue.routeToFallback(qitemId, OTHER, "test reroute");
|
||||
const firstGeneration = gen(SEAT);
|
||||
sessionRegistry.updateStatus(sessionRegistry.registerSession(seatNodeId, SEAT).id, "running");
|
||||
expect(gen(SEAT)).not.toBe(firstGeneration);
|
||||
|
||||
swap();
|
||||
|
||||
expect(jobState(wakeRef)).not.toBe("active");
|
||||
expect(wakeLive(qitemId)).toBe(false);
|
||||
expect(parkedWhenIdle(OTHER).obligations.unhealthyHeldCount).toBe(1);
|
||||
expect((await fireAllActive()).filter((d) => d.targetSession === SEAT)).toEqual([]);
|
||||
});
|
||||
|
||||
it.each([false, true])("a repeating wait survives the reroute and wakes only the new owner (swaps=%s)", async (withSwaps) => {
|
||||
const { qitemId, wakeRef } = await park({ repeating: true });
|
||||
if (withSwaps) swap();
|
||||
queue.routeToFallback(qitemId, OTHER, "test reroute");
|
||||
if (withSwaps) {
|
||||
sessionRegistry.updateStatus(sessionRegistry.registerSession(seatNodeId, SEAT).id, "running");
|
||||
swap();
|
||||
}
|
||||
|
||||
expect(jobState(wakeRef)).toBe("active");
|
||||
expect(wakeLive(qitemId)).toBe(true);
|
||||
const deliveries = await fireAllActive();
|
||||
expect(deliveries.map((d) => d.targetSession)).toEqual([OTHER]);
|
||||
});
|
||||
|
||||
it("control: a stale timer that is no longer the row's current wake still stops", async () => {
|
||||
const { qitemId, wakeRef: first } = await park();
|
||||
queue.update({ qitemId, actorSession: SEAT, state: "blocked", blockedOn: "external:ci", transitionNote: "continuation: re-parked", wakeAfterSeconds: 600 } as never);
|
||||
// Legacy residue: a superseded timer left active by a path that did not retire it.
|
||||
db.prepare("UPDATE watchdog_jobs SET state = 'active', terminal_reason = NULL WHERE job_id = ?").run(first);
|
||||
|
||||
swap();
|
||||
|
||||
expect(jobState(first)).toBe("stopped");
|
||||
expect(wakeLive(qitemId)).toBe(true);
|
||||
});
|
||||
});
|
||||
@@ -48,7 +48,8 @@ describe("GHOST-STAGE (e) DefaultOccupantInvalidator", () => {
|
||||
watchdog: { dropArmedByRegisteringGeneration },
|
||||
log: (m) => logs.push(m),
|
||||
}).invalidateRetiringOccupant({ retiringSessionName: "seat@rig", successorSessionName: "seat@rig", retiringGeneration: "gen-retired" });
|
||||
expect(dropArmedByRegisteringGeneration).toHaveBeenCalledWith("gen-retired");
|
||||
// No queue dep wired, so no park timers are kept.
|
||||
expect(dropArmedByRegisteringGeneration).toHaveBeenCalledWith("gen-retired", []);
|
||||
expect(logs.some((l) => /stopped 2 armed watchdog/.test(l))).toBe(true);
|
||||
});
|
||||
|
||||
|
||||
@@ -17,6 +17,8 @@ import { TmuxAdapter } from "../src/adapters/tmux.js";
|
||||
import type { RuntimeAdapter } from "../src/domain/runtime-adapter.js";
|
||||
import { observeCodexSandbox } from "../src/domain/permission-drift.js";
|
||||
import { AppliedLaunchObservationStore } from "../src/domain/applied-launch-observation-store.js";
|
||||
import { QueueRepository } from "../src/domain/queue-repository.js";
|
||||
import { DefaultOccupantInvalidator, type OccupantInvalidator } from "../src/domain/occupant-invalidator.js";
|
||||
|
||||
describe("SeatHandoverService", () => {
|
||||
let db: Database.Database;
|
||||
@@ -96,7 +98,7 @@ describe("SeatHandoverService", () => {
|
||||
return { runtime: "codex", launchHarness, checkReady } as unknown as RuntimeAdapter;
|
||||
}
|
||||
|
||||
function newService(adapter: TmuxAdapter = tmux()): SeatHandoverService {
|
||||
function newService(adapter: TmuxAdapter = tmux(), occupantInvalidator: OccupantInvalidator = { invalidateRetiringOccupant }): SeatHandoverService {
|
||||
return new SeatHandoverService({
|
||||
db,
|
||||
rigRepo,
|
||||
@@ -109,7 +111,7 @@ describe("SeatHandoverService", () => {
|
||||
runtimeAdapters: { codex: codexAdapter() },
|
||||
contextUsageStore: { readSidecar } as never,
|
||||
resumeTokenCapturer: { captureCodexThreadId } as never,
|
||||
occupantInvalidator: { invalidateRetiringOccupant },
|
||||
occupantInvalidator,
|
||||
activityOracle: { declareOccupantSwap },
|
||||
predecessorRecapResolver: resolvePredecessorRecap as never,
|
||||
readinessTimeoutMs: 50,
|
||||
@@ -1166,6 +1168,26 @@ describe("SeatHandoverService", () => {
|
||||
expect(durableRows()).toBe(before);
|
||||
});
|
||||
|
||||
it("a successful fresh handover keeps the seat's blocked row's park timer live (real invalidator and repositories)", async () => {
|
||||
seedSeat();
|
||||
const gen = (session: string) => sessionRegistry.currentOccupantGenerationForSession(session);
|
||||
const queue = new QueueRepository(db, eventBus, { validateRig: () => true, resolveOccupantGeneration: gen });
|
||||
const jobs = new WatchdogJobsRepository(db, undefined, gen);
|
||||
queue.attachWatchdogJobsRepository(jobs);
|
||||
const row = await queue.create({ sourceSession: "orch-lead@seat-rig", destinationSession: "dev-impl@seat-rig", body: "work" } as never);
|
||||
queue.update({ qitemId: row.qitemId, actorSession: "dev-impl@seat-rig", state: "blocked", blockedOn: "external:ci", transitionNote: "continuation: check CI", wakeAfterSeconds: 1800 } as never);
|
||||
const timer = (queue.getParkWakeStatus(row.qitemId) as { ref: string }).ref;
|
||||
const handover = newService(tmux(), new DefaultOccupantInvalidator({
|
||||
enforcer: { invalidateOccupant() {} }, contextUsage: { invalidateOccupantSidecar() {} }, watchdog: jobs, queue,
|
||||
}));
|
||||
|
||||
const result = await handover.handover({ seatRef: "dev-impl@seat-rig", reason: "context-wall", source: "fresh" });
|
||||
|
||||
expect(result.ok).toBe(true);
|
||||
expect(jobs.getById(timer)?.state).toBe("active");
|
||||
expect((queue.getParkWakeStatus(row.qitemId) as { live: boolean }).live).toBe(true);
|
||||
});
|
||||
|
||||
it("fails before mutation when the seat has no current occupant", async () => {
|
||||
seedSeat({ withSession: false });
|
||||
const discovered = seedDiscovery();
|
||||
|
||||
Reference in New Issue
Block a user