Skip to content

Commit ed9f165

Browse files
committed
persist schedule phase
1 parent 2f5ac02 commit ed9f165

6 files changed

Lines changed: 151 additions & 1 deletion

File tree

apps/webapp/app/v3/scheduleEngine.server.ts

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -72,6 +72,7 @@ function createScheduleEngine() {
7272
distributionWindow: {
7373
seconds: env.SCHEDULE_WORKER_DISTRIBUTION_WINDOW_SECONDS,
7474
},
75+
schedulePhaseSecret: env.ENCRYPTION_KEY,
7576
tracer,
7677
meter,
7778
onTriggerScheduledTask: async ({

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

Lines changed: 27 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,7 @@ import type {
1515
TriggerScheduledTaskCallback,
1616
TriggerScheduleParams,
1717
} from "./types.js";
18+
import { calculateSchedulePhase } from "./scheduleTiming.js";
1819
import { scheduleWorkerCatalog } from "./workerCatalog.js";
1920
import { tryCatch } from "@trigger.dev/core/utils";
2021

@@ -168,6 +169,32 @@ export class ScheduleEngine {
168169
instance.taskSchedule.generatorExpression
169170
);
170171

172+
const hasScheduleWindow =
173+
instance.taskSchedule.windowDurationSeconds !== null ||
174+
instance.taskSchedule.windowPercentage !== null;
175+
if (hasScheduleWindow && instance.schedulePhase === null) {
176+
const schedulePhase = calculateSchedulePhase({
177+
secret: this.options.schedulePhaseSecret,
178+
environmentId: instance.environmentId,
179+
deduplicationKey: instance.taskSchedule.deduplicationKey,
180+
});
181+
const result = await this.prisma.taskScheduleInstance.updateMany({
182+
where: {
183+
id: instance.id,
184+
schedulePhase: null,
185+
},
186+
data: {
187+
schedulePhase,
188+
},
189+
});
190+
191+
span.setAttribute("schedule_phase", schedulePhase);
192+
span.setAttribute("schedule_phase_persisted", result.count === 1);
193+
} else if (instance.schedulePhase !== null) {
194+
span.setAttribute("schedule_phase", instance.schedulePhase);
195+
span.setAttribute("schedule_phase_persisted", false);
196+
}
197+
171198
const fromTimestamp = params.fromTimestamp ?? new Date();
172199
span.setAttribute("from_timestamp", fromTimestamp.toISOString());
173200

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

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -50,6 +50,7 @@ export interface ScheduleEngineOptions {
5050
distributionWindow?: {
5151
seconds: number;
5252
};
53+
schedulePhaseSecret: string | Buffer;
5354
tracer?: Tracer;
5455
meter?: Meter;
5556
onTriggerScheduledTask: TriggerScheduledTaskCallback;

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

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@ describe("ScheduleEngine Integration", () => {
2222
prisma,
2323
redis: redisOptions,
2424
distributionWindow: { seconds: 10 },
25+
schedulePhaseSecret: "test-schedule-phase-secret",
2526
worker: {
2627
concurrency: 1,
2728
disabled: false, // Enable worker for full integration test

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

Lines changed: 115 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,7 @@ import { containerTest } from "@internal/testcontainers";
22
import { trace } from "@internal/tracing";
33
import { describe, expect, vi } from "vitest";
44
import type { TriggerScheduledTaskParams } from "../src/engine/types.js";
5-
import { ScheduleEngine } from "../src/index.js";
5+
import { calculateSchedulePhase, ScheduleEngine } from "../src/index.js";
66

77
describe("ScheduleEngine Integration (part 2)", () => {
88
// Deploy-moment backward compatibility. At deploy time, in-flight Redis jobs
@@ -20,6 +20,7 @@ describe("ScheduleEngine Integration (part 2)", () => {
2020
prisma,
2121
redis: redisOptions,
2222
distributionWindow: { seconds: 10 },
23+
schedulePhaseSecret: "test-schedule-phase-secret",
2324
worker: {
2425
concurrency: 1,
2526
disabled: true, // Don't actually run the worker — calling triggerScheduledTask directly
@@ -108,4 +109,117 @@ describe("ScheduleEngine Integration (part 2)", () => {
108109
}
109110
}
110111
);
112+
113+
containerTest(
114+
"should assign a stable schedule phase once when a window is configured",
115+
{ timeout: 30_000 },
116+
async ({ prisma, redisOptions }) => {
117+
const schedulePhaseSecret = "test-schedule-phase-secret";
118+
const engine = new ScheduleEngine({
119+
prisma,
120+
redis: redisOptions,
121+
distributionWindow: { seconds: 10 },
122+
schedulePhaseSecret,
123+
worker: {
124+
concurrency: 1,
125+
disabled: true,
126+
pollIntervalMs: 1000,
127+
},
128+
tracer: trace.getTracer("test", "0.0.0"),
129+
onTriggerScheduledTask: async () => ({ success: true }),
130+
isDevEnvironmentConnectedHandler: vi.fn().mockResolvedValue(true),
131+
});
132+
133+
try {
134+
const organization = await prisma.organization.create({
135+
data: { title: "Schedule Phase Org", slug: "schedule-phase-org" },
136+
});
137+
const project = await prisma.project.create({
138+
data: {
139+
name: "Schedule Phase Project",
140+
slug: "schedule-phase-project",
141+
externalRef: "schedule-phase-ref",
142+
organizationId: organization.id,
143+
},
144+
});
145+
const environment = await prisma.runtimeEnvironment.create({
146+
data: {
147+
slug: "schedule-phase-env",
148+
type: "PRODUCTION",
149+
projectId: project.id,
150+
organizationId: organization.id,
151+
apiKey: "tr_schedule_phase",
152+
pkApiKey: "pk_schedule_phase",
153+
shortcode: "phase",
154+
},
155+
});
156+
const taskSchedule = await prisma.taskSchedule.create({
157+
data: {
158+
friendlyId: "sched_phase",
159+
taskIdentifier: "schedule-phase-task",
160+
projectId: project.id,
161+
deduplicationKey: "schedule-phase-dedup",
162+
generatorExpression: "*/5 * * * *",
163+
generatorDescription: "Every 5 minutes",
164+
timezone: "UTC",
165+
type: "DECLARATIVE",
166+
},
167+
});
168+
const scheduleInstance = await prisma.taskScheduleInstance.create({
169+
data: {
170+
taskScheduleId: taskSchedule.id,
171+
environmentId: environment.id,
172+
projectId: project.id,
173+
},
174+
});
175+
176+
await engine.registerNextTaskScheduleInstance({ instanceId: scheduleInstance.id });
177+
178+
const unwindowedInstance = await prisma.taskScheduleInstance.findUniqueOrThrow({
179+
where: { id: scheduleInstance.id },
180+
select: { schedulePhase: true },
181+
});
182+
expect(unwindowedInstance.schedulePhase).toBeNull();
183+
184+
await prisma.taskSchedule.update({
185+
where: { id: taskSchedule.id },
186+
data: { windowDurationSeconds: 60 },
187+
});
188+
189+
const expectedPhase = calculateSchedulePhase({
190+
secret: schedulePhaseSecret,
191+
environmentId: environment.id,
192+
deduplicationKey: taskSchedule.deduplicationKey,
193+
});
194+
195+
await Promise.all([
196+
engine.registerNextTaskScheduleInstance({ instanceId: scheduleInstance.id }),
197+
engine.registerNextTaskScheduleInstance({ instanceId: scheduleInstance.id }),
198+
engine.registerNextTaskScheduleInstance({ instanceId: scheduleInstance.id }),
199+
]);
200+
201+
const assignedInstance = await prisma.taskScheduleInstance.findUniqueOrThrow({
202+
where: { id: scheduleInstance.id },
203+
select: { schedulePhase: true },
204+
});
205+
expect(assignedInstance.schedulePhase).toBe(expectedPhase);
206+
207+
const pinnedPhase = 1_234_567_890;
208+
await prisma.taskScheduleInstance.update({
209+
where: { id: scheduleInstance.id },
210+
data: { schedulePhase: pinnedPhase },
211+
});
212+
213+
await engine.registerNextTaskScheduleInstance({ instanceId: scheduleInstance.id });
214+
215+
const preservedInstance = await prisma.taskScheduleInstance.findUniqueOrThrow({
216+
where: { id: scheduleInstance.id },
217+
select: { schedulePhase: true },
218+
});
219+
expect(preservedInstance.schedulePhase).toBe(pinnedPhase);
220+
} finally {
221+
await engine.quit();
222+
}
223+
}
224+
);
111225
});

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

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@ describe("Schedule Recovery", () => {
1616
prisma,
1717
redis: redisOptions,
1818
distributionWindow: { seconds: 10 },
19+
schedulePhaseSecret: "test-schedule-phase-secret",
1920
worker: {
2021
concurrency: 1,
2122
disabled: true, // Disable worker to prevent automatic execution
@@ -118,6 +119,7 @@ describe("Schedule Recovery", () => {
118119
prisma,
119120
redis: redisOptions,
120121
distributionWindow: { seconds: 10 },
122+
schedulePhaseSecret: "test-schedule-phase-secret",
121123
worker: {
122124
concurrency: 1,
123125
disabled: true, // Disable worker to prevent automatic execution
@@ -223,6 +225,7 @@ describe("Schedule Recovery", () => {
223225
prisma,
224226
redis: redisOptions,
225227
distributionWindow: { seconds: 10 },
228+
schedulePhaseSecret: "test-schedule-phase-secret",
226229
worker: {
227230
concurrency: 1,
228231
disabled: true, // Disable worker to prevent automatic execution
@@ -334,6 +337,7 @@ describe("Schedule Recovery", () => {
334337
prisma,
335338
redis: redisOptions,
336339
distributionWindow: { seconds: 10 },
340+
schedulePhaseSecret: "test-schedule-phase-secret",
337341
worker: {
338342
concurrency: 1,
339343
disabled: true, // Disable worker to prevent automatic execution
@@ -404,6 +408,7 @@ describe("Schedule Recovery", () => {
404408
prisma,
405409
redis: redisOptions,
406410
distributionWindow: { seconds: 10 },
411+
schedulePhaseSecret: "test-schedule-phase-secret",
407412
worker: { concurrency: 1, disabled: true, pollIntervalMs: 1000 },
408413
tracer: trace.getTracer("test", "0.0.0"),
409414
onTriggerScheduledTask: async () => ({ success: true }),
@@ -505,6 +510,7 @@ describe("Schedule Recovery", () => {
505510
prisma,
506511
redis: redisOptions,
507512
distributionWindow: { seconds: 10 },
513+
schedulePhaseSecret: "test-schedule-phase-secret",
508514
worker: { concurrency: 1, disabled: true, pollIntervalMs: 1000 },
509515
tracer: trace.getTracer("test", "0.0.0"),
510516
onTriggerScheduledTask: async () => ({ success: true }),

0 commit comments

Comments
 (0)