Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
34 changes: 30 additions & 4 deletions services/api/src/agents/scheduler.ts
Original file line number Diff line number Diff line change
Expand Up @@ -39,9 +39,12 @@ export class AgentScheduler {
.from(agentSchedules)
.where(and(eq(agentSchedules.enabled, true), lte(agentSchedules.nextRunAt, now)));
let scheduled = 0;
let coalesced = 0;
for (const schedule of due) {
const claimed = await this.claim(schedule, now);
const advance = scheduleAdvance(schedule.cron, schedule.timezone, schedule.nextRunAt, now);
const claimed = await this.claim(schedule, now, advance.nextRunAt);
if (!claimed) continue;
if (advance.missed) coalesced += 1;
try {
const [projectManifest, projection] = await Promise.all([
this.projectManifests.load(schedule.orgId, schedule.projectId),
Expand Down Expand Up @@ -75,7 +78,7 @@ export class AgentScheduler {
});
}
}
return { projects: activeProjects.length, due: due.length, scheduled, failures };
return { projects: activeProjects.length, due: due.length, scheduled, coalesced, failures };
}

async status(orgId: string, projectId: string) {
Expand Down Expand Up @@ -200,8 +203,7 @@ export class AgentScheduler {
});
}

private async claim(schedule: typeof agentSchedules.$inferSelect, now: Date) {
const nextRunAt = nextOccurrence(schedule.cron, schedule.timezone, schedule.nextRunAt);
private async claim(schedule: typeof agentSchedules.$inferSelect, now: Date, nextRunAt: Date) {
return (
await this.db
.update(agentSchedules)
Expand All @@ -226,3 +228,27 @@ export class AgentScheduler {
export function nextOccurrence(cron: string, timezone: string, from: Date) {
return cronParser.parseExpression(cron, { currentDate: from, tz: timezone }).next().toDate();
}
/**
* Where a due schedule should point once its occurrence is claimed.
*
* The next occurrence is computed from `now`, not from the occurrence being
* claimed. Advancing by one cron step from the stored value makes the loop
* edge-triggered: a schedule that fell behind while the worker was down stays
* due after each claim and dispatches again on the following tick, one paid run
* per missed occurrence — 24 turns in 24 minutes for an hourly schedule after a
* day of downtime. Advancing from `now` makes it level-triggered: the loop
* converges on "this schedule is due" and runs it once, however long the gap.
*
* A schedule claimed on time is unaffected, because the next occurrence after
* `now` and the next occurrence after its own due instant are the same one.
*/
export function scheduleAdvance(cron: string, timezone: string, dueAt: Date, now: Date) {
const from = now > dueAt ? now : dueAt;
return {
nextRunAt: nextOccurrence(cron, timezone, from),
// Whether this claim absorbs more than the occurrence it satisfies: the
// following occurrence had already come due as well. One extra cron step,
// never a walk over the backlog, so an outage of any length costs the same.
missed: nextOccurrence(cron, timezone, dueAt) <= now,
};
}
54 changes: 54 additions & 0 deletions services/api/test/agent-automation.integration.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -939,6 +939,60 @@ describe("agent automations use persistent story workspaces", async () => {
await db.select().from(storyMessages).where(eq(storyMessages.storyId, scheduledStory.id)),
).toHaveLength(1);
});

it("runs a schedule that fell behind once, not once per missed occurrence", async () => {
// security-audit:nightly is `0 2 * * *` UTC. Put it a week in arrears, the
// shape of a worker that was down: seven occurrences have come due.
const dueAt = new Date("2026-01-03T02:00:00.000Z");
const now = new Date("2026-01-10T02:00:00.000Z");
const scheduleRow = and(
eq(agentSchedules.projectId, projectId),
eq(agentSchedules.agentName, "security-audit"),
eq(agentSchedules.triggerName, "nightly"),
);
await db
.update(agentSchedules)
.set({ nextRunAt: dueAt, lastScheduledAt: null })
.where(scheduleRow);

const scheduledStory = (
await db
.select()
.from(stories)
.where(
and(
eq(stories.projectId, projectId),
eq(stories.provider, "schedule"),
eq(stories.externalId, "security-audit:nightly"),
),
)
)[0];
if (!scheduledStory) throw new Error("expected scheduled story");
const messagesBefore = (
await db.select().from(storyMessages).where(eq(storyMessages.storyId, scheduledStory.id))
).length;

// Three ticks at the same instant: the worker runs `* * * * *`, so the
// backlog would drain a turn per minute until it caught up.
const results = [];
for (let tick = 0; tick < 3; tick += 1) results.push(await scheduler.tick(now));

// `scheduled` is attributable to this fixture: the catalog and manifest
// sources throw for any project but this one, so a schedule left in the
// shared test database by another suite becomes a failure, never a run.
// The tick counters are global for the same reason, so the rest of the
// assertions read this project's own rows.
expect(results.reduce((sum, result) => sum + result.scheduled, 0)).toBe(1);
expect(
await db.select().from(storyMessages).where(eq(storyMessages.storyId, scheduledStory.id)),
).toHaveLength(messagesBefore + 1);

// The claim satisfies the occurrence it observed and leaves the schedule
// ahead of the clock, so the catch-up cannot restart on the next tick.
const [after] = await db.select().from(agentSchedules).where(scheduleRow);
expect(after?.lastScheduledAt).toEqual(dueAt);
expect(after?.nextRunAt).toEqual(new Date("2026-01-11T02:00:00.000Z"));
});
});

function render(agent: ReturnType<typeof manifest>) {
Expand Down
74 changes: 73 additions & 1 deletion services/api/test/agent-scheduler.test.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
import { describe, expect, it } from "vitest";
import { nextOccurrence } from "../src/agents/scheduler.js";
import { nextOccurrence, scheduleAdvance } from "../src/agents/scheduler.js";

describe("agent scheduler clock", () => {
it("uses the trigger timezone across daylight-saving boundaries", () => {
Expand All @@ -23,3 +23,75 @@ describe("agent scheduler clock", () => {
expect(first).toEqual(new Date("2026-07-15T07:15:00.000Z"));
});
});

describe("agent scheduler catch-up", () => {
const hourly = ["0 * * * *", "UTC"] as const;

it("leaves an on-time claim exactly where it was", () => {
const dueAt = new Date("2026-09-07T06:00:00.000Z");
// The worker ticks every minute, so a claim lands within the same period as
// the occurrence it satisfies rather than precisely on it.
for (const now of [dueAt, new Date("2026-09-07T06:00:07.000Z")]) {
const advance = scheduleAdvance(...hourly, dueAt, now);
expect(advance.nextRunAt).toEqual(new Date("2026-09-07T07:00:00.000Z"));
expect(advance.nextRunAt).toEqual(nextOccurrence(...hourly, dueAt));
expect(advance.missed).toBe(false);
}
});

it("collapses a backlog of any size into a single claim", () => {
const dueAt = new Date("2026-09-06T06:00:00.000Z");
const now = new Date("2026-09-07T06:00:00.000Z");

const advance = scheduleAdvance(...hourly, dueAt, now);
expect(advance.nextRunAt).toEqual(new Date("2026-09-07T07:00:00.000Z"));
expect(advance.missed).toBe(true);

// The loop the tick performs: claim while due. Advancing one cron step from
// the stored occurrence dispatches once per missed hour; advancing from now
// leaves the schedule ahead of the clock after a single claim.
let claims = 0;
let nextRunAt = dueAt;
while (nextRunAt <= now) {
nextRunAt = scheduleAdvance(...hourly, nextRunAt, now).nextRunAt;
claims += 1;
}
expect(claims).toBe(1);
expect(nextRunAt.getTime()).toBeGreaterThan(now.getTime());
});

it("reports a single missed occurrence, not only a long outage", () => {
// One tick late by more than a period: still a coalesced claim, and the
// difference between this and a normal minute is what the counter carries.
const advance = scheduleAdvance(
"*/5 * * * *",
"UTC",
new Date("2026-09-07T06:00:00.000Z"),
new Date("2026-09-07T06:07:00.000Z"),
);
expect(advance.nextRunAt).toEqual(new Date("2026-09-07T06:10:00.000Z"));
expect(advance.missed).toBe(true);
});

it("keeps the trigger timezone when catching up across a daylight-saving shift", () => {
// The outage spans spring-forward. 02:00 local does not exist on 2026-03-08
// in New York, so a UTC-arithmetic catch-up would land an hour off.
const advance = scheduleAdvance(
"0 2 * * *",
"America/New_York",
new Date("2026-03-06T07:00:00.000Z"),
new Date("2026-03-09T12:00:00.000Z"),
);
expect(advance.nextRunAt).toEqual(new Date("2026-03-10T06:00:00.000Z"));
expect(advance.missed).toBe(true);
});

it("never moves a schedule backwards when the clock is behind its due instant", () => {
// Defensive: the tick only selects nextRunAt <= now, but the function is
// total and must not hand back an occurrence the claim would re-fire.
const dueAt = new Date("2026-09-07T06:00:00.000Z");
const advance = scheduleAdvance(...hourly, dueAt, new Date("2026-09-07T05:30:00.000Z"));
expect(advance.nextRunAt).toEqual(new Date("2026-09-07T07:00:00.000Z"));
expect(advance.missed).toBe(false);
});
});
Loading