test: expand coverage across jobs and tools (#380)

This commit is contained in:
SnapOtter
2026-06-29 22:06:24 +08:00
committed by GitHub
parent ef342c268b
commit fd6ebe77b5
22 changed files with 3692 additions and 0 deletions
@@ -0,0 +1,153 @@
import { afterEach, describe, expect, it, vi } from "vitest";
const statfsMock = vi.hoisted(() => vi.fn());
const getSettingStringMock = vi.hoisted(() => vi.fn());
const deliverWebhookMock = vi.hoisted(() => vi.fn());
const decryptMock = vi.hoisted(() => vi.fn());
const isEncryptedMock = vi.hoisted(() => vi.fn());
const getActiveLicenseMock = vi.hoisted(() => vi.fn());
const selectMock = vi.hoisted(() => vi.fn());
function queryChain<T>(result: T) {
const chain = {
from: vi.fn(() => chain),
where: vi.fn(() => Promise.resolve(result)),
};
return chain;
}
async function loadAlertEvaluator() {
vi.resetModules();
statfsMock.mockReset();
getSettingStringMock.mockReset();
deliverWebhookMock.mockReset();
decryptMock.mockReset();
isEncryptedMock.mockReset();
getActiveLicenseMock.mockReset();
selectMock.mockReset();
vi.doMock("node:fs/promises", () => ({
statfs: statfsMock,
}));
vi.doMock("drizzle-orm", () => ({
and: vi.fn(() => "and"),
eq: vi.fn(() => "eq"),
gte: vi.fn(() => "gte"),
sql: vi.fn(() => "sql"),
}));
vi.doMock("../../../../apps/api/src/config.js", () => ({
env: {
WORKSPACE_PATH: "/workspace",
DATA_ENCRYPTION_KEY: "test-key",
},
}));
vi.doMock("../../../../apps/api/src/db/index.js", () => ({
db: {
select: selectMock,
},
schema: {
auditLog: {
action: "action",
createdAt: "createdAt",
},
},
}));
vi.doMock("../../../../apps/api/src/lib/settings-helpers.js", () => ({
getSettingString: getSettingStringMock,
}));
vi.doMock("../../../../apps/api/src/lib/encryption.js", () => ({
decrypt: decryptMock,
isEncrypted: isEncryptedMock,
}));
vi.doMock("../../../../apps/api/src/lib/webhook-delivery.js", () => ({
deliverWebhook: deliverWebhookMock,
}));
vi.doMock("@snapotter/enterprise", () => ({
getActiveLicense: getActiveLicenseMock,
}));
return import("../../../../apps/api/src/jobs/alert-evaluator.js");
}
describe("alert evaluator behavior", () => {
afterEach(() => {
vi.restoreAllMocks();
});
it("returns early when webhook destination settings are invalid JSON", async () => {
const { evaluateAlerts } = await loadAlertEvaluator();
getSettingStringMock.mockResolvedValueOnce("{invalid");
await evaluateAlerts();
expect(statfsMock).not.toHaveBeenCalled();
expect(deliverWebhookMock).not.toHaveBeenCalled();
});
it("returns early when no enabled alert destinations are configured", async () => {
const { evaluateAlerts } = await loadAlertEvaluator();
getSettingStringMock.mockResolvedValueOnce(
JSON.stringify([
{ url: "https://example.test/siem", authHeader: "", enabled: true, type: "siem" },
{ url: "https://example.test/alerts", authHeader: "", enabled: false, type: "alerts" },
]),
);
await evaluateAlerts();
expect(statfsMock).not.toHaveBeenCalled();
expect(deliverWebhookMock).not.toHaveBeenCalled();
});
it("delivers triggered alerts to enabled alert webhooks with decrypted auth", async () => {
vi.spyOn(Date, "now").mockReturnValue(new Date("2026-06-29T12:00:00.000Z").getTime());
const { evaluateAlerts } = await loadAlertEvaluator();
getSettingStringMock
.mockResolvedValueOnce(
JSON.stringify([
{
url: "https://example.test/alerts",
authHeader: "enc:token",
enabled: true,
type: "alerts",
},
{
url: "https://example.test/ignored",
authHeader: "",
enabled: true,
type: "siem",
},
]),
)
.mockResolvedValueOnce(JSON.stringify({ timestamp: "2026-06-27T11:59:00.000Z" }));
statfsMock.mockResolvedValue({ bfree: 100, bsize: 1024 });
selectMock.mockReturnValue(queryChain([{ count: 21 }]));
getActiveLicenseMock.mockReturnValue({ expiresAt: "2026-07-05T12:00:00.000Z" });
isEncryptedMock.mockReturnValue(true);
decryptMock.mockResolvedValue("Bearer decrypted");
deliverWebhookMock.mockResolvedValue({ success: true });
await evaluateAlerts();
expect(deliverWebhookMock).toHaveBeenCalledTimes(1);
expect(deliverWebhookMock).toHaveBeenCalledWith(
"https://example.test/alerts",
"Bearer decrypted",
expect.arrayContaining([
expect.objectContaining({ condition: "disk_space_low" }),
expect.objectContaining({ condition: "auth_anomaly", failedLogins: 21 }),
expect.objectContaining({ condition: "backup_stale" }),
expect.objectContaining({ condition: "license_expiring", daysLeft: 6 }),
]),
{ maxRetries: 1 },
);
});
});
@@ -0,0 +1,119 @@
import { afterEach, describe, expect, it, vi } from "vitest";
const isFeatureEnabledMock = vi.hoisted(() => vi.fn());
const selectMock = vi.hoisted(() => vi.fn());
const deleteMock = vi.hoisted(() => vi.fn());
const upsertSettingMock = vi.hoisted(() => vi.fn());
const mkdirMock = vi.hoisted(() => vi.fn());
function queryChain<T>(result: T) {
const chain = {
from: vi.fn(() => chain),
where: vi.fn(() => Promise.resolve(result)),
};
return chain;
}
async function loadAuditArchive() {
vi.resetModules();
isFeatureEnabledMock.mockReset();
selectMock.mockReset();
deleteMock.mockReset();
upsertSettingMock.mockReset();
mkdirMock.mockReset();
vi.doMock("node:fs/promises", () => ({
mkdir: mkdirMock,
stat: vi.fn(),
}));
vi.doMock("node:fs", () => ({
createWriteStream: vi.fn(),
}));
vi.doMock("node:stream/promises", () => ({
pipeline: vi.fn(),
}));
vi.doMock("drizzle-orm", () => ({
eq: vi.fn(() => "eq"),
lt: vi.fn(() => "lt"),
}));
vi.doMock("@snapotter/enterprise", () => ({
isFeatureEnabled: isFeatureEnabledMock,
}));
vi.doMock("../../../../apps/api/src/config.js", () => ({
env: { FILES_STORAGE_PATH: "/data/files" },
}));
vi.doMock("../../../../apps/api/src/db/index.js", () => ({
db: {
select: selectMock,
delete: deleteMock,
},
schema: {
settings: {
key: "settings.key",
value: "settings.value",
},
auditLog: {
createdAt: "auditLog.createdAt",
},
},
}));
vi.doMock("../../../../apps/api/src/lib/settings-helpers.js", () => ({
upsertSetting: upsertSettingMock,
}));
return import("../../../../apps/api/src/jobs/audit-archive.js");
}
describe("audit archive job behavior", () => {
afterEach(() => {
vi.restoreAllMocks();
});
it("returns before reading archive settings when the enterprise feature is disabled", async () => {
const { runAuditArchive } = await loadAuditArchive();
isFeatureEnabledMock.mockReturnValue(false);
await runAuditArchive();
expect(selectMock).not.toHaveBeenCalled();
expect(upsertSettingMock).not.toHaveBeenCalled();
});
it("returns when archive months is missing or disabled", async () => {
const { runAuditArchive } = await loadAuditArchive();
isFeatureEnabledMock.mockReturnValue(true);
selectMock.mockReturnValueOnce(queryChain([{ value: "0" }]));
await runAuditArchive();
expect(mkdirMock).not.toHaveBeenCalled();
expect(upsertSettingMock).not.toHaveBeenCalled();
});
it("clears archival state when there are no rows older than the boundary", async () => {
const deleteWhere = vi.fn().mockResolvedValue(undefined);
const { runAuditArchive } = await loadAuditArchive();
deleteMock.mockReturnValue({ where: deleteWhere });
isFeatureEnabledMock.mockReturnValue(true);
selectMock
.mockReturnValueOnce(queryChain([{ value: "1" }]))
.mockReturnValueOnce(queryChain([]))
.mockReturnValueOnce(queryChain([]));
await runAuditArchive();
expect(upsertSettingMock).toHaveBeenCalledWith(
"audit_archival_state",
expect.stringContaining('"state":"EXPORTING"'),
);
expect(mkdirMock).toHaveBeenCalledWith("/data/audit-archives", { recursive: true });
expect(deleteWhere).toHaveBeenCalled();
});
});
@@ -0,0 +1,168 @@
import { afterEach, describe, expect, it, vi } from "vitest";
const insertedValues = vi.hoisted(() => vi.fn());
const queueAdd = vi.hoisted(() => vi.fn());
const getJob = vi.hoisted(() => vi.fn());
const queueEventClose = vi.hoisted(() => vi.fn());
const flowProducerClose = vi.hoisted(() => vi.fn());
async function loadEnqueueModule() {
vi.resetModules();
insertedValues.mockReset();
queueAdd.mockReset();
getJob.mockReset();
queueEventClose.mockReset();
flowProducerClose.mockReset();
queueAdd.mockResolvedValue({ id: "job-1" });
queueEventClose.mockResolvedValue(undefined);
flowProducerClose.mockResolvedValue(undefined);
vi.doMock("bullmq", () => ({
QueueEvents: vi.fn(() => ({
close: queueEventClose,
waitUntilReady: vi.fn().mockResolvedValue(undefined),
})),
FlowProducer: vi.fn(() => ({
close: flowProducerClose,
})),
}));
vi.doMock("../../../../apps/api/src/config.js", () => ({
env: { SYNC_WAIT_MS: 50 },
}));
vi.doMock("../../../../apps/api/src/db/index.js", () => ({
db: {
insert: vi.fn(() => ({
values: insertedValues.mockResolvedValue(undefined),
})),
},
schema: {
jobs: {},
},
}));
vi.doMock("../../../../apps/api/src/jobs/connection.js", () => ({
createBullMQConnection: vi.fn(() => ({ mocked: "connection" })),
}));
vi.doMock("../../../../apps/api/src/jobs/queues.js", () => ({
getQueue: vi.fn(() => ({
add: queueAdd,
getJob,
})),
}));
return import("../../../../apps/api/src/jobs/enqueue.js");
}
describe("job enqueue helpers", () => {
afterEach(() => {
vi.restoreAllMocks();
});
it("strips NUL bytes recursively before persisting settings but keeps queue data intact", async () => {
const { enqueueToolJob } = await loadEnqueueModule();
const data = {
jobId: "job-1",
userId: null,
toolId: "tool-a",
pool: "image",
kind: "single",
inputRefs: ["uploads/job-1/input.png"],
filename: "input.png",
settings: {
title: "a\0b",
nested: { value: "c\0d" },
list: ["e\0f", 1],
},
} as never;
await enqueueToolJob(data);
expect(insertedValues).toHaveBeenCalledWith(
expect.objectContaining({
id: "job-1",
settings: {
title: "ab",
nested: { value: "cd" },
list: ["ef", 1],
},
}),
);
expect(queueAdd).toHaveBeenCalledWith(
"tool-a",
expect.objectContaining({
settings: {
title: "a\0b",
nested: { value: "c\0d" },
list: ["e\0f", 1],
},
}),
{ jobId: "job-1" },
);
});
it("persists redacted dbSettings while enqueueing real settings", async () => {
const { enqueueToolJob } = await loadEnqueueModule();
await enqueueToolJob({
jobId: "job-2",
userId: null,
toolId: "ftp-upload",
pool: "system",
kind: "single",
inputRefs: [],
filename: "file.txt",
settings: { password: "secret" },
dbSettings: { password: "[redacted]" },
} as never);
expect(insertedValues).toHaveBeenCalledWith(
expect.objectContaining({ settings: { password: "[redacted]" } }),
);
expect(queueAdd).toHaveBeenCalledWith(
"ftp-upload",
expect.objectContaining({ settings: { password: "secret" } }),
{ jobId: "job-2" },
);
});
it("waitForJob returns null when the job is missing or the sync window times out", async () => {
const { waitForJob } = await loadEnqueueModule();
getJob.mockResolvedValueOnce(undefined);
await expect(waitForJob("image", "missing")).resolves.toBeNull();
getJob.mockResolvedValueOnce({
waitUntilFinished: vi.fn().mockRejectedValue(new Error("job timed out before finishing")),
});
await expect(waitForJob("image", "slow", 25)).resolves.toBeNull();
});
it("waitForJob rethrows real job failures", async () => {
const { waitForJob } = await loadEnqueueModule();
getJob.mockResolvedValueOnce({
waitUntilFinished: vi.fn().mockRejectedValue(new Error("processor failed")),
});
await expect(waitForJob("image", "failed")).rejects.toThrow("processor failed");
});
it("closes lazy QueueEvents and FlowProducer singletons", async () => {
const { closeFlowProducer, closeQueueEvents, getFlowProducer, warmQueueEvents, waitForJob } =
await loadEnqueueModule();
await warmQueueEvents();
getFlowProducer();
getJob.mockResolvedValueOnce(undefined);
await waitForJob("image", "job-1");
await closeQueueEvents();
await closeFlowProducer();
expect(queueEventClose).toHaveBeenCalled();
expect(flowProducerClose).toHaveBeenCalledTimes(1);
});
});
@@ -0,0 +1,144 @@
import AdmZip from "adm-zip";
import { afterEach, describe, expect, it, vi } from "vitest";
const selectMock = vi.hoisted(() => vi.fn());
const readStoredFileMock = vi.hoisted(() => vi.fn());
const putObjectMock = vi.hoisted(() => vi.fn());
function queryChain<T>(result: T) {
const chain = {
from: vi.fn(() => chain),
where: vi.fn(() => Promise.resolve(result)),
};
return chain;
}
async function loadGdprExport() {
vi.resetModules();
selectMock.mockReset();
readStoredFileMock.mockReset();
putObjectMock.mockReset();
vi.doMock("drizzle-orm", () => ({
eq: vi.fn(() => "eq"),
}));
vi.doMock("../../../../apps/api/src/db/index.js", () => ({
db: {
select: selectMock,
},
schema: {
users: { id: "users.id" },
userFiles: { userId: "userFiles.userId" },
jobs: { userId: "jobs.userId" },
auditLog: { actorId: "auditLog.actorId" },
},
}));
vi.doMock("../../../../apps/api/src/lib/file-storage.js", () => ({
readStoredFile: readStoredFileMock,
}));
vi.doMock("../../../../apps/api/src/lib/object-storage.js", () => ({
putObject: putObjectMock,
}));
return import("../../../../apps/api/src/jobs/gdpr-export.js");
}
describe("GDPR export job behavior", () => {
afterEach(() => {
vi.restoreAllMocks();
});
it("throws before writing output when the user does not exist", async () => {
const { gdprExportJob } = await loadGdprExport();
selectMock.mockReturnValueOnce(queryChain([]));
await expect(gdprExportJob("missing-user", "job-1")).rejects.toThrow(
"User missing-user not found",
);
expect(putObjectMock).not.toHaveBeenCalled();
});
it("writes a ZIP without passwordHash and skips missing library file contents", async () => {
const { gdprExportJob } = await loadGdprExport();
selectMock
.mockReturnValueOnce(
queryChain([
{
id: "user-1",
email: "ada@example.test",
passwordHash: "do-not-export",
createdAt: new Date("2026-06-01T00:00:00.000Z"),
},
]),
)
.mockReturnValueOnce(
queryChain([
{
id: "file-1",
userId: "user-1",
storedName: "stored/a",
originalName: "a.txt",
createdAt: new Date("2026-06-02T00:00:00.000Z"),
},
{
id: "file-2",
userId: "user-1",
storedName: "stored/missing",
originalName: "missing.txt",
createdAt: new Date("2026-06-03T00:00:00.000Z"),
},
]),
)
.mockReturnValueOnce(
queryChain([
{
id: "job-a",
userId: "user-1",
createdAt: new Date("2026-06-04T00:00:00.000Z"),
startedAt: null,
completedAt: new Date("2026-06-04T00:01:00.000Z"),
deleteAfter: null,
},
]),
)
.mockReturnValueOnce(
queryChain([
{
id: "audit-1",
actorId: "user-1",
action: "LOGIN",
createdAt: new Date("2026-06-05T00:00:00.000Z"),
},
]),
);
readStoredFileMock
.mockResolvedValueOnce(Buffer.from("file contents"))
.mockRejectedValueOnce(new Error("missing"));
await expect(gdprExportJob("user-1", "export-job")).resolves.toEqual({
outputRef: "outputs/export-job/gdpr-export.zip",
});
expect(putObjectMock).toHaveBeenCalledTimes(1);
const [outputRef, zipBuffer] = putObjectMock.mock.calls[0];
expect(outputRef).toBe("outputs/export-job/gdpr-export.zip");
const zip = new AdmZip(zipBuffer);
const profile = JSON.parse(zip.readAsText("profile.json"));
const files = JSON.parse(zip.readAsText("files.json"));
const jobs = JSON.parse(zip.readAsText("jobs.json"));
const audit = JSON.parse(zip.readAsText("audit-log.json"));
expect(profile).toMatchObject({ id: "user-1", email: "ada@example.test" });
expect(profile).not.toHaveProperty("passwordHash");
expect(files[0].createdAt).toBe("2026-06-02T00:00:00.000Z");
expect(jobs[0].completedAt).toBe("2026-06-04T00:01:00.000Z");
expect(audit[0].createdAt).toBe("2026-06-05T00:00:00.000Z");
expect(zip.readAsText("library-files/file-1_a.txt")).toBe("file contents");
expect(zip.getEntry("library-files/file-2_missing.txt")).toBeNull();
});
});
+115
View File
@@ -0,0 +1,115 @@
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
const queueInstances: Array<{
name: string;
options: Record<string, unknown>;
close: ReturnType<typeof vi.fn>;
getJobCounts: ReturnType<typeof vi.fn>;
getJobs: ReturnType<typeof vi.fn>;
}> = [];
async function loadQueuesModule() {
vi.resetModules();
queueInstances.length = 0;
vi.doMock("bullmq", () => ({
Queue: vi.fn((name: string, options: Record<string, unknown>) => {
const queue = {
name,
options,
close: vi.fn().mockResolvedValue(undefined),
getJobCounts: vi.fn().mockResolvedValue({ active: 0, waiting: 0, delayed: 0, failed: 0 }),
getJobs: vi.fn().mockResolvedValue([]),
};
queueInstances.push(queue);
return queue;
}),
}));
vi.doMock("../../../../apps/api/src/jobs/connection.js", () => ({
createBullMQConnection: vi.fn(() => ({ mocked: "connection" })),
}));
return import("../../../../apps/api/src/jobs/queues.js");
}
describe("job queues", () => {
beforeEach(() => {
vi.useRealTimers();
});
afterEach(() => {
vi.restoreAllMocks();
});
it("creates one cached queue per pool with pool-specific retry attempts", async () => {
const { getQueue } = await loadQueuesModule();
const imageQueue = getQueue("image");
const sameImageQueue = getQueue("image");
const aiQueue = getQueue("ai");
expect(sameImageQueue).toBe(imageQueue);
expect(aiQueue).not.toBe(imageQueue);
expect(queueInstances).toHaveLength(2);
expect(queueInstances[0].name).toContain("image");
expect(queueInstances[1].name).toContain("ai");
expect(queueInstances[0].options).toMatchObject({
defaultJobOptions: { attempts: 2 },
});
expect(queueInstances[1].options).toMatchObject({
defaultJobOptions: { attempts: 1 },
});
});
it("aggregates counts only from queues that have been created", async () => {
const { getQueue, queueCounts, perPoolCounts } = await loadQueuesModule();
getQueue("image");
getQueue("docs");
queueInstances[0].getJobCounts.mockResolvedValueOnce({ active: 2, waiting: 3, delayed: 4 });
queueInstances[1].getJobCounts.mockResolvedValueOnce({ active: 5, waiting: 7, delayed: 11 });
await expect(queueCounts()).resolves.toEqual({ active: 7, waiting: 10, delayed: 15 });
queueInstances[0].getJobCounts.mockResolvedValueOnce({ active: 13, waiting: 17 });
queueInstances[1].getJobCounts.mockResolvedValueOnce({ active: 19, waiting: 23 });
await expect(perPoolCounts()).resolves.toMatchObject({
image: { active: 13, waiting: 17 },
docs: { active: 19, waiting: 23 },
ai: { active: 0, waiting: 0 },
media: { active: 0, waiting: 0 },
system: { active: 0, waiting: 0 },
});
});
it("reports oldest waiting age for per-pool health when waiting jobs exist", async () => {
vi.useFakeTimers();
vi.setSystemTime(new Date("2026-06-29T12:00:00.000Z"));
const { getQueue, perPoolHealth } = await loadQueuesModule();
getQueue("media");
queueInstances[0].getJobCounts.mockResolvedValueOnce({ active: 1, waiting: 1, failed: 2 });
queueInstances[0].getJobs.mockResolvedValueOnce([
{ timestamp: new Date("2026-06-29T11:59:45.000Z").getTime() },
]);
await expect(perPoolHealth()).resolves.toMatchObject({
media: { active: 1, waiting: 1, failed: 2, oldestWaitingMs: 15_000 },
image: { active: 0, waiting: 0, failed: 0, oldestWaitingMs: null },
});
});
it("closes cached queues and clears counts", async () => {
const { getQueue, closeQueues, queueCounts } = await loadQueuesModule();
getQueue("image");
getQueue("ai");
await closeQueues();
expect(queueInstances[0].close).toHaveBeenCalledTimes(1);
expect(queueInstances[1].close).toHaveBeenCalledTimes(1);
await expect(queueCounts()).resolves.toEqual({ active: 0, waiting: 0, delayed: 0 });
});
});
@@ -0,0 +1,181 @@
import { afterEach, describe, expect, it, vi } from "vitest";
const readSiemConfigMock = vi.hoisted(() => vi.fn());
const deliverWebhookMock = vi.hoisted(() => vi.fn());
const upsertSettingMock = vi.hoisted(() => vi.fn());
const decryptMock = vi.hoisted(() => vi.fn());
const isEncryptedMock = vi.hoisted(() => vi.fn());
const selectMock = vi.hoisted(() => vi.fn());
function queryChain<T>(result: T, terminalWhere = false) {
const chain = {
from: vi.fn(() => chain),
where: vi.fn(() => (terminalWhere ? Promise.resolve(result) : chain)),
orderBy: vi.fn(() => chain),
limit: vi.fn(() => Promise.resolve(result)),
};
return chain;
}
async function loadSiemForward() {
vi.resetModules();
readSiemConfigMock.mockReset();
deliverWebhookMock.mockReset();
upsertSettingMock.mockReset();
decryptMock.mockReset();
isEncryptedMock.mockReset();
selectMock.mockReset();
vi.doMock("drizzle-orm", () => ({
asc: vi.fn(() => "asc"),
eq: vi.fn(() => "eq"),
gte: vi.fn(() => "gte"),
}));
vi.doMock("../../../../apps/api/src/config.js", () => ({
env: { DATA_ENCRYPTION_KEY: "test-key" },
}));
vi.doMock("../../../../apps/api/src/db/index.js", () => ({
db: {
select: selectMock,
},
schema: {
auditLog: {
action: "action",
actorId: "actorId",
actorUsername: "actorUsername",
targetType: "targetType",
targetId: "targetId",
ipAddress: "ipAddress",
details: "details",
createdAt: "createdAt",
},
settings: {
key: "key",
value: "value",
},
},
}));
vi.doMock("../../../../apps/api/src/lib/encryption.js", () => ({
decrypt: decryptMock,
isEncrypted: isEncryptedMock,
}));
vi.doMock("../../../../apps/api/src/lib/settings-helpers.js", () => ({
upsertSetting: upsertSettingMock,
}));
vi.doMock("../../../../apps/api/src/lib/webhook-delivery.js", () => ({
deliverWebhook: deliverWebhookMock,
}));
vi.doMock("../../../../apps/api/src/routes/enterprise/siem.js", () => ({
readSiemConfig: readSiemConfigMock,
}));
return import("../../../../apps/api/src/jobs/siem-forward.js");
}
describe("SIEM forwarding behavior", () => {
afterEach(() => {
vi.restoreAllMocks();
});
it("returns without querying audit rows when SIEM is disabled", async () => {
const { runSiemForward } = await loadSiemForward();
readSiemConfigMock.mockResolvedValue({ enabled: false, webhookUrl: "https://siem.test" });
await expect(runSiemForward()).resolves.toBeUndefined();
expect(selectMock).not.toHaveBeenCalled();
expect(deliverWebhookMock).not.toHaveBeenCalled();
});
it("opens the circuit breaker at five consecutive failures", async () => {
const { runSiemForward } = await loadSiemForward();
readSiemConfigMock.mockResolvedValue({ enabled: true, webhookUrl: "https://siem.test" });
selectMock.mockReturnValueOnce(queryChain([{ value: "5" }], true));
await expect(runSiemForward()).resolves.toBeUndefined();
expect(selectMock).toHaveBeenCalledTimes(1);
expect(deliverWebhookMock).not.toHaveBeenCalled();
});
it("maps audit rows, decrypts auth, advances cursor, and resets failures after success", async () => {
const { runSiemForward } = await loadSiemForward();
const createdAt = new Date("2026-06-29T12:00:00.000Z");
readSiemConfigMock.mockResolvedValue({
enabled: true,
webhookUrl: "https://siem.test/events",
authHeader: "enc:auth",
});
selectMock
.mockReturnValueOnce(queryChain([{ value: "2" }], true))
.mockReturnValueOnce(queryChain([{ value: "2026-06-29T11:00:00.000Z" }], true))
.mockReturnValueOnce(
queryChain([
{
createdAt,
action: "LOGIN_FAILED",
actorId: "user-1",
actorUsername: "ada",
targetType: "session",
targetId: "session-1",
ipAddress: "203.0.113.10",
details: { reason: "bad_password" },
},
]),
);
isEncryptedMock.mockReturnValue(true);
decryptMock.mockResolvedValue("Bearer clear");
deliverWebhookMock.mockResolvedValue({ success: true });
await expect(runSiemForward()).resolves.toEqual({ forwarded: 1 });
expect(deliverWebhookMock).toHaveBeenCalledWith("https://siem.test/events", "Bearer clear", [
{
timestamp: "2026-06-29T12:00:00.000Z",
event: "LOGIN_FAILED",
actorId: "user-1",
actorUsername: "ada",
targetType: "session",
targetId: "session-1",
ip: "203.0.113.10",
details: { reason: "bad_password" },
},
]);
expect(upsertSettingMock).toHaveBeenCalledWith(
"siem_last_forwarded_at",
"2026-06-29T12:00:00.000Z",
);
expect(upsertSettingMock).toHaveBeenCalledWith("siem_consecutive_failures", "0");
});
it("increments failure counter when delivery fails", async () => {
const { runSiemForward } = await loadSiemForward();
readSiemConfigMock.mockResolvedValue({
enabled: true,
webhookUrl: "https://siem.test/events",
authHeader: "",
});
selectMock
.mockReturnValueOnce(queryChain([{ value: "4" }], true))
.mockReturnValueOnce(queryChain([], true))
.mockReturnValueOnce(
queryChain([
{
createdAt: new Date("2026-06-29T12:00:00.000Z"),
action: "FILE_DELETED",
},
]),
);
deliverWebhookMock.mockResolvedValue({ success: false, error: "downstream 500" });
await expect(runSiemForward()).resolves.toBeUndefined();
expect(upsertSettingMock).toHaveBeenCalledWith("siem_consecutive_failures", "5");
});
});
@@ -0,0 +1,195 @@
import { afterEach, describe, expect, it, vi } from "vitest";
const getQueueMock = vi.hoisted(() => vi.fn());
const runSiemForwardMock = vi.hoisted(() => vi.fn());
const runAuditArchiveMock = vi.hoisted(() => vi.fn());
const dbExecuteMock = vi.hoisted(() => vi.fn());
const dbUpdateMock = vi.hoisted(() => vi.fn());
const storageReconciliationJobMock = vi.hoisted(() => vi.fn());
const gdprExportJobMock = vi.hoisted(() => vi.fn());
const evaluateAlertsMock = vi.hoisted(() => vi.fn());
async function loadSystemJobs(cleanupIntervalMinutes = 15) {
vi.resetModules();
getQueueMock.mockReset();
runSiemForwardMock.mockReset();
runAuditArchiveMock.mockReset();
dbExecuteMock.mockReset();
dbUpdateMock.mockReset();
storageReconciliationJobMock.mockReset();
gdprExportJobMock.mockReset();
evaluateAlertsMock.mockReset();
vi.doMock("drizzle-orm", () => ({
and: vi.fn(() => "and"),
eq: vi.fn(() => "eq"),
inArray: vi.fn(() => "inArray"),
isNotNull: vi.fn(() => "isNotNull"),
lt: vi.fn(() => "lt"),
sql: vi.fn(() => "sql"),
}));
vi.doMock("../../../../apps/api/src/config.js", () => ({
env: {
CLEANUP_INTERVAL_MINUTES: cleanupIntervalMinutes,
JOBS_RETENTION_DAYS: 30,
AUDIT_RETENTION_DAYS: 90,
},
}));
vi.doMock("../../../../apps/api/src/db/index.js", () => ({
db: {
execute: dbExecuteMock.mockResolvedValue(undefined),
update: dbUpdateMock,
},
schema: {
jobs: {
id: "jobs.id",
},
},
}));
vi.doMock("../../../../apps/api/src/lib/cleanup.js", () => ({
getMaxAgeMs: vi.fn().mockResolvedValue(0),
}));
vi.doMock("../../../../apps/api/src/lib/object-storage.js", () => ({
deletePrefix: vi.fn(),
listJobDirs: vi.fn().mockResolvedValue([]),
}));
vi.doMock("../../../../apps/api/src/lib/settings-helpers.js", () => ({
getSettingNumber: vi.fn(),
}));
vi.doMock("../../../../apps/api/src/jobs/audit-archive.js", () => ({
runAuditArchive: runAuditArchiveMock,
}));
vi.doMock("../../../../apps/api/src/jobs/queues.js", () => ({
getQueue: getQueueMock,
}));
vi.doMock("../../../../apps/api/src/jobs/siem-forward.js", () => ({
runSiemForward: runSiemForwardMock,
}));
vi.doMock("../../../../apps/api/src/jobs/storage-reconciliation.js", () => ({
storageReconciliationJob: storageReconciliationJobMock,
}));
vi.doMock("../../../../apps/api/src/jobs/gdpr-export.js", () => ({
gdprExportJob: gdprExportJobMock,
}));
vi.doMock("../../../../apps/api/src/jobs/alert-evaluator.js", () => ({
evaluateAlerts: evaluateAlertsMock,
}));
return import("../../../../apps/api/src/jobs/system-jobs.js");
}
describe("system jobs behavior", () => {
afterEach(() => {
vi.restoreAllMocks();
});
it("decides expiry from local mtimes, S3 job rows, and rowless S3 directories", async () => {
const { decideExpiry } = await loadSystemJobs();
const cutoffMs = new Date("2026-06-29T12:00:00.000Z").getTime();
const rowsById = new Map([
[
"job-old",
{
createdAt: new Date("2026-06-20T12:00:00.000Z"),
completedAt: null,
},
],
[
"job-new",
{
createdAt: new Date("2026-06-20T12:00:00.000Z"),
completedAt: new Date("2026-06-29T12:01:00.000Z"),
},
],
]);
expect(decideExpiry({ key: "uploads/job-a", mtimeMs: cutoffMs - 1 }, cutoffMs, rowsById)).toBe(
"expired",
);
expect(decideExpiry({ key: "outputs/job-b", mtimeMs: cutoffMs }, cutoffMs, rowsById)).toBe(
"keep",
);
expect(decideExpiry({ key: "uploads/job-old", mtimeMs: 0 }, cutoffMs, rowsById)).toBe(
"expired",
);
expect(decideExpiry({ key: "outputs/job-new", mtimeMs: 0 }, cutoffMs, rowsById)).toBe("keep");
expect(decideExpiry({ key: "uploads/orphan", mtimeMs: 0 }, cutoffMs, rowsById)).toBe("skip");
});
it("schedules repeatable jobs and removes storage TTL scheduler when cleanup is disabled", async () => {
const queue = {
upsertJobScheduler: vi.fn().mockResolvedValue(undefined),
removeJobScheduler: vi.fn().mockResolvedValue(undefined),
};
const { SYSTEM_JOBS, scheduleSystemJobs } = await loadSystemJobs(0);
getQueueMock.mockReturnValue(queue);
await scheduleSystemJobs();
expect(queue.removeJobScheduler).toHaveBeenCalledWith(SYSTEM_JOBS.storageTtl);
expect(queue.upsertJobScheduler).toHaveBeenCalledWith(SYSTEM_JOBS.sessionPurge, {
every: 60 * 60_000,
});
expect(queue.upsertJobScheduler).toHaveBeenCalledWith(SYSTEM_JOBS.retention, {
every: 6 * 60 * 60_000,
});
expect(queue.upsertJobScheduler).toHaveBeenCalledWith(SYSTEM_JOBS.auditArchive, {
pattern: "0 2 1 * *",
});
expect(queue.upsertJobScheduler).toHaveBeenCalledWith(SYSTEM_JOBS.storageReconciliation, {
pattern: "0 3 * * 0",
});
expect(queue.upsertJobScheduler).toHaveBeenCalledWith(SYSTEM_JOBS.alertEvaluator, {
every: 60_000,
});
});
it("dispatches one-shot system jobs and updates GDPR export job rows", async () => {
const updateWhere = vi.fn().mockResolvedValue(undefined);
const updateSet = vi.fn(() => ({ where: updateWhere }));
const { SYSTEM_JOBS, runSystemJob } = await loadSystemJobs();
dbUpdateMock.mockReturnValue({ set: updateSet });
runSiemForwardMock.mockResolvedValue({ forwarded: 2 });
gdprExportJobMock.mockResolvedValue({ outputRef: "outputs/export-job/gdpr-export.zip" });
storageReconciliationJobMock.mockResolvedValue(undefined);
evaluateAlertsMock.mockResolvedValue(undefined);
await expect(runSystemJob({ name: SYSTEM_JOBS.siemForward } as never)).resolves.toEqual({
forwarded: 2,
});
await expect(runSystemJob({ name: SYSTEM_JOBS.storageReconciliation } as never)).resolves.toBe(
undefined,
);
await expect(
runSystemJob({
name: SYSTEM_JOBS.gdprExport,
data: { userId: "user-1", jobId: "export-job" },
} as never),
).resolves.toEqual({ outputRef: "outputs/export-job/gdpr-export.zip" });
await expect(runSystemJob({ name: SYSTEM_JOBS.alertEvaluator } as never)).resolves.toBe(
undefined,
);
expect(gdprExportJobMock).toHaveBeenCalledWith("user-1", "export-job");
expect(updateSet).toHaveBeenCalledWith(
expect.objectContaining({
status: "completed",
outputRefs: ["outputs/export-job/gdpr-export.zip"],
}),
);
await expect(runSystemJob({ name: "system:unknown" } as never)).rejects.toThrow(
"Unknown system job: system:unknown",
);
});
});
+171
View File
@@ -0,0 +1,171 @@
import { afterEach, describe, expect, it, vi } from "vitest";
async function loadWorker() {
vi.resetModules();
vi.doMock("node:fs/promises", () => ({
mkdir: vi.fn(),
readFile: vi.fn(),
rm: vi.fn(),
}));
vi.doMock("@snapotter/shared", () => ({
ANALYTICS_EVENTS: {},
TOOLS: [],
getBundleForTool: vi.fn(() => null),
}));
vi.doMock("bullmq", () => ({
UnrecoverableError: class UnrecoverableError extends Error {},
Worker: vi.fn(() => ({
on: vi.fn(),
close: vi.fn().mockResolvedValue(undefined),
})),
}));
vi.doMock("drizzle-orm", () => ({
eq: vi.fn(() => "eq"),
}));
vi.doMock("../../../../apps/api/src/config.js", () => ({
env: {
SCRATCH_PATH: "",
JOB_TIMEOUT_LONG_S: 60,
JOB_TIMEOUT_FAST_S: 15,
},
}));
vi.doMock("../../../../apps/api/src/db/index.js", () => ({
db: {},
schema: { jobs: {} },
}));
vi.doMock("../../../../apps/api/src/lib/analytics.js", () => ({
captureException: vi.fn(),
trackEvent: vi.fn(),
}));
vi.doMock("../../../../apps/api/src/lib/analytics-gate.js", () => ({
analyticsEnabled: vi.fn(() => false),
}));
vi.doMock("../../../../apps/api/src/lib/env.js", () => ({
resolveConcurrency: vi.fn(() => 2),
}));
vi.doMock("../../../../apps/api/src/lib/errors.js", () => ({
friendlyError: vi.fn((message: string) => message),
}));
vi.doMock("../../../../apps/api/src/lib/logger.js", () => ({
logger: {
error: vi.fn(),
info: vi.fn(),
},
}));
vi.doMock("../../../../apps/api/src/lib/metrics.js", () => ({
jobDuration: { observe: vi.fn() },
jobsTotal: { inc: vi.fn() },
}));
vi.doMock("../../../../apps/api/src/lib/object-storage.js", () => ({
getObjectBuffer: vi.fn(),
putObject: vi.fn(),
}));
vi.doMock("../../../../apps/api/src/routes/progress.js", () => ({
publishEphemeral: vi.fn(),
updateSingleFileProgress: vi.fn(),
}));
vi.doMock("../../../../apps/api/src/routes/tool-factory.js", () => ({
getToolConfig: vi.fn(),
}));
vi.doMock("../../../../apps/api/src/jobs/ai-handlers.js", () => ({
hasAiJobHandler: vi.fn(() => false),
runAiToolJob: vi.fn(),
}));
vi.doMock("../../../../apps/api/src/jobs/batch-progress.js", () => ({
recordChildOutcome: vi.fn(),
}));
vi.doMock("../../../../apps/api/src/jobs/cancel.js", () => ({
registerCancelable: vi.fn(() => new AbortController()),
unregisterCancelable: vi.fn(),
}));
vi.doMock("../../../../apps/api/src/jobs/connection.js", () => ({
createBullMQConnection: vi.fn(() => ({})),
}));
vi.doMock("../../../../apps/api/src/jobs/postprocess.js", () => ({
autoSaveToLibrary: vi.fn(),
buildOutputName: vi.fn(),
generatePreview: vi.fn(),
}));
vi.doMock("../../../../apps/api/src/jobs/system-jobs.js", () => ({
runSystemJob: vi.fn(),
}));
return import("../../../../apps/api/src/jobs/worker.js");
}
describe("worker result payload behavior", () => {
afterEach(() => {
vi.restoreAllMocks();
});
it("builds legacy download, preview, saved-file, and tool payload fields", async () => {
const { buildLegacyResultPayload } = await loadWorker();
expect(
buildLegacyResultPayload(
{
outputRefs: ["outputs/job-1/report final.pdf"],
filename: "report final.pdf",
contentType: "application/pdf",
originalSize: 100,
processedSize: 80,
previewRef: "outputs/job-1/preview.png",
savedFileId: "file-2",
resultPayload: { pageCount: 3 },
},
"job-1",
),
).toEqual({
jobId: "job-1",
downloadUrl: "/api/v1/download/job-1/report%20final.pdf",
previewUrl: "/api/v1/download/job-1/preview.png",
originalSize: 100,
processedSize: 80,
savedFileId: "file-2",
pageCount: 3,
});
});
it("omits optional legacy payload fields when the job result does not include them", async () => {
const { buildLegacyResultPayload } = await loadWorker();
expect(
buildLegacyResultPayload(
{
outputRefs: ["outputs/job-2/out.png"],
filename: "out.png",
contentType: "image/png",
originalSize: 10,
processedSize: 8,
},
"job-2",
),
).toEqual({
jobId: "job-2",
downloadUrl: "/api/v1/download/job-2/out.png",
originalSize: 10,
processedSize: 8,
});
});
});