Skip to content

Commit 6a50ab7

Browse files
committed
fix(run-engine): anchor scheduling delay to the current queue stint on re-enqueue
The wait metric measured dequeue time minus the original trigger time even when a run re-entered the queue after a waitpoint, checkpoint, or pending version, so the whole wait or checkpoint duration showed up as scheduling delay. Re-enqueues now anchor to the re-enqueue time while first enqueues keep the trigger or delay anchor; queue ordering is unchanged, so re-enqueued runs keep their original position. Nacked runs never left the queue stint and keep the original anchor. Also replaces the :ck: suffix regexes on user-controlled queue names with indexOf slicing (identical semantics) to remove a polynomial regex flagged by code scanning.
1 parent 1daa52a commit 6a50ab7

4 files changed

Lines changed: 33 additions & 12 deletions

File tree

internal-packages/run-engine/src/engine/systems/enqueueSystem.ts

Lines changed: 9 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -98,11 +98,16 @@ export class EnqueueSystem {
9898
// Force development runs to use the environment id as the worker queue.
9999
const workerQueue = env.type === "DEVELOPMENT" ? env.id : run.workerQueue;
100100

101-
const eligibleAtMs = (run.queueTimestamp ?? run.createdAt).getTime();
102-
const timestamp = eligibleAtMs - run.priorityMs;
101+
// Ordering keeps the run's original position; the scheduling-delay anchor is the
102+
// trigger/delay time only on first enqueue (includeTtl). Re-enqueues anchor to now,
103+
// else the wait metric absorbs the whole waitpoint/checkpoint duration.
104+
const queuePositionMs = (run.queueTimestamp ?? run.createdAt).getTime();
105+
const timestamp = queuePositionMs - run.priorityMs;
106+
const eligibleAtMs = includeTtl ? queuePositionMs : Date.now();
103107

104-
// Include TTL only when explicitly requested (first enqueue from trigger).
105-
// Re-enqueues (waitpoint, checkpoint, delayed, pending version) must not add TTL.
108+
// Include TTL only when explicitly requested (first enqueue from trigger or the
109+
// delayed-run system). Re-enqueues (waitpoint, checkpoint, pending version) must
110+
// not add TTL.
106111
let ttlExpiresAt: number | undefined;
107112
if (includeTtl && run.ttl) {
108113
const expireAt = parseNaturalLanguageDuration(run.ttl);

internal-packages/run-engine/src/engine/tests/ttl.test.ts

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -293,7 +293,12 @@ describe("RunEngine ttl", () => {
293293
);
294294
assertNonNullable(messageAfterTrigger);
295295
expect(messageAfterTrigger.ttlExpiresAt).toBeDefined();
296+
// First enqueue anchors the scheduling-delay clock at the trigger time.
297+
expect(messageAfterTrigger.eligibleAtMs).toBe(
298+
(run.queueTimestamp ?? run.createdAt).getTime()
299+
);
296300

301+
const beforeReenqueue = Date.now();
297302
await engine.enqueueSystem.enqueueRun({
298303
run,
299304
env: authenticatedEnvironment,
@@ -308,6 +313,10 @@ describe("RunEngine ttl", () => {
308313
);
309314
assertNonNullable(messageAfterReenqueue);
310315
expect(messageAfterReenqueue.ttlExpiresAt).toBeUndefined();
316+
// Re-enqueues anchor to now so the wait metric measures only this queue stint,
317+
// while the ordering timestamp keeps the run's original position.
318+
expect(messageAfterReenqueue.eligibleAtMs).toBeGreaterThanOrEqual(beforeReenqueue);
319+
expect(messageAfterReenqueue.timestamp).toBe(messageAfterTrigger.timestamp);
311320
} finally {
312321
await engine.quit();
313322
}

internal-packages/run-engine/src/run-queue/index.ts

Lines changed: 7 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -603,7 +603,7 @@ export class RunQueue {
603603
const queuedResult = stats?.[i * 2];
604604
const runningResult = stats?.[i * 2 + 1];
605605
return {
606-
concurrencyKey: member.match(/:ck:(.+)$/)?.[1] ?? "",
606+
concurrencyKey: this.#concurrencyKeyFromQueue(member) ?? "",
607607
queued: queuedResult && !queuedResult[0] ? ((queuedResult[1] as number) ?? 0) : 0,
608608
running: runningResult && !runningResult[0] ? ((runningResult[1] as number) ?? 0) : 0,
609609
oldestEnqueuedAt: score,
@@ -2033,6 +2033,11 @@ export class RunQueue {
20332033
this.options.queueMetrics?.emitGauge(queue, fields);
20342034
}
20352035

2036+
#concurrencyKeyFromQueue(queue: string): string | undefined {
2037+
const idx = queue.indexOf(":ck:");
2038+
return idx === -1 || idx + 4 >= queue.length ? undefined : queue.slice(idx + 4);
2039+
}
2040+
20362041
#emitQueueMetric(shardKey: string, fields: Record<string, string | number>): void {
20372042
// Counters roll up per BASE queue: normalize the CK-qualified queue to its base so all
20382043
// concurrency keys share one monotonic odometer (and one shard/order key), matching the
@@ -2042,7 +2047,7 @@ export class RunQueue {
20422047
let baseFields = fields;
20432048
if (typeof fields.q === "string") {
20442049
baseFields = { ...fields, q: this.keys.baseQueueKeyFromQueue(fields.q) };
2045-
const ck = fields.q.match(/:ck:(.+)$/)?.[1];
2050+
const ck = this.#concurrencyKeyFromQueue(fields.q);
20462051
if (ck && ck !== "*") baseFields.ck = ck;
20472052
}
20482053
this.options.queueMetrics?.emit(baseQueue, baseFields);

internal-packages/run-engine/src/run-queue/keyProducer.ts

Lines changed: 8 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -141,8 +141,7 @@ export class RunQueueFullKeyProducer implements RunQueueKeyProducer {
141141
}
142142

143143
queueConcurrencyLimitKeyFromQueue(queue: string) {
144-
const concurrencyQueueName = queue.replace(/:ck:.+$/, "");
145-
return `${concurrencyQueueName}:${constants.CONCURRENCY_LIMIT_PART}`;
144+
return `${this.baseQueueKeyFromQueue(queue)}:${constants.CONCURRENCY_LIMIT_PART}`;
146145
}
147146

148147
queueCurrentConcurrencyKeyFromQueue(queue: string) {
@@ -313,12 +312,14 @@ export class RunQueueFullKeyProducer implements RunQueueKeyProducer {
313312
}
314313

315314
ckIndexKeyFromQueue(queue: string): string {
316-
const baseQueue = queue.replace(/:ck:.+$/, "");
317-
return `${baseQueue}:${constants.CK_INDEX_PART}`;
315+
return `${this.baseQueueKeyFromQueue(queue)}:${constants.CK_INDEX_PART}`;
318316
}
319317

318+
// indexOf instead of /:ck:.+$/ (queue names are user-controlled; polynomial regex).
319+
// Only strips when at least one character follows ":ck:", matching the old semantics.
320320
baseQueueKeyFromQueue(queue: string): string {
321-
return queue.replace(/:ck:.+$/, "");
321+
const idx = queue.indexOf(":ck:");
322+
return idx === -1 || idx + 4 >= queue.length ? queue : queue.slice(0, idx);
322323
}
323324

324325
queueLengthCounterKey(env: RunQueueKeyProducerEnvironment, queue: string): string {
@@ -342,7 +343,8 @@ export class RunQueueFullKeyProducer implements RunQueueKeyProducer {
342343
}
343344

344345
toCkWildcard(queue: string): string {
345-
return queue.replace(/:ck:.+$/, ":ck:*");
346+
const base = this.baseQueueKeyFromQueue(queue);
347+
return base === queue ? queue : `${base}:ck:*`;
346348
}
347349

348350
descriptorFromQueue(queue: string): QueueDescriptor {

0 commit comments

Comments
 (0)