summaryrefslogtreecommitdiffhomepage
path: root/packages/provider-concurrency/src/provider-wrapper.ts
diff options
context:
space:
mode:
authorAdam Malczewski <[email protected]>2026-06-27 18:49:50 +0900
committerAdam Malczewski <[email protected]>2026-06-27 18:49:50 +0900
commitd19f75caca6f6ad49e95b83021cb4cf00a39297f (patch)
treebc106b082cc492095d502be9702a65365acbbd2a /packages/provider-concurrency/src/provider-wrapper.ts
parent87e85e026e54b1dc25b0648af298ab0a8a715701 (diff)
downloaddispatch-d19f75caca6f6ad49e95b83021cb4cf00a39297f.tar.gz
dispatch-d19f75caca6f6ad49e95b83021cb4cf00a39297f.zip
feat(concurrency): add "queued" ConversationStatus — emit when request blocks on acquire, re-emit "active" when slot granted
Diffstat (limited to 'packages/provider-concurrency/src/provider-wrapper.ts')
-rw-r--r--packages/provider-concurrency/src/provider-wrapper.ts11
1 files changed, 10 insertions, 1 deletions
diff --git a/packages/provider-concurrency/src/provider-wrapper.ts b/packages/provider-concurrency/src/provider-wrapper.ts
index ee3ca85..aa08e5b 100644
--- a/packages/provider-concurrency/src/provider-wrapper.ts
+++ b/packages/provider-concurrency/src/provider-wrapper.ts
@@ -25,12 +25,20 @@ import type { ConcurrencyLimiter } from "./concurrency-manager.js";
* @param conversationId The agent requesting the stream (for slot attribution).
* @param promptStartedAt When the agent's current prompt (turn) started
* (epoch-ms, for oldest-agent-first scheduling).
+ * @param onQueued Called synchronously when `acquire()` decides to
+ * queue the request (cannot grant immediately).
+ * Lets the caller emit a "queued" status signal.
+ * @param onAcquired Called when `acquire()` resolves (slot granted,
+ * whether immediately or after queueing). Lets the
+ * caller emit an "active" status signal.
*/
export function wrapProviderWithConcurrency(
provider: ProviderContract,
limiter: ConcurrencyLimiter,
conversationId: string,
promptStartedAt: number,
+ onQueued?: () => void,
+ onAcquired?: () => void,
): ProviderContract {
const innerStream = provider.stream;
const providerId = provider.id;
@@ -42,7 +50,8 @@ export function wrapProviderWithConcurrency(
tools: readonly ToolContract[],
opts?: ProviderStreamOptions,
): AsyncIterable<ProviderEvent> {
- const release = await limiter.acquire(providerId, conversationId, promptStartedAt);
+ const release = await limiter.acquire(providerId, conversationId, promptStartedAt, onQueued);
+ onAcquired?.();
try {
for await (const event of innerStream(messages, tools, opts)) {
if (event.type === "error" && event.code === "429") {