Skip to content

Commit e5d23bf

Browse files
committed
[refactor] Simplify queue creation in AiProvider
- Replaced the manual queue setup process in the createQueue method with a call to createDefaultQueue for improved readability and maintainability. - This change reduces complexity by leveraging a dedicated function for queue creation, streamlining the overall implementation.
1 parent 8f45111 commit e5d23bf

2 files changed

Lines changed: 47 additions & 26 deletions

File tree

packages/ai/src/provider/AiProvider.ts

Lines changed: 2 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,7 @@
77
import { TaskInput, TaskOutput } from "@workglow/task-graph";
88
import { globalServiceRegistry, WORKER_MANAGER, type WorkerServer } from "@workglow/util";
99
import type { ModelConfig } from "../model/ModelSchema";
10+
import { createDefaultQueue } from "../queue/createDefaultQueue";
1011
import { type AiProviderRunFn, getAiProviderRegistry } from "./AiProviderRegistry";
1112

1213
/**
@@ -183,31 +184,6 @@ export abstract class AiProvider<TModelConfig extends ModelConfig = ModelConfig>
183184
* Uses InMemoryQueueStorage with a ConcurrencyLimiter.
184185
*/
185186
protected async createQueue(concurrency: number): Promise<void> {
186-
// Lazy imports to avoid circular dependencies -- these packages
187-
// are always available at runtime in any environment that uses AiProvider.
188-
const { InMemoryQueueStorage } = await import("@workglow/storage");
189-
const { ConcurrencyLimiter, JobQueueClient, JobQueueServer } =
190-
await import("@workglow/job-queue");
191-
const { getTaskQueueRegistry } = await import("@workglow/task-graph");
192-
const { AiJob } = await import("../job/AiJob");
193-
194-
const storage = new InMemoryQueueStorage(this.name);
195-
await storage.setupDatabase();
196-
197-
const server = new JobQueueServer(AiJob as any, {
198-
storage,
199-
queueName: this.name,
200-
limiter: new ConcurrencyLimiter(concurrency, 100),
201-
});
202-
203-
const client = new JobQueueClient({
204-
storage,
205-
queueName: this.name,
206-
});
207-
208-
client.attach(server);
209-
210-
getTaskQueueRegistry().registerQueue({ server, client, storage });
211-
await server.start();
187+
await createDefaultQueue(this.name, concurrency);
212188
}
213189
}
Lines changed: 45 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,45 @@
1+
/**
2+
* @license
3+
* Copyright 2025 Steven Roussey <sroussey@gmail.com>
4+
* SPDX-License-Identifier: Apache-2.0
5+
*/
6+
7+
import { ConcurrencyLimiter, JobQueueClient, JobQueueServer } from "@workglow/job-queue";
8+
import { getTaskQueueRegistry } from "@workglow/task-graph";
9+
import { InMemoryQueueStorage } from "@workglow/storage";
10+
11+
import { AiJob } from "../job/AiJob";
12+
13+
/**
14+
* Create and register a default job queue for an AI provider.
15+
* Uses InMemoryQueueStorage with a ConcurrencyLimiter.
16+
*
17+
* Extracted to a separate module to avoid circular dependencies between
18+
* AiProvider, AiJob, and the storage/job-queue/task-graph packages.
19+
*
20+
* @param providerName - Unique provider identifier (used as queue name)
21+
* @param concurrency - Maximum number of concurrent jobs
22+
*/
23+
export async function createDefaultQueue(
24+
providerName: string,
25+
concurrency: number
26+
): Promise<void> {
27+
const storage = new InMemoryQueueStorage(providerName);
28+
await storage.setupDatabase();
29+
30+
const server = new JobQueueServer(AiJob as any, {
31+
storage,
32+
queueName: providerName,
33+
limiter: new ConcurrencyLimiter(concurrency, 100),
34+
});
35+
36+
const client = new JobQueueClient({
37+
storage,
38+
queueName: providerName,
39+
});
40+
41+
client.attach(server);
42+
43+
getTaskQueueRegistry().registerQueue({ server, client, storage });
44+
await server.start();
45+
}

0 commit comments

Comments
 (0)