fix(run-engine): pass through engine fair dequeue selection strategy options instead of using defaults (#2565)

This commit is contained in:
Eric Allam
2025-09-26 17:12:36 +01:00
committed by James Ritchie
parent 94d96b85b5
commit 0b7ea8ade0
2 changed files with 27 additions and 9 deletions
@@ -129,15 +129,22 @@ export class RunEngine {
const keys = new RunQueueFullKeyProducer();
const queueSelectionStrategyOptions = {
keys,
redis: { ...options.queue.redis, keyPrefix: `${options.queue.redis.keyPrefix}runqueue:` },
defaultEnvConcurrencyLimit: options.queue?.defaultEnvConcurrency ?? 10,
...options.queue?.queueSelectionStrategyOptions,
};
this.logger.log("RunEngine FairQueueSelectionStrategy queueSelectionStrategyOptions", {
options: queueSelectionStrategyOptions,
});
this.runQueue = new RunQueue({
name: "rq",
tracer: trace.getTracer("rq"),
keys,
queueSelectionStrategy: new FairQueueSelectionStrategy({
keys,
redis: { ...options.queue.redis, keyPrefix: `${options.queue.redis.keyPrefix}runqueue:` },
defaultEnvConcurrencyLimit: options.queue?.defaultEnvConcurrency ?? 10,
}),
queueSelectionStrategy: new FairQueueSelectionStrategy(queueSelectionStrategyOptions),
defaultEnvConcurrency: options.queue?.defaultEnvConcurrency ?? 10,
defaultEnvConcurrencyBurstFactor: options.queue?.defaultEnvConcurrencyBurstFactor,
logger: new Logger("RunQueue", options.queue?.logLevel ?? "info"),
@@ -1730,10 +1737,13 @@ export class RunEngine {
});
if (!taskRun) {
this.logger.error("RunEngine.handleRepairSnapshot SUSPENDED/FINISHED task run not found", {
runId,
snapshotId,
});
this.logger.error(
"RunEngine.handleRepairSnapshot SUSPENDED/FINISHED task run not found",
{
runId,
snapshotId,
}
);
return;
}
@@ -1706,6 +1706,14 @@ export class RunQueue {
const message = await this.#dequeueMessageFromKey(messageKey);
if (!message) {
this.logger.error("Failed to dequeue message from worker queue", {
messageKey,
workerQueue,
workerQueueKey,
workerQueueLength,
service: this.name,
});
return;
}