Files

376 lines
24 KiB
TypeScript

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<void>((resolve) => { resolveSaveStarted = resolve; });
const saveRelease = new Promise<void>((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<void>((resolve) => { resolveLoadStarted = resolve; });
const loadRelease = new Promise<void>((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, channelGroupId: "image-default", routeSnapshotVersion: 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, channelGroupId: "image-default", routeSnapshotVersion: 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, channelGroupId: "image-default", routeSnapshotVersion: 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, 1);
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);
});