diff --git a/server/src/__tests__/routines-service.test.ts b/server/src/__tests__/routines-service.test.ts index a68955595a..9dccbabc0c 100644 --- a/server/src/__tests__/routines-service.test.ts +++ b/server/src/__tests__/routines-service.test.ts @@ -2059,6 +2059,100 @@ describeEmbeddedPostgres("routine service live-execution coalescing", () => { expect(newRuns).toMatchObject([{ status: "issue_created" }]); }); + it("coalesces multiple missed sub-hourly ticks into one catch-up run", async () => { + const { routine, svc } = await seedFixture(); + await db.update(routines).set({ + catchUpPolicy: "enqueue_missed_with_cap", + }).where(eq(routines.id, routine.id)); + const { trigger } = await svc.createTrigger(routine.id, { + kind: "schedule", + cronExpression: "*/10 * * * *", + timezone: "UTC", + }, {}); + await db.update(routineTriggers).set({ + nextRunAt: new Date("2026-07-16T00:00:00.000Z"), + }).where(eq(routineTriggers.id, trigger.id)); + + expect(await svc.tickScheduledTriggers(new Date("2026-07-16T01:05:00.000Z"))).toEqual({ triggered: 1 }); + + const runs = await db.select().from(routineRuns).where(eq(routineRuns.routineId, routine.id)); + expect(runs).toHaveLength(1); + expect(runs[0]?.status).toBe("issue_created"); + const updatedTrigger = await db.select().from(routineTriggers).where(eq(routineTriggers.id, trigger.id)).then((rows) => rows[0]); + expect(updatedTrigger?.nextRunAt).toEqual(new Date("2026-07-16T01:10:00.000Z")); + }); + + it("continues replaying each missed hourly tick", async () => { + const { routine, svc } = await seedFixture(); + await db.update(routines).set({ + catchUpPolicy: "enqueue_missed_with_cap", + }).where(eq(routines.id, routine.id)); + const { trigger } = await svc.createTrigger(routine.id, { + kind: "schedule", + cronExpression: "0 * * * *", + timezone: "UTC", + }, {}); + await db.update(routineTriggers).set({ + nextRunAt: new Date("2026-07-16T00:00:00.000Z"), + }).where(eq(routineTriggers.id, trigger.id)); + + expect(await svc.tickScheduledTriggers(new Date("2026-07-16T02:30:00.000Z"))).toEqual({ triggered: 3 }); + + const runs = await db.select().from(routineRuns).where(eq(routineRuns.routineId, routine.id)); + expect(runs).toHaveLength(3); + expect(runs.filter((run) => run.status === "issue_created")).toHaveLength(1); + expect(runs.filter((run) => run.status === "coalesced")).toHaveLength(2); + const updatedTrigger = await db.select().from(routineTriggers).where(eq(routineTriggers.id, trigger.id)).then((rows) => rows[0]); + expect(updatedTrigger?.nextRunAt).toEqual(new Date("2026-07-16T03:00:00.000Z")); + }); + + it("continues replaying missed ticks for daily schedules with multiple minute values", async () => { + const { routine, svc } = await seedFixture(); + await db.update(routines).set({ + catchUpPolicy: "enqueue_missed_with_cap", + }).where(eq(routines.id, routine.id)); + const { trigger } = await svc.createTrigger(routine.id, { + kind: "schedule", + cronExpression: "0,30 9 * * *", + timezone: "UTC", + }, {}); + await db.update(routineTriggers).set({ + nextRunAt: new Date("2026-07-14T09:00:00.000Z"), + }).where(eq(routineTriggers.id, trigger.id)); + + expect(await svc.tickScheduledTriggers(new Date("2026-07-15T10:00:00.000Z"))).toEqual({ triggered: 4 }); + + const runs = await db.select().from(routineRuns).where(eq(routineRuns.routineId, routine.id)); + expect(runs).toHaveLength(4); + expect(runs.filter((run) => run.status === "issue_created")).toHaveLength(1); + expect(runs.filter((run) => run.status === "coalesced")).toHaveLength(3); + const updatedTrigger = await db.select().from(routineTriggers).where(eq(routineTriggers.id, trigger.id)).then((rows) => rows[0]); + expect(updatedTrigger?.nextRunAt).toEqual(new Date("2026-07-16T09:00:00.000Z")); + }); + + it("coalesces sub-hourly schedules restricted to weekdays", async () => { + const { routine, svc } = await seedFixture(); + await db.update(routines).set({ + catchUpPolicy: "enqueue_missed_with_cap", + }).where(eq(routines.id, routine.id)); + const { trigger } = await svc.createTrigger(routine.id, { + kind: "schedule", + cronExpression: "*/10 * * * 1-5", + timezone: "UTC", + }, {}); + await db.update(routineTriggers).set({ + nextRunAt: new Date("2026-07-13T00:00:00.000Z"), + }).where(eq(routineTriggers.id, trigger.id)); + + expect(await svc.tickScheduledTriggers(new Date("2026-07-13T01:05:00.000Z"))).toEqual({ triggered: 1 }); + + const runs = await db.select().from(routineRuns).where(eq(routineRuns.routineId, routine.id)); + expect(runs).toHaveLength(1); + expect(runs[0]?.status).toBe("issue_created"); + const updatedTrigger = await db.select().from(routineTriggers).where(eq(routineTriggers.id, trigger.id)).then((rows) => rows[0]); + expect(updatedTrigger?.nextRunAt).toEqual(new Date("2026-07-13T01:10:00.000Z")); + }); + it("applies the armed cutoff to webhook dispatch but not manual API runs", async () => { const runtimeEnv = { PAPERCLIP_IN_WORKTREE: "true", PAPERCLIP_INSTANCE_ID: "worktree-routines-test" }; const { routine, svc } = await seedFixture({ runtimeEnv }); diff --git a/server/src/services/routines.ts b/server/src/services/routines.ts index 4808decfcc..e64da63290 100644 --- a/server/src/services/routines.ts +++ b/server/src/services/routines.ts @@ -239,6 +239,22 @@ export function nextCronTickInTimeZone(expression: string, timeZone: string, aft return null; } +function isSubHourlyCronExpression(expression: string, timeZone: string, after: Date) { + const firstTick = nextCronTickInTimeZone(expression, timeZone, after); + if (!firstTick) return false; + + const windowEnd = firstTick.getTime() + 24 * 60 * 60 * 1000; + let occurrenceCount = 1; + let cursor = firstTick; + while (occurrenceCount <= 24) { + const nextTick = nextCronTickInTimeZone(expression, timeZone, cursor); + if (!nextTick || nextTick.getTime() >= windowEnd) return false; + occurrenceCount += 1; + cursor = nextTick; + } + return true; +} + function nextResultText(status: string, issueId?: string | null) { if (status === "issue_created" && issueId) return `Created execution issue ${issueId}`; if (status === "coalesced") return "Coalesced into an existing live execution issue"; @@ -1594,6 +1610,7 @@ export function routineService( executionWorkspacePreference?: string | null; executionWorkspaceSettings?: Record | null; descriptionAppendix?: string | null; + nextRunAtOverride?: Date | null; actor?: Actor; }) { const projectId = input.projectId ?? input.routine.projectId ?? null; @@ -1718,9 +1735,11 @@ export function routineService( }) .returning(); - const nextRunAt = input.trigger?.kind === "schedule" && input.trigger.cronExpression && input.trigger.timezone - ? nextCronTickInTimeZone(input.trigger.cronExpression, input.trigger.timezone, triggeredAt) - : undefined; + const nextRunAt = input.nextRunAtOverride !== undefined + ? input.nextRunAtOverride + : input.trigger?.kind === "schedule" && input.trigger.cronExpression && input.trigger.timezone + ? nextCronTickInTimeZone(input.trigger.cronExpression, input.trigger.timezone, triggeredAt) + : undefined; let createdIssue: Awaited> | null = null; try { @@ -2936,12 +2955,16 @@ export function routineService( let claimedNextRunAt = nextCronTickInTimeZone(row.trigger.cronExpression, row.trigger.timezone, now); if (!projectPaused && !worktreeSuppressed && row.routine.catchUpPolicy === "enqueue_missed_with_cap") { - let cursor: Date | null = row.trigger.nextRunAt; - runCount = 0; - while (cursor && cursor <= now && runCount < MAX_CATCH_UP_RUNS) { - runCount += 1; - claimedNextRunAt = nextCronTickInTimeZone(row.trigger.cronExpression, row.trigger.timezone, cursor); - cursor = claimedNextRunAt; + if (isSubHourlyCronExpression(row.trigger.cronExpression, row.trigger.timezone, now)) { + claimedNextRunAt = nextCronTickInTimeZone(row.trigger.cronExpression, row.trigger.timezone, now); + } else { + let cursor: Date | null = row.trigger.nextRunAt; + runCount = 0; + while (cursor && cursor <= now && runCount < MAX_CATCH_UP_RUNS) { + runCount += 1; + claimedNextRunAt = nextCronTickInTimeZone(row.trigger.cronExpression, row.trigger.timezone, cursor); + cursor = claimedNextRunAt; + } } } @@ -2999,6 +3022,7 @@ export function routineService( routine: row.routine, trigger: row.trigger, source: "schedule", + nextRunAtOverride: claimedNextRunAt, }); triggered += 1; } diff --git a/ui/src/components/routine-sections/editable-sections.tsx b/ui/src/components/routine-sections/editable-sections.tsx index 12f6220d74..d00d6742a3 100644 --- a/ui/src/components/routine-sections/editable-sections.tsx +++ b/ui/src/components/routine-sections/editable-sections.tsx @@ -64,7 +64,7 @@ const catchUpPolicyOptions = [ { value: "enqueue_missed_with_cap", title: "Enqueue missed with cap", - description: "Catch up missed schedule windows in capped batches after recovery.", + description: "Catch up missed schedule windows after recovery; sub-hourly schedules are combined into one catch-up run, slower schedules replay each missed window up to a cap.", }, ]; diff --git a/ui/src/pages/Routines.tsx b/ui/src/pages/Routines.tsx index 837b6aac4f..0b7e33d68e 100644 --- a/ui/src/pages/Routines.tsx +++ b/ui/src/pages/Routines.tsx @@ -56,7 +56,7 @@ const concurrencyPolicyDescriptions: Record = { }; const catchUpPolicyDescriptions: Record = { skip_missed: "Ignore windows that were missed while the scheduler or routine was paused.", - enqueue_missed_with_cap: "Catch up missed schedule windows in capped batches after recovery.", + enqueue_missed_with_cap: "Catch up missed schedule windows after recovery; sub-hourly schedules are combined into one catch-up run, slower schedules replay each missed window up to a cap.", }; function autoResizeTextarea(element: HTMLTextAreaElement | null) {