Skip to content

Commit 466cfcd

Browse files
committed
preserve existing jobs when schedule unchanged
1 parent 0044bbf commit 466cfcd

5 files changed

Lines changed: 155 additions & 11 deletions

File tree

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

Lines changed: 11 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -705,6 +705,12 @@ export async function syncDeclarativeSchedules(
705705
);
706706

707707
if (existingSchedule) {
708+
const normalizedWindow = normalizeScheduleWindow(task.schedule.window);
709+
const timingChanged =
710+
existingSchedule.generatorExpression !== task.schedule.cron ||
711+
existingSchedule.timezone !== task.schedule.timezone ||
712+
existingSchedule.windowDurationSeconds !== normalizedWindow.windowDurationSeconds ||
713+
existingSchedule.windowPercentage !== normalizedWindow.windowPercentage;
708714
const schedule = await prisma.taskSchedule.update({
709715
where: {
710716
id: existingSchedule.id,
@@ -713,7 +719,7 @@ export async function syncDeclarativeSchedules(
713719
generatorExpression: task.schedule.cron,
714720
generatorDescription: cronstrue.toString(task.schedule.cron),
715721
timezone: task.schedule.timezone,
716-
...normalizeScheduleWindow(task.schedule.window),
722+
...normalizedWindow,
717723
},
718724
include: {
719725
instances: true,
@@ -723,7 +729,10 @@ export async function syncDeclarativeSchedules(
723729
missingSchedules.delete(existingSchedule.id);
724730
const instance = schedule.instances.at(0);
725731
if (instance) {
726-
await scheduleEngine.registerNextTaskScheduleInstance({ instanceId: instance.id });
732+
await scheduleEngine.registerNextTaskScheduleInstance({
733+
instanceId: instance.id,
734+
preserveExistingJob: !timingChanged,
735+
});
727736
} else {
728737
throw new CreateDeclarativeScheduleError(
729738
`Missing instance for declarative schedule ${schedule.id}`

apps/webapp/test/syncDeclarativeSchedules.test.ts

Lines changed: 85 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4,8 +4,17 @@ import { describe, expect, vi } from "vitest";
44
import type { AuthenticatedEnvironment } from "~/services/apiAuth.server";
55
import { syncDeclarativeSchedules } from "~/v3/services/createBackgroundWorker.server";
66

7+
const { registerNextTaskScheduleInstance } = vi.hoisted(() => ({
8+
registerNextTaskScheduleInstance: vi.fn().mockResolvedValue(undefined),
9+
}));
10+
11+
vi.mock("~/v3/scheduleEngine.server", () => ({
12+
scheduleEngine: { registerNextTaskScheduleInstance },
13+
}));
14+
715
vi.setConfig({ testTimeout: 60_000 });
816

17+
type TasksArg = Parameters<typeof syncDeclarativeSchedules>[0];
918
type WorkerArg = Parameters<typeof syncDeclarativeSchedules>[1];
1019
const noWorker = {} as unknown as WorkerArg;
1120

@@ -82,6 +91,82 @@ function countingPrisma(prisma: PrismaClient) {
8291
const asEnv = (env: { id: string; projectId: string; type: string }) =>
8392
env as unknown as AuthenticatedEnvironment;
8493

94+
function declarativeTasks(schedule: { cron: string; timezone: string; window?: string }): TasksArg {
95+
return [{ id: "my-task", schedule }] as TasksArg;
96+
}
97+
98+
async function seedScheduledTask(
99+
prisma: PrismaClient,
100+
projectId: string,
101+
runtimeEnvironmentId: string
102+
) {
103+
const worker = await prisma.backgroundWorker.create({
104+
data: {
105+
friendlyId: `worker_${runtimeEnvironmentId}`,
106+
contentHash: `hash_${runtimeEnvironmentId}`,
107+
version: "20260811.1",
108+
metadata: {},
109+
projectId,
110+
runtimeEnvironmentId,
111+
},
112+
});
113+
114+
await prisma.backgroundWorkerTask.create({
115+
data: {
116+
friendlyId: `task_${runtimeEnvironmentId}`,
117+
slug: "my-task",
118+
filePath: "src/trigger/my-task.ts",
119+
workerId: worker.id,
120+
projectId,
121+
runtimeEnvironmentId,
122+
triggerSource: "SCHEDULED",
123+
},
124+
});
125+
}
126+
127+
describe("syncDeclarativeSchedules registration", () => {
128+
containerTest(
129+
"preserves an existing Redis job when declarative timing is unchanged",
130+
async ({ prisma }) => {
131+
registerNextTaskScheduleInstance.mockClear();
132+
const { project, prodEnv } = await seedProjectWithEnvs(prisma);
133+
const schedule = await makeDeclarativeSchedule(prisma, project.id, [prodEnv.id]);
134+
await seedScheduledTask(prisma, project.id, prodEnv.id);
135+
136+
await syncDeclarativeSchedules(
137+
declarativeTasks({ cron: "0 * * * *", timezone: "UTC" }),
138+
noWorker,
139+
asEnv(prodEnv),
140+
prisma
141+
);
142+
143+
expect(registerNextTaskScheduleInstance).toHaveBeenCalledWith({
144+
instanceId: schedule.instances[0].id,
145+
preserveExistingJob: true,
146+
});
147+
}
148+
);
149+
150+
containerTest("replaces the Redis job when declarative timing changes", async ({ prisma }) => {
151+
registerNextTaskScheduleInstance.mockClear();
152+
const { project, prodEnv } = await seedProjectWithEnvs(prisma);
153+
const schedule = await makeDeclarativeSchedule(prisma, project.id, [prodEnv.id]);
154+
await seedScheduledTask(prisma, project.id, prodEnv.id);
155+
156+
await syncDeclarativeSchedules(
157+
declarativeTasks({ cron: "30 * * * *", timezone: "UTC", window: "30m" }),
158+
noWorker,
159+
asEnv(prodEnv),
160+
prisma
161+
);
162+
163+
expect(registerNextTaskScheduleInstance).toHaveBeenCalledWith({
164+
instanceId: schedule.instances[0].id,
165+
preserveExistingJob: false,
166+
});
167+
});
168+
});
169+
85170
describe("syncDeclarativeSchedules deletion path", () => {
86171
containerTest(
87172
"does not issue any instance delete when the env owns no instance of the missing schedules",

internal-packages/schedule-engine/src/engine/index.ts

Lines changed: 22 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -327,6 +327,7 @@ export class ScheduleEngine {
327327
exactScheduleTime: nominalAt,
328328
effectiveScheduleTime: effectiveAt,
329329
lastScheduleTime,
330+
preserveExistingJob: params.preserveExistingJob,
330331
});
331332

332333
// Record metrics
@@ -752,16 +753,19 @@ export class ScheduleEngine {
752753
exactScheduleTime,
753754
effectiveScheduleTime,
754755
lastScheduleTime,
756+
preserveExistingJob = false,
755757
}: {
756758
instanceId: string;
757759
exactScheduleTime: Date;
758760
effectiveScheduleTime: Date;
759761
lastScheduleTime?: Date;
762+
preserveExistingJob?: boolean;
760763
}) {
761764
return startSpan(this.tracer, "enqueueScheduledTask", async (span) => {
762765
span.setAttribute("instanceId", instanceId);
763766
span.setAttribute("exactScheduleTime", exactScheduleTime.toISOString());
764767
span.setAttribute("effectiveScheduleTime", effectiveScheduleTime.toISOString());
768+
span.setAttribute("preserveExistingJob", preserveExistingJob);
765769
if (lastScheduleTime) {
766770
span.setAttribute("lastScheduleTime", lastScheduleTime.toISOString());
767771
}
@@ -790,27 +794,38 @@ export class ScheduleEngine {
790794
distributedExecutionTime: distributedExecutionTime.toISOString(),
791795
distributionOffsetMs,
792796
distributionWindowSeconds: this.distributionWindowSeconds,
797+
preserveExistingJob,
793798
});
794799

795800
try {
796-
await this.worker.enqueue({
801+
const job = {
797802
id: `scheduled-task-instance:${instanceId}`,
798-
job: "schedule.triggerScheduledTask",
803+
job: "schedule.triggerScheduledTask" as const,
799804
payload: {
800805
instanceId,
801806
exactScheduleTime,
802807
effectiveScheduleTime,
803808
lastScheduleTime,
804809
},
805810
availableAt: distributedExecutionTime,
806-
});
811+
};
812+
let enqueued = true;
813+
if (preserveExistingJob) {
814+
enqueued = await this.worker.enqueueOnce(job);
815+
} else {
816+
await this.worker.enqueue(job);
817+
}
807818

808819
span.setAttribute("enqueue_success", true);
820+
span.setAttribute("existing_job_preserved", !enqueued);
809821

810-
this.logger.debug("Successfully enqueued scheduled task", {
811-
instanceId,
812-
jobId: `scheduled-task-instance:${instanceId}`,
813-
});
822+
this.logger.debug(
823+
enqueued ? "Successfully enqueued scheduled task" : "Preserved existing scheduled task",
824+
{
825+
instanceId,
826+
jobId: job.id,
827+
}
828+
);
814829
} catch (error) {
815830
this.logger.error("Failed to enqueue scheduled task", {
816831
instanceId,

internal-packages/schedule-engine/src/engine/types.ts

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -96,4 +96,9 @@ export interface RegisterScheduleInstanceParams {
9696
* disconnected, etc.) do NOT advance this — only real fires do.
9797
*/
9898
lastScheduleTime?: Date;
99+
/**
100+
* Keep an existing stable-ID Redis job unchanged, while still creating it
101+
* when missing. Intended for no-op reconciliation of unchanged schedules.
102+
*/
103+
preserveExistingJob?: boolean;
99104
}

internal-packages/schedule-engine/test/scheduleEngine2.test.ts

Lines changed: 32 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -220,7 +220,11 @@ describe("ScheduleEngine Integration (part 2)", () => {
220220
},
221221
});
222222

223-
await engine.registerNextTaskScheduleInstance({ instanceId: scheduleInstance.id });
223+
// Atomic preserve mode still creates the stable-ID job when it is missing.
224+
await engine.registerNextTaskScheduleInstance({
225+
instanceId: scheduleInstance.id,
226+
preserveExistingJob: true,
227+
});
224228

225229
const unwindowedInstance = await prisma.taskScheduleInstance.findUniqueOrThrow({
226230
where: { id: scheduleInstance.id },
@@ -291,6 +295,32 @@ describe("ScheduleEngine Integration (part 2)", () => {
291295
});
292296
expect(preservedInstance.schedulePhase).toBe(pinnedPhase);
293297

298+
const pendingBeforeNoop = await engine.getJob(
299+
`scheduled-task-instance:${scheduleInstance.id}`
300+
);
301+
302+
// No-op reconciliation preserves the existing payload and score atomically.
303+
await engine.registerNextTaskScheduleInstance({
304+
instanceId: scheduleInstance.id,
305+
preserveExistingJob: true,
306+
});
307+
const pendingAfterNoop = await engine.getJob(
308+
`scheduled-task-instance:${scheduleInstance.id}`
309+
);
310+
expect(pendingAfterNoop).toEqual(pendingBeforeNoop);
311+
312+
await prisma.taskSchedule.update({
313+
where: { id: taskSchedule.id },
314+
data: { windowDurationSeconds: 120 },
315+
});
316+
317+
// Normal registration still replaces the job when timing changed.
318+
await engine.registerNextTaskScheduleInstance({ instanceId: scheduleInstance.id });
319+
const pendingAfterTimingChange = await engine.getJob(
320+
`scheduled-task-instance:${scheduleInstance.id}`
321+
);
322+
expect(pendingAfterTimingChange).not.toEqual(pendingBeforeNoop);
323+
294324
const intervalMs = 5 * 60_000;
295325
const exactScheduleTime = new Date(Math.floor(Date.now() / intervalMs) * intervalMs);
296326
const effectiveScheduleTime = new Date(exactScheduleTime.getTime() + 45_000);
@@ -317,7 +347,7 @@ describe("ScheduleEngine Integration (part 2)", () => {
317347
nominalAt: nextNominalAt,
318348
nextNominalAt: followingNominalAt,
319349
schedulePhase: pinnedPhase,
320-
window: { type: "duration", durationSeconds: 60 },
350+
window: { type: "duration", durationSeconds: 120 },
321351
});
322352

323353
expect(new Date(nextJobPayload.exactScheduleTime)).toEqual(nextNominalAt);

0 commit comments

Comments
 (0)