import test from "node:test"; import assert from "node:assert/strict"; import { mkdtemp, readFile, rm, writeFile } from "node:fs/promises"; import { tmpdir } from "node:os"; import path from "node:path"; import { FileRepository } from "../src/infra/repository.ts"; import { MemoryTaskQueue } from "../src/infra/queue.ts"; import { ensureStagingDir, gcStaging, putStagingObject, safeStagingPath } from "../src/infra/staging.ts"; import { appendEvent, createStore, refundBalance, releaseBalance, reserveBalance, settleBalance } from "../src/store.ts"; import { recoverExpiredReservations, recoverExpiredTaskLeases } from "../src/jobs/task-worker.ts"; import { enqueueMessage, processMessageOutbox } from "../src/adapters/messaging.ts"; import { enqueueExpiredWebDavRetention } from "../src/adapters/webdav.ts"; test("file repository writes an atomic, private snapshot", async () => { const root = await mkdtemp(path.join(tmpdir(), "miragenflow-repo-")); try { const repository = new FileRepository<{ marker: string }>(path.join(root, "nested", "snapshot.json")); await repository.save({ marker: "persisted" }); assert.deepEqual(await repository.load(), { marker: "persisted" }); assert.equal((await repository.health()).status, "ready"); assert.equal(JSON.parse(await readFile(path.join(root, "nested", "snapshot.json"), "utf8")).marker, "persisted"); } finally { await rm(root, { recursive: true, force: true }); } }); test("file repository tolerates concurrent atomic saves", async () => { const root = await mkdtemp(path.join(tmpdir(), "miragenflow-repo-race-")); try { const repository = new FileRepository<{ marker: string }>(path.join(root, "snapshot.json")); await Promise.all(Array.from({ length: 8 }, (_, index) => repository.save({ marker: String(index) }))); assert.match((await repository.load())?.marker || "", /^[0-7]$/); } finally { await rm(root, { recursive: true, force: true }); } }); test("postgres reload waits for an in-flight snapshot write", async () => { const store = createStore(); let persisted: unknown = { users: [] }; let loadCalled = false; let releaseSave!: () => void; let resolveSaveStarted!: () => void; const saveStarted = new Promise((resolve) => { resolveSaveStarted = resolve; }); const saveRelease = new Promise((resolve) => { releaseSave = resolve; }); let firstSave = true; store.repository = { adapter: "postgres", async load() { loadCalled = true; return persisted; }, async save(value: unknown) { if (firstSave) { firstSave = false; resolveSaveStarted(); await saveRelease; } persisted = value; }, getRevision() { return 0; }, async health() { return { adapter: "postgres", status: "ready" }; }, }; store.users.set("reload-race-user", { id: "reload-race-user", version: 1, email: "reload@example.com", password: "hash", verified: true, status: "active", mfaRequired: false, mfaEnabled: false, roles: ["user"], failedLoginCount: 0, createdAt: new Date().toISOString() }); store.persist(); await saveStarted; const reload = store.loadPersisted(); await new Promise((resolve) => setTimeout(resolve, 0)); assert.equal(loadCalled, false); assert.equal(store.users.has("reload-race-user"), true); releaseSave(); await reload; assert.equal(loadCalled, true); assert.equal(store.users.has("reload-race-user"), true); await store.persistAsync(); const users = (persisted as { users?: Array<[string, unknown]> }).users || []; assert.equal(users.some(([id]) => id === "reload-race-user"), true); }); test("postgres reload retries when a write starts during the database read", async () => { const store = createStore(); let persisted: unknown = { users: [] }; let firstLoad = true; let releaseLoad!: () => void; let resolveLoadStarted!: () => void; const loadStarted = new Promise((resolve) => { resolveLoadStarted = resolve; }); const loadRelease = new Promise((resolve) => { releaseLoad = resolve; }); store.repository = { adapter: "postgres", async load() { if (firstLoad) { firstLoad = false; resolveLoadStarted(); await loadRelease; } return persisted; }, async save(value: unknown) { persisted = value; }, getRevision() { return 0; }, async health() { return { adapter: "postgres", status: "ready" }; }, }; const reload = store.loadPersisted(); await loadStarted; store.users.set("reload-during-read-user", { id: "reload-during-read-user", version: 1, email: "reload-read@example.com", password: "hash", verified: true, status: "active", mfaRequired: false, mfaEnabled: false, roles: ["user"], failedLoginCount: 0, createdAt: new Date().toISOString() }); store.persist(); releaseLoad(); await reload; assert.equal(store.users.has("reload-during-read-user"), true); }); test("store transaction rolls back in-memory mutations and defers events until commit", async () => { const store = createStore(); const taskId = "transaction-task"; store.tasks.set(taskId, { id: taskId, ownerId: "transaction-user", taskType: "image", modelProductId: "basic-image-v1", status: "queued", estimatedCost: 1, reservedCost: 0, createdAt: new Date().toISOString(), updatedAt: new Date().toISOString(), eventSequence: 0, routeSnapshot: { groupId: "image-default", groupVersion: 1, publicModelId: "basic-image-v1", routeVersion: 1, orderedCandidates: [], strategy: "strict", totalRetryBudget: 0 }, modelSnapshot: { publicModelId: "basic-image-v1", modelVersion: 1, name: "基础模型", tier: "basic", capabilities: ["image"] }, pricingSnapshot: { ruleId: "price-basic-image", ruleVersion: 1, publicModelId: "basic-image-v1", unitPrice: 1, quantity: 1, billingUnit: "output", multiplier: 1, balanceUnitVersion: 1 }, planSnapshot: { queuePriority: 0, maxConcurrent: 1 }, attempts: [], outputs: [] }); let observed = 0; store.taskSubscribers.set(taskId, new Set([() => { observed += 1; }])); await assert.rejects(store.transact(() => { const task = store.tasks.get(taskId)!; task.status = "running"; appendEvent(store, taskId, { taskId, type: "task.running", payload: {} }); throw new Error("rollback"); }), /rollback/); assert.equal(store.tasks.get(taskId)?.status, "queued"); assert.equal(store.events.get(taskId)?.length || 0, 0); assert.equal(observed, 0); await store.transact(() => { const task = store.tasks.get(taskId)!; task.status = "running"; appendEvent(store, taskId, { taskId, type: "task.running", payload: {} }); }); assert.equal(observed, 1); }); test("memory queue leases a task once and allows ack", async () => { const queue = new MemoryTaskQueue(); await queue.enqueue("task-1"); const claim = await queue.claim("worker-1", 1_000); assert.equal(claim?.taskId, "task-1"); assert.equal((await queue.claim("worker-2", 1_000)), undefined); await queue.ack("task-1", claim!.leaseToken); assert.equal((await queue.health()).status, "ready"); }); test("staging objects stay below the private root and expired files are collected", async () => { const root = await mkdtemp(path.join(tmpdir(), "miragenflow-staging-")); try { assert.equal((await ensureStagingDir(root)).status, "ready"); assert.throws(() => safeStagingPath(root, "../escape"), /escapes root|invalid/); const object = await putStagingObject(root, "owner-1", new Uint8Array([1, 2, 3])); assert.equal(object.path.startsWith(path.resolve(root)), true); const protectedCount = await gcStaging(root, 0, Date.now() + 1_000, (key) => key === object.key); assert.equal(protectedCount, 0); const removed = await gcStaging(root, 0, Date.now() + 1_000); assert.equal(removed, 1); } finally { await rm(root, { recursive: true, force: true }); } }); test("expired worker lease moves a running task to unknown without releasing reserve", () => { const store = createStore(); const now = new Date(Date.now() - 10_000).toISOString(); store.tasks.set("lease-task", { id: "lease-task", ownerId: "lease-user", taskType: "image", modelProductId: "basic-image", status: "running", estimatedCost: 10, reservedCost: 10, createdAt: now, updatedAt: now, eventSequence: 0, routeSnapshot: { groupId: "image-default", groupVersion: 1, publicModelId: "basic-image-v1", routeVersion: 1, orderedCandidates: [], strategy: "strict", totalRetryBudget: 0 }, modelSnapshot: { publicModelId: "basic-image-v1", modelVersion: 1, name: "基础模型", tier: "basic", capabilities: ["image"] }, pricingSnapshot: { ruleId: "price-basic-image", ruleVersion: 1, publicModelId: "basic-image-v1", unitPrice: 1, quantity: 1, billingUnit: "output", multiplier: 1, balanceUnitVersion: 1 }, planSnapshot: { queuePriority: 0, maxConcurrent: 1 }, leaseExpiresAt: now, attempts: [{ id: "attempt-1", channelId: "channel-a", sequence: 1, status: "started", startedAt: now, leaseExpiresAt: now }], outputs: [] }); store.balances.set("lease-user", { available: 0, reserved: 10 }); assert.equal(recoverExpiredTaskLeases(store), 1); const task = store.tasks.get("lease-task")!; assert.equal(task.status, "unknown"); assert.equal(task.reservedCost, 10); assert.equal(task.attempts[0].reconciliationStatus, "pending"); assert.equal(store.balances.get("lease-user")?.reserved, 10); }); test("bucket allocations release only the requested amount and settled refunds create a new bucket", () => { const store = createStore(); const userId = "bucket-user"; store.balances.set(userId, { available: 20, reserved: 0 }); store.buckets.set(userId, [ { id: "bucket-a", userId, source: "recharge", remaining: 10, priority: 0 }, { id: "bucket-b", userId, source: "plan", remaining: 10, priority: 1 }, ]); assert.equal(reserveBalance(store, userId, 10, "task-bucket"), true); assert.deepEqual(store.buckets.get(userId)?.map((bucket) => bucket.remaining), [0, 10]); assert.equal(releaseBalance(store, userId, 4, "task-bucket"), true); assert.equal(store.balances.get(userId)?.available, 14); assert.equal(store.balances.get(userId)?.reserved, 6); assert.deepEqual(store.buckets.get(userId)?.map((bucket) => bucket.remaining), [4, 10]); assert.equal(settleBalance(store, userId, 6, 6, "task-bucket"), true); assert.equal(store.balances.get(userId)?.reserved, 0); assert.equal(refundBalance(store, userId, 3, "task-bucket-refund"), true); assert.equal(store.balances.get(userId)?.available, 17); assert.equal(store.ledger.filter((entry) => entry.userId === userId).length, 4); }); test("partial settlement records the unused reserve as an explicit release", () => { const store = createStore(); const userId = "partial-settle-user"; store.balances.set(userId, { available: 10, reserved: 0 }); store.buckets.set(userId, [{ id: "partial-bucket", userId, source: "recharge", remaining: 10, priority: 0 }]); assert.equal(reserveBalance(store, userId, 10, "partial-task"), true); assert.equal(settleBalance(store, userId, 10, 6, "partial-task"), true); assert.equal(store.ledger.filter((entry) => entry.referenceId === "partial-task" && entry.type === "release").length, 1); assert.equal(store.ledger.find((entry) => entry.referenceId === "partial-task" && entry.type === "release")?.amount, 4); assert.equal(store.ledger.find((entry) => entry.referenceId === "partial-task" && entry.type === "release")?.bucketId, "partial-bucket"); assert.equal(store.ledger.find((entry) => entry.referenceId === "partial-task" && entry.type === "settle")?.bucketId, "partial-bucket"); assert.equal(store.balances.get(userId)?.available, 4); }); test("legacy aggregate gaps materialize as UUID buckets even with a preferred bucket", () => { const store = createStore(); const userId = "legacy-gap-user"; store.balances.set(userId, { available: 20, reserved: 0 }); store.buckets.set(userId, [{ id: "existing-bucket", userId, source: "recharge", remaining: 5, priority: 0 }]); assert.equal(reserveBalance(store, userId, 10, "legacy-gap-task", { preferredBucketIds: new Set(["existing-bucket"]) }), true); const buckets = store.buckets.get(userId) || []; assert.equal(buckets.length, 2); assert.match(buckets[1].id, /^[0-9a-f]{8}-[0-9a-f]{4}-4[0-9a-f]{3}-[89ab][0-9a-f]{3}-[0-9a-f]{12}$/i); assert.equal(buckets.reduce((total, bucket) => total + bucket.remaining, 0), 10); assert.deepEqual({ available: store.balances.get(userId)?.available, reserved: store.balances.get(userId)?.reserved }, { available: 10, reserved: 10 }); assert.equal(store.balances.get(userId)?.version, 2); assert.equal((store.reserveAllocations.get(`${userId}:legacy-gap-task`) || []).reduce((total, allocation) => total + allocation.amount, 0), 10); }); test("cross-bucket settlement keeps per-bucket release and charge attribution", () => { const store = createStore(); const userId = "cross-bucket-settle-user"; store.balances.set(userId, { available: 20, reserved: 0 }); store.buckets.set(userId, [ { id: "cross-bucket-a", userId, source: "recharge", remaining: 6, priority: 0 }, { id: "cross-bucket-b", userId, source: "plan", remaining: 14, priority: 1 }, ]); assert.equal(reserveBalance(store, userId, 15, "cross-bucket-task"), true); assert.equal(settleBalance(store, userId, 15, 8, "cross-bucket-task"), true); const entries = store.ledger.filter((entry) => entry.referenceId === "cross-bucket-task"); assert.equal(entries.filter((entry) => entry.type === "settle").reduce((total, entry) => total - entry.amount, 0), 8); assert.equal(entries.filter((entry) => entry.type === "release").reduce((total, entry) => total + entry.amount, 0), 7); assert.ok(entries.filter((entry) => entry.type === "settle" || entry.type === "release").every((entry) => Boolean(entry.bucketId))); assert.deepEqual({ available: store.balances.get(userId)?.available, reserved: store.balances.get(userId)?.reserved }, { available: 12, reserved: 0 }); assert.equal((store.buckets.get(userId) || []).reduce((total, bucket) => total + bucket.remaining, 0), 12); assert.equal(store.reserveAllocations.has(`${userId}:cross-bucket-task`), false); }); test("release and settlement validation failures leave all balance state untouched", () => { const releaseStore = createStore(); releaseStore.balances.set("missing-release-user", { available: 0, reserved: 5 }); releaseStore.reserveAllocations.set("missing-release-user:missing-release-task", [{ bucketId: "missing-bucket", amount: 5 }]); const releaseBefore = JSON.stringify({ balances: [...releaseStore.balances], buckets: [...releaseStore.buckets], allocations: [...releaseStore.reserveAllocations], ledger: releaseStore.ledger }); assert.equal(releaseBalance(releaseStore, "missing-release-user", 5, "missing-release-task"), false); assert.equal(JSON.stringify({ balances: [...releaseStore.balances], buckets: [...releaseStore.buckets], allocations: [...releaseStore.reserveAllocations], ledger: releaseStore.ledger }), releaseBefore); const settleStore = createStore(); settleStore.balances.set("missing-settle-user", { available: 0, reserved: 5 }); settleStore.reserveAllocations.set("missing-settle-user:missing-settle-task", [{ bucketId: "missing-bucket", amount: 5 }]); const settleBefore = JSON.stringify({ balances: [...settleStore.balances], buckets: [...settleStore.buckets], allocations: [...settleStore.reserveAllocations], ledger: settleStore.ledger }); assert.equal(settleBalance(settleStore, "missing-settle-user", 5, 3, "missing-settle-task"), false); assert.equal(JSON.stringify({ balances: [...settleStore.balances], buckets: [...settleStore.buckets], allocations: [...settleStore.reserveAllocations], ledger: settleStore.ledger }), settleBefore); }); test("zero-charge settlement writes an idempotent marker without losing the bucket", () => { const store = createStore(); const userId = "zero-charge-user"; store.balances.set(userId, { available: 5, reserved: 0 }); store.buckets.set(userId, [{ id: "zero-charge-bucket", userId, source: "recharge", remaining: 5, priority: 0 }]); assert.equal(reserveBalance(store, userId, 5, "zero-charge-task"), true); assert.equal(settleBalance(store, userId, 5, 0, "zero-charge-task"), true); assert.equal(store.ledger.filter((entry) => entry.idempotencyKey === `settle:${userId}:zero-charge-task` && entry.amount === 0).length, 1); assert.deepEqual({ available: store.balances.get(userId)?.available, reserved: store.balances.get(userId)?.reserved }, { available: 5, reserved: 0 }); assert.equal(store.buckets.get(userId)?.[0].remaining, 5); }); test("reserved plan balance that expires before release is not restored", () => { const store = createStore(); const userId = "expired-reserved-release-user"; const referenceId = "expired-reserved-release-task"; store.balances.set(userId, { available: 10, reserved: 0 }); store.buckets.set(userId, [{ id: "expired-reserved-release-bucket", userId, source: "plan", remaining: 10, priority: 0, expiresAt: new Date(Date.now() + 60_000).toISOString() }]); assert.equal(reserveBalance(store, userId, 10, referenceId), true); store.buckets.get(userId)![0].expiresAt = new Date(Date.now() - 1_000).toISOString(); assert.equal(releaseBalance(store, userId, 4, referenceId), true); assert.deepEqual({ available: store.balances.get(userId)?.available, reserved: store.balances.get(userId)?.reserved }, { available: 0, reserved: 6 }); assert.equal(store.buckets.get(userId)?.[0].remaining, 0); assert.deepEqual(store.reserveAllocations.get(`${userId}:${referenceId}`), [{ bucketId: "expired-reserved-release-bucket", amount: 6 }]); assert.equal(store.ledger.find((entry) => entry.idempotencyKey === `release:${userId}:${referenceId}:4:expire:expired-reserved-release-bucket`)?.amount, -4); assert.equal(settleBalance(store, userId, 6, 6, referenceId), true); assert.deepEqual({ available: store.balances.get(userId)?.available, reserved: store.balances.get(userId)?.reserved }, { available: 0, reserved: 0 }); }); test("unused part of an expired reservation is expired during settlement", () => { const store = createStore(); const userId = "expired-reserved-settle-user"; const referenceId = "expired-reserved-settle-task"; store.balances.set(userId, { available: 10, reserved: 0 }); store.buckets.set(userId, [{ id: "expired-reserved-settle-bucket", userId, source: "plan", remaining: 10, priority: 0, expiresAt: new Date(Date.now() + 60_000).toISOString() }]); assert.equal(reserveBalance(store, userId, 10, referenceId), true); store.buckets.get(userId)![0].expiresAt = new Date(Date.now() - 1_000).toISOString(); assert.equal(settleBalance(store, userId, 10, 6, referenceId), true); assert.deepEqual({ available: store.balances.get(userId)?.available, reserved: store.balances.get(userId)?.reserved }, { available: 0, reserved: 0 }); assert.equal(store.buckets.get(userId)?.[0].remaining, 0); assert.equal(store.reserveAllocations.has(`${userId}:${referenceId}`), false); assert.equal(store.ledger.find((entry) => entry.idempotencyKey === `settle:${userId}:${referenceId}:expire:expired-reserved-settle-bucket`)?.amount, -4); assert.equal(store.ledger.find((entry) => entry.idempotencyKey === `settle:${userId}:${referenceId}:expired-reserved-settle-bucket`)?.amount, -6); assert.equal(store.ledger.some((entry) => entry.referenceId === referenceId && entry.type === "release"), false); }); test("expired unknown tasks without a provider request id remain held for reconciliation", () => { const store = createStore(); const userId = "unknown-expiry-user"; store.balances.set(userId, { available: 10, reserved: 0 }); store.buckets.set(userId, [{ id: "bucket", userId, source: "recharge", remaining: 10, priority: 0 }]); assert.equal(reserveBalance(store, userId, 10, "unknown-expiry-task"), true); const expired = new Date(Date.now() - 1_000).toISOString(); store.tasks.set("unknown-expiry-task", { id: "unknown-expiry-task", ownerId: userId, taskType: "image", modelProductId: "basic-image", status: "unknown", estimatedCost: 10, reservedCost: 10, reserveExpiresAt: expired, createdAt: expired, updatedAt: expired, eventSequence: 0, routeSnapshot: { groupId: "image-default", groupVersion: 1, publicModelId: "basic-image-v1", routeVersion: 1, orderedCandidates: [], strategy: "strict", totalRetryBudget: 0 }, modelSnapshot: { publicModelId: "basic-image-v1", modelVersion: 1, name: "基础模型", tier: "basic", capabilities: ["image"] }, pricingSnapshot: { ruleId: "price-basic-image", ruleVersion: 1, publicModelId: "basic-image-v1", unitPrice: 1, quantity: 1, billingUnit: "output", multiplier: 1, balanceUnitVersion: 1 }, planSnapshot: { queuePriority: 0, maxConcurrent: 1 }, attempts: [{ id: "unknown-attempt", channelId: "channel-a", sequence: 1, status: "unknown", startedAt: expired, finishedAt: expired, reconciliationStatus: "pending" }], outputs: [] }); assert.equal(recoverExpiredReservations(store), 0); assert.equal(store.tasks.get("unknown-expiry-task")?.status, "unknown"); assert.equal(store.tasks.get("unknown-expiry-task")?.reservedCost, 10); assert.deepEqual({ available: store.balances.get(userId)?.available, reserved: store.balances.get(userId)?.reserved }, { available: 0, reserved: 10 }); }); test("file snapshots use a version marker and reject corrupt state", async () => { const root = await mkdtemp(path.join(tmpdir(), "miragenflow-snapshot-")); try { const file = path.join(root, "store.json"); const store = createStore(file); store.persist(); const persisted = JSON.parse(await readFile(file, "utf8")) as { snapshotVersion?: number }; assert.equal(persisted.snapshotVersion, 2); await writeFile(file, "{not-json", "utf8"); assert.throws(() => createStore(file), /persistence snapshot rejected/); } finally { await rm(root, { recursive: true, force: true }); } }); test("message outbox fails over to the next enabled provider", async () => { const previous = process.env.MIRAGENFLOW_MESSAGE_FAIL; process.env.MIRAGENFLOW_MESSAGE_FAIL = "email-primary"; const store = createStore(); store.messageProviders.set("email-primary", { id: "email-primary", channel: "email", enabled: true, priority: 0 }); store.messageProviders.set("email-secondary", { id: "email-secondary", channel: "email", enabled: true, priority: 1 }); const record = enqueueMessage(store, { channel: "email", target: "test@example.com", purpose: "register", templateData: { code: "123456" } }); await processMessageOutbox(store); assert.equal(record.status, "sent"); assert.ok(record.providerMessageId?.startsWith("email-secondary-")); if (previous === undefined) delete process.env.MIRAGENFLOW_MESSAGE_FAIL; else process.env.MIRAGENFLOW_MESSAGE_FAIL = previous; }); test("message outbox reuses a supplied idempotency key", () => { const store = createStore(); const first = enqueueMessage(store, { channel: "email", target: "same@example.com", purpose: "notification", templateData: { message: "hello" }, idempotencyKey: "welcome:user-1" }); const second = enqueueMessage(store, { channel: "email", target: "same@example.com", purpose: "notification", templateData: { message: "changed" }, idempotencyKey: "welcome:user-1" }); assert.equal(second.id, first.id); assert.equal(store.messageOutbox.size, 1); }); test("expired WebDAV retention queues only tracked remote files", () => { const store = createStore(); const userId = "webdav-retention-user"; store.webdav.set(userId, { userId, configured: true, directory: "miragenflow", state: "ready", retentionDays: 30, manifestRetentionExpiresAt: new Date(Date.now() - 1_000).toISOString() }); store.webdavFiles.set(`${userId}:canvas/manifest.json`, { userId, path: "canvas/manifest.json", mimeType: "application/json", data: "eA==", checksum: "a".repeat(64), etag: '"v1"', version: 1, updatedAt: new Date().toISOString() }); store.webdavFiles.set(`${userId}:assets/orphan.png`, { userId, path: "assets/orphan.png", mimeType: "image/png", data: "eA==", checksum: "b".repeat(64), etag: '"v1"', version: 1, updatedAt: new Date().toISOString() }); assert.equal(enqueueExpiredWebDavRetention(store), 2); assert.equal(store.webdav.get(userId)?.retentionState, "deleting"); assert.deepEqual([...store.webdavJobs.values()].map((job) => job.intent), ["retention-delete", "retention-delete"]); assert.equal(enqueueExpiredWebDavRetention(store), 0); }); test("memory queue renews a lease and moves exhausted work to dead letters", async () => { const queue = new MemoryTaskQueue(); await queue.enqueue("lease-task", { maxAttempts: 1 }); const lease = await queue.claim("worker", 20); assert.ok(lease); assert.equal(await queue.renew!("lease-task", lease!.leaseToken, 1000), true); await queue.nack!("lease-task", lease!.leaseToken, { retry: true, maxAttempts: 1 }); assert.deepEqual(await queue.deadLetters!(), ["lease-task"]); assert.equal(await queue.renew!("lease-task", lease!.leaseToken, 1000), false); }); test("memory queue rejects stale ack/nack and makes the expired lease claimable again", async () => { const queue = new MemoryTaskQueue(); await queue.enqueue("stale-lease-task"); const lease = await queue.claim("worker", 10); assert.ok(lease); await new Promise((resolve) => setTimeout(resolve, 25)); await queue.ack("stale-lease-task", lease!.leaseToken); const reclaimed = await queue.claim("worker-2", 1_000); assert.equal(reclaimed?.taskId, "stale-lease-task"); await queue.enqueue("stale-nack-task"); const staleNack = await queue.claim("worker-3", 10); assert.ok(staleNack); await new Promise((resolve) => setTimeout(resolve, 25)); assert.equal(await queue.nack!("stale-nack-task", staleNack!.leaseToken, { retry: true }), false); });