improve the dead letter queue stuff
This commit is contained in:
Vendored
+1
-1
@@ -146,7 +146,7 @@
|
||||
"type": "node-terminal",
|
||||
"request": "launch",
|
||||
"name": "Debug RunQueue tests",
|
||||
"command": "pnpm run test ./src/run-queue/tests/dequeueMessageFromMasterQueue.test.ts",
|
||||
"command": "pnpm run test ./src/run-queue/index.test.ts",
|
||||
"cwd": "${workspaceFolder}/internal-packages/run-engine",
|
||||
"sourceMaps": true
|
||||
}
|
||||
|
||||
@@ -826,7 +826,7 @@ describe("RunQueue", () => {
|
||||
}
|
||||
);
|
||||
|
||||
redisTest("Dead Letter Queue", { timeout: 8_000 }, async ({ redisContainer, redisOptions }) => {
|
||||
redisTest("Dead Letter Queue", async ({ redisContainer, redisOptions }) => {
|
||||
const queue = new RunQueue({
|
||||
...testOptions,
|
||||
retryOptions: {
|
||||
@@ -891,36 +891,30 @@ describe("RunQueue", () => {
|
||||
expect(taskConcurrency2).toBe(0);
|
||||
|
||||
//check the message is still there
|
||||
const exists2 = await redis.exists(key);
|
||||
expect(exists2).toBe(1);
|
||||
const message = await queue.readMessage(messages[0].message.orgId, messages[0].messageId);
|
||||
expect(message).toBeDefined();
|
||||
|
||||
//check it's in the dlq
|
||||
const dlqKey = "dlq";
|
||||
const dlqExists = await redis.exists(dlqKey);
|
||||
expect(dlqExists).toBe(1);
|
||||
const dlqMembers = await redis.zrange(dlqKey, 0, -1);
|
||||
expect(dlqMembers).toContain(messageProd.runId);
|
||||
const deadLetterQueueLengthBefore = await queue.lengthOfDeadLetterQueue(authenticatedEnvProd);
|
||||
expect(deadLetterQueueLengthBefore).toBe(1);
|
||||
|
||||
const existsInDlq = await queue.messageInDeadLetterQueue(
|
||||
authenticatedEnvProd,
|
||||
messageProd.runId
|
||||
);
|
||||
expect(existsInDlq).toBe(true);
|
||||
|
||||
//redrive
|
||||
const redisClient = createRedisClient({
|
||||
host: redisContainer.getHost(),
|
||||
port: redisContainer.getPort(),
|
||||
password: redisContainer.getPassword(),
|
||||
});
|
||||
|
||||
// Publish redrive message
|
||||
await redisClient.publish(
|
||||
"rq:redrive",
|
||||
JSON.stringify({ runId: messageProd.runId, orgId: messageProd.orgId })
|
||||
);
|
||||
await queue.redriveMessage(authenticatedEnvProd, messageProd.runId);
|
||||
|
||||
// Wait for the item to be redrived and processed
|
||||
await setTimeout(5_000);
|
||||
await redisClient.quit();
|
||||
|
||||
//shouldn't be in the dlq now
|
||||
const dlqMembersAfter = await redis.zrange(dlqKey, 0, -1);
|
||||
expect(dlqMembersAfter).not.toContain(messageProd.runId);
|
||||
const existsInDlqAfter = await queue.messageInDeadLetterQueue(
|
||||
authenticatedEnvProd,
|
||||
messageProd.runId
|
||||
);
|
||||
expect(existsInDlqAfter).toBe(false);
|
||||
|
||||
//dequeue
|
||||
const messages3 = await queue.dequeueMessageFromMasterQueue("test_12345", "main", 10);
|
||||
|
||||
@@ -161,6 +161,28 @@ export class RunQueue {
|
||||
return this.redis.zcard(this.keys.envQueueKey(env));
|
||||
}
|
||||
|
||||
public async lengthOfDeadLetterQueue(env: MinimalAuthenticatedEnvironment) {
|
||||
return this.redis.zcard(this.keys.deadLetterQueueKey(env));
|
||||
}
|
||||
|
||||
public async messageInDeadLetterQueue(env: MinimalAuthenticatedEnvironment, messageId: string) {
|
||||
const result = await this.redis.zscore(this.keys.deadLetterQueueKey(env), messageId);
|
||||
return !!result;
|
||||
}
|
||||
|
||||
public async redriveMessage(env: MinimalAuthenticatedEnvironment, messageId: string) {
|
||||
// Publish redrive message
|
||||
await this.redis.publish(
|
||||
"rq:redrive",
|
||||
JSON.stringify({
|
||||
runId: messageId,
|
||||
orgId: env.organization.id,
|
||||
envId: env.id,
|
||||
projectId: env.project.id,
|
||||
})
|
||||
);
|
||||
}
|
||||
|
||||
public async oldestMessageInQueue(
|
||||
env: MinimalAuthenticatedEnvironment,
|
||||
queue: string,
|
||||
@@ -275,6 +297,10 @@ export class RunQueue {
|
||||
return this.redis.scard(this.keys.projectCurrentConcurrencyKey(env));
|
||||
}
|
||||
|
||||
public async messageExists(orgId: string, messageId: string) {
|
||||
return this.redis.exists(this.keys.messageKey(orgId, messageId));
|
||||
}
|
||||
|
||||
public async currentConcurrencyOfTask(
|
||||
env: MinimalAuthenticatedEnvironment,
|
||||
taskIdentifier: string
|
||||
@@ -282,6 +308,41 @@ export class RunQueue {
|
||||
return this.redis.scard(this.keys.taskIdentifierCurrentConcurrencyKey(env, taskIdentifier));
|
||||
}
|
||||
|
||||
public async readMessage(orgId: string, messageId: string) {
|
||||
return this.#trace(
|
||||
"readMessage",
|
||||
async (span) => {
|
||||
const rawMessage = await this.redis.get(this.keys.messageKey(orgId, messageId));
|
||||
|
||||
if (!rawMessage) {
|
||||
return;
|
||||
}
|
||||
|
||||
const message = OutputPayload.safeParse(JSON.parse(rawMessage));
|
||||
|
||||
if (!message.success) {
|
||||
this.logger.error(`[${this.name}] Failed to parse message`, {
|
||||
messageId,
|
||||
error: message.error,
|
||||
service: this.name,
|
||||
});
|
||||
|
||||
return;
|
||||
}
|
||||
|
||||
return message.data;
|
||||
},
|
||||
{
|
||||
attributes: {
|
||||
[SEMATTRS_MESSAGING_OPERATION]: "receive",
|
||||
[SEMATTRS_MESSAGE_ID]: messageId,
|
||||
[SEMATTRS_MESSAGING_SYSTEM]: "marqs",
|
||||
[SemanticAttributes.RUN_ID]: messageId,
|
||||
},
|
||||
}
|
||||
);
|
||||
}
|
||||
|
||||
public async enqueueMessage({
|
||||
env,
|
||||
message,
|
||||
@@ -434,7 +495,7 @@ export class RunQueue {
|
||||
return this.#trace(
|
||||
"acknowledgeMessage",
|
||||
async (span) => {
|
||||
const message = await this.#readMessage(orgId, messageId);
|
||||
const message = await this.readMessage(orgId, messageId);
|
||||
|
||||
if (!message) {
|
||||
this.logger.log(`[${this.name}].acknowledgeMessage() message not found`, {
|
||||
@@ -452,18 +513,7 @@ export class RunQueue {
|
||||
});
|
||||
|
||||
await this.#callAcknowledgeMessage({
|
||||
messageId,
|
||||
messageQueue: message.queue,
|
||||
masterQueues: message.masterQueues,
|
||||
messageKey: this.keys.messageKey(orgId, messageId),
|
||||
concurrencyKey: this.keys.currentConcurrencyKeyFromQueue(message.queue),
|
||||
envConcurrencyKey: this.keys.envCurrentConcurrencyKeyFromQueue(message.queue),
|
||||
taskConcurrencyKey: this.keys.taskIdentifierCurrentConcurrencyKeyFromQueue(
|
||||
message.queue,
|
||||
message.taskIdentifier
|
||||
),
|
||||
envQueueKey: this.keys.envQueueKeyFromQueue(message.queue),
|
||||
projectConcurrencyKey: this.keys.projectCurrentConcurrencyKeyFromQueue(message.queue),
|
||||
message,
|
||||
});
|
||||
},
|
||||
{
|
||||
@@ -497,7 +547,7 @@ export class RunQueue {
|
||||
async (span) => {
|
||||
const maxAttempts = this.retryOptions.maxAttempts ?? defaultRetrySettings.maxAttempts;
|
||||
|
||||
const message = await this.#readMessage(orgId, messageId);
|
||||
const message = await this.readMessage(orgId, messageId);
|
||||
if (!message) {
|
||||
this.logger.log(`[${this.name}].nackMessage() message not found`, {
|
||||
orgId,
|
||||
@@ -516,75 +566,16 @@ export class RunQueue {
|
||||
[SemanticAttributes.MASTER_QUEUES]: message.masterQueues.join(","),
|
||||
});
|
||||
|
||||
const messageKey = this.keys.messageKey(orgId, messageId);
|
||||
const messageQueue = message.queue;
|
||||
const concurrencyKey = this.keys.currentConcurrencyKeyFromQueue(message.queue);
|
||||
const envConcurrencyKey = this.keys.envCurrentConcurrencyKeyFromQueue(message.queue);
|
||||
const taskConcurrencyKey = this.keys.taskIdentifierCurrentConcurrencyKeyFromQueue(
|
||||
message.queue,
|
||||
message.taskIdentifier
|
||||
);
|
||||
const projectConcurrencyKey = this.keys.projectCurrentConcurrencyKeyFromQueue(
|
||||
message.queue
|
||||
);
|
||||
const envQueueKey = this.keys.envQueueKeyFromQueue(message.queue);
|
||||
|
||||
if (incrementAttemptCount) {
|
||||
message.attempt = message.attempt + 1;
|
||||
if (message.attempt >= maxAttempts) {
|
||||
await this.redis.moveToDeadLetterQueue(
|
||||
messageKey,
|
||||
messageQueue,
|
||||
concurrencyKey,
|
||||
envConcurrencyKey,
|
||||
projectConcurrencyKey,
|
||||
envQueueKey,
|
||||
taskConcurrencyKey,
|
||||
"dlq",
|
||||
messageId,
|
||||
messageQueue,
|
||||
JSON.stringify(message.masterQueues),
|
||||
this.options.redis.keyPrefix ?? ""
|
||||
);
|
||||
await this.#callMoveToDeadLetterQueue({ message });
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
const nextRetryDelay = calculateNextRetryDelay(this.retryOptions, message.attempt);
|
||||
const messageScore = retryAt ?? (nextRetryDelay ? Date.now() + nextRetryDelay : Date.now());
|
||||
await this.#callNackMessage({ message });
|
||||
|
||||
this.logger.debug("Calling nackMessage", {
|
||||
messageKey,
|
||||
messageQueue,
|
||||
masterQueues: message.masterQueues,
|
||||
concurrencyKey,
|
||||
envConcurrencyKey,
|
||||
projectConcurrencyKey,
|
||||
envQueueKey,
|
||||
taskConcurrencyKey,
|
||||
messageId,
|
||||
messageScore,
|
||||
attempt: message.attempt,
|
||||
service: this.name,
|
||||
});
|
||||
|
||||
await this.redis.nackMessage(
|
||||
//keys
|
||||
messageKey,
|
||||
messageQueue,
|
||||
concurrencyKey,
|
||||
envConcurrencyKey,
|
||||
projectConcurrencyKey,
|
||||
envQueueKey,
|
||||
taskConcurrencyKey,
|
||||
//args
|
||||
messageId,
|
||||
messageQueue,
|
||||
JSON.stringify(message),
|
||||
String(messageScore),
|
||||
JSON.stringify(message.masterQueues),
|
||||
this.options.redis.keyPrefix ?? ""
|
||||
);
|
||||
return true;
|
||||
},
|
||||
{
|
||||
@@ -606,7 +597,7 @@ export class RunQueue {
|
||||
return this.#trace(
|
||||
"releaseConcurrency",
|
||||
async (span) => {
|
||||
const message = await this.#readMessage(orgId, messageId);
|
||||
const message = await this.readMessage(orgId, messageId);
|
||||
|
||||
if (!message) {
|
||||
this.logger.log(`[${this.name}].acknowledgeMessage() message not found`, {
|
||||
@@ -652,7 +643,7 @@ export class RunQueue {
|
||||
return this.#trace(
|
||||
"reacquireConcurrency",
|
||||
async (span) => {
|
||||
const message = await this.#readMessage(orgId, messageId);
|
||||
const message = await this.readMessage(orgId, messageId);
|
||||
|
||||
if (!message) {
|
||||
this.logger.log(`[${this.name}].acknowledgeMessage() message not found`, {
|
||||
@@ -702,16 +693,21 @@ export class RunQueue {
|
||||
|
||||
private async handleRedriveMessage(channel: string, message: string) {
|
||||
try {
|
||||
const { runId, orgId } = JSON.parse(message) as any;
|
||||
if (typeof orgId !== "string" || typeof runId !== "string") {
|
||||
const { runId, envId, projectId, orgId } = JSON.parse(message) as any;
|
||||
if (
|
||||
typeof orgId !== "string" ||
|
||||
typeof runId !== "string" ||
|
||||
typeof envId !== "string" ||
|
||||
typeof projectId !== "string"
|
||||
) {
|
||||
this.logger.error(
|
||||
"handleRedriveMessage: invalid message format: runId and orgId must be strings",
|
||||
"handleRedriveMessage: invalid message format: runId, envId, projectId and orgId must be strings",
|
||||
{ message, channel }
|
||||
);
|
||||
return;
|
||||
}
|
||||
|
||||
const data = await this.#readMessage(orgId, runId);
|
||||
const data = await this.readMessage(orgId, runId);
|
||||
|
||||
if (!data) {
|
||||
this.logger.error(`handleRedriveMessage: couldn't read message`, { orgId, runId, channel });
|
||||
@@ -739,7 +735,10 @@ export class RunQueue {
|
||||
});
|
||||
|
||||
//remove from the dlq
|
||||
const result = await this.redis.zrem("dlq", runId);
|
||||
const result = await this.redis.zrem(
|
||||
this.keys.deadLetterQueueKey({ envId, orgId, projectId }),
|
||||
runId
|
||||
);
|
||||
|
||||
if (result === 0) {
|
||||
this.logger.error(`handleRedriveMessage: couldn't remove message from dlq`, {
|
||||
@@ -800,41 +799,6 @@ export class RunQueue {
|
||||
this.subscriber.on("message", this.handleRedriveMessage.bind(this));
|
||||
}
|
||||
|
||||
async #readMessage(orgId: string, messageId: string) {
|
||||
return this.#trace(
|
||||
"readMessage",
|
||||
async (span) => {
|
||||
const rawMessage = await this.redis.get(this.keys.messageKey(orgId, messageId));
|
||||
|
||||
if (!rawMessage) {
|
||||
return;
|
||||
}
|
||||
|
||||
const message = OutputPayload.safeParse(JSON.parse(rawMessage));
|
||||
|
||||
if (!message.success) {
|
||||
this.logger.error(`[${this.name}] Failed to parse message`, {
|
||||
messageId,
|
||||
error: message.error,
|
||||
service: this.name,
|
||||
});
|
||||
|
||||
return;
|
||||
}
|
||||
|
||||
return message.data;
|
||||
},
|
||||
{
|
||||
attributes: {
|
||||
[SEMATTRS_MESSAGING_OPERATION]: "receive",
|
||||
[SEMATTRS_MESSAGE_ID]: messageId,
|
||||
[SEMATTRS_MESSAGING_SYSTEM]: "marqs",
|
||||
[SemanticAttributes.RUN_ID]: messageId,
|
||||
},
|
||||
}
|
||||
);
|
||||
}
|
||||
|
||||
async #callEnqueueMessage(
|
||||
message: OutputPayload,
|
||||
masterQueues: string[],
|
||||
@@ -1085,35 +1049,30 @@ export class RunQueue {
|
||||
};
|
||||
}
|
||||
|
||||
async #callAcknowledgeMessage({
|
||||
messageId,
|
||||
masterQueues,
|
||||
messageKey,
|
||||
messageQueue,
|
||||
concurrencyKey,
|
||||
envConcurrencyKey,
|
||||
taskConcurrencyKey,
|
||||
envQueueKey,
|
||||
projectConcurrencyKey,
|
||||
}: {
|
||||
masterQueues: string[];
|
||||
messageKey: string;
|
||||
messageQueue: string;
|
||||
concurrencyKey: string;
|
||||
envConcurrencyKey: string;
|
||||
taskConcurrencyKey: string;
|
||||
envQueueKey: string;
|
||||
projectConcurrencyKey: string;
|
||||
messageId: string;
|
||||
}) {
|
||||
async #callAcknowledgeMessage({ message }: { message: OutputPayload }) {
|
||||
const messageId = message.runId;
|
||||
const messageKey = this.keys.messageKey(message.orgId, messageId);
|
||||
const messageQueue = message.queue;
|
||||
const queueCurrentConcurrencyKey = this.keys.currentConcurrencyKeyFromQueue(message.queue);
|
||||
const envCurrentConcurrencyKey = this.keys.envCurrentConcurrencyKeyFromQueue(message.queue);
|
||||
const projectCurrentConcurrencyKey = this.keys.projectCurrentConcurrencyKeyFromQueue(
|
||||
message.queue
|
||||
);
|
||||
const envQueueKey = this.keys.envQueueKeyFromQueue(message.queue);
|
||||
const taskCurrentConcurrencyKey = this.keys.taskIdentifierCurrentConcurrencyKeyFromQueue(
|
||||
message.queue,
|
||||
message.taskIdentifier
|
||||
);
|
||||
const masterQueues = message.masterQueues;
|
||||
|
||||
this.logger.debug("Calling acknowledgeMessage", {
|
||||
messageKey,
|
||||
messageQueue,
|
||||
concurrencyKey,
|
||||
envConcurrencyKey,
|
||||
projectConcurrencyKey,
|
||||
queueCurrentConcurrencyKey,
|
||||
envCurrentConcurrencyKey,
|
||||
projectCurrentConcurrencyKey,
|
||||
envQueueKey,
|
||||
taskConcurrencyKey,
|
||||
taskCurrentConcurrencyKey,
|
||||
messageId,
|
||||
masterQueues,
|
||||
service: this.name,
|
||||
@@ -1125,11 +1084,11 @@ export class RunQueue {
|
||||
return this.redis.acknowledgeMessage(
|
||||
messageKey,
|
||||
messageQueue,
|
||||
concurrencyKey,
|
||||
envConcurrencyKey,
|
||||
projectConcurrencyKey,
|
||||
queueCurrentConcurrencyKey,
|
||||
envCurrentConcurrencyKey,
|
||||
projectCurrentConcurrencyKey,
|
||||
envQueueKey,
|
||||
taskConcurrencyKey,
|
||||
taskCurrentConcurrencyKey,
|
||||
queueReserveConcurrencyKey,
|
||||
envReserveConcurrencyKey,
|
||||
messageId,
|
||||
@@ -1139,6 +1098,90 @@ export class RunQueue {
|
||||
);
|
||||
}
|
||||
|
||||
async #callNackMessage({ message, retryAt }: { message: OutputPayload; retryAt?: number }) {
|
||||
const messageId = message.runId;
|
||||
const messageKey = this.keys.messageKey(message.orgId, message.runId);
|
||||
const messageQueue = message.queue;
|
||||
const queueCurrentConcurrencyKey = this.keys.currentConcurrencyKeyFromQueue(message.queue);
|
||||
const envCurrentConcurrencyKey = this.keys.envCurrentConcurrencyKeyFromQueue(message.queue);
|
||||
const projectCurrentConcurrencyKey = this.keys.projectCurrentConcurrencyKeyFromQueue(
|
||||
message.queue
|
||||
);
|
||||
const envQueueKey = this.keys.envQueueKeyFromQueue(message.queue);
|
||||
const taskCurrentConcurrencyKey = this.keys.taskIdentifierCurrentConcurrencyKeyFromQueue(
|
||||
message.queue,
|
||||
message.taskIdentifier
|
||||
);
|
||||
|
||||
const nextRetryDelay = calculateNextRetryDelay(this.retryOptions, message.attempt);
|
||||
const messageScore = retryAt ?? (nextRetryDelay ? Date.now() + nextRetryDelay : Date.now());
|
||||
|
||||
this.logger.debug("Calling nackMessage", {
|
||||
messageKey,
|
||||
messageQueue,
|
||||
masterQueues: message.masterQueues,
|
||||
queueCurrentConcurrencyKey,
|
||||
envCurrentConcurrencyKey,
|
||||
projectCurrentConcurrencyKey,
|
||||
envQueueKey,
|
||||
taskCurrentConcurrencyKey,
|
||||
messageId,
|
||||
messageScore,
|
||||
attempt: message.attempt,
|
||||
service: this.name,
|
||||
});
|
||||
|
||||
await this.redis.nackMessage(
|
||||
//keys
|
||||
messageKey,
|
||||
messageQueue,
|
||||
queueCurrentConcurrencyKey,
|
||||
envCurrentConcurrencyKey,
|
||||
projectCurrentConcurrencyKey,
|
||||
envQueueKey,
|
||||
taskCurrentConcurrencyKey,
|
||||
//args
|
||||
messageId,
|
||||
messageQueue,
|
||||
JSON.stringify(message),
|
||||
String(messageScore),
|
||||
JSON.stringify(message.masterQueues),
|
||||
this.options.redis.keyPrefix ?? ""
|
||||
);
|
||||
}
|
||||
|
||||
async #callMoveToDeadLetterQueue({ message }: { message: OutputPayload }) {
|
||||
const messageId = message.runId;
|
||||
const messageKey = this.keys.messageKey(message.orgId, message.runId);
|
||||
const messageQueue = message.queue;
|
||||
const queueCurrentConcurrencyKey = this.keys.currentConcurrencyKeyFromQueue(message.queue);
|
||||
const envCurrentConcurrencyKey = this.keys.envCurrentConcurrencyKeyFromQueue(message.queue);
|
||||
const projectCurrentConcurrencyKey = this.keys.projectCurrentConcurrencyKeyFromQueue(
|
||||
message.queue
|
||||
);
|
||||
const envQueueKey = this.keys.envQueueKeyFromQueue(message.queue);
|
||||
const taskCurrentConcurrencyKey = this.keys.taskIdentifierCurrentConcurrencyKeyFromQueue(
|
||||
message.queue,
|
||||
message.taskIdentifier
|
||||
);
|
||||
const deadLetterQueueKey = this.keys.deadLetterQueueKeyFromQueue(message.queue);
|
||||
|
||||
await this.redis.moveToDeadLetterQueue(
|
||||
messageKey,
|
||||
messageQueue,
|
||||
queueCurrentConcurrencyKey,
|
||||
envCurrentConcurrencyKey,
|
||||
projectCurrentConcurrencyKey,
|
||||
envQueueKey,
|
||||
taskCurrentConcurrencyKey,
|
||||
deadLetterQueueKey,
|
||||
messageId,
|
||||
messageQueue,
|
||||
JSON.stringify(message.masterQueues),
|
||||
this.options.redis.keyPrefix ?? ""
|
||||
);
|
||||
}
|
||||
|
||||
#callUpdateGlobalConcurrencyLimits({
|
||||
envConcurrencyLimitKey,
|
||||
envConcurrencyLimit,
|
||||
@@ -1501,11 +1544,11 @@ redis.call('SREM', envReserveConcurrencyKey, messageId)
|
||||
-- Keys:
|
||||
local messageKey = KEYS[1]
|
||||
local messageQueueKey = KEYS[2]
|
||||
local concurrencyKey = KEYS[3]
|
||||
local envConcurrencyKey = KEYS[4]
|
||||
local projectConcurrencyKey = KEYS[5]
|
||||
local queueCurrentConcurrencyKey = KEYS[3]
|
||||
local envCurrentConcurrencyKey = KEYS[4]
|
||||
local projectCurrentConcurrencyKey = KEYS[5]
|
||||
local envQueueKey = KEYS[6]
|
||||
local taskConcurrencyKey = KEYS[7]
|
||||
local taskCurrentConcurrencyKey = KEYS[7]
|
||||
|
||||
-- Args:
|
||||
local messageId = ARGV[1]
|
||||
@@ -1519,10 +1562,10 @@ local keyPrefix = ARGV[6]
|
||||
redis.call('SET', messageKey, messageData)
|
||||
|
||||
-- Update the concurrency keys
|
||||
redis.call('SREM', concurrencyKey, messageId)
|
||||
redis.call('SREM', envConcurrencyKey, messageId)
|
||||
redis.call('SREM', projectConcurrencyKey, messageId)
|
||||
redis.call('SREM', taskConcurrencyKey, messageId)
|
||||
redis.call('SREM', queueCurrentConcurrencyKey, messageId)
|
||||
redis.call('SREM', envCurrentConcurrencyKey, messageId)
|
||||
redis.call('SREM', projectCurrentConcurrencyKey, messageId)
|
||||
redis.call('SREM', taskCurrentConcurrencyKey, messageId)
|
||||
|
||||
-- Enqueue the message into the queue
|
||||
redis.call('ZADD', messageQueueKey, messageScore, messageId)
|
||||
@@ -1547,11 +1590,11 @@ end
|
||||
-- Keys:
|
||||
local messageKey = KEYS[1]
|
||||
local messageQueue = KEYS[2]
|
||||
local concurrencyKey = KEYS[3]
|
||||
local queueCurrentConcurrencyKey = KEYS[3]
|
||||
local envCurrentConcurrencyKey = KEYS[4]
|
||||
local projectCurrentConcurrencyKey = KEYS[5]
|
||||
local envQueueKey = KEYS[6]
|
||||
local taskCurrentConcurrencyKey = KEYS[7]
|
||||
local taskCurrentConcurrencyKeyPrefix = KEYS[7]
|
||||
local deadLetterQueueKey = KEYS[8]
|
||||
|
||||
-- Args:
|
||||
@@ -1579,10 +1622,10 @@ end
|
||||
redis.call('ZADD', deadLetterQueueKey, tonumber(redis.call('TIME')[1]), messageId)
|
||||
|
||||
-- Update the concurrency keys
|
||||
redis.call('SREM', concurrencyKey, messageId)
|
||||
redis.call('SREM', queueCurrentConcurrencyKey, messageId)
|
||||
redis.call('SREM', envCurrentConcurrencyKey, messageId)
|
||||
redis.call('SREM', projectCurrentConcurrencyKey, messageId)
|
||||
redis.call('SREM', taskCurrentConcurrencyKey, messageId)
|
||||
redis.call('SREM', taskCurrentConcurrencyKeyPrefix, messageId)
|
||||
`,
|
||||
});
|
||||
|
||||
@@ -1760,11 +1803,11 @@ declare module "@internal/redis" {
|
||||
nackMessage(
|
||||
messageKey: string,
|
||||
messageQueue: string,
|
||||
concurrencyKey: string,
|
||||
envConcurrencyKey: string,
|
||||
projectConcurrencyKey: string,
|
||||
queueCurrentConcurrencyKey: string,
|
||||
envCurrentConcurrencyKey: string,
|
||||
projectCurrentConcurrencyKey: string,
|
||||
envQueueKey: string,
|
||||
taskConcurrencyKey: string,
|
||||
taskCurrentConcurrencyKey: string,
|
||||
messageId: string,
|
||||
messageQueueName: string,
|
||||
messageData: string,
|
||||
@@ -1777,11 +1820,11 @@ declare module "@internal/redis" {
|
||||
moveToDeadLetterQueue(
|
||||
messageKey: string,
|
||||
messageQueue: string,
|
||||
concurrencyKey: string,
|
||||
envConcurrencyKey: string,
|
||||
projectConcurrencyKey: string,
|
||||
queueCurrentConcurrencyKey: string,
|
||||
envCurrentConcurrencyKey: string,
|
||||
projectCurrentConcurrencyKey: string,
|
||||
envQueueKey: string,
|
||||
taskConcurrencyKey: string,
|
||||
taskCurrentConcurrencyKey: string,
|
||||
deadLetterQueueKey: string,
|
||||
messageId: string,
|
||||
messageQueueName: string,
|
||||
|
||||
@@ -13,6 +13,7 @@ const constants = {
|
||||
TASK_PART: "task",
|
||||
MESSAGE_PART: "message",
|
||||
RESERVE_CONCURRENCY_PART: "reserveConcurrency",
|
||||
DEAD_LETTER_QUEUE_PART: "deadLetter",
|
||||
} as const;
|
||||
|
||||
export class RunQueueFullKeyProducer implements RunQueueKeyProducer {
|
||||
@@ -256,6 +257,30 @@ export class RunQueueFullKeyProducer implements RunQueueKeyProducer {
|
||||
|
||||
return this.envReserveConcurrencyKey(descriptor);
|
||||
}
|
||||
deadLetterQueueKey(env: MinimalAuthenticatedEnvironment): string;
|
||||
deadLetterQueueKey(env: EnvDescriptor): string;
|
||||
deadLetterQueueKey(envOrDescriptor: EnvDescriptor | MinimalAuthenticatedEnvironment): string {
|
||||
if ("id" in envOrDescriptor) {
|
||||
return [
|
||||
this.orgKeySection(envOrDescriptor.organization.id),
|
||||
this.projKeySection(envOrDescriptor.project.id),
|
||||
this.envKeySection(envOrDescriptor.id),
|
||||
constants.DEAD_LETTER_QUEUE_PART,
|
||||
].join(":");
|
||||
} else {
|
||||
return [
|
||||
this.orgKeySection(envOrDescriptor.orgId),
|
||||
this.projKeySection(envOrDescriptor.projectId),
|
||||
this.envKeySection(envOrDescriptor.envId),
|
||||
constants.DEAD_LETTER_QUEUE_PART,
|
||||
].join(":");
|
||||
}
|
||||
}
|
||||
deadLetterQueueKeyFromQueue(queue: string): string {
|
||||
const descriptor = this.descriptorFromQueue(queue);
|
||||
|
||||
return this.deadLetterQueueKey(descriptor);
|
||||
}
|
||||
|
||||
private queueReserveConcurrencyKeyFromDescriptor(descriptor: QueueDescriptor) {
|
||||
return [
|
||||
|
||||
@@ -0,0 +1,244 @@
|
||||
import { redisTest } from "@internal/testcontainers";
|
||||
import { trace } from "@internal/tracing";
|
||||
import { Logger } from "@trigger.dev/core/logger";
|
||||
import { describe } from "node:test";
|
||||
import { FairQueueSelectionStrategy } from "../fairQueueSelectionStrategy.js";
|
||||
import { RunQueue } from "../index.js";
|
||||
import { RunQueueFullKeyProducer } from "../keyProducer.js";
|
||||
import { InputPayload } from "../types.js";
|
||||
import { setTimeout } from "node:timers/promises";
|
||||
|
||||
const testOptions = {
|
||||
name: "rq",
|
||||
tracer: trace.getTracer("rq"),
|
||||
workers: 1,
|
||||
defaultEnvConcurrency: 25,
|
||||
logger: new Logger("RunQueue", "warn"),
|
||||
retryOptions: {
|
||||
maxAttempts: 5,
|
||||
factor: 1.1,
|
||||
minTimeoutInMs: 100,
|
||||
maxTimeoutInMs: 1_000,
|
||||
randomize: true,
|
||||
},
|
||||
keys: new RunQueueFullKeyProducer(),
|
||||
};
|
||||
|
||||
const authenticatedEnvDev = {
|
||||
id: "e1234",
|
||||
type: "DEVELOPMENT" as const,
|
||||
maximumConcurrencyLimit: 10,
|
||||
project: { id: "p1234" },
|
||||
organization: { id: "o1234" },
|
||||
};
|
||||
|
||||
const messageDev: InputPayload = {
|
||||
runId: "r4321",
|
||||
taskIdentifier: "task/my-task",
|
||||
orgId: "o1234",
|
||||
projectId: "p1234",
|
||||
environmentId: "e4321",
|
||||
environmentType: "DEVELOPMENT",
|
||||
queue: "task/my-task",
|
||||
timestamp: Date.now(),
|
||||
attempt: 0,
|
||||
};
|
||||
|
||||
vi.setConfig({ testTimeout: 60_000 });
|
||||
|
||||
describe("RunQueue.nackMessage", () => {
|
||||
redisTest("nacking a message clears all concurrency", async ({ redisContainer }) => {
|
||||
const queue = new RunQueue({
|
||||
...testOptions,
|
||||
queueSelectionStrategy: new FairQueueSelectionStrategy({
|
||||
redis: {
|
||||
keyPrefix: "runqueue:test:",
|
||||
host: redisContainer.getHost(),
|
||||
port: redisContainer.getPort(),
|
||||
},
|
||||
keys: testOptions.keys,
|
||||
}),
|
||||
redis: {
|
||||
keyPrefix: "runqueue:test:",
|
||||
host: redisContainer.getHost(),
|
||||
port: redisContainer.getPort(),
|
||||
},
|
||||
});
|
||||
|
||||
try {
|
||||
const envMasterQueue = `env:${authenticatedEnvDev.id}`;
|
||||
|
||||
// Enqueue message with reserve concurrency
|
||||
await queue.enqueueMessage({
|
||||
env: authenticatedEnvDev,
|
||||
message: messageDev,
|
||||
masterQueues: ["main", envMasterQueue],
|
||||
reserveConcurrency: {
|
||||
messageId: messageDev.runId,
|
||||
recursiveQueue: true,
|
||||
},
|
||||
});
|
||||
|
||||
// Verify reserve concurrency is set
|
||||
const queueReserveConcurrency = await queue.reserveConcurrencyOfQueue(
|
||||
authenticatedEnvDev,
|
||||
messageDev.queue
|
||||
);
|
||||
expect(queueReserveConcurrency).toBe(1);
|
||||
|
||||
const envReserveConcurrency = await queue.reserveConcurrencyOfEnvironment(
|
||||
authenticatedEnvDev
|
||||
);
|
||||
expect(envReserveConcurrency).toBe(1);
|
||||
|
||||
// Dequeue message
|
||||
const dequeued = await queue.dequeueMessageFromMasterQueue("test_12345", envMasterQueue, 10);
|
||||
expect(dequeued.length).toBe(1);
|
||||
|
||||
// Verify current concurrency is set and reserve is cleared
|
||||
const queueCurrentConcurrency = await queue.currentConcurrencyOfQueue(
|
||||
authenticatedEnvDev,
|
||||
messageDev.queue
|
||||
);
|
||||
expect(queueCurrentConcurrency).toBe(1);
|
||||
|
||||
const envCurrentConcurrency = await queue.currentConcurrencyOfEnvironment(
|
||||
authenticatedEnvDev
|
||||
);
|
||||
expect(envCurrentConcurrency).toBe(1);
|
||||
|
||||
// Nack the message
|
||||
await queue.nackMessage({
|
||||
orgId: messageDev.orgId,
|
||||
messageId: messageDev.runId,
|
||||
});
|
||||
|
||||
// Verify all concurrency is cleared
|
||||
const queueCurrentConcurrencyAfterNack = await queue.currentConcurrencyOfQueue(
|
||||
authenticatedEnvDev,
|
||||
messageDev.queue
|
||||
);
|
||||
expect(queueCurrentConcurrencyAfterNack).toBe(0);
|
||||
|
||||
const envCurrentConcurrencyAfterNack = await queue.currentConcurrencyOfEnvironment(
|
||||
authenticatedEnvDev
|
||||
);
|
||||
expect(envCurrentConcurrencyAfterNack).toBe(0);
|
||||
|
||||
const projectCurrentConcurrencyAfterNack = await queue.currentConcurrencyOfProject(
|
||||
authenticatedEnvDev
|
||||
);
|
||||
expect(projectCurrentConcurrencyAfterNack).toBe(0);
|
||||
|
||||
const taskCurrentConcurrencyAfterNack = await queue.currentConcurrencyOfTask(
|
||||
authenticatedEnvDev,
|
||||
messageDev.taskIdentifier
|
||||
);
|
||||
expect(taskCurrentConcurrencyAfterNack).toBe(0);
|
||||
|
||||
const envQueueLength = await queue.lengthOfEnvQueue(authenticatedEnvDev);
|
||||
expect(envQueueLength).toBe(1);
|
||||
|
||||
const message = await queue.readMessage(messageDev.orgId, messageDev.runId);
|
||||
expect(message?.attempt).toBe(1);
|
||||
|
||||
//we need to wait because the default wait is 1 second
|
||||
await setTimeout(300);
|
||||
|
||||
// Now we should be able to dequeue it again
|
||||
const dequeued2 = await queue.dequeueMessageFromMasterQueue("test_12345", envMasterQueue, 10);
|
||||
expect(dequeued2.length).toBe(1);
|
||||
} finally {
|
||||
await queue.quit();
|
||||
}
|
||||
});
|
||||
|
||||
redisTest(
|
||||
"nacking a message with maxAttempts reached should be moved to dead letter queue",
|
||||
async ({ redisContainer }) => {
|
||||
const queue = new RunQueue({
|
||||
...testOptions,
|
||||
retryOptions: {
|
||||
...testOptions.retryOptions,
|
||||
maxAttempts: 2, // Set lower for testing
|
||||
},
|
||||
queueSelectionStrategy: new FairQueueSelectionStrategy({
|
||||
redis: {
|
||||
keyPrefix: "runqueue:test:",
|
||||
host: redisContainer.getHost(),
|
||||
port: redisContainer.getPort(),
|
||||
},
|
||||
keys: testOptions.keys,
|
||||
}),
|
||||
redis: {
|
||||
keyPrefix: "runqueue:test:",
|
||||
host: redisContainer.getHost(),
|
||||
port: redisContainer.getPort(),
|
||||
},
|
||||
});
|
||||
|
||||
try {
|
||||
const envMasterQueue = `env:${authenticatedEnvDev.id}`;
|
||||
|
||||
await queue.enqueueMessage({
|
||||
env: authenticatedEnvDev,
|
||||
message: messageDev,
|
||||
masterQueues: ["main", envMasterQueue],
|
||||
});
|
||||
|
||||
const dequeued = await queue.dequeueMessageFromMasterQueue(
|
||||
"test_12345",
|
||||
envMasterQueue,
|
||||
10
|
||||
);
|
||||
expect(dequeued.length).toBe(1);
|
||||
|
||||
await queue.nackMessage({
|
||||
orgId: messageDev.orgId,
|
||||
messageId: messageDev.runId,
|
||||
});
|
||||
|
||||
// Wait for any requeue delay
|
||||
await setTimeout(300);
|
||||
|
||||
// Message should not be requeued as max attempts reached
|
||||
const envQueueLength = await queue.lengthOfEnvQueue(authenticatedEnvDev);
|
||||
expect(envQueueLength).toBe(1);
|
||||
|
||||
const message = await queue.readMessage(messageDev.orgId, messageDev.runId);
|
||||
expect(message?.attempt).toBe(1);
|
||||
|
||||
// Now we dequeue and nack again, and it should be moved to dead letter queue
|
||||
const dequeued3 = await queue.dequeueMessageFromMasterQueue(
|
||||
"test_12345",
|
||||
envMasterQueue,
|
||||
10
|
||||
);
|
||||
expect(dequeued3.length).toBe(1);
|
||||
|
||||
const envQueueLengthDequeue = await queue.lengthOfEnvQueue(authenticatedEnvDev);
|
||||
expect(envQueueLengthDequeue).toBe(0);
|
||||
|
||||
const deadLetterQueueLengthBefore = await queue.lengthOfDeadLetterQueue(
|
||||
authenticatedEnvDev
|
||||
);
|
||||
expect(deadLetterQueueLengthBefore).toBe(0);
|
||||
|
||||
await queue.nackMessage({
|
||||
orgId: messageDev.orgId,
|
||||
messageId: messageDev.runId,
|
||||
});
|
||||
|
||||
const envQueueLengthAfterNack = await queue.lengthOfEnvQueue(authenticatedEnvDev);
|
||||
expect(envQueueLengthAfterNack).toBe(0);
|
||||
|
||||
const deadLetterQueueLengthAfterNack = await queue.lengthOfDeadLetterQueue(
|
||||
authenticatedEnvDev
|
||||
);
|
||||
expect(deadLetterQueueLengthAfterNack).toBe(1);
|
||||
} finally {
|
||||
await queue.quit();
|
||||
}
|
||||
}
|
||||
);
|
||||
});
|
||||
@@ -92,6 +92,10 @@ export interface RunQueueKeyProducer {
|
||||
|
||||
reserveConcurrencyKeyFromQueue(queue: string): string;
|
||||
envReserveConcurrencyKeyFromQueue(queue: string): string;
|
||||
|
||||
deadLetterQueueKey(env: MinimalAuthenticatedEnvironment): string;
|
||||
deadLetterQueueKey(env: EnvDescriptor): string;
|
||||
deadLetterQueueKeyFromQueue(queue: string): string;
|
||||
}
|
||||
|
||||
export type EnvQueues = {
|
||||
|
||||
@@ -12,5 +12,6 @@ export default defineConfig({
|
||||
singleThread: true,
|
||||
},
|
||||
},
|
||||
testTimeout: 60_000,
|
||||
},
|
||||
});
|
||||
|
||||
@@ -0,0 +1,11 @@
|
||||
import { redisTest } from "@internal/testcontainers";
|
||||
import { trace } from "@internal/tracing";
|
||||
import { Logger } from "@trigger.dev/core/logger";
|
||||
import { describe } from "node:test";
|
||||
import { FairQueueSelectionStrategy } from "../fairQueueSelectionStrategy.js";
|
||||
import { RunQueue } from "../index.js";
|
||||
import { RunQueueFullKeyProducer } from "../keyProducer.js";
|
||||
import { InputPayload } from "../types.js";
|
||||
import { createRedisClient } from "@internal/redis";
|
||||
|
||||
// ... existing code ...
|
||||
Reference in New Issue
Block a user