Skip to content

Commit 8fe554a

Browse files
carderneTrigger.dev RepoOps
authored andcommitted
fix(webapp): keep scheduled last run times accurate
Mono-RevId: b455565ff705ccad4d1bba3fe3e02b25fbcd063e
1 parent 8f15fa5 commit 8fe554a

5 files changed

Lines changed: 188 additions & 51 deletions

File tree

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,6 @@
1+
---
2+
area: webapp
3+
type: fix
4+
---
5+
6+
Keep scheduled task "Last run" times stable across unchanged deployments and align them with configured schedule windows.

‎apps/webapp/app/v3/scheduleTimings.server.ts‎

Lines changed: 53 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -88,7 +88,10 @@ export function resolveScheduleTimings(
8888
{ phaseSecret, includeLastRun, now = new Date() }: ResolveScheduleTimingsOptions
8989
): ScheduleTiming[] {
9090
const nominalCache = new Map<string, Date[]>();
91-
const previousCache = new Map<string, Date | undefined>();
91+
const previousCache = new Map<
92+
string,
93+
{ latestNominal: Date; previousNominal?: Date } | undefined
94+
>();
9295

9396
return inputs.map((input) => {
9497
const window: NormalizedScheduleWindow | undefined = resolveScheduleWindow({
@@ -134,51 +137,80 @@ export function resolveScheduleTimings(
134137
return {
135138
nextRun: nominalAt,
136139
nextRunEffectiveAt: effectiveAt,
137-
lastRun: includeLastRun ? resolveLastRun(input, now, previousCache) : undefined,
140+
lastRun: includeLastRun
141+
? resolveLastRun(input, now, nominalAt, phase, window, previousCache)
142+
: undefined,
138143
};
139144
});
140145
}
141146

142147
/**
143-
* Approximates "last run" from the cron's previous slot.
148+
* Approximates "last run" from the most recent effective schedule time.
144149
*
145-
* Skips inactive schedules — the previous slot reflects what *would* have
146-
* fired. Skips slots that predate `updatedAt`: any config change (cron edited,
147-
* timezone changed, deactivate/reactivate) bumps `updatedAt`, and a slot from
148-
* before the most recent change didn't fire under the current configuration.
149-
*
150-
* `cron-parser` throws on malformed expressions, so this degrades to undefined
151-
* per row rather than failing the whole list. Best-effort by design; the runs
152-
* page is the source of truth.
150+
* Skips inactive schedules and effective times that predate `updatedAt`. Best-effort by design;
151+
* the runs page is the source of truth.
153152
*/
154153
function resolveLastRun(
155154
input: ScheduleTimingInput,
156155
now: Date,
157-
cache: Map<string, Date | undefined>
156+
nextNominal: Date,
157+
phase: number,
158+
window: NormalizedScheduleWindow | undefined,
159+
cache: Map<string, { latestNominal: Date; previousNominal?: Date } | undefined>
158160
): Date | undefined {
159161
if (!input.active) {
160162
return undefined;
161163
}
162164

163165
const key = cacheKey(input.cron, input.timezone);
166+
let nominalTimes = cache.get(key);
164167

165-
let previous: Date | undefined;
166-
if (cache.has(key)) {
167-
previous = cache.get(key);
168-
} else {
168+
if (!cache.has(key)) {
169169
try {
170-
previous = previousScheduledTimestamp(input.cron, input.timezone, now);
170+
nominalTimes = {
171+
latestNominal: previousScheduledTimestamp(
172+
input.cron,
173+
input.timezone,
174+
new Date(now.getTime() + 1)
175+
),
176+
};
171177
} catch {
172-
previous = undefined;
178+
nominalTimes = undefined;
173179
}
174-
cache.set(key, previous);
180+
cache.set(key, nominalTimes);
175181
}
176182

177-
if (!previous) {
183+
if (!nominalTimes) {
178184
return undefined;
179185
}
180186

181-
return previous.getTime() > input.updatedAt.getTime() ? previous : undefined;
187+
const latestEffective = calculateEffectiveScheduleTime({
188+
nominalAt: nominalTimes.latestNominal,
189+
nextNominalAt: nextNominal,
190+
schedulePhase: phase,
191+
window,
192+
minimumWindowDurationSeconds: input.minimumWindowDurationSeconds,
193+
}).effectiveAt;
194+
if (latestEffective.getTime() <= now.getTime()) {
195+
return latestEffective.getTime() > input.updatedAt.getTime() ? latestEffective : undefined;
196+
}
197+
198+
if (!nominalTimes.previousNominal) {
199+
nominalTimes.previousNominal = previousScheduledTimestamp(
200+
input.cron,
201+
input.timezone,
202+
nominalTimes.latestNominal
203+
);
204+
}
205+
const previousEffective = calculateEffectiveScheduleTime({
206+
nominalAt: nominalTimes.previousNominal,
207+
nextNominalAt: nominalTimes.latestNominal,
208+
schedulePhase: phase,
209+
window,
210+
minimumWindowDurationSeconds: input.minimumWindowDurationSeconds,
211+
}).effectiveAt;
212+
213+
return previousEffective.getTime() > input.updatedAt.getTime() ? previousEffective : undefined;
182214
}
183215

184216
/**

‎apps/webapp/app/v3/services/createBackgroundWorker.server.test.ts‎

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -394,6 +394,7 @@ describe("syncDeclarativeSchedules registration", () => {
394394
async ({ prisma, redisOptions }) => {
395395
const { project, prodEnv } = await seedProjectWithEnvs(prisma);
396396
const schedule = await makeDeclarativeSchedule(prisma, project.id, [prodEnv.id]);
397+
const updatedAt = schedule.updatedAt;
397398
await seedScheduledTask(prisma, project.id, prodEnv.id);
398399
const engine = createTestScheduleEngine(prisma, redisOptions);
399400

@@ -412,6 +413,12 @@ describe("syncDeclarativeSchedules registration", () => {
412413

413414
expect(before).toBeDefined();
414415
expect(await engine.getJob(jobId)).toEqual(before);
416+
expect(
417+
await prisma.taskSchedule.findUniqueOrThrow({
418+
where: { id: schedule.id },
419+
select: { updatedAt: true },
420+
})
421+
).toEqual({ updatedAt });
415422
} finally {
416423
await engine.quit();
417424
}

‎apps/webapp/app/v3/services/createBackgroundWorker.server.ts‎

Lines changed: 24 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -1399,6 +1399,7 @@ async function prepareDeclarativeSchedules(
13991399
minimumWindowDurationSeconds: true,
14001400
instances: {
14011401
select: {
1402+
id: true,
14021403
environmentId: true,
14031404
},
14041405
},
@@ -1537,21 +1538,29 @@ export async function syncDeclarativeSchedules(
15371538
!scheduleWindowsEqual(previousWindow, nextWindow) ||
15381539
// Clearing the minimum changes the effective range.
15391540
existingSchedule.minimumWindowDurationSeconds !== minimumWindowDurationSeconds;
1540-
const schedule = await prisma.taskSchedule.update({
1541-
where: {
1542-
id: existingSchedule.id,
1543-
},
1544-
data: {
1545-
generatorExpression: task.schedule.cron,
1546-
generatorDescription: cronstrue.toString(task.schedule.cron),
1547-
timezone: task.schedule.timezone,
1548-
minimumWindowDurationSeconds,
1549-
...normalizedWindow,
1550-
},
1551-
include: {
1552-
instances: true,
1553-
},
1554-
});
1541+
const persistedValuesChanged =
1542+
existingSchedule.generatorExpression !== task.schedule.cron ||
1543+
existingSchedule.timezone !== task.schedule.timezone ||
1544+
existingSchedule.windowDurationSeconds !== normalizedWindow.windowDurationSeconds ||
1545+
existingSchedule.windowPercentage !== normalizedWindow.windowPercentage ||
1546+
existingSchedule.minimumWindowDurationSeconds !== minimumWindowDurationSeconds;
1547+
const schedule = persistedValuesChanged
1548+
? await prisma.taskSchedule.update({
1549+
where: {
1550+
id: existingSchedule.id,
1551+
},
1552+
data: {
1553+
generatorExpression: task.schedule.cron,
1554+
generatorDescription: cronstrue.toString(task.schedule.cron),
1555+
timezone: task.schedule.timezone,
1556+
minimumWindowDurationSeconds,
1557+
...normalizedWindow,
1558+
},
1559+
include: {
1560+
instances: true,
1561+
},
1562+
})
1563+
: existingSchedule;
15551564

15561565
missingSchedules.delete(existingSchedule.id);
15571566
const instances = timingChanged

‎apps/webapp/test/scheduleTimings.test.ts‎

Lines changed: 98 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -4,9 +4,9 @@ import {
44
SCHEDULE_PHASE_DENOMINATOR,
55
calculateEffectiveScheduleTime,
66
calculateSchedulePhase,
7+
resolveScheduleWindow,
78
} from "@internal/schedule-engine";
89
import { CronPattern } from "~/v3/schedules";
9-
import { type NormalizedScheduleWindow } from "@trigger.dev/core/v3";
1010
import { parseExpression } from "cron-parser";
1111
import { describe, expect, it } from "vitest";
1212
import {
@@ -77,24 +77,48 @@ function referenceResolve(
7777
deduplicationKey: input.deduplicationKey,
7878
});
7979

80-
const window: NormalizedScheduleWindow | undefined =
81-
input.windowPercentage !== null
82-
? { type: "percentage", percentage: input.windowPercentage }
83-
: input.windowDurationSeconds !== null
84-
? { type: "duration", durationSeconds: input.windowDurationSeconds }
85-
: undefined;
80+
const window = resolveScheduleWindow({
81+
windowDurationSeconds: input.windowDurationSeconds,
82+
windowPercentage: input.windowPercentage,
83+
defaultWindowDurationSeconds: input.defaultWindowDurationSeconds,
84+
}).window;
8685

8786
const { effectiveAt } = calculateEffectiveScheduleTime({
8887
nominalAt: nominalTimes[0],
8988
nextNominalAt: nominalTimes[1],
9089
schedulePhase: phase,
9190
window,
91+
minimumWindowDurationSeconds: input.minimumWindowDurationSeconds,
9292
});
9393

9494
let lastRun: Date | undefined;
9595
if (includeLastRun && input.active) {
9696
try {
97-
const previous = previousScheduledTimestamp(input.cron, input.timezone, now);
97+
const latestNominal = previousScheduledTimestamp(
98+
input.cron,
99+
input.timezone,
100+
new Date(now.getTime() + 1)
101+
);
102+
const previousNominal = previousScheduledTimestamp(
103+
input.cron,
104+
input.timezone,
105+
latestNominal
106+
);
107+
const latestEffective = calculateEffectiveScheduleTime({
108+
nominalAt: latestNominal,
109+
nextNominalAt: nominalTimes[0],
110+
schedulePhase: phase,
111+
window,
112+
minimumWindowDurationSeconds: input.minimumWindowDurationSeconds,
113+
}).effectiveAt;
114+
const previousEffective = calculateEffectiveScheduleTime({
115+
nominalAt: previousNominal,
116+
nextNominalAt: latestNominal,
117+
schedulePhase: phase,
118+
window,
119+
minimumWindowDurationSeconds: input.minimumWindowDurationSeconds,
120+
}).effectiveAt;
121+
const previous = latestEffective <= now ? latestEffective : previousEffective;
98122
lastRun = previous.getTime() > input.updatedAt.getTime() ? previous : undefined;
99123
} catch {
100124
lastRun = undefined;
@@ -255,15 +279,24 @@ describe("resolveScheduleTimings", () => {
255279
});
256280

257281
it("skips lastRun when the previous slot predates the last config change", () => {
258-
const [stale] = resolveScheduleTimings([input({ cron: "0 0 * * *", updatedAt: now })], {
259-
phaseSecret: PHASE_SECRET,
260-
includeLastRun: true,
261-
now,
262-
});
282+
const [stale] = resolveScheduleTimings(
283+
[input({ cron: "0 0 * * *", schedulePhase: 0, updatedAt: now })],
284+
{
285+
phaseSecret: PHASE_SECRET,
286+
includeLastRun: true,
287+
now,
288+
}
289+
);
263290
expect(stale.lastRun).toBeUndefined();
264291

265292
const [fresh] = resolveScheduleTimings(
266-
[input({ cron: "0 0 * * *", updatedAt: new Date("2020-01-01T00:00:00.000Z") })],
293+
[
294+
input({
295+
cron: "0 0 * * *",
296+
schedulePhase: 0,
297+
updatedAt: new Date("2020-01-01T00:00:00.000Z"),
298+
}),
299+
],
267300
{ phaseSecret: PHASE_SECRET, includeLastRun: true, now }
268301
);
269302
expect(fresh.lastRun).toEqual(new Date("2024-06-15T00:00:00.000Z"));
@@ -282,7 +315,7 @@ describe("resolveScheduleTimings", () => {
282315
});
283316

284317
it("resolves lastRun for a valid expression", () => {
285-
const [valid] = resolveScheduleTimings([input({ cron: "0 0 * * *" })], {
318+
const [valid] = resolveScheduleTimings([input({ cron: "0 0 * * *", schedulePhase: 0 })], {
286319
phaseSecret: PHASE_SECRET,
287320
includeLastRun: true,
288321
now,
@@ -291,6 +324,56 @@ describe("resolveScheduleTimings", () => {
291324
expect(valid.lastRun).toEqual(new Date("2024-06-15T00:00:00.000Z"));
292325
});
293326

327+
it("resolves lastRun from the most recent effective window time", () => {
328+
const [beforeCurrentWindow] = resolveScheduleTimings(
329+
[
330+
input({
331+
cron: "0 * * * *",
332+
schedulePhase: SCHEDULE_PHASE_DENOMINATOR / 2,
333+
windowPercentage: 100,
334+
}),
335+
],
336+
{
337+
phaseSecret: PHASE_SECRET,
338+
includeLastRun: true,
339+
now,
340+
}
341+
);
342+
const [afterCurrentWindow] = resolveScheduleTimings(
343+
[
344+
input({
345+
cron: "0 * * * *",
346+
schedulePhase: SCHEDULE_PHASE_DENOMINATOR / 2,
347+
windowPercentage: 100,
348+
}),
349+
],
350+
{
351+
phaseSecret: PHASE_SECRET,
352+
includeLastRun: true,
353+
now: new Date("2024-06-15T09:40:00.000Z"),
354+
}
355+
);
356+
357+
expect(beforeCurrentWindow.lastRun).toEqual(new Date("2024-06-15T08:30:00.000Z"));
358+
expect(afterCurrentWindow.lastRun).toEqual(new Date("2024-06-15T09:30:00.000Z"));
359+
});
360+
361+
it("compares updatedAt with the effective window time", () => {
362+
const [timing] = resolveScheduleTimings(
363+
[
364+
input({
365+
cron: "0 * * * *",
366+
schedulePhase: SCHEDULE_PHASE_DENOMINATOR / 2,
367+
windowPercentage: 100,
368+
updatedAt: new Date("2024-06-15T08:15:00.000Z"),
369+
}),
370+
],
371+
{ phaseSecret: PHASE_SECRET, includeLastRun: true, now }
372+
);
373+
374+
expect(timing.lastRun).toEqual(new Date("2024-06-15T08:30:00.000Z"));
375+
});
376+
294377
it("honours a caller-supplied schedulePhase over the derived one", () => {
295378
const rows = [input({ schedulePhase: 0, windowPercentage: 100, cron: "0 * * * *" })];
296379

0 commit comments

Comments
 (0)