summaryrefslogtreecommitdiffhomepage
path: root/packages/provider-concurrency/src/provider-wrapper.test.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.test.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.test.ts')
-rw-r--r--packages/provider-concurrency/src/provider-wrapper.test.ts73
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:
| {