diff options
| author | Adam Malczewski <[email protected]> | 2026-06-27 18:49:50 +0900 |
|---|---|---|
| committer | Adam Malczewski <[email protected]> | 2026-06-27 18:49:50 +0900 |
| commit | d19f75caca6f6ad49e95b83021cb4cf00a39297f (patch) | |
| tree | bc106b082cc492095d502be9702a65365acbbd2a /packages/provider-concurrency/src/provider-wrapper.ts | |
| parent | 87e85e026e54b1dc25b0648af298ab0a8a715701 (diff) | |
| download | dispatch-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.ts | 11 |
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") { |
