Correctly use the new release concurrency queue in the run engine

This commit is contained in:
Eric Allam
2025-03-13 22:45:58 +00:00
parent bf41703dfd
commit b7ddf20d62
6 changed files with 131 additions and 118 deletions
@@ -71,7 +71,7 @@ import {
isPendingExecuting,
} from "./statuses.js";
import { HeartbeatTimeouts, RunEngineOptions, TriggerParams } from "./types.js";
import { ReleaseConcurrencyQueue } from "./releaseConcurrencyQueue.js";
import { ReleaseConcurrencyTokenBucketQueue } from "./releaseConcurrencyTokenBucketQueue.js";
const workerCatalog = {
finishWaitpoint: {
@@ -139,7 +139,11 @@ export class RunEngine {
private logger = new Logger("RunEngine", "debug");
private tracer: Tracer;
private heartbeatTimeouts: HeartbeatTimeouts;
private releaseConcurrencyQueue: ReleaseConcurrencyQueue;
private releaseConcurrencyQueue: ReleaseConcurrencyTokenBucketQueue<{
orgId: string;
projectId: string;
envId: string;
}>;
eventBus = new EventEmitter<EventBusEvents>();
constructor(private readonly options: RunEngineOptions) {
@@ -244,15 +248,43 @@ export class RunEngine {
};
// Initialize the ReleaseConcurrencyQueue
this.releaseConcurrencyQueue = new ReleaseConcurrencyQueue({
this.releaseConcurrencyQueue = new ReleaseConcurrencyTokenBucketQueue({
redis: {
...options.queue.redis, // Use base queue redis options
...options.releaseConcurrency?.redis, // Allow overrides
keyPrefix: `${options.queue.redis.keyPrefix}release-concurrency:`,
},
maxTokens: options.releaseConcurrency?.maxTokens ?? 10, // Default to 10 tokens
executor: async (releaseQueue, runId) => {
await this.#executeReleasedConcurrencyFromQueue(releaseQueue, runId);
retry: {
maxRetries: 5, // TODO: Make this configurable
backoff: {
minDelay: 1000, // TODO: Make this configurable
maxDelay: 10000, // TODO: Make this configurable
factor: 2, // TODO: Make this configurable
},
},
executor: async (descriptor, runId) => {
await this.#executeReleasedConcurrencyFromQueue(descriptor, runId);
},
maxTokens: async (descriptor) => {
const environment = await this.prisma.runtimeEnvironment.findFirstOrThrow({
where: { id: descriptor.envId },
select: {
maximumConcurrencyLimit: true,
},
});
return (
environment.maximumConcurrencyLimit * (options.releaseConcurrency?.maxTokensRatio ?? 1.0)
);
},
keys: {
fromDescriptor: (descriptor) =>
`org:${descriptor.orgId}:proj:${descriptor.projectId}:env:${descriptor.envId}`,
toDescriptor: (name) => ({
orgId: name.split(":")[1],
projectId: name.split(":")[3],
envId: name.split(":")[5],
}),
},
tracer: this.tracer,
});
@@ -2149,25 +2181,24 @@ export class RunEngine {
}
await this.releaseConcurrencyQueue.attemptToRelease(
this.runQueue.keys.releaseConcurrencyKey({
{
orgId: run.runtimeEnvironment.organizationId,
projectId: run.runtimeEnvironment.projectId,
envId: run.runtimeEnvironment.id,
}),
},
snapshot.runId
);
return;
}
async #executeReleasedConcurrencyFromQueue(releaseQueue: string, runId: string) {
const releaseQueueDescriptor =
this.runQueue.keys.releaseConcurrencyDescriptorFromQueue(releaseQueue);
async #executeReleasedConcurrencyFromQueue(
descriptor: { orgId: string; projectId: string; envId: string },
runId: string
) {
this.logger.debug("Executing released concurrency", {
releaseQueue,
descriptor,
runId,
releaseQueueDescriptor,
});
// - Runlock the run
@@ -2186,7 +2217,7 @@ export class RunEngine {
return;
}
return await this.runQueue.releaseConcurrency(releaseQueueDescriptor.orgId, snapshot.runId);
return await this.runQueue.releaseConcurrency(descriptor.orgId, snapshot.runId);
});
}
@@ -2418,11 +2449,11 @@ export class RunEngine {
// Refill the token bucket for the release concurrency queue
await this.releaseConcurrencyQueue.refillTokens(
this.runQueue.keys.releaseConcurrencyKey({
{
orgId: run.runtimeEnvironment.organizationId,
projectId: run.runtimeEnvironment.projectId,
envId: run.runtimeEnvironment.id,
}),
},
1
);
@@ -20,6 +20,7 @@ export type ReleaseConcurrencyQueueOptions<T> = {
fromDescriptor: (releaseQueue: T) => string;
toDescriptor: (releaseQueue: string) => T;
};
maxTokens: (descriptor: T) => Promise<number>;
consumersCount?: number;
masterQueuesKey?: string;
tracer?: Tracer;
@@ -36,7 +37,7 @@ const QueueItemMetadata = z.object({
type QueueItemMetadata = z.infer<typeof QueueItemMetadata>;
export class ReleaseConcurrencyQueue<T> {
export class ReleaseConcurrencyTokenBucketQueue<T> {
private redis: Redis;
private logger: Logger;
private abortController: AbortController;
@@ -47,6 +48,7 @@ export class ReleaseConcurrencyQueue<T> {
private consumersCount: number;
private pollInterval: number;
private keys: ReleaseConcurrencyQueueOptions<T>["keys"];
private maxTokens: ReleaseConcurrencyQueueOptions<T>["maxTokens"];
private batchSize: number;
private maxRetries: number;
private backoff: NonNullable<Required<ReleaseConcurrencyQueueRetryOptions["backoff"]>>;
@@ -62,6 +64,7 @@ export class ReleaseConcurrencyQueue<T> {
this.consumersCount = options.consumersCount ?? 1;
this.pollInterval = options.pollInterval ?? 1000;
this.keys = options.keys;
this.maxTokens = options.maxTokens;
this.batchSize = options.batchSize ?? 5;
this.maxRetries = options.retry?.maxRetries ?? 3;
this.backoff = {
@@ -86,7 +89,8 @@ export class ReleaseConcurrencyQueue<T> {
* If there is no token available, then we'll add the operation to a queue
* and wait until the token is available.
*/
public async attemptToRelease(releaseQueueDescriptor: T, runId: string, maxTokens: number) {
public async attemptToRelease(releaseQueueDescriptor: T, runId: string) {
const maxTokens = await this.maxTokens(releaseQueueDescriptor);
const releaseQueue = this.keys.fromDescriptor(releaseQueueDescriptor);
const result = await this.redis.consumeToken(
@@ -113,7 +117,8 @@ export class ReleaseConcurrencyQueue<T> {
*
* This will add the amount of tokens to the token bucket.
*/
public async refillTokens(releaseQueueDescriptor: T, maxTokens: number, amount: number = 1) {
public async refillTokens(releaseQueueDescriptor: T, amount: number = 1) {
const maxTokens = await this.maxTokens(releaseQueueDescriptor);
const releaseQueue = this.keys.fromDescriptor(releaseQueueDescriptor);
if (amount < 0) {
@@ -533,7 +538,7 @@ class ReleaseConcurrencyQueueConsumer<T> {
private logger: Logger;
constructor(
private readonly queue: ReleaseConcurrencyQueue<T>,
private readonly queue: ReleaseConcurrencyTokenBucketQueue<T>,
private readonly pollInterval: number,
private readonly signal: AbortSignal,
logger?: Logger
@@ -1,15 +1,18 @@
import { redisTest, StartedRedisContainer } from "@internal/testcontainers";
import { ReleaseConcurrencyQueue } from "../releaseConcurrencyQueue.js";
import { ReleaseConcurrencyTokenBucketQueue } from "../releaseConcurrencyTokenBucketQueue.js";
import { setTimeout } from "node:timers/promises";
type TestQueueDescriptor = {
name: string;
};
function createReleaseConcurrencyQueue(redisContainer: StartedRedisContainer) {
function createReleaseConcurrencyQueue(
redisContainer: StartedRedisContainer,
maxTokens: number = 2
) {
const executedRuns: { releaseQueue: string; runId: string }[] = [];
const queue = new ReleaseConcurrencyQueue<TestQueueDescriptor>({
const queue = new ReleaseConcurrencyTokenBucketQueue<TestQueueDescriptor>({
redis: {
keyPrefix: "release-queue:test:",
host: redisContainer.getHost(),
@@ -18,6 +21,7 @@ function createReleaseConcurrencyQueue(redisContainer: StartedRedisContainer) {
executor: async (releaseQueue, runId) => {
executedRuns.push({ releaseQueue: releaseQueue.name, runId });
},
maxTokens: async (_) => maxTokens,
keys: {
fromDescriptor: (descriptor) => descriptor.name,
toDescriptor: (name) => ({ name }),
@@ -33,12 +37,12 @@ function createReleaseConcurrencyQueue(redisContainer: StartedRedisContainer) {
describe("ReleaseConcurrencyQueue", () => {
redisTest("Should manage token bucket and queue correctly", async ({ redisContainer }) => {
const { queue, executedRuns } = createReleaseConcurrencyQueue(redisContainer);
const { queue, executedRuns } = createReleaseConcurrencyQueue(redisContainer, 2);
try {
// First two attempts should execute immediately (we have 2 tokens)
await queue.attemptToRelease({ name: "test-queue" }, "run1", 2);
await queue.attemptToRelease({ name: "test-queue" }, "run2", 2);
await queue.attemptToRelease({ name: "test-queue" }, "run1");
await queue.attemptToRelease({ name: "test-queue" }, "run2");
// Verify first two runs were executed
expect(executedRuns).toHaveLength(2);
@@ -46,11 +50,11 @@ describe("ReleaseConcurrencyQueue", () => {
expect(executedRuns).toContainEqual({ releaseQueue: "test-queue", runId: "run2" });
// Third attempt should be queued (no tokens left)
await queue.attemptToRelease({ name: "test-queue" }, "run3", 2);
await queue.attemptToRelease({ name: "test-queue" }, "run3");
expect(executedRuns).toHaveLength(2); // Still 2, run3 is queued
// Refill one token, should execute run3
await queue.refillTokens({ name: "test-queue" }, 2, 1);
await queue.refillTokens({ name: "test-queue" }, 1);
// Now we need to wait for the queue to be processed
await setTimeout(1000);
@@ -63,15 +67,15 @@ describe("ReleaseConcurrencyQueue", () => {
});
redisTest("Should handle multiple refills correctly", async ({ redisContainer }) => {
const { queue, executedRuns } = createReleaseConcurrencyQueue(redisContainer);
const { queue, executedRuns } = createReleaseConcurrencyQueue(redisContainer, 3);
try {
// Queue up 5 runs (more than maxTokens)
await queue.attemptToRelease({ name: "test-queue" }, "run1", 3);
await queue.attemptToRelease({ name: "test-queue" }, "run2", 3);
await queue.attemptToRelease({ name: "test-queue" }, "run3", 3);
await queue.attemptToRelease({ name: "test-queue" }, "run4", 3);
await queue.attemptToRelease({ name: "test-queue" }, "run5", 3);
await queue.attemptToRelease({ name: "test-queue" }, "run1");
await queue.attemptToRelease({ name: "test-queue" }, "run2");
await queue.attemptToRelease({ name: "test-queue" }, "run3");
await queue.attemptToRelease({ name: "test-queue" }, "run4");
await queue.attemptToRelease({ name: "test-queue" }, "run5");
// First 3 should be executed immediately (maxTokens = 3)
expect(executedRuns).toHaveLength(3);
@@ -80,7 +84,7 @@ describe("ReleaseConcurrencyQueue", () => {
expect(executedRuns).toContainEqual({ releaseQueue: "test-queue", runId: "run3" });
// Refill 2 tokens
await queue.refillTokens({ name: "test-queue" }, 3, 2);
await queue.refillTokens({ name: "test-queue" }, 2);
await setTimeout(1000);
@@ -94,14 +98,14 @@ describe("ReleaseConcurrencyQueue", () => {
});
redisTest("Should handle multiple queues independently", async ({ redisContainer }) => {
const { queue, executedRuns } = createReleaseConcurrencyQueue(redisContainer);
const { queue, executedRuns } = createReleaseConcurrencyQueue(redisContainer, 1);
try {
// Add runs to different queues
await queue.attemptToRelease({ name: "queue1" }, "run1", 1);
await queue.attemptToRelease({ name: "queue1" }, "run2", 1);
await queue.attemptToRelease({ name: "queue2" }, "run3", 1);
await queue.attemptToRelease({ name: "queue2" }, "run4", 1);
await queue.attemptToRelease({ name: "queue1" }, "run1");
await queue.attemptToRelease({ name: "queue1" }, "run2");
await queue.attemptToRelease({ name: "queue2" }, "run3");
await queue.attemptToRelease({ name: "queue2" }, "run4");
// Only first run from each queue should be executed
expect(executedRuns).toHaveLength(2);
@@ -109,7 +113,7 @@ describe("ReleaseConcurrencyQueue", () => {
expect(executedRuns).toContainEqual({ releaseQueue: "queue2", runId: "run3" });
// Refill tokens for queue1
await queue.refillTokens({ name: "queue1" }, 1, 1);
await queue.refillTokens({ name: "queue1" }, 1);
await setTimeout(1000);
@@ -118,7 +122,7 @@ describe("ReleaseConcurrencyQueue", () => {
expect(executedRuns).toContainEqual({ releaseQueue: "queue1", runId: "run2" });
// Refill tokens for queue2
await queue.refillTokens({ name: "queue2" }, 1, 1);
await queue.refillTokens({ name: "queue2" }, 1);
await setTimeout(1000);
@@ -131,19 +135,19 @@ describe("ReleaseConcurrencyQueue", () => {
});
redisTest("Should not allow refilling more than maxTokens", async ({ redisContainer }) => {
const { queue, executedRuns } = createReleaseConcurrencyQueue(redisContainer);
const { queue, executedRuns } = createReleaseConcurrencyQueue(redisContainer, 1);
try {
// Add two runs
await queue.attemptToRelease({ name: "test-queue" }, "run1", 1);
await queue.attemptToRelease({ name: "test-queue" }, "run2", 1);
await queue.attemptToRelease({ name: "test-queue" }, "run1");
await queue.attemptToRelease({ name: "test-queue" }, "run2");
// First run should be executed immediately
expect(executedRuns).toHaveLength(1);
expect(executedRuns).toContainEqual({ releaseQueue: "test-queue", runId: "run1" });
// Refill with more tokens than needed
await queue.refillTokens({ name: "test-queue" }, 1, 5);
await queue.refillTokens({ name: "test-queue" }, 5);
await setTimeout(1000);
@@ -152,7 +156,7 @@ describe("ReleaseConcurrencyQueue", () => {
expect(executedRuns).toContainEqual({ releaseQueue: "test-queue", runId: "run2" });
// Add another run - should NOT execute immediately because we don't have excess tokens
await queue.attemptToRelease({ name: "test-queue" }, "run3", 1);
await queue.attemptToRelease({ name: "test-queue" }, "run3");
expect(executedRuns).toHaveLength(2);
} finally {
await queue.quit();
@@ -160,35 +164,35 @@ describe("ReleaseConcurrencyQueue", () => {
});
redisTest("Should maintain FIFO order when releasing", async ({ redisContainer }) => {
const { queue, executedRuns } = createReleaseConcurrencyQueue(redisContainer);
const { queue, executedRuns } = createReleaseConcurrencyQueue(redisContainer, 1);
try {
// Queue up multiple runs
await queue.attemptToRelease({ name: "test-queue" }, "run1", 1);
await queue.attemptToRelease({ name: "test-queue" }, "run2", 1);
await queue.attemptToRelease({ name: "test-queue" }, "run3", 1);
await queue.attemptToRelease({ name: "test-queue" }, "run4", 1);
await queue.attemptToRelease({ name: "test-queue" }, "run1");
await queue.attemptToRelease({ name: "test-queue" }, "run2");
await queue.attemptToRelease({ name: "test-queue" }, "run3");
await queue.attemptToRelease({ name: "test-queue" }, "run4");
// First run should be executed immediately
expect(executedRuns).toHaveLength(1);
expect(executedRuns[0]).toEqual({ releaseQueue: "test-queue", runId: "run1" });
// Refill tokens one at a time and verify order
await queue.refillTokens({ name: "test-queue" }, 1, 1);
await queue.refillTokens({ name: "test-queue" }, 1);
await setTimeout(1000);
expect(executedRuns).toHaveLength(2);
expect(executedRuns[1]).toEqual({ releaseQueue: "test-queue", runId: "run2" });
await queue.refillTokens({ name: "test-queue" }, 1, 1);
await queue.refillTokens({ name: "test-queue" }, 1);
await setTimeout(1000);
expect(executedRuns).toHaveLength(3);
expect(executedRuns[2]).toEqual({ releaseQueue: "test-queue", runId: "run3" });
await queue.refillTokens({ name: "test-queue" }, 1, 1);
await queue.refillTokens({ name: "test-queue" }, 1);
await setTimeout(1000);
@@ -206,7 +210,7 @@ describe("ReleaseConcurrencyQueue", () => {
const executedRuns: { releaseQueue: string; runId: string }[] = [];
const queue = new ReleaseConcurrencyQueue<string>({
const queue = new ReleaseConcurrencyTokenBucketQueue<string>({
redis: {
keyPrefix: "release-queue:test:",
host: redisContainer.getHost(),
@@ -218,6 +222,7 @@ describe("ReleaseConcurrencyQueue", () => {
}
executedRuns.push({ releaseQueue, runId });
},
maxTokens: async (_) => 2,
keys: {
fromDescriptor: (descriptor) => descriptor,
toDescriptor: (name) => name,
@@ -236,12 +241,12 @@ describe("ReleaseConcurrencyQueue", () => {
try {
// Attempt to release with failing executor
await queue.attemptToRelease("test-queue", "run1", 2);
await queue.attemptToRelease("test-queue", "run1");
// Does not execute because the executor throws an error
expect(executedRuns).toHaveLength(0);
// Token should have been returned to the bucket so this should try to execute immediately and fail again
await queue.attemptToRelease("test-queue", "run2", 2);
await queue.attemptToRelease("test-queue", "run2");
expect(executedRuns).toHaveLength(0);
// Allow executor to succeed
@@ -260,21 +265,21 @@ describe("ReleaseConcurrencyQueue", () => {
);
redisTest("Should handle invalid token amounts", async ({ redisContainer }) => {
const { queue, executedRuns } = createReleaseConcurrencyQueue(redisContainer);
const { queue, executedRuns } = createReleaseConcurrencyQueue(redisContainer, 1);
try {
// Try to refill with negative tokens
await expect(queue.refillTokens({ name: "test-queue" }, 1, -1)).rejects.toThrow();
await expect(queue.refillTokens({ name: "test-queue" }, -1)).rejects.toThrow();
// Try to refill with zero tokens
await queue.refillTokens({ name: "test-queue" }, 1, 0);
await queue.refillTokens({ name: "test-queue" }, 0);
await setTimeout(1000);
expect(executedRuns).toHaveLength(0);
// Verify normal operation still works
await queue.attemptToRelease({ name: "test-queue" }, "run1", 1);
await queue.attemptToRelease({ name: "test-queue" }, "run1");
expect(executedRuns).toHaveLength(1);
} finally {
await queue.quit();
@@ -284,7 +289,7 @@ describe("ReleaseConcurrencyQueue", () => {
redisTest("Should handle concurrent operations correctly", async ({ redisContainer }) => {
const executedRuns: { releaseQueue: string; runId: string }[] = [];
const queue = new ReleaseConcurrencyQueue<string>({
const queue = new ReleaseConcurrencyTokenBucketQueue<string>({
redis: {
keyPrefix: "release-queue:test:",
host: redisContainer.getHost(),
@@ -299,6 +304,7 @@ describe("ReleaseConcurrencyQueue", () => {
fromDescriptor: (descriptor) => descriptor,
toDescriptor: (name) => name,
},
maxTokens: async (_) => 2,
batchSize: 5,
pollInterval: 50,
});
@@ -306,17 +312,17 @@ describe("ReleaseConcurrencyQueue", () => {
try {
// Attempt multiple concurrent releases
await Promise.all([
queue.attemptToRelease("test-queue", "run1", 2),
queue.attemptToRelease("test-queue", "run2", 2),
queue.attemptToRelease("test-queue", "run3", 2),
queue.attemptToRelease("test-queue", "run4", 2),
queue.attemptToRelease("test-queue", "run1"),
queue.attemptToRelease("test-queue", "run2"),
queue.attemptToRelease("test-queue", "run3"),
queue.attemptToRelease("test-queue", "run4"),
]);
// Should only execute maxTokens (2) runs
expect(executedRuns).toHaveLength(2);
// Attempt concurrent refills
queue.refillTokens("test-queue", 2, 2);
await queue.refillTokens("test-queue", 2);
await setTimeout(1000);
@@ -341,25 +347,25 @@ describe("ReleaseConcurrencyQueue", () => {
});
redisTest("Should clean up Redis resources on quit", async ({ redisContainer }) => {
const { queue } = createReleaseConcurrencyQueue(redisContainer);
const { queue } = createReleaseConcurrencyQueue(redisContainer, 1);
// Add some data
await queue.attemptToRelease({ name: "test-queue" }, "run1", 1);
await queue.attemptToRelease({ name: "test-queue" }, "run2", 1);
await queue.attemptToRelease({ name: "test-queue" }, "run1");
await queue.attemptToRelease({ name: "test-queue" }, "run2");
// Quit the queue
await queue.quit();
// Verify we can't perform operations after quit
await expect(queue.attemptToRelease({ name: "test-queue" }, "run3", 1)).rejects.toThrow();
await expect(queue.refillTokens({ name: "test-queue" }, 1, 1)).rejects.toThrow();
await expect(queue.attemptToRelease({ name: "test-queue" }, "run3")).rejects.toThrow();
await expect(queue.refillTokens({ name: "test-queue" }, 1)).rejects.toThrow();
});
redisTest("Should stop retrying after max retries is reached", async ({ redisContainer }) => {
let failCount = 0;
const executedRuns: { releaseQueue: string; runId: string; attempt: number }[] = [];
const queue = new ReleaseConcurrencyQueue<string>({
const queue = new ReleaseConcurrencyTokenBucketQueue<string>({
redis: {
keyPrefix: "release-queue:test:",
host: redisContainer.getHost(),
@@ -374,6 +380,7 @@ describe("ReleaseConcurrencyQueue", () => {
fromDescriptor: (descriptor) => descriptor,
toDescriptor: (name) => name,
},
maxTokens: async (_) => 1,
retry: {
maxRetries: 2, // Set max retries to 2 (will attempt 3 times total: initial + 2 retries)
backoff: {
@@ -387,7 +394,7 @@ describe("ReleaseConcurrencyQueue", () => {
try {
// Attempt to release - this will fail and retry
await queue.attemptToRelease("test-queue", "run1", 1);
await queue.attemptToRelease("test-queue", "run1");
// Wait for retries to occur
await setTimeout(2000);
@@ -404,7 +411,7 @@ describe("ReleaseConcurrencyQueue", () => {
// Attempt a new release to verify the token was returned
let secondRunAttempted = false;
const queue2 = new ReleaseConcurrencyQueue<string>({
const queue2 = new ReleaseConcurrencyTokenBucketQueue<string>({
redis: {
keyPrefix: "release-queue:test:",
host: redisContainer.getHost(),
@@ -417,6 +424,7 @@ describe("ReleaseConcurrencyQueue", () => {
fromDescriptor: (descriptor) => descriptor,
toDescriptor: (name) => name,
},
maxTokens: async (_) => 1,
retry: {
maxRetries: 2,
backoff: {
@@ -428,7 +436,7 @@ describe("ReleaseConcurrencyQueue", () => {
pollInterval: 50,
});
await queue2.attemptToRelease("test-queue", "run2", 1);
await queue2.attemptToRelease("test-queue", "run2");
expect(secondRunAttempted).toBe(true); // Should execute immediately because token was returned
await queue2.quit();
@@ -441,7 +449,7 @@ describe("ReleaseConcurrencyQueue", () => {
const executedRuns: { releaseQueue: string; runId: string; attempt: number }[] = [];
const runAttempts: Record<string, number> = {};
const queue = new ReleaseConcurrencyQueue<string>({
const queue = new ReleaseConcurrencyTokenBucketQueue<string>({
redis: {
keyPrefix: "release-queue:test:",
host: redisContainer.getHost(),
@@ -456,6 +464,7 @@ describe("ReleaseConcurrencyQueue", () => {
fromDescriptor: (descriptor) => descriptor,
toDescriptor: (name) => name,
},
maxTokens: async (_) => 3,
retry: {
maxRetries: 2,
backoff: {
@@ -471,9 +480,9 @@ describe("ReleaseConcurrencyQueue", () => {
try {
// Queue up multiple runs
await Promise.all([
queue.attemptToRelease("test-queue", "run1", 3),
queue.attemptToRelease("test-queue", "run2", 3),
queue.attemptToRelease("test-queue", "run3", 3),
queue.attemptToRelease("test-queue", "run1"),
queue.attemptToRelease("test-queue", "run2"),
queue.attemptToRelease("test-queue", "run3"),
]);
// Wait for all retries to complete
@@ -514,7 +523,7 @@ describe("ReleaseConcurrencyQueue", () => {
const minDelay = 100;
const factor = 2;
const queue = new ReleaseConcurrencyQueue<string>({
const queue = new ReleaseConcurrencyTokenBucketQueue<string>({
redis: {
keyPrefix: "release-queue:test:",
host: redisContainer.getHost(),
@@ -530,6 +539,7 @@ describe("ReleaseConcurrencyQueue", () => {
fromDescriptor: (descriptor) => descriptor,
toDescriptor: (name) => name,
},
maxTokens: async (_) => 1,
retry: {
maxRetries: 2,
backoff: {
@@ -543,7 +553,7 @@ describe("ReleaseConcurrencyQueue", () => {
try {
startTime = Date.now();
await queue.attemptToRelease("test-queue", "run1", 1);
await queue.attemptToRelease("test-queue", "run1");
// Wait for all retries to complete
await setTimeout(1000);
@@ -36,7 +36,7 @@ export type RunEngineOptions = {
queueRunsWaitingForWorkerBatchSize?: number;
tracer: Tracer;
releaseConcurrency?: {
maxTokens?: number;
maxTokensRatio?: number;
redis?: Partial<RedisOptions>;
};
};
@@ -299,35 +299,7 @@ export class RunQueueFullKeyProducer implements RunQueueKeyProducer {
concurrencyKey: parts.at(9),
};
}
releaseConcurrencyKey(env: EnvDescriptor): string;
releaseConcurrencyKey(env: MinimalAuthenticatedEnvironment): string;
releaseConcurrencyKey(envOrDescriptor: EnvDescriptor | MinimalAuthenticatedEnvironment): string {
if ("id" in envOrDescriptor) {
return [
this.orgKeySection(envOrDescriptor.organization.id),
this.projKeySection(envOrDescriptor.project.id),
this.envKeySection(envOrDescriptor.id),
"release-concurrency",
].join(":");
} else {
return [
this.orgKeySection(envOrDescriptor.orgId),
this.projKeySection(envOrDescriptor.projectId),
this.envKeySection(envOrDescriptor.envId),
"release-concurrency",
].join(":");
}
}
releaseConcurrencyDescriptorFromQueue(queue: string): EnvDescriptor {
const parts = queue.split(":");
return {
orgId: parts[1],
projectId: parts[3],
envId: parts[5],
};
}
private envKeySection(envId: string) {
return `${constants.ENV_PART}:${envId}`;
}
@@ -96,11 +96,6 @@ export interface RunQueueKeyProducer {
deadLetterQueueKey(env: MinimalAuthenticatedEnvironment): string;
deadLetterQueueKey(env: EnvDescriptor): string;
deadLetterQueueKeyFromQueue(queue: string): string;
releaseConcurrencyKey(env: MinimalAuthenticatedEnvironment): string;
releaseConcurrencyKey(env: EnvDescriptor): string;
releaseConcurrencyDescriptorFromQueue(queue: string): EnvDescriptor;
}
export type EnvQueues = {