diff --git a/services/api/src/agents/scheduler.ts b/services/api/src/agents/scheduler.ts index 6fbdd211..be903073 100644 --- a/services/api/src/agents/scheduler.ts +++ b/services/api/src/agents/scheduler.ts @@ -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), @@ -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) { @@ -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) @@ -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, + }; +} diff --git a/services/api/test/agent-automation.integration.test.ts b/services/api/test/agent-automation.integration.test.ts index dd79f30b..b67310c5 100644 --- a/services/api/test/agent-automation.integration.test.ts +++ b/services/api/test/agent-automation.integration.test.ts @@ -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) { diff --git a/services/api/test/agent-scheduler.test.ts b/services/api/test/agent-scheduler.test.ts index 0f9dd33f..9959f08c 100644 --- a/services/api/test/agent-scheduler.test.ts +++ b/services/api/test/agent-scheduler.test.ts @@ -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", () => { @@ -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); + }); +});