Skip to content
228 changes: 214 additions & 14 deletions apps/scheduler/src/request-scheduler.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -85,6 +85,7 @@ function makeCollections(
queuedQueueName?: string;
};
count: number;
requestIds: string[];
}
>();
for (const request of requests.filter(
Expand All @@ -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,
Expand All @@ -108,6 +111,7 @@ function makeCollections(
...(queuedQueueName ? { queuedQueueName } : {}),
},
count: 1,
requestIds: [request._id],
});
}
}
Expand Down Expand Up @@ -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<string, any>;
$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;
Comment on lines +185 to +186

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Addressed on 739d92f. The test collection mock now treats an omitted run.queuedQueueName filter as unrestricted, matching MongoDB production behavior. Added coverage for a pending request retaining an old queuedQueueName so the terminalization path is exercised correctly.

if (filter._id?.$in && !filter._id.$in.includes(candidate._id)) return false;
if (
typeof filter.workerType === "string" &&
candidate.workerType !== filter.workerType
Expand All @@ -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 };
Expand Down Expand Up @@ -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", [
Expand All @@ -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(() => {});
Expand All @@ -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 () => {
Expand Down Expand Up @@ -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 () => {
Expand Down
Loading