diff --git a/apps/scheduler/src/request-scheduler.test.ts b/apps/scheduler/src/request-scheduler.test.ts index fe25923e..f017f3b1 100644 --- a/apps/scheduler/src/request-scheduler.test.ts +++ b/apps/scheduler/src/request-scheduler.test.ts @@ -85,6 +85,7 @@ function makeCollections( queuedQueueName?: string; }; count: number; + requestIds: string[]; } >(); for (const request of requests.filter( @@ -97,8 +98,10 @@ function makeCollections( `${request.workerType}\0${request.agentVersion ?? ""}` + `\0${queuedQueueName ?? ""}`; const current = grouped.get(key); - if (current) current.count++; - else { + if (current) { + current.count++; + current.requestIds.push(request._id); + } else { grouped.set(key, { _id: { workerType: request.workerType, @@ -108,6 +111,7 @@ function makeCollections( ...(queuedQueueName ? { queuedQueueName } : {}), }, count: 1, + requestIds: [request._id], }); } } @@ -166,17 +170,21 @@ function makeCollections( updateMany: vi.fn().mockImplementation( async ( filter: { + _id?: { $in: string[] }; workerType?: string; agentVersion?: string; + "run.status"?: string; "run.queuedQueueName"?: string | { $exists: false }; }, update: { - $set: { "run.status": string }; + $set: Record; $unset?: { "run.queuedQueueName"?: string }; }, ) => { const matches = requests.filter((candidate) => { - if (candidate.run?.status !== "queued" || candidate.deletedAt) return false; + const expectedStatus = filter["run.status"] ?? "queued"; + if (candidate.run?.status !== expectedStatus || candidate.deletedAt) return false; + if (filter._id?.$in && !filter._id.$in.includes(candidate._id)) return false; if ( typeof filter.workerType === "string" && candidate.workerType !== filter.workerType @@ -186,15 +194,23 @@ function makeCollections( candidate.agentVersion !== filter.agentVersion ) return false; const queueFilter = filter["run.queuedQueueName"]; + if (queueFilter === undefined) return true; if (typeof queueFilter === "string") { return candidate.run.queuedQueueName === queueFilter; } return candidate.run.queuedQueueName === undefined; }); for (const request of matches) { - request.run!.status = update.$set["run.status"] as "pending"; + if (!request.run) continue; + for (const [key, value] of Object.entries(update.$set)) { + if (key === "updatedAt") { + request.updatedAt = value as Date; + } else if (key.startsWith("run.")) { + (request.run as any)[key.slice(4)] = value; + } + } if (update.$unset?.["run.queuedQueueName"] !== undefined) { - delete request.run!.queuedQueueName; + delete request.run.queuedQueueName; } } return { matchedCount: matches.length, modifiedCount: matches.length }; @@ -437,11 +453,13 @@ describe("RequestScheduler", () => { ); }); - it("leaves invalid targets pending and emits actionable telemetry", async () => { + it("terminalizes permanently invalid pinned targets but keeps recoverable targets pending", async () => { const requests = [ makeRequest("unknown", "missing-worker", "v1"), makeRequest("versionless", "known-worker", undefined), makeRequest("inactive", "known-worker", "retired"), + makeRequest("temporarily-unavailable", "paused-worker", "v1"), + makeRequest("deleted", "deleted-worker", "v1"), ]; const { requestCollection, agentCollection } = makeCollections(requests, [ makeAgent("known-worker", [ @@ -452,6 +470,16 @@ describe("RequestScheduler", () => { status: "retired", }, ]), + makeAgent( + "paused-worker", + [{ agentVersion: "v1", queueName: "paused-queue" }], + { available: false }, + ), + makeAgent( + "deleted-worker", + [{ agentVersion: "v1", queueName: "deleted-queue" }], + { deletedAt: new Date() }, + ), ]); const queue = makeQueueClient(); const warn = vi.spyOn(console, "warn").mockImplementation(() => {}); @@ -464,16 +492,184 @@ describe("RequestScheduler", () => { await (scheduler as any).dispatch(); expect(queue.sendMessage).not.toHaveBeenCalled(); - expect(requests.every((request) => request.run?.status === "pending")).toBe( - true, - ); + expect(requests[0].run?.status).toBe("pending"); + expect(requests[1].run?.status).toBe("pending"); + expect(requests[2].run).toMatchObject({ + status: "done", + outcome: "failed", + errorCode: "agent_version_unavailable", + }); + expect(requests[2].run?.error).toContain("retired"); + expect(requests[3].run?.status).toBe("pending"); + expect(requests[4].run).toMatchObject({ + status: "done", + outcome: "failed", + errorCode: "agent_deleted", + }); + expect(requests[2].run?.finishedAt).toBeInstanceOf(Date); + expect(requests[4].run?.finishedAt).toBeInstanceOf(Date); expect(warn.mock.calls.flat().join("\n")).toContain("agent_not_found"); expect(warn.mock.calls.flat().join("\n")).toContain( "agent_version_missing", ); expect(warn.mock.calls.flat().join("\n")).toContain( - "agent_version_unavailable", + "Terminalized 1 request(s)", ); + expect(warn.mock.calls.flat().join("\n")).toContain("agent_unavailable"); + }); + + it("terminalizes missing or retired exact versions even when the agent is unavailable", async () => { + const requests = [ + makeRequest("active", "paused-worker", "v1"), + makeRequest("retired", "paused-worker", "v0"), + makeRequest("missing", "paused-worker", "v9"), + ]; + const { requestCollection, agentCollection } = makeCollections(requests, [ + makeAgent( + "paused-worker", + [ + { agentVersion: "v0", queueName: "paused-queue", status: "retired" }, + { agentVersion: "v1", queueName: "paused-queue" }, + ], + { available: false }, + ), + ]); + const scheduler = new RequestScheduler( + requestCollection, + agentCollection, + () => makeQueueClient(), + ); + vi.spyOn(console, "warn").mockImplementation(() => {}); + + await (scheduler as any).dispatch(); + + expect(requests[0].run?.status).toBe("pending"); + for (const request of requests.slice(1)) { + expect(request.run).toMatchObject({ + status: "done", + outcome: "failed", + errorCode: "agent_version_unavailable", + }); + } + }); + + it("revalidates a target that registers after pending candidates are observed", async () => { + const request = makeRequest("late-registration", "worker", "v2"); + const initialAgent = makeAgent("worker", [ + { agentVersion: "v1", queueName: "queue-v1" }, + ]); + const refreshedAgent = makeAgent("worker", [ + { agentVersion: "v1", queueName: "queue-v1" }, + { agentVersion: "v2", queueName: "queue-v2" }, + ]); + const { requestCollection, agentCollection } = makeCollections( + [request], + [initialAgent], + ); + agentCollection.find + .mockImplementationOnce(() => ({ toArray: async () => [initialAgent] })) + .mockImplementationOnce(() => ({ toArray: async () => [refreshedAgent] })); + const scheduler = new RequestScheduler( + requestCollection, + agentCollection, + () => makeQueueClient(), + ); + + await (scheduler as any).dispatch(); + + expect(request.run?.status).toBe("pending"); + }); + + it("terminalizes only requests present in the inspected pending snapshot", async () => { + const requests = [makeRequest("observed", "worker", "retired")]; + const { requestCollection, agentCollection } = makeCollections(requests, [ + makeAgent("worker", [ + { + agentVersion: "retired", + queueName: "old-queue", + status: "retired", + }, + ]), + ]); + const updateMany = requestCollection.updateMany.getMockImplementation(); + requestCollection.updateMany.mockImplementationOnce(async (...args: any[]) => { + requests.push(makeRequest("arrived-late", "worker", "retired")); + return updateMany!(...args); + }); + const scheduler = new RequestScheduler( + requestCollection, + agentCollection, + () => makeQueueClient(), + ); + vi.spyOn(console, "warn").mockImplementation(() => {}); + + await (scheduler as any).dispatch(); + + expect(requests[0].run?.status).toBe("done"); + expect(requests[1].run?.status).toBe("pending"); + }); + + it("terminalization ignores a stale queuedQueueName when no queue filter is specified", async () => { + const request = makeRequest("resumed", "worker", "retired"); + request.run!.queuedQueueName = "old-queue"; + const { requestCollection, agentCollection } = makeCollections([request], [ + makeAgent("worker", [ + { + agentVersion: "retired", + queueName: "old-queue", + status: "retired", + }, + ]), + ]); + const scheduler = new RequestScheduler( + requestCollection, + agentCollection, + () => makeQueueClient(), + ); + vi.spyOn(console, "warn").mockImplementation(() => {}); + + await (scheduler as any).dispatch(); + + expect(request.run).toMatchObject({ + status: "done", + outcome: "failed", + errorCode: "agent_version_unavailable", + }); + }); + + it("dispatches an exact active version to its current queue without substituting versions", async () => { + const moved = makeRequest("moved", "worker", "v1"); + const retired = makeRequest("retired", "worker", "v0"); + const { requestCollection, agentCollection } = makeCollections( + [moved, retired], + [ + makeAgent("worker", [ + { agentVersion: "v0", queueName: "old-queue", status: "retired" }, + { agentVersion: "v1", queueName: "new-queue" }, + { agentVersion: "v2", queueName: "newest-queue" }, + ]), + ], + ); + const newQueue = makeQueueClient(); + const newestQueue = makeQueueClient(); + const scheduler = new RequestScheduler( + requestCollection, + agentCollection, + (queueName) => (queueName === "new-queue" ? newQueue : newestQueue), + { targetQueueDepth: 1 }, + ); + + await (scheduler as any).dispatch(); + + expect(moved.run?.status).toBe("queued"); + expect(moved.run?.queuedQueueName).toBe("new-queue"); + expect(retired.run).toMatchObject({ + status: "done", + outcome: "failed", + errorCode: "agent_version_unavailable", + }); + expect(newQueue.sendMessage).toHaveBeenCalledTimes(1); + expect(newestQueue.sendMessage).not.toHaveBeenCalled(); }); it("does not route unavailable, deleted, or queue-less registry targets", async () => { @@ -508,9 +704,13 @@ describe("RequestScheduler", () => { await (scheduler as any).dispatch(); expect(factory).not.toHaveBeenCalled(); - expect(requests.every((request) => request.run?.status === "pending")).toBe( - true, - ); + expect(requests[0].run?.status).toBe("pending"); + expect(requests[1].run).toMatchObject({ + status: "done", + outcome: "failed", + errorCode: "agent_deleted", + }); + expect(requests[2].run?.status).toBe("pending"); }); it("returns a claim to pending when queue send fails", async () => { diff --git a/apps/scheduler/src/request-scheduler.ts b/apps/scheduler/src/request-scheduler.ts index 8d253424..cbd33df9 100644 --- a/apps/scheduler/src/request-scheduler.ts +++ b/apps/scheduler/src/request-scheduler.ts @@ -113,7 +113,7 @@ export class RequestScheduler { } } - await this.maybeReportInvalidPendingTargets(agents, queueTargets); + await this.maybeReportInvalidPendingTargets(); } catch (error) { console.error("[Scheduler] Failed to refresh agent registry:", error); trackEvent({ @@ -232,10 +232,7 @@ export class RequestScheduler { })); } - private async maybeReportInvalidPendingTargets( - agents: CodingAgentDocument[], - queueTargets: QueueTarget[], - ): Promise { + private async maybeReportInvalidPendingTargets(): Promise { const now = Date.now(); if ( now - this.lastInvalidTargetReportAt < @@ -246,7 +243,7 @@ export class RequestScheduler { this.lastInvalidTargetReportAt = now; try { - await this.reportInvalidPendingTargets(agents, queueTargets); + await this.reportInvalidPendingTargets(); } catch (error) { console.error("[Scheduler] Failed to inspect invalid pending targets:", error); trackEvent({ @@ -522,21 +519,12 @@ export class RequestScheduler { } } - private async reportInvalidPendingTargets( - agents: CodingAgentDocument[], - queueTargets: QueueTarget[], - ): Promise { - const validTargets = new Set( - queueTargets.flatMap((queue) => - queue.targets.map((target) => - this.targetKey(target.workerType, target.agentVersion), - ), - ), - ); + private async reportInvalidPendingTargets(): Promise { const pendingTargets = await this.requestCollection .aggregate<{ _id: { workerType?: string; agentVersion?: string }; count: number; + requestIds: string[]; }>([ { $match: { @@ -551,11 +539,25 @@ export class RequestScheduler { agentVersion: "$agentVersion", }, count: { $sum: 1 }, + requestIds: { $push: "$_id" }, }, }, ]) .toArray(); + // Re-read the registry only after collecting the pending candidates. The + // dispatch cycle may have spent time doing queue I/O, during which a new + // agent/version can register. Using a fresh registry here avoids + // terminalizing requests based on the cycle's older snapshot. + const agents = await this.agentCollection.find({}).toArray(); + const queueTargets = this.buildQueueTargets(agents); + const validTargets = new Set( + queueTargets.flatMap((queue) => + queue.targets.map((target) => + this.targetKey(target.workerType, target.agentVersion), + ), + ), + ); const agentsById = new Map(agents.map((agent) => [agent._id, agent])); const invalid = pendingTargets .filter( @@ -564,7 +566,7 @@ export class RequestScheduler { this.targetKey(_id.workerType ?? "", _id.agentVersion ?? ""), ), ) - .map(({ _id, count }) => { + .map(({ _id, count, requestIds }) => { const workerType = _id.workerType ?? "(missing)"; const agentVersion = _id.agentVersion ?? "(missing)"; const agent = agentsById.get(_id.workerType ?? ""); @@ -576,19 +578,30 @@ export class RequestScheduler { ) { reason = "agent_queue_conflict"; } else if (agent?.deletedAt) reason = "agent_deleted"; - else if (agent && agent.available !== true) reason = "agent_unavailable"; else if (agent && !_id.agentVersion) reason = "agent_version_missing"; else if (agent) { const version = (agent.versions ?? []).find( (candidate) => candidate.agentVersion === _id.agentVersion, ); if (!version || version.status !== "active") { + // Missing/retired exact versions are permanently invalid even when + // the agent itself is temporarily unavailable. reason = "agent_version_unavailable"; + } else if (agent.available !== true) { + reason = "agent_unavailable"; } else if (!version.queueName?.trim()) { reason = "agent_queue_missing"; } } - return { workerType, agentVersion, count, reason }; + return { + workerType, + agentVersion, + rawWorkerType: _id.workerType, + rawAgentVersion: _id.agentVersion, + count, + requestIds, + reason, + }; }) .sort( (left, right) => @@ -596,6 +609,51 @@ export class RequestScheduler { left.agentVersion.localeCompare(right.agentVersion), ); + // Requests pinned to targets that can never become runnable again should + // not remain pending forever. Keep recoverable registry/configuration + // states pending, but terminalize deleted agents and missing/retired exact + // versions. The status filter makes this safe against concurrent dispatch. + const terminalizedCounts = new Map(); + for (const target of invalid) { + if ( + target.reason !== "agent_deleted" && + target.reason !== "agent_version_unavailable" + ) { + continue; + } + if (!target.rawWorkerType || !target.rawAgentVersion) continue; + + const finishedAt = new Date(); + const error = + target.reason === "agent_deleted" + ? `Agent "${target.rawWorkerType}" is deleted and cannot run this request` + : `Agent version "${target.rawAgentVersion}" for "${target.rawWorkerType}" is no longer available`; + const result = await this.requestCollection.updateMany( + { + _id: { $in: target.requestIds }, + workerType: target.rawWorkerType, + agentVersion: target.rawAgentVersion, + "run.status": "pending", + deletedAt: { $exists: false }, + } as any, + { + $set: { + "run.status": "done", + "run.outcome": "failed", + "run.error": error, + "run.errorCode": target.reason, + "run.finishedAt": finishedAt, + "run.updatedAt": finishedAt, + updatedAt: finishedAt, + }, + } as any, + ); + terminalizedCounts.set( + this.targetKey(target.workerType, target.agentVersion), + result.modifiedCount, + ); + } + const signature = JSON.stringify( invalid.map(({ workerType, agentVersion, reason }) => ({ workerType, @@ -607,6 +665,32 @@ export class RequestScheduler { this.invalidTargetSignature = signature; for (const target of invalid) { + const terminalized = + terminalizedCounts.get( + this.targetKey(target.workerType, target.agentVersion), + ) ?? 0; + if (terminalized > 0) { + console.warn( + `[Scheduler] Terminalized ${terminalized} request(s) for permanently invalid target ` + + `worker=${target.workerType}, version=${target.agentVersion}: ${target.reason}`, + ); + trackEvent({ + name: "scheduler.invalid_pending_target_terminalized", + properties: { + workerType: target.workerType, + agentVersion: target.agentVersion, + reason: target.reason, + terminalizedCount: String(terminalized), + }, + }); + trackMetric({ + name: "scheduler.invalid_pending_requests_terminalized", + value: terminalized, + properties: { reason: target.reason }, + }); + continue; + } + console.warn( `[Scheduler] Leaving ${target.count} request(s) pending for invalid target ` + `worker=${target.workerType}, version=${target.agentVersion}: ${target.reason}`,