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.test.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.test.ts')
| -rw-r--r-- | packages/provider-concurrency/src/provider-wrapper.test.ts | 73 |
1 files changed, 73 insertions, 0 deletions
diff --git a/packages/provider-concurrency/src/provider-wrapper.test.ts b/packages/provider-concurrency/src/provider-wrapper.test.ts index e024d59..e59ab39 100644 --- a/packages/provider-concurrency/src/provider-wrapper.test.ts +++ b/packages/provider-concurrency/src/provider-wrapper.test.ts @@ -139,6 +139,79 @@ describe("wrapProviderWithConcurrency", () => { expect(models).toEqual([{ id: "model-1" }]); }); + it("calls onQueued when the request blocks and onAcquired when the slot is granted", async () => { + let queuedCalled = false; + let acquiredCalled = false; + + const blockingLimiter: ConcurrencyLimiter = { + acquire(_providerId, _convId, _promptAt, onQueued) { + // Simulate a queued request: call onQueued, then resolve on next tick. + onQueued?.(); + return new Promise((resolve) => { + setTimeout(() => { + resolve(() => {}); + }, 0); + }); + }, + reportRateLimit() {}, + }; + + const provider = fakeProvider([{ type: "finish", reason: "stop" }]); + const wrapped = wrapProviderWithConcurrency( + provider, + blockingLimiter, + "conv1", + 0, + () => { + queuedCalled = true; + }, + () => { + acquiredCalled = true; + }, + ); + + for await (const _e of wrapped.stream([], [])) { + // consume + } + + expect(queuedCalled).toBe(true); + expect(acquiredCalled).toBe(true); + }); + + it("does NOT call onQueued when the slot is granted immediately", async () => { + let queuedCalled = false; + let acquiredCalled = false; + + const immediateLimiter: ConcurrencyLimiter = { + acquire(_providerId, _convId, _promptAt, _onQueued) { + // Grant immediately — do NOT call onQueued. + return Promise.resolve(() => {}); + }, + reportRateLimit() {}, + }; + + const provider = fakeProvider([{ type: "finish", reason: "stop" }]); + const wrapped = wrapProviderWithConcurrency( + provider, + immediateLimiter, + "conv1", + 0, + () => { + queuedCalled = true; + }, + () => { + acquiredCalled = true; + }, + ); + + for await (const _e of wrapped.stream([], [])) { + // consume + } + + expect(queuedCalled).toBe(false); + expect(acquiredCalled).toBe(true); + }); + it("passes through messages, tools, and opts to the inner stream", async () => { let receivedArgs: | { |
