From b7ddf20d62888eddff6189aaea3455ba9cd4451d Mon Sep 17 00:00:00 2001 From: Eric Allam Date: Thu, 13 Mar 2025 22:45:58 +0000 Subject: [PATCH] Correctly use the new release concurrency queue in the run engine --- .../run-engine/src/engine/index.ts | 65 ++++++--- ... => releaseConcurrencyTokenBucketQueue.ts} | 13 +- ...eleaseConcurrencyTokenBucketQueue.test.ts} | 136 ++++++++++-------- .../run-engine/src/engine/types.ts | 2 +- .../run-engine/src/run-queue/keyProducer.ts | 28 ---- .../run-engine/src/run-queue/types.ts | 5 - 6 files changed, 131 insertions(+), 118 deletions(-) rename internal-packages/run-engine/src/engine/{releaseConcurrencyQueue.ts => releaseConcurrencyTokenBucketQueue.ts} (96%) rename internal-packages/run-engine/src/engine/tests/{releaseConcurrencyQueue.test.ts => releaseConcurrencyTokenBucketQueue.test.ts} (85%) diff --git a/internal-packages/run-engine/src/engine/index.ts b/internal-packages/run-engine/src/engine/index.ts index 4bab97fa0..6b24768a7 100644 --- a/internal-packages/run-engine/src/engine/index.ts +++ b/internal-packages/run-engine/src/engine/index.ts @@ -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(); 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 ); diff --git a/internal-packages/run-engine/src/engine/releaseConcurrencyQueue.ts b/internal-packages/run-engine/src/engine/releaseConcurrencyTokenBucketQueue.ts similarity index 96% rename from internal-packages/run-engine/src/engine/releaseConcurrencyQueue.ts rename to internal-packages/run-engine/src/engine/releaseConcurrencyTokenBucketQueue.ts index 378ced1c5..c4cb99e17 100644 --- a/internal-packages/run-engine/src/engine/releaseConcurrencyQueue.ts +++ b/internal-packages/run-engine/src/engine/releaseConcurrencyTokenBucketQueue.ts @@ -20,6 +20,7 @@ export type ReleaseConcurrencyQueueOptions = { fromDescriptor: (releaseQueue: T) => string; toDescriptor: (releaseQueue: string) => T; }; + maxTokens: (descriptor: T) => Promise; consumersCount?: number; masterQueuesKey?: string; tracer?: Tracer; @@ -36,7 +37,7 @@ const QueueItemMetadata = z.object({ type QueueItemMetadata = z.infer; -export class ReleaseConcurrencyQueue { +export class ReleaseConcurrencyTokenBucketQueue { private redis: Redis; private logger: Logger; private abortController: AbortController; @@ -47,6 +48,7 @@ export class ReleaseConcurrencyQueue { private consumersCount: number; private pollInterval: number; private keys: ReleaseConcurrencyQueueOptions["keys"]; + private maxTokens: ReleaseConcurrencyQueueOptions["maxTokens"]; private batchSize: number; private maxRetries: number; private backoff: NonNullable>; @@ -62,6 +64,7 @@ export class ReleaseConcurrencyQueue { 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 { * 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 { * * 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 { private logger: Logger; constructor( - private readonly queue: ReleaseConcurrencyQueue, + private readonly queue: ReleaseConcurrencyTokenBucketQueue, private readonly pollInterval: number, private readonly signal: AbortSignal, logger?: Logger diff --git a/internal-packages/run-engine/src/engine/tests/releaseConcurrencyQueue.test.ts b/internal-packages/run-engine/src/engine/tests/releaseConcurrencyTokenBucketQueue.test.ts similarity index 85% rename from internal-packages/run-engine/src/engine/tests/releaseConcurrencyQueue.test.ts rename to internal-packages/run-engine/src/engine/tests/releaseConcurrencyTokenBucketQueue.test.ts index 91c7f069e..9cce28ec1 100644 --- a/internal-packages/run-engine/src/engine/tests/releaseConcurrencyQueue.test.ts +++ b/internal-packages/run-engine/src/engine/tests/releaseConcurrencyTokenBucketQueue.test.ts @@ -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({ + const queue = new ReleaseConcurrencyTokenBucketQueue({ 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({ + const queue = new ReleaseConcurrencyTokenBucketQueue({ 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({ + const queue = new ReleaseConcurrencyTokenBucketQueue({ 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({ + const queue = new ReleaseConcurrencyTokenBucketQueue({ 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({ + const queue2 = new ReleaseConcurrencyTokenBucketQueue({ 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 = {}; - const queue = new ReleaseConcurrencyQueue({ + const queue = new ReleaseConcurrencyTokenBucketQueue({ 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({ + const queue = new ReleaseConcurrencyTokenBucketQueue({ 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); diff --git a/internal-packages/run-engine/src/engine/types.ts b/internal-packages/run-engine/src/engine/types.ts index c5f4fc478..9be634151 100644 --- a/internal-packages/run-engine/src/engine/types.ts +++ b/internal-packages/run-engine/src/engine/types.ts @@ -36,7 +36,7 @@ export type RunEngineOptions = { queueRunsWaitingForWorkerBatchSize?: number; tracer: Tracer; releaseConcurrency?: { - maxTokens?: number; + maxTokensRatio?: number; redis?: Partial; }; }; diff --git a/internal-packages/run-engine/src/run-queue/keyProducer.ts b/internal-packages/run-engine/src/run-queue/keyProducer.ts index 033da935b..a24bf8dbd 100644 --- a/internal-packages/run-engine/src/run-queue/keyProducer.ts +++ b/internal-packages/run-engine/src/run-queue/keyProducer.ts @@ -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}`; } diff --git a/internal-packages/run-engine/src/run-queue/types.ts b/internal-packages/run-engine/src/run-queue/types.ts index cbf59cced..563cececa 100644 --- a/internal-packages/run-engine/src/run-queue/types.ts +++ b/internal-packages/run-engine/src/run-queue/types.ts @@ -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 = {