Skip to content

Commit d12e958

Browse files
committed
feat(sdk,core,cli): the concurrency option and named concurrency limits
Tasks declare concurrency directly: an inline { perKey, total } shape caps the task, and concurrencyLimit() declares a named, shareable limit that tasks hold via the same option (at most one inline plus two named). Trigger calls switch a run's named limits with their own concurrency option, strings only like queue. The queue tuple syntax and combinedConcurrencyLimit never ship: queue() is a line again, its concurrencyLimit deprecated in place, and the manifest carries the new declarations for the server to compile.
1 parent 3b04574 commit d12e958

16 files changed

Lines changed: 287 additions & 147 deletions

File tree

.changeset/queue-combined-concurrency-limit.md

Lines changed: 0 additions & 18 deletions
This file was deleted.

.changeset/queue-gates.md

Lines changed: 0 additions & 22 deletions
This file was deleted.

.changeset/task-concurrency.md

Lines changed: 20 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,20 @@
1+
---
2+
"@trigger.dev/sdk": patch
3+
"@trigger.dev/core": patch
4+
---
5+
6+
Control a task's concurrency with the new `concurrency` option, and share limits across tasks with named concurrency limits. An inline shape caps the task itself; `concurrencyLimit()` declares a limit any task can hold (up to two named limits per task), and a trigger call can switch a run's named limits with its own `concurrency` option.
7+
8+
```ts
9+
import { concurrencyLimit, task } from "@trigger.dev/sdk";
10+
11+
export const openaiLimit = concurrencyLimit({ name: "openai", total: 25 });
12+
13+
export const generateSummary = task({
14+
id: "generate-summary",
15+
concurrency: [{ perKey: 1, total: 5 }, openaiLimit],
16+
run: async (payload) => {},
17+
});
18+
```
19+
20+
`perKey` caps each `concurrencyKey` pool and `total` caps across everything, keys or not. The queue-level `concurrencyLimit` option keeps working unchanged and is deprecated in favor of `concurrency`. Enforcement happens server-side; servers without support accept the option but do not enforce it yet.

packages/cli-v3/src/entryPoints/dev-index-worker.ts

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -196,6 +196,7 @@ await sendMessageInCatalog(
196196
prompts: convertPromptSchemasToJsonSchemas(resourceCatalog.listPromptManifests()),
197197
skills: resourceCatalog.listSkillManifests(),
198198
queues: resourceCatalog.listQueueManifests(),
199+
concurrencyLimits: resourceCatalog.listConcurrencyLimitManifests(),
199200
configPath: buildManifest.configPath,
200201
runtime: buildManifest.runtime,
201202
runtimeVersion: detectRuntimeVersion(),

packages/cli-v3/src/entryPoints/managed-index-worker.ts

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -192,6 +192,7 @@ await sendMessageInCatalog(
192192
prompts: convertPromptSchemasToJsonSchemas(resourceCatalog.listPromptManifests()),
193193
skills: resourceCatalog.listSkillManifests(),
194194
queues: resourceCatalog.listQueueManifests(),
195+
concurrencyLimits: resourceCatalog.listConcurrencyLimitManifests(),
195196
configPath: buildManifest.configPath,
196197
runtime: buildManifest.runtime,
197198
runtimeVersion: detectRuntimeVersion(),

packages/core/src/v3/resource-catalog/catalog.ts

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,7 @@ import type {
77
WebhookManifest,
88
WebhookMetadata,
99
WorkerManifest,
10+
ConcurrencyLimitManifest,
1011
} from "../schemas/index.js";
1112
import type {
1213
PromptMetadataWithFunctions,
@@ -27,6 +28,8 @@ export interface ResourceCatalog {
2728
registerWorkerManifest(workerManifest: WorkerManifest): void;
2829
registerQueueMetadata(queue: QueueManifest): void;
2930
listQueueManifests(): Array<QueueManifest>;
31+
registerConcurrencyLimitMetadata(limit: ConcurrencyLimitManifest): void;
32+
listConcurrencyLimitManifests(): Array<ConcurrencyLimitManifest>;
3033
getTaskSchema(id: string): TaskSchema | undefined;
3134
registerPromptMetadata(prompt: PromptMetadataWithFunctions): void;
3235
listPromptManifests(): Array<PromptManifest>;

packages/core/src/v3/resource-catalog/index.ts

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,7 @@ import type {
99
WebhookManifest,
1010
WebhookMetadata,
1111
WorkerManifest,
12+
ConcurrencyLimitManifest,
1213
} from "../schemas/index.js";
1314
import type {
1415
PromptMetadataWithFunctions,
@@ -94,6 +95,14 @@ export class ResourceCatalogAPI {
9495
return this.#getCatalog().listQueueManifests();
9596
}
9697

98+
public registerConcurrencyLimitMetadata(limit: ConcurrencyLimitManifest): void {
99+
this.#getCatalog().registerConcurrencyLimitMetadata(limit);
100+
}
101+
102+
public listConcurrencyLimitManifests(): Array<ConcurrencyLimitManifest> {
103+
return this.#getCatalog().listConcurrencyLimitManifests();
104+
}
105+
97106
public registerPromptMetadata(prompt: PromptMetadataWithFunctions): void {
98107
this.#getCatalog().registerPromptMetadata(prompt);
99108
}

packages/core/src/v3/resource-catalog/noopResourceCatalog.ts

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,7 @@ import type {
77
WebhookManifest,
88
WebhookMetadata,
99
WorkerManifest,
10+
ConcurrencyLimitManifest,
1011
} from "../schemas/index.js";
1112
import {
1213
type PromptMetadataWithFunctions,
@@ -72,6 +73,12 @@ export class NoopResourceCatalog implements ResourceCatalog {
7273
return [];
7374
}
7475

76+
registerConcurrencyLimitMetadata(limit: ConcurrencyLimitManifest): void {}
77+
78+
listConcurrencyLimitManifests(): Array<ConcurrencyLimitManifest> {
79+
return [];
80+
}
81+
7582
registerPromptMetadata(prompt: PromptMetadataWithFunctions): void {
7683
// noop
7784
}

packages/core/src/v3/resource-catalog/standardResourceCatalog.ts

Lines changed: 25 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,7 @@ import type {
1010
WebhookMetadata,
1111
WorkerManifest,
1212
QueueManifest,
13+
ConcurrencyLimitManifest,
1314
} from "../schemas/index.js";
1415
import type {
1516
PromptMetadataWithFunctions,
@@ -41,6 +42,7 @@ export class StandardResourceCatalog implements ResourceCatalog {
4142
private _promptSchemas: Map<string, TaskSchema> = new Map();
4243
private _currentFileContext?: Omit<TaskFileMetadata, "exportName">;
4344
private _queueMetadata: Map<string, QueueManifest> = new Map();
45+
private _concurrencyLimitMetadata: Map<string, ConcurrencyLimitManifest> = new Map();
4446
private _skillMetadata: Map<string, SkillMetadata> = new Map();
4547
private _skillFileMetadata: Map<string, TaskFileMetadata> = new Map();
4648
private _webhookMetadata: Map<string, WebhookMetadata> = new Map();
@@ -99,6 +101,25 @@ export class StandardResourceCatalog implements ResourceCatalog {
99101
this._queueMetadata.set(queue.name, queue);
100102
}
101103

104+
registerConcurrencyLimitMetadata(limit: ConcurrencyLimitManifest): void {
105+
const existing = this._concurrencyLimitMetadata.get(limit.name);
106+
107+
//if it exists already with different settings, log a warning and keep the first definition
108+
if (existing) {
109+
if (existing.perKey !== limit.perKey || existing.total !== limit.total) {
110+
console.warn(
111+
`Concurrency limit "${limit.name}" is defined twice, with different settings.` +
112+
`\n - perKey: ${existing.perKey} vs ${limit.perKey}` +
113+
`\n - total: ${existing.total} vs ${limit.total}` +
114+
`\n Keeping the first definition.`
115+
);
116+
return;
117+
}
118+
}
119+
120+
this._concurrencyLimitMetadata.set(limit.name, limit);
121+
}
122+
102123
registerWorkerManifest(workerManifest: WorkerManifest): void {
103124
for (const task of workerManifest.tasks) {
104125
this._taskFileMetadata.set(task.id, {
@@ -222,6 +243,10 @@ export class StandardResourceCatalog implements ResourceCatalog {
222243
return Array.from(this._queueMetadata.values());
223244
}
224245

246+
listConcurrencyLimitManifests(): Array<ConcurrencyLimitManifest> {
247+
return Array.from(this._concurrencyLimitMetadata.values());
248+
}
249+
225250
getTaskManifest(id: string): TaskManifest | undefined {
226251
const metadata = this._taskMetadata.get(id);
227252
const fileMetadata = this._taskFileMetadata.get(id);

packages/core/src/v3/schemas/api.ts

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -336,6 +336,7 @@ export const TriggerTaskRequestBody = z
336336
)
337337
.max(2)
338338
.optional(),
339+
concurrency: z.string().min(1).max(128).array().max(2).optional(),
339340
concurrencyKey: ConcurrencyKeySchema.optional(),
340341
delay: z.string().or(z.coerce.date()).optional(),
341342
idempotencyKey: z
@@ -451,6 +452,7 @@ export const BatchTriggerTaskItem = z.object({
451452
)
452453
.max(2)
453454
.optional(),
455+
concurrency: z.string().min(1).max(128).array().max(2).optional(),
454456
tags: RunTags.optional(),
455457
test: z.boolean().optional(),
456458
ttl: z.string().or(z.number().nonnegative().int()).optional(),

0 commit comments

Comments
 (0)