Fix issues with special characters in queue/task names causing runs to get stuck in queued

This commit is contained in:
Eric Allam
2024-05-16 17:35:44 +01:00
parent ba61bfe3b9
commit c9733f357f
3 changed files with 29 additions and 5 deletions
@@ -12,7 +12,7 @@ import { prisma } from "~/db.server";
import { createNewSession, disconnectSession } from "~/models/runtimeEnvironment.server";
import { AuthenticatedEnvironment } from "~/services/apiAuth.server";
import { logger } from "~/services/logger.server";
import { marqs } from "~/v3/marqs/index.server";
import { marqs, sanitizeQueueName } from "~/v3/marqs/index.server";
import { EnvironmentVariablesRepository } from "../environmentVariables/environmentVariablesRepository.server";
import { generateFriendlyId } from "../friendlyIdentifiers";
import { CancelAttemptService } from "../services/cancelAttempt.server";
@@ -452,11 +452,21 @@ export class DevQueueConsumer {
const queue = await prisma.taskQueue.findUnique({
where: {
runtimeEnvironmentId_name: { runtimeEnvironmentId: this.env.id, name: lockedTaskRun.queue },
runtimeEnvironmentId_name: {
runtimeEnvironmentId: this.env.id,
name: sanitizeQueueName(lockedTaskRun.queue),
},
},
});
if (!queue) {
logger.debug("[DevQueueConsumer] Failed to find queue", {
queueName: lockedTaskRun.queue,
sanitizedName: sanitizeQueueName(lockedTaskRun.queue),
taskRun: lockedTaskRun.id,
messageId: message.messageId,
});
await marqs?.nackMessage(message.messageId);
setTimeout(() => this.#doWork(), 1000);
return;
@@ -21,7 +21,7 @@ import { z } from "zod";
import { prisma } from "~/db.server";
import { logger } from "~/services/logger.server";
import { singleton } from "~/utils/singleton";
import { marqs } from "~/v3/marqs/index.server";
import { marqs, sanitizeQueueName } from "~/v3/marqs/index.server";
import { EnvironmentVariablesRepository } from "../environmentVariables/environmentVariablesRepository.server";
import { generateFriendlyId } from "../friendlyIdentifiers";
import { socketIo } from "../handleSocketIo.server";
@@ -408,7 +408,7 @@ export class SharedQueueConsumer {
where: {
runtimeEnvironmentId_name: {
runtimeEnvironmentId: lockedTaskRun.runtimeEnvironmentId,
name: lockedTaskRun.queue,
name: sanitizeQueueName(lockedTaskRun.queue),
},
},
});
@@ -635,12 +635,17 @@ export class SharedQueueConsumer {
where: {
runtimeEnvironmentId_name: {
runtimeEnvironmentId: resumableAttempt.runtimeEnvironmentId,
name: resumableRun.queue,
name: sanitizeQueueName(resumableRun.queue),
},
},
});
if (!queue) {
logger.debug("SharedQueueConsumer queue not found, so nacking message", {
queueName: sanitizeQueueName(resumableRun.queue),
attempt: resumableAttempt,
});
await this.#nackAndDoMoreWork(message.messageId, this._options.nextTickInterval);
return;
}
@@ -17,6 +17,15 @@ export const simplestTask = task({
},
});
export const taskWithSpecialCharacters = task({
id: "admin:special-characters",
run: async (payload: { url: string }) => {
return {
message: "This task has special characters in its ID",
};
},
});
export const createJsonHeroDoc = task({
id: "create-jsonhero-doc",
run: async (payload: { title: string; content: any }, { ctx }) => {