diff --git a/apps/webapp/app/v3/marqs/devQueueConsumer.server.ts b/apps/webapp/app/v3/marqs/devQueueConsumer.server.ts index 7fc7224b7..90ffcd92d 100644 --- a/apps/webapp/app/v3/marqs/devQueueConsumer.server.ts +++ b/apps/webapp/app/v3/marqs/devQueueConsumer.server.ts @@ -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; diff --git a/apps/webapp/app/v3/marqs/sharedQueueConsumer.server.ts b/apps/webapp/app/v3/marqs/sharedQueueConsumer.server.ts index 399634ae9..d56af33aa 100644 --- a/apps/webapp/app/v3/marqs/sharedQueueConsumer.server.ts +++ b/apps/webapp/app/v3/marqs/sharedQueueConsumer.server.ts @@ -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; } diff --git a/references/v3-catalog/src/trigger/simple.ts b/references/v3-catalog/src/trigger/simple.ts index ff767e284..da727e750 100644 --- a/references/v3-catalog/src/trigger/simple.ts +++ b/references/v3-catalog/src/trigger/simple.ts @@ -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 }) => {