Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
2.3 kB · 66 lines
TypeScript
12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667import { describe, expect, test } from "vitest";import { ConsumerScheduler } from "../src/agents/scheduler.js";
describe("consumer scheduler", () => { test("runs independent keys concurrently while preserving order within each key", async () => { const scheduler = new ConsumerScheduler({ concurrency: 2 }); const events: string[] = []; let releaseSlow: (() => void) | undefined; const slowGate = new Promise<void>((resolve) => { releaseSlow = resolve; });
scheduler.enqueue("slow", async () => { events.push("slow-one:start"); await slowGate; events.push("slow-one:end"); }); scheduler.enqueue("slow", async () => { events.push("slow-two:start"); events.push("slow-two:end"); }); scheduler.enqueue("fast", async () => { events.push("fast:start"); events.push("fast:end"); });
await waitUntil(() => events.includes("fast:end")); expect(events).toContain("slow-one:start"); expect(events).not.toContain("slow-two:start");
releaseSlow?.(); await scheduler.drain(); expect(events.indexOf("slow-one:end")).toBeLessThan(events.indexOf("slow-two:start")); });
test("honors the configured global concurrency bound", async () => { const scheduler = new ConsumerScheduler({ concurrency: 1 }); const events: string[] = []; let releaseFirst: (() => void) | undefined; const firstGate = new Promise<void>((resolve) => { releaseFirst = resolve; });
scheduler.enqueue("first", async () => { events.push("first:start"); await firstGate; events.push("first:end"); }); scheduler.enqueue("second", async () => { events.push("second:start"); });
await waitUntil(() => events.includes("first:start")); await new Promise((resolve) => setTimeout(resolve, 10)); expect(events).not.toContain("second:start");
releaseFirst?.(); await scheduler.drain(); expect(events).toEqual(["first:start", "first:end", "second:start"]); });});
async function waitUntil(predicate: () => boolean, timeoutMs = 1_000): Promise<void> { const deadline = Date.now() + timeoutMs; while (Date.now() < deadline) { if (predicate()) return; await new Promise((resolve) => setTimeout(resolve, 5)); } throw new Error("Timed out waiting for scheduler state");}