feat(catalog): persist idempotent bulk operations and retryable deliveries
CI / check (push) Successful in 3m7s
CI / deploy (push) Successful in 1m37s
CI / publish-container (push) Successful in 44s

This commit is contained in:
Simo committed 2026-09-13 18:05:03 +02:00
1 parent 76e46d295c
commit 13665e0c3c
53 files changed
+1316 -76

No files matched your search

@@ -82,6 +82,7 @@ export function BulkOfferEditor({
const [preview, setPreview] = useState<{
input: BulkInput;
data: Preview;
requestKey: string;
} | null>(null);
const [busy, setBusy] = useState(false);
const busyRef = useRef(false);
@@ -149,7 +150,7 @@ export function BulkOfferEditor({
setError(result.error);
return;
}
setPreview({ input, data: result.data });
setPreview({ input, data: result.data, requestKey: crypto.randomUUID() });
} catch {
setError(t("requestFailed"));
} finally {
@@ -172,12 +173,14 @@ export function BulkOfferEditor({
const result = await applyBulkOffers(
preview.input,
preview.data.fingerprint,
preview.requestKey,
);
if (!result.ok) {
setError(result.error);
setPreview(null);
return;
}
const undoKey = crypto.randomUUID();
toast.success(t("success", { count: result.data.changedCount }), {
duration: 15000,
description: t("undoHint"),
@@ -188,7 +191,10 @@ export function BulkOfferEditor({
if (busyRef.current || (beforeEdit && !beforeEdit())) return;
busyRef.current = true;
try {
const restored = await undoBulkOffers(result.data.historyIds);
const restored = await undoBulkOffers(
result.data.historyIds,
undoKey,
);
if (!restored.ok) {
toast.error(restored.error);
return;
@@ -210,7 +216,6 @@ export function BulkOfferEditor({
await onApplied();
} catch {
setError(t("requestFailed"));
setPreview(null);
} finally {
busyRef.current = false;
setBusy(false);
@@ -57,10 +57,8 @@ it("gates edits before database and sends one update after commit", async () =>
expect((await applyBulkOffers(input, "x")).ok).toBe(true);
expect(state.permission).toHaveBeenCalledWith("edit");
expect(state.export).toHaveBeenCalledTimes(1);
expect(state.send).toHaveBeenCalledTimes(1);
expect(state.apply.mock.invocationCallOrder[0]).toBeLessThan(
state.send.mock.invocationCallOrder[0],
);
expect(state.send).not.toHaveBeenCalled();
expect(state.refresh).toHaveBeenCalled();
});
it("stops denied permission before export or data access", async () => {
state.permission.mockRejectedValueOnce(Error("Denied"));
@@ -76,7 +74,7 @@ it("does not notify hotel when a transaction fails", async () => {
it("does not report a committed edit as failed on audit failure", async () => {
state.audit.mockRejectedValueOnce(Error("log failed"));
expect((await applyBulkOffers(input, "x")).ok).toBe(true);
expect(state.send).toHaveBeenCalledTimes(1);
expect(state.send).not.toHaveBeenCalled();
});
it("loads destination choices with view permission only", async () => {
@@ -95,9 +93,9 @@ it("undo enforces edit access and exports only successful restores", async () =>
state.undo.mockResolvedValueOnce({ changedCount: 2 });
expect((await undoBulkOffers([10, 11])).ok).toBe(true);
expect(state.permission).toHaveBeenCalledWith("edit");
expect(state.undo).toHaveBeenCalledWith([10, 11], 1);
expect(state.undo).toHaveBeenCalledWith([10, 11], 1, undefined);
expect(state.export).toHaveBeenCalledTimes(1);
expect(state.send).toHaveBeenCalledTimes(1);
expect(state.send).not.toHaveBeenCalled();
});
it("undo denial stops before export and a conflict does not notify", async () => {
state.permission.mockRejectedValueOnce(Error("Denied"));
@@ -107,3 +105,10 @@ it("undo denial stops before export and a conflict does not notify", async () =>
expect((await undoBulkOffers([10])).ok).toBe(false);
expect(state.send).not.toHaveBeenCalled();
});
it("passes a stable operation key and does not synchronously notify after commit", async () => {
const key = "123e4567-e89b-42d3-a456-426614174000";
expect((await applyBulkOffers(input, "x", key)).ok).toBe(true);
expect(state.apply).toHaveBeenCalledWith(input, "x", 1, key);
expect(state.send).not.toHaveBeenCalled();
});
@@ -115,6 +115,17 @@ vi.mock("@/lib/db", () => ({
},
}));
vi.mock("@/features/operations/server", async () => {
const { db } = await import("@/lib/db");
return {
runOperation: (
_input: unknown,
work: (tx: unknown, id: string) => Promise<unknown>,
) => db.transaction((tx) => work(tx, "operation-fixture")),
enqueueEffect: vi.fn(),
};
});
import {
applyBulkOffersCommand,
listBulkOfferDestinationsCommand,
+60 -26
View File
@@ -1,5 +1,5 @@
import "server-only";
import { createHash } from "node:crypto";
import { createHash, randomUUID } from "node:crypto";
import { sql } from "drizzle-orm";
import {
applyHistory,
@@ -7,6 +7,7 @@ import {
readHistory,
recordHistory,
} from "@/features/history/server";
import { enqueueEffect, runOperation } from "@/features/operations/server";
import { db } from "@/lib/db";
import {
type BulkOfferInput,
@@ -115,11 +116,12 @@ export async function applyBulkOffersCommand(
value: BulkOfferInput,
fingerprint: string,
userId?: number,
requestKey: string = randomUUID(),
) {
const input = bulkOfferInputSchema.parse(value);
if (!/^[a-f0-9]{64}$/.test(fingerprint))
throw new CatalogInputError("A valid preview is required");
return db.transaction(async (tx) => {
const work = async (tx: Tx, operationId?: string) => {
// Read references, then lock pages before offers to match category structural commands.
const initial = await readOffers(tx, input.ids);
const pages = await readPages(tx, initial, input, true);
@@ -167,8 +169,27 @@ export async function applyBulkOffersCommand(
historyIds.push(historyId);
}
}
if (operationId && result.changedCount > 0) {
await enqueueEffect(tx, operationId, "catalog.refresh");
await enqueueEffect(tx, operationId, "catalog.export.request");
}
return { changedCount: result.changedCount, historyIds };
});
};
return userId
? runOperation(
{
actorId: userId,
kind: "catalog.bulk.apply",
key: requestKey,
input: {
...input,
ids: [...input.ids].sort((a, b) => a - b),
fingerprint,
},
},
work,
)
: db.transaction((tx) => work(tx));
}
export async function listBulkOfferDestinationsCommand() {
@@ -188,6 +209,7 @@ export async function listBulkOfferDestinationsCommand() {
export async function undoBulkOffersCommand(
historyIds: number[],
userId: number,
requestKey: string = randomUUID(),
) {
if (
!Array.isArray(historyIds) ||
@@ -197,27 +219,39 @@ export async function undoBulkOffersCommand(
historyIds.some((id) => !Number.isSafeInteger(id) || id < 1)
)
throw new CatalogInputError("Invalid undo selection");
return db.transaction(async (tx) => {
const entries = [];
for (const id of historyIds) {
const entry = await readHistory(tx, id);
if (entry.kind !== "catalog_offer")
throw new CatalogInputError("Invalid undo selection");
entries.push(entry);
}
entries.sort((a, b) => Number(a.targetId) - Number(b.targetId));
if (new Set(entries.map((entry) => entry.targetId)).size !== entries.length)
throw new CatalogInputError("Duplicate undo target");
await lockOfferHistoryPages(tx, entries);
try {
for (const entry of entries) await applyHistory(tx, entry, userId);
} catch (error) {
if (error instanceof Error && error.message === "conflict")
throw new CatalogConflict(
"Offers changed after this update. Undo was not applied.",
);
throw error;
}
return { changedCount: entries.length };
});
return runOperation(
{
actorId: userId,
kind: "catalog.bulk.undo",
key: requestKey,
input: { historyIds: [...historyIds].sort((a, b) => a - b) },
},
async (tx, operationId) => {
const entries = [];
for (const id of historyIds) {
const entry = await readHistory(tx, id);
if (entry.kind !== "catalog_offer")
throw new CatalogInputError("Invalid undo selection");
entries.push(entry);
}
entries.sort((a, b) => Number(a.targetId) - Number(b.targetId));
if (
new Set(entries.map((entry) => entry.targetId)).size !== entries.length
)
throw new CatalogInputError("Duplicate undo target");
await lockOfferHistoryPages(tx, entries);
try {
for (const entry of entries) await applyHistory(tx, entry, userId);
} catch (error) {
if (error instanceof Error && error.message === "conflict")
throw new CatalogConflict(
"Offers changed after this update. Undo was not applied.",
);
throw error;
}
await enqueueEffect(tx, operationId, "catalog.refresh");
await enqueueEffect(tx, operationId, "catalog.export.request");
return { changedCount: entries.length };
},
);
}
@@ -79,6 +79,17 @@ vi.mock("@/lib/db", () => ({
},
}));
vi.mock("@/features/operations/server", async () => {
const { db } = await import("@/lib/db");
return {
runOperation: (
_input: unknown,
work: (tx: unknown, id: string) => Promise<unknown>,
) => db.transaction((tx) => work(tx, "operation-fixture")),
enqueueEffect: vi.fn(),
};
});
import { undoBulkOffersCommand } from "./bulk-offers";
beforeEach(() => {
@@ -605,15 +605,15 @@ describe("housekeeping foundation completion contracts", () => {
expect(html).toContain(`>${sentinel}</button>`);
});
it("validates the complete 142-row migration matrix without issues", () => {
it("validates the complete 143-row migration matrix without issues", () => {
const discovered = discoverLegacyPages();
const issues = validateMigrationEntries(
discovered,
HOUSEKEEPING_MIGRATION_MATRIX,
);
expect(HOUSEKEEPING_MIGRATION_MATRIX).toHaveLength(142);
expect(discovered).toHaveLength(142);
expect(HOUSEKEEPING_MIGRATION_MATRIX).toHaveLength(143);
expect(discovered).toHaveLength(143);
expect(issues).toEqual([]);
});
});
@@ -29,7 +29,7 @@ describe("discoverLegacyPages", () => {
it("discovers the exact legacy administration inventory", () => {
const pages = discoverLegacyPages();
expect(pages).toHaveLength(142);
expect(pages).toHaveLength(143);
expect(pages).toContainEqual({
surface: "admin",
legacyPath: "/admin/users/:id/edit",
@@ -4,10 +4,10 @@ import { HOUSEKEEPING_MIGRATION_MATRIX } from "./matrix";
import { validateMigrationEntries } from "./validate-matrix";
describe("HOUSEKEEPING_MIGRATION_MATRIX", () => {
it("covers all 142 legacy pages exactly once", () => {
it("covers all 143 legacy pages exactly once", () => {
const discovered = discoverLegacyPages();
expect(HOUSEKEEPING_MIGRATION_MATRIX).toHaveLength(142);
expect(HOUSEKEEPING_MIGRATION_MATRIX).toHaveLength(143);
expect(
validateMigrationEntries(discovered, HOUSEKEEPING_MIGRATION_MATRIX),
).toEqual([]);
@@ -18,8 +18,8 @@ const SYSTEM_PREFIXES = [
] as const;
describe("systemMigrationEntries", () => {
it("covers all 22 System pages exactly once", () => {
expect(systemMigrationEntries).toHaveLength(22);
it("covers all 23 System pages exactly once", () => {
expect(systemMigrationEntries).toHaveLength(23);
expect(
validateMigrationEntries(
ownedLegacyPages(SYSTEM_PREFIXES),
@@ -58,6 +58,24 @@ export const systemMigrationEntries: readonly MigrationEntry[] = [
localization: "PARTIAL",
accessibility: "PARTIAL",
}),
plannedSystemEntry({
surface: "admin",
legacyPath: "/admin/devops/deliveries",
sourceFile: "src/app/admin/devops/deliveries/page.tsx",
targetPath: "/admin/system/operations/deliveries",
decision: "REHOST",
capabilities: { read: [PERMS.DEVOPS_VIEW], mutate: [PERMS.DEVOPS_EDIT] },
dependencies: {
queries: ["cms_operations", "cms_outbox"],
mutations: ["retryDelivery"],
},
auditRequirement: "PRIVILEGED_MUTATION",
localization: "PARTIAL",
accessibility: "PARTIAL",
notes: [
"Retries failed effects only; does not replay committed catalog mutations",
],
}),
// Operator alerts and read-only analytics.
plannedSystemEntry({
surface: "admin",
+9
View File
@@ -0,0 +1,9 @@
# Durable catalog operations
Bulk offer apply and undo accept a request UUID. The authenticated actor, operation kind and UUID identify one mutation. A canonical payload hash rejects reuse for different changes. The result and outbox entries commit in the same database transaction as the offers and audit history; replay returns the stored result.
The jobs worker claims pending effects under a database lock with a 120-second lease. A claim token protects completion from stale workers. Failed delivery uses exponential retry, stopping after eight attempts. `/admin/devops/deliveries` shows the latest 100 effects to DEVOPS_VIEW; DEVOPS_EDIT can retry failed effects without replaying the catalog mutation.
Delivery is at least once: a crash after sending and before acknowledgement may repeat an effect. Catalog refresh is repeatable; export delivery means a request entered the existing export queue, not that Git publication or client refresh completed. Export disabled by configuration remains a no-op. This first adapter covers bulk offer apply/undo, not every CMS mutation.
Migration0031 is required before the updated worker starts. Unit tests cover payload identity, dispatch limits, retries and authorization; the separate Docker integration suite covers concurrent SQL requests, transaction rollback and claim ownership. No production database was used for local tests.
@@ -0,0 +1,51 @@
import { expect, it, vi } from "vitest";
import { dispatchEffects } from "./dispatcher";
it("settles only after successful delivery and preserves the claim token", async () => {
const effect = {
id: "a",
token: "lease1",
topic: "catalog.refresh" as const,
attempts: 1,
};
const repo = {
claim: vi.fn().mockResolvedValueOnce(effect).mockResolvedValue(null),
complete: vi.fn(),
fail: vi.fn(),
};
const deliver = vi.fn();
await dispatchEffects(repo, deliver);
expect(deliver).toHaveBeenCalledWith(effect);
expect(repo.complete).toHaveBeenCalledWith(effect);
expect(repo.fail).not.toHaveBeenCalled();
});
it("persists failures and continues to other pending work", async () => {
const effect = {
id: "a",
token: "lease1",
topic: "catalog.refresh" as const,
attempts: 1,
};
const repo = {
claim: vi.fn().mockResolvedValueOnce(effect).mockResolvedValue(null),
complete: vi.fn(),
fail: vi.fn(),
};
await dispatchEffects(repo, vi.fn().mockRejectedValue(Error("private")));
expect(repo.fail).toHaveBeenCalledWith(effect);
expect(repo.complete).not.toHaveBeenCalled();
});
it("bounds work per tick instead of draining indefinitely", async () => {
const repo = {
claim: vi.fn().mockResolvedValue({
id: "a",
token: "t",
topic: "catalog.refresh",
attempts: 1,
}),
complete: vi.fn(),
fail: vi.fn(),
};
await dispatchEffects(repo, async () => {});
expect(repo.claim).toHaveBeenCalledTimes(20);
});
+23
View File
@@ -0,0 +1,23 @@
import type { EffectClaim } from "./model";
export interface EffectRepository {
claim(): Promise<EffectClaim | null>;
complete(claim: EffectClaim): Promise<unknown>;
fail(claim: EffectClaim): Promise<unknown>;
}
/** At-least-once delivery: handlers must tolerate replay after an expired claim. */
export async function dispatchEffects(
repo: EffectRepository,
deliver: (claim: EffectClaim) => Promise<unknown>,
) {
for (let i = 0; i < 20; i++) {
const claim = await repo.claim();
if (!claim) return;
try {
await deliver(claim);
} catch {
await repo.fail(claim);
continue;
}
await repo.complete(claim);
}
}
+31
View File
@@ -0,0 +1,31 @@
import { expect, it } from "vitest";
import { operationHash, retryDelay, validateOperation } from "./model";
it("hashes equivalent object payloads equally without ignoring array order", () => {
expect(operationHash({ b: 2, a: { y: 1, x: 0 } })).toBe(
operationHash({ a: { x: 0, y: 1 }, b: 2 }),
);
expect(operationHash([1, 2])).not.toBe(operationHash([2, 1]));
});
it("rejects unsupported or oversized input rather than hashing ambiguous values", () => {
expect(() => operationHash({ amount: NaN })).toThrow();
expect(() => operationHash({ value: undefined })).toThrow();
expect(() => operationHash("x".repeat(100000))).toThrow();
});
it("validates actor and request identity before database access", () => {
const valid = {
actorId: 1,
kind: "catalog.bulk.apply",
key: "123e4567-e89b-42d3-a456-426614174000",
input: {},
};
expect(() => validateOperation(valid)).not.toThrow();
expect(() => validateOperation({ ...valid, actorId: 0 })).toThrow();
expect(() => validateOperation({ ...valid, key: "shared" })).toThrow();
expect(() => validateOperation({ ...valid, kind: "bad kind" })).toThrow();
});
it("backs off boundedly for repeated delivery failures", () => {
expect(retryDelay(1)).toBe(60);
expect(retryDelay(2)).toBe(120);
expect(retryDelay(20)).toBe(3600);
});
+62
View File
@@ -0,0 +1,62 @@
import { createHash } from "node:crypto";
export interface OperationInput {
actorId: number;
kind: string;
key: string;
input: unknown;
}
export class OperationConflict extends Error {
constructor() {
super(
"This request key was already used for different changes. Refresh the preview.",
);
}
}
export function validateOperation(value: OperationInput) {
if (
!Number.isSafeInteger(value.actorId) ||
value.actorId < 1 ||
!/^[a-z][a-z0-9.-]{2,63}$/.test(value.kind) ||
!/^[a-f0-9]{8}-[a-f0-9]{4}-[a-f0-9]{4}-[a-f0-9]{4}-[a-f0-9]{12}$/i.test(
value.key,
)
)
throw new Error("Invalid operation identity");
}
function canonical(value: unknown, depth = 0): unknown {
if (depth > 30) throw new Error("Operation input is too deep");
if (value === null || typeof value === "string" || typeof value === "boolean")
return value;
if (typeof value === "number" && Number.isFinite(value)) return value;
if (Array.isArray(value)) return value.map((v) => canonical(v, depth + 1));
if (
value &&
typeof value === "object" &&
Object.getPrototypeOf(value) === Object.prototype
)
return Object.fromEntries(
Object.keys(value)
.sort()
.map((key) => [
key,
canonical((value as Record<string, unknown>)[key], depth + 1),
]),
);
throw new Error("Unsupported operation input");
}
export function operationHash(value: unknown) {
const serialized = JSON.stringify(canonical(value));
if (Buffer.byteLength(serialized) > 65536)
throw new Error("Operation input is too large");
return createHash("sha256").update(serialized).digest("hex");
}
export function retryDelay(attempt: number) {
return Math.min(3600, 60 * 2 ** Math.max(0, attempt - 1));
}
export type EffectTopic = "catalog.refresh" | "catalog.export.request";
export interface EffectClaim {
id: string;
token: string;
topic: EffectTopic;
attempts: number;
}
@@ -0,0 +1,51 @@
import { beforeEach, expect, it, vi } from "vitest";
const mocks = vi.hoisted(() => ({
permission: vi.fn(),
retry: vi.fn(),
audit: vi.fn(),
refresh: vi.fn(),
error: vi.fn(),
}));
vi.mock("next/cache", () => ({ revalidatePath: mocks.refresh }));
vi.mock("@/lib/admin/guard", () => ({ requirePermission: mocks.permission }));
vi.mock("@/lib/logger", () => ({ logger: { error: mocks.error } }));
vi.mock("@/lib/services/audit", () => ({ logAudit: mocks.audit }));
vi.mock("./server", () => ({ retryEffect: mocks.retry }));
import { retryDelivery } from "./retry-action";
const id = "840f4023-04cd-4d7b-92c5-047aa620029c";
beforeEach(() => {
vi.resetAllMocks();
mocks.permission.mockResolvedValue({ id: 7 });
mocks.retry.mockResolvedValue(undefined);
mocks.audit.mockResolvedValue(undefined);
});
it("rejects unauthorized access before touching the delivery", async () => {
mocks.permission.mockRejectedValue(new Error("Forbidden"));
await expect(retryDelivery(id)).rejects.toThrow("Forbidden");
expect(mocks.retry).not.toHaveBeenCalled();
expect(mocks.audit).not.toHaveBeenCalled();
});
it("reports a rejected retry without recording a successful action", async () => {
mocks.retry.mockRejectedValue(new Error("Database unavailable"));
expect(await retryDelivery(id)).toMatchObject({ ok: false });
expect(mocks.audit).not.toHaveBeenCalled();
expect(mocks.refresh).not.toHaveBeenCalled();
});
it("records the delivery UUID as audit metadata", async () => {
expect(await retryDelivery(id)).toMatchObject({ ok: true });
expect(mocks.audit).toHaveBeenCalledWith(
expect.objectContaining({ userId: 7, after: { deliveryId: id } }),
);
});
it("does not report a committed retry as failed when audit and refresh fail", async () => {
mocks.audit.mockRejectedValue(new Error("Audit unavailable"));
mocks.refresh.mockImplementation(() => {
throw new Error("Refresh unavailable");
});
expect(await retryDelivery(id)).toMatchObject({ ok: true });
expect(mocks.retry).toHaveBeenCalledTimes(1);
expect(mocks.error).toHaveBeenCalled();
});
+35
View File
@@ -0,0 +1,35 @@
"use server";
import { revalidatePath } from "next/cache";
import { requirePermission } from "@/lib/admin/guard";
import { logger } from "@/lib/logger";
import { PERMS } from "@/lib/permission-slugs";
import { logAudit } from "@/lib/services/audit";
import { retryEffect } from "./server";
export async function retryDelivery(id: string) {
const staff = await requirePermission(PERMS.DEVOPS_EDIT);
try {
await retryEffect(id);
} catch {
return {
ok: false as const,
error: "Delivery could not be retried. Refresh and try again.",
};
}
await logAudit({
userId: staff.id,
action: "operation_delivery_retry",
target: "operations",
after: { deliveryId: id },
}).catch((error) =>
logger.error("Delivery retried; audit failed", {
module: "operations",
error,
}),
);
try {
revalidatePath("/admin/devops/deliveries");
} catch {
/* Delivery retry remains saved. */
}
return { ok: true as const, data: {} };
}
+23
View File
@@ -0,0 +1,23 @@
"use client";
import { useTranslations } from "next-intl";
import { Button } from "@/components/ui/button";
import { useServerAction } from "@/hooks/use-server-action";
import { retryDelivery } from "./retry-action";
export function DeliveryRetry({ id }: { id: string }) {
const t = useTranslations("pages.admin.deliveries");
const { run, isPending } = useServerAction();
return (
<Button
variant="outline"
disabled={isPending}
onClick={() =>
run(() => retryDelivery(id), {
successMessage: t("retried"),
errorMessage: t("retryError"),
})
}
>
{t("retry")}
</Button>
);
}
@@ -0,0 +1,23 @@
import { beforeEach, expect, it, vi } from "vitest";
const execute = vi.hoisted(() => vi.fn());
vi.mock("@/lib/db", () => ({ db: { execute } }));
import { retryEffect } from "./server";
const id = "840f4023-04cd-4d7b-92c5-047aa620029c";
beforeEach(() => vi.resetAllMocks());
it("does not report a missing or no-longer-failed delivery as retried", async () => {
execute.mockResolvedValue([{ affectedRows: 0 }]);
await expect(retryEffect(id)).rejects.toThrow("no longer available");
});
it("accepts a delivery that was actually moved back to pending", async () => {
execute.mockResolvedValue([{ affectedRows: 1 }]);
await expect(retryEffect(id)).resolves.toBeUndefined();
});
it("rejects malformed identifiers before a database query", async () => {
await expect(
retryEffect("------------------------------------"),
).rejects.toThrow("Invalid delivery");
expect(execute).not.toHaveBeenCalled();
});
+117
View File
@@ -0,0 +1,117 @@
import "server-only";
import { randomUUID } from "node:crypto";
import { sql } from "drizzle-orm";
import { db } from "@/lib/db";
import {
type EffectClaim,
type EffectTopic,
OperationConflict,
type OperationInput,
operationHash,
retryDelay,
validateOperation,
} from "./model";
export type OperationTransaction = Parameters<
Parameters<typeof db.transaction>[0]
>[0];
export async function runOperation<T>(
input: OperationInput,
work: (tx: OperationTransaction, operationId: string) => Promise<T>,
): Promise<T> {
validateOperation(input);
const hash = operationHash(input.input);
return db.transaction(async (tx) => {
await tx.execute(
sql`INSERT INTO cms_operations (id,actor_id,kind,request_key,request_hash) VALUES (${randomUUID()},${input.actorId},${input.kind},${input.key},${hash}) ON DUPLICATE KEY UPDATE id=id`,
);
const [rows] = await tx.execute(
sql`SELECT id,request_hash AS requestHash,result_json AS resultJson FROM cms_operations WHERE actor_id=${input.actorId} AND kind=${input.kind} AND request_key=${input.key} FOR UPDATE`,
);
const row = (
rows as unknown as Array<{
id: string;
requestHash: string;
resultJson: string | null;
}>
)[0];
if (!row) throw new Error("Operation unavailable");
if (row.requestHash !== hash) throw new OperationConflict();
if (row.resultJson !== null) return JSON.parse(row.resultJson) as T;
const result = await work(tx, row.id);
const serialized = JSON.stringify(result);
if (!serialized || Buffer.byteLength(serialized) > 1048576)
throw new Error("Operation result exceeds storage limit");
await tx.execute(
sql`UPDATE cms_operations SET result_json=${serialized} WHERE id=${row.id}`,
);
return result;
});
}
export async function enqueueEffect(
tx: OperationTransaction,
operationId: string,
topic: EffectTopic,
) {
await tx.execute(
sql`INSERT INTO cms_outbox (id,operation_id,topic) VALUES (${randomUUID()},${operationId},${topic}) ON DUPLICATE KEY UPDATE id=id`,
);
}
export const effectRepository = {
async claim(): Promise<EffectClaim | null> {
return db.transaction(async (tx) => {
const [rows] = await tx.execute(
sql`SELECT id,topic,attempts FROM cms_outbox WHERE (status='pending' AND available_at<=UTC_TIMESTAMP(3)) OR (status='running' AND lease_until<UTC_TIMESTAMP(3)) ORDER BY available_at,id LIMIT 1 FOR UPDATE`,
);
const row = (
rows as unknown as Array<{
id: string;
topic: EffectTopic;
attempts: number;
}>
)[0];
if (!row) return null;
const token = randomUUID();
await tx.execute(
sql`UPDATE cms_outbox SET status='running',lease_token=${token},lease_until=DATE_ADD(UTC_TIMESTAMP(3),INTERVAL 120 SECOND),attempts=attempts+1 WHERE id=${row.id}`,
);
return { ...row, token, attempts: Number(row.attempts) + 1 };
});
},
async complete(claim: EffectClaim) {
await db.execute(
sql`UPDATE cms_outbox SET status='done',lease_token=NULL,lease_until=NULL,last_error=NULL WHERE id=${claim.id} AND status='running' AND lease_token=${claim.token}`,
);
},
async fail(claim: EffectClaim) {
await db.execute(
sql`UPDATE cms_outbox SET status=${claim.attempts >= 8 ? "failed" : "pending"},available_at=DATE_ADD(UTC_TIMESTAMP(3),INTERVAL ${retryDelay(claim.attempts)} SECOND),lease_token=NULL,lease_until=NULL,last_error='Delivery failed; inspect the linked operation and service diagnostics.' WHERE id=${claim.id} AND status='running' AND lease_token=${claim.token}`,
);
},
};
export async function listEffects() {
const [rows] = await db.execute(
sql`SELECT e.id,e.operation_id AS operationId,e.topic,e.status,e.attempts,e.last_error AS lastError,e.created_at AS createdAt,o.actor_id AS actorId,o.kind FROM cms_outbox e JOIN cms_operations o ON o.id=e.operation_id ORDER BY e.created_at DESC,e.id DESC LIMIT 100`,
);
return rows as unknown as Array<{
id: string;
operationId: string;
topic: EffectTopic;
status: string;
attempts: number;
lastError: string | null;
createdAt: string;
actorId: number;
kind: string;
}>;
}
export async function retryEffect(id: string) {
if (
!/^[a-f0-9]{8}-[a-f0-9]{4}-[a-f0-9]{4}-[a-f0-9]{4}-[a-f0-9]{12}$/i.test(id)
)
throw new Error("Invalid delivery");
const [result] = await db.execute(
sql`UPDATE cms_outbox SET status='pending',attempts=0,available_at=UTC_TIMESTAMP(3),last_error=NULL WHERE id=${id} AND status='failed'`,
);
if ((result as unknown as { affectedRows: number }).affectedRows !== 1)
throw new Error("Delivery is no longer available for retry");
}
+32
View File
@@ -0,0 +1,32 @@
import "server-only";
import { sendCatalogUpdate } from "@/features/catalog/server/sync-status";
import { logger } from "@/lib/logger";
import {
catalogExportEnabled,
catalogExportQueue,
} from "@/lib/services/catalog-git-queue";
import { dispatchEffects } from "./dispatcher";
import { effectRepository } from "./server";
let running = false;
export async function drainOperationEffects() {
if (running) return;
running = true;
try {
await dispatchEffects(effectRepository, async (claim) => {
if (claim.topic === "catalog.refresh") {
if (!(await sendCatalogUpdate()).sent)
throw new Error("Hotel update not delivered");
} else if (claim.topic === "catalog.export.request") {
if (catalogExportEnabled()) await catalogExportQueue().request();
} else throw new Error("Unknown delivery topic");
});
} catch (error) {
logger.error("Operation delivery tick failed", {
module: "operations",
error,
});
} finally {
running = false;
}
}