Skip to content

Commit 52ac1e1

Browse files
committed
add feature flag
1 parent 78f6764 commit 52ac1e1

7 files changed

Lines changed: 38 additions & 15 deletions

File tree

apps/webapp/app/env.server.ts

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1672,6 +1672,7 @@ const EnvironmentSchema = z
16721672
SCHEDULE_WORKER_CONCURRENCY_LIMIT: z.coerce.number().int().default(50),
16731673
SCHEDULE_WORKER_SHUTDOWN_TIMEOUT_MS: z.coerce.number().int().default(30_000),
16741674
SCHEDULE_WORKER_DISTRIBUTION_WINDOW_SECONDS: z.coerce.number().int().default(30),
1675+
SCHEDULE_WORKER_CRON_SPREAD_ENABLED: BoolEnv.default(false),
16751676

16761677
SCHEDULE_WORKER_REDIS_HOST: z
16771678
.string()

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

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -73,6 +73,7 @@ function createScheduleEngine() {
7373
seconds: env.SCHEDULE_WORKER_DISTRIBUTION_WINDOW_SECONDS,
7474
},
7575
schedulePhaseSecret: env.ENCRYPTION_KEY,
76+
cronSpreadEnabled: env.SCHEDULE_WORKER_CRON_SPREAD_ENABLED,
7677
tracer,
7778
meter,
7879
onTriggerScheduledTask: async ({

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

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -225,14 +225,17 @@ export class ScheduleEngine {
225225
instance.taskSchedule.timezone,
226226
nominalAt
227227
);
228-
const { effectiveAt } = calculateEffectiveScheduleTime({
228+
const { effectiveAt: candidateEffectiveAt } = calculateEffectiveScheduleTime({
229229
nominalAt,
230230
nextNominalAt,
231231
schedulePhase,
232232
window: scheduleWindow,
233233
});
234+
const effectiveAt = this.options.cronSpreadEnabled ? candidateEffectiveAt : nominalAt;
234235

236+
span.setAttribute("cron_spread_enabled", this.options.cronSpreadEnabled);
235237
span.setAttribute("next_scheduled_timestamp", nominalAt.toISOString());
238+
span.setAttribute("candidate_effective_schedule_time", candidateEffectiveAt.toISOString());
236239
span.setAttribute("effective_schedule_time", effectiveAt.toISOString());
237240

238241
const schedulingDelayMs = effectiveAt.getTime() - Date.now();
@@ -242,7 +245,9 @@ export class ScheduleEngine {
242245
instanceId: params.instanceId,
243246
taskIdentifier: instance.taskSchedule.taskIdentifier,
244247
nominalAt: nominalAt.toISOString(),
248+
candidateEffectiveAt: candidateEffectiveAt.toISOString(),
245249
effectiveAt: effectiveAt.toISOString(),
250+
cronSpreadEnabled: this.options.cronSpreadEnabled,
246251
schedulingDelayMs,
247252
generatorExpression: instance.taskSchedule.generatorExpression,
248253
timezone: instance.taskSchedule.timezone,

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

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -52,6 +52,7 @@ export interface ScheduleEngineOptions {
5252
seconds: number;
5353
};
5454
schedulePhaseSecret: string | Buffer;
55+
cronSpreadEnabled: boolean;
5556
tracer?: Tracer;
5657
meter?: Meter;
5758
onTriggerScheduledTask: TriggerScheduledTaskCallback;

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

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

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

Lines changed: 22 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,7 @@ describe("ScheduleEngine Integration (part 2)", () => {
2727
redis: redisOptions,
2828
distributionWindow: { seconds: 10 },
2929
schedulePhaseSecret: "test-schedule-phase-secret",
30+
cronSpreadEnabled: false,
3031
worker: {
3132
concurrency: 1,
3233
disabled: true, // Don't actually run the worker — calling triggerScheduledTask directly
@@ -79,6 +80,7 @@ describe("ScheduleEngine Integration (part 2)", () => {
7980
type: "DECLARATIVE",
8081
active: true,
8182
externalId: "legacy-ext",
83+
windowDurationSeconds: 60,
8284
},
8385
});
8486

@@ -119,22 +121,27 @@ describe("ScheduleEngine Integration (part 2)", () => {
119121
effectiveScheduleTime: string;
120122
};
121123
const nextNominalAt = new Date("2026-04-30T10:10:00.000Z");
122-
const followingNominalAt = new Date("2026-04-30T10:15:00.000Z");
123-
const schedulePhase = calculateSchedulePhase({
124-
secret: "test-schedule-phase-secret",
125-
environmentId: environment.id,
126-
deduplicationKey: taskSchedule.deduplicationKey,
127-
});
128-
const { effectiveAt: nextEffectiveAt } = calculateEffectiveScheduleTime({
129-
nominalAt: nextNominalAt,
130-
nextNominalAt: followingNominalAt,
131-
schedulePhase,
132-
});
133124

134-
// The next job advances from the legacy job's nominal T, not from E
135-
// or the current wall clock, and newly enqueued jobs carry both times.
125+
// The next job advances from the legacy job's nominal T, not from the
126+
// current wall clock. With cron spread disabled, actual eligibility
127+
// remains nominal even though registration still calculates candidate E.
136128
expect(new Date(nextJobPayload.exactScheduleTime)).toEqual(nextNominalAt);
137-
expect(new Date(nextJobPayload.effectiveScheduleTime)).toEqual(nextEffectiveAt);
129+
expect(new Date(nextJobPayload.effectiveScheduleTime)).toEqual(nextNominalAt);
130+
expect(nextJob!.timestamp).toEqual(
131+
calculateDistributedExecutionTime(nextNominalAt, 10, scheduleInstance.id)
132+
);
133+
134+
const updatedInstance = await prisma.taskScheduleInstance.findUniqueOrThrow({
135+
where: { id: scheduleInstance.id },
136+
select: { schedulePhase: true },
137+
});
138+
expect(updatedInstance.schedulePhase).toBe(
139+
calculateSchedulePhase({
140+
secret: "test-schedule-phase-secret",
141+
environmentId: environment.id,
142+
deduplicationKey: taskSchedule.deduplicationKey,
143+
})
144+
);
138145
} finally {
139146
await engine.quit();
140147
}
@@ -152,6 +159,7 @@ describe("ScheduleEngine Integration (part 2)", () => {
152159
redis: redisOptions,
153160
distributionWindow: { seconds: 10 },
154161
schedulePhaseSecret,
162+
cronSpreadEnabled: true,
155163
worker: {
156164
concurrency: 1,
157165
disabled: true,

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

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@ describe("Schedule Recovery", () => {
1717
redis: redisOptions,
1818
distributionWindow: { seconds: 10 },
1919
schedulePhaseSecret: "test-schedule-phase-secret",
20+
cronSpreadEnabled: true,
2021
worker: {
2122
concurrency: 1,
2223
disabled: true, // Disable worker to prevent automatic execution
@@ -120,6 +121,7 @@ describe("Schedule Recovery", () => {
120121
redis: redisOptions,
121122
distributionWindow: { seconds: 10 },
122123
schedulePhaseSecret: "test-schedule-phase-secret",
124+
cronSpreadEnabled: true,
123125
worker: {
124126
concurrency: 1,
125127
disabled: true, // Disable worker to prevent automatic execution
@@ -226,6 +228,7 @@ describe("Schedule Recovery", () => {
226228
redis: redisOptions,
227229
distributionWindow: { seconds: 10 },
228230
schedulePhaseSecret: "test-schedule-phase-secret",
231+
cronSpreadEnabled: true,
229232
worker: {
230233
concurrency: 1,
231234
disabled: true, // Disable worker to prevent automatic execution
@@ -338,6 +341,7 @@ describe("Schedule Recovery", () => {
338341
redis: redisOptions,
339342
distributionWindow: { seconds: 10 },
340343
schedulePhaseSecret: "test-schedule-phase-secret",
344+
cronSpreadEnabled: true,
341345
worker: {
342346
concurrency: 1,
343347
disabled: true, // Disable worker to prevent automatic execution
@@ -409,6 +413,7 @@ describe("Schedule Recovery", () => {
409413
redis: redisOptions,
410414
distributionWindow: { seconds: 10 },
411415
schedulePhaseSecret: "test-schedule-phase-secret",
416+
cronSpreadEnabled: true,
412417
worker: { concurrency: 1, disabled: true, pollIntervalMs: 1000 },
413418
tracer: trace.getTracer("test", "0.0.0"),
414419
onTriggerScheduledTask: async () => ({ success: true }),
@@ -511,6 +516,7 @@ describe("Schedule Recovery", () => {
511516
redis: redisOptions,
512517
distributionWindow: { seconds: 10 },
513518
schedulePhaseSecret: "test-schedule-phase-secret",
519+
cronSpreadEnabled: true,
514520
worker: { concurrency: 1, disabled: true, pollIntervalMs: 1000 },
515521
tracer: trace.getTracer("test", "0.0.0"),
516522
onTriggerScheduledTask: async () => ({ success: true }),

0 commit comments

Comments
 (0)