chore(run-engine): add additional logging around dequeueing and worker queues (#2562)

This commit is contained in:
Eric Allam
2025-09-26 11:50:21 +01:00
committed by James Ritchie
parent 8dfeb0ab93
commit 5c06ba19ea
3 changed files with 23 additions and 8 deletions
@@ -1414,6 +1414,11 @@ export class RunEngine {
throw new NotImplementedError("There shouldn't be a heartbeat for QUEUED_EXECUTING");
}
case "PENDING_EXECUTING": {
this.logger.log("RunEngine stalled snapshot PENDING_EXECUTING", {
runId,
snapshotId: latestSnapshot.id,
});
//the run didn't start executing, we need to requeue it
const run = await prisma.taskRun.findFirst({
where: { id: runId },
@@ -143,6 +143,15 @@ export class DequeueSystem {
const orgId = message.message.orgId;
const runId = message.messageId;
this.$.logger.info("DequeueSystem.dequeueFromWorkerQueue dequeued message", {
runId,
orgId,
environmentId: message.message.environmentId,
environmentType: message.message.environmentType,
workerQueueLength: message.workerQueueLength ?? 0,
workerQueue,
});
span.setAttribute("run_id", runId);
span.setAttribute("org_id", orgId);
span.setAttribute("environment_id", message.message.environmentId);
@@ -1417,27 +1417,28 @@ export class RunQueue {
const pipeline = this.redis.pipeline();
const workerQueueKeys = new Set<string>();
const operations = [];
for (const message of messages) {
const workerQueueKey = this.keys.workerQueueKey(
this.#getWorkerQueueFromMessage(message.message)
);
workerQueueKeys.add(workerQueueKey);
const messageKeyValue = this.keys.messageKey(message.message.orgId, message.messageId);
operations.push({
workerQueueKey: workerQueueKey,
messageId: message.messageId,
});
pipeline.rpush(workerQueueKey, messageKeyValue);
}
span.setAttribute("worker_queue_count", workerQueueKeys.size);
span.setAttribute("worker_queue_keys", Array.from(workerQueueKeys));
span.setAttribute("operations_count", operations.length);
this.logger.debug("enqueueMessagesToWorkerQueues pipeline", {
this.logger.info("enqueueMessagesToWorkerQueues", {
service: this.name,
messages,
workerQueueKeys: Array.from(workerQueueKeys),
operations,
});
await pipeline.exec();