diff --git a/apps/webapp/app/routes/api.v1.projects.$projectRef.background-workers.ts b/apps/webapp/app/routes/api.v1.projects.$projectRef.background-workers.ts index 7e0a1329a..872ca6f2f 100644 --- a/apps/webapp/app/routes/api.v1.projects.$projectRef.background-workers.ts +++ b/apps/webapp/app/routes/api.v1.projects.$projectRef.background-workers.ts @@ -58,7 +58,7 @@ export async function action({ request, params }: ActionFunctionArgs) { { status: 200 } ); } catch (e) { - logger.error("Failed to create background worker", { error: e }); + logger.error("Failed to create background worker", { error: JSON.stringify(e) }); if (e instanceof ServiceValidationError) { return json({ error: e.message }, { status: 400 }); diff --git a/apps/webapp/app/v3/marqs/index.server.ts b/apps/webapp/app/v3/marqs/index.server.ts index 1464e351a..d273a0c72 100644 --- a/apps/webapp/app/v3/marqs/index.server.ts +++ b/apps/webapp/app/v3/marqs/index.server.ts @@ -986,8 +986,15 @@ export class MarQS { } async #callEnqueueMessage(message: MessagePayload) { + const concurrencyKey = this.keys.currentConcurrencyKeyFromQueue(message.queue); + const envConcurrencyKey = this.keys.envCurrentConcurrencyKeyFromQueue(message.queue); + const orgConcurrencyKey = this.keys.orgCurrentConcurrencyKeyFromQueue(message.queue); + logger.debug("Calling enqueueMessage", { messagePayload: message, + concurrencyKey, + envConcurrencyKey, + orgConcurrencyKey, service: this.name, }); @@ -995,6 +1002,9 @@ export class MarQS { message.queue, message.parentQueue, this.keys.messageKey(message.messageId), + concurrencyKey, + envConcurrencyKey, + orgConcurrencyKey, message.queue, message.messageId, JSON.stringify(message), @@ -1268,11 +1278,14 @@ export class MarQS { #registerCommands() { this.redis.defineCommand("enqueueMessage", { - numberOfKeys: 3, + numberOfKeys: 6, lua: ` local queue = KEYS[1] local parentQueue = KEYS[2] local messageKey = KEYS[3] +local concurrencyKey = KEYS[4] +local envCurrentConcurrencyKey = KEYS[5] +local orgCurrentConcurrencyKey = KEYS[6] local queueName = ARGV[1] local messageId = ARGV[2] @@ -1292,6 +1305,11 @@ if #earliestMessage == 0 then else redis.call('ZADD', parentQueue, earliestMessage[2], queueName) end + +-- Update the concurrency keys +redis.call('SREM', concurrencyKey, messageId) +redis.call('SREM', envCurrentConcurrencyKey, messageId) +redis.call('SREM', orgCurrentConcurrencyKey, messageId) `, }); @@ -1621,6 +1639,9 @@ declare module "ioredis" { queue: string, parentQueue: string, messageKey: string, + concurrencyKey: string, + envConcurrencyKey: string, + orgConcurrencyKey: string, queueName: string, messageId: string, messageData: string,