Files
EpicNext-Cms/src/features/housekeeping/foundation/commands/dispatcher.ts
T

522 lines
12 KiB
TypeScript

import { isIP } from "node:net";
import { z } from "zod";
import type { AuditEntry, HousekeepingAuditWriter } from "@/lib/services/audit";
import { satisfiesCapability } from "../capability-context";
import {
fail,
type HousekeepingCapabilityContext,
type HousekeepingPartialCompletion,
type HousekeepingResult,
mapUnknownError,
} from "../contracts";
import { createCorrelationId } from "../correlation";
import { writeIntent, writeOutcome } from "./audit-envelope";
import { confirmHousekeepingCommand } from "./confirmation";
import {
getHousekeepingCommand,
type HousekeepingCommand,
type HousekeepingCommandContext,
isHousekeepingCommandId,
} from "./registry";
const normalizedCommandId = z
.string()
.transform((value) => value.normalize("NFC").trim())
.pipe(
z
.string()
.min(1)
.max(160)
.refine(isHousekeepingCommandId, "invalid command id"),
);
const normalizedReason = z
.string()
.transform((value) => value.normalize("NFC").trim())
.pipe(z.string().max(1000));
const commandRequestSchema = z.strictObject({
commandId: normalizedCommandId,
input: z.unknown(),
reason: normalizedReason.optional(),
});
const housekeepingErrorCodeSchema = z.enum([
"UNAUTHENTICATED",
"FORBIDDEN",
"VALIDATION",
"NOT_FOUND",
"CONFLICT",
"RATE_LIMITED",
"DEPENDENCY_UNAVAILABLE",
"TIMEOUT",
"INTERNAL",
]);
const housekeepingErrorSchema = z.strictObject({
code: housekeepingErrorCodeSchema,
messageKey: z.string().min(1).max(160),
fieldErrors: z.record(z.string(), z.array(z.string())).optional(),
});
const housekeepingResultSchema = z.discriminatedUnion("ok", [
z.strictObject({
ok: z.literal(true),
data: z.unknown(),
correlationId: z.string().min(1).max(160),
completion: z
.strictObject({
status: z.literal("partial"),
external: z.enum(["not-required", "completed", "failed"]),
audit: z.enum(["persisted", "unavailable"]),
})
.optional(),
}),
z.strictObject({
ok: z.literal(false),
error: housekeepingErrorSchema,
correlationId: z.string().min(1).max(160),
completion: z
.strictObject({
status: z.literal("partial"),
external: z.enum(["not-required", "completed", "failed"]),
audit: z.enum(["persisted", "unavailable"]),
})
.optional(),
}),
]);
export interface HousekeepingCommandRequest {
commandId: string;
input: unknown;
reason?: string;
}
export interface HousekeepingCommandDependencies {
context: HousekeepingCapabilityContext;
ipAddress: string;
audit: HousekeepingAuditWriter;
rateLimit(key: string, attempts: number, windowMs: number): Promise<boolean>;
}
export async function dispatchHousekeepingCommand(
request: unknown,
dependencies: HousekeepingCommandDependencies,
): Promise<HousekeepingResult<unknown>> {
const correlationId = createCorrelationId();
const ipAddress = canonicalizeServerIp(dependencies.ipAddress);
if (ipAddress === undefined) {
return persistReturnedOutcome(
dependencies.audit,
createDispatchAuditEntry(
"server-context",
dependencies.context,
correlationId,
),
"failure",
mapUnknownError(new Error("invalid server IP"), correlationId),
);
}
if (!isPlainRecord(request)) {
return persistMalformedRequestOutcome(
dependencies,
ipAddress,
correlationId,
);
}
let parsedRequest: ReturnType<typeof commandRequestSchema.safeParse>;
try {
parsedRequest = commandRequestSchema.safeParse(request);
} catch (error) {
return persistReturnedOutcome(
dependencies.audit,
createDispatchAuditEntry(
"request-envelope",
dependencies.context,
correlationId,
ipAddress,
),
"failure",
mapUnknownError(error, correlationId),
);
}
if (!parsedRequest.success) {
return persistMalformedRequestOutcome(
dependencies,
ipAddress,
correlationId,
);
}
const commandRequest = parsedRequest.data;
const command = getHousekeepingCommand(commandRequest.commandId);
if (!command) {
const result = fail(
"NOT_FOUND",
"errors.housekeeping.notFound",
correlationId,
);
return persistReturnedOutcome(
dependencies.audit,
createUnknownCommandAuditEntry(
commandRequest.commandId,
dependencies.context,
ipAddress,
correlationId,
),
"denied",
result,
);
}
const auditEntry = createAuditEntry(
command,
dependencies.context,
ipAddress,
correlationId,
commandRequest.reason || undefined,
);
let authorized: boolean;
try {
authorized = satisfiesCapability(dependencies.context, command.capability);
} catch (error) {
return persistReturnedOutcome(
dependencies.audit,
auditEntry,
"failure",
mapUnknownError(error, correlationId),
);
}
if (!authorized) {
return persistReturnedOutcome(
dependencies.audit,
auditEntry,
"denied",
fail("FORBIDDEN", "errors.housekeeping.forbidden", correlationId),
);
}
let parsedInput: ReturnType<typeof command.input.safeParse>;
try {
parsedInput = command.input.safeParse(commandRequest.input);
} catch (error) {
return persistReturnedOutcome(
dependencies.audit,
auditEntry,
"failure",
mapUnknownError(error, correlationId),
);
}
if (!parsedInput.success) {
const fieldErrors: Record<string, readonly string[]> = {};
for (const issue of parsedInput.error.issues) {
const field = String(issue.path[0] ?? "input");
fieldErrors[field] = ["errors.validation.invalid"];
}
return persistReturnedOutcome(
dependencies.audit,
auditEntry,
"denied",
fail(
"VALIDATION",
"errors.housekeeping.validation",
correlationId,
fieldErrors,
),
);
}
let confirmation: HousekeepingResult<unknown>;
try {
confirmation = confirmHousekeepingCommand(
command,
commandRequest.reason || undefined,
correlationId,
);
} catch (error) {
return persistReturnedOutcome(
dependencies.audit,
auditEntry,
"failure",
mapUnknownError(error, correlationId),
);
}
if (!confirmation.ok) {
return persistReturnedOutcome(
dependencies.audit,
auditEntry,
"denied",
confirmation,
);
}
let allowed: boolean;
try {
allowed = await dependencies.rateLimit(
createRateLimitKey(command.id, dependencies.context, ipAddress),
command.rateLimit.attempts,
command.rateLimit.windowMs,
);
} catch (error) {
return persistReturnedOutcome(
dependencies.audit,
auditEntry,
"failure",
mapUnknownError(error, correlationId),
);
}
if (!allowed) {
return persistReturnedOutcome(
dependencies.audit,
auditEntry,
"denied",
fail("RATE_LIMITED", "errors.housekeeping.rateLimited", correlationId),
);
}
if (command.risk === "sensitive") {
try {
await writeIntent(dependencies.audit, auditEntry);
} catch (error) {
return mapUnknownError(error, correlationId);
}
}
const commandContext: HousekeepingCommandContext = {
capability: dependencies.context,
correlationId,
ipAddress,
};
let correlatedResult: HousekeepingResult<unknown>;
try {
const rawResult = await command.execute(commandContext, parsedInput.data);
correlatedResult = validateAndCorrelateResult(rawResult, correlationId);
} catch (error) {
return persistReturnedOutcome(
dependencies.audit,
auditEntry,
"failure",
mapUnknownError(error, correlationId),
);
}
if (!correlatedResult.ok) {
return persistReturnedOutcome(
dependencies.audit,
auditEntry,
outcomeForFailure(correlatedResult),
correlatedResult,
);
}
if (correlatedResult.completion?.status === "partial") {
return persistPartialCompletionOutcome(
dependencies.audit,
auditEntry,
correlatedResult,
);
}
return persistSuccessfulOutcome(
dependencies.audit,
auditEntry,
correlatedResult,
);
}
function createAuditEntry(
command: HousekeepingCommand<unknown, unknown>,
context: HousekeepingCapabilityContext,
ipAddress: string,
correlationId: string,
reason: string | undefined,
): AuditEntry {
return {
userId: context.actor.id,
action: command.id,
target: command.owner,
correlationId,
domain: command.owner,
...(reason === undefined ? {} : { reason }),
ipAddress,
};
}
function createDispatchAuditEntry(
target: "request-envelope" | "server-context",
context: HousekeepingCapabilityContext,
correlationId: string,
ipAddress?: string,
): AuditEntry {
return {
userId: context.actor.id,
action: "housekeeping.command.dispatch",
target,
correlationId,
...(ipAddress === undefined ? {} : { ipAddress }),
};
}
function createUnknownCommandAuditEntry(
commandId: string,
context: HousekeepingCapabilityContext,
ipAddress: string,
correlationId: string,
): AuditEntry {
return {
userId: context.actor.id,
action: "housekeeping.command.dispatch",
target: commandId,
correlationId,
ipAddress,
};
}
function createRateLimitKey(
commandId: string,
context: HousekeepingCapabilityContext,
ipAddress: string,
): string {
return `housekeeping-command:${context.actor.id}:${ipAddress}:${commandId}`;
}
function canonicalizeServerIp(value: unknown): string | undefined {
if (typeof value !== "string") return undefined;
const candidate = value.trim();
const version = isIP(candidate);
if (version === 4) {
return candidate
.split(".")
.map((part) => String(Number(part)))
.join(".");
}
if (version === 6) {
try {
return new URL(`http://[${candidate}]/`).hostname.slice(1, -1);
} catch {
return undefined;
}
}
return undefined;
}
function isPlainRecord(value: unknown): value is Record<string, unknown> {
if (typeof value !== "object" || value === null || Array.isArray(value)) {
return false;
}
try {
const prototype = Object.getPrototypeOf(value);
return prototype === Object.prototype || prototype === null;
} catch {
return false;
}
}
function validateAndCorrelateResult(
value: unknown,
correlationId: string,
): HousekeepingResult<unknown> {
if (!isPlainRecord(value) || !Object.hasOwn(value, "ok")) {
throw new Error("invalid housekeeping command result");
}
if (value.ok === true && !Object.hasOwn(value, "data")) {
throw new Error("invalid housekeeping command result");
}
if (value.ok === false && !Object.hasOwn(value, "error")) {
throw new Error("invalid housekeeping command result");
}
const parsed = housekeepingResultSchema.safeParse(value);
if (!parsed.success) {
throw new Error("invalid housekeeping command result");
}
return { ...parsed.data, correlationId };
}
function outcomeForFailure(
result: Extract<HousekeepingResult<unknown>, { ok: false }>,
): "failure" | "denied" {
return result.error.code === "FORBIDDEN" ||
result.error.code === "VALIDATION" ||
result.error.code === "RATE_LIMITED"
? "denied"
: "failure";
}
function persistMalformedRequestOutcome(
dependencies: HousekeepingCommandDependencies,
ipAddress: string,
correlationId: string,
): Promise<HousekeepingResult<never>> {
return persistReturnedOutcome(
dependencies.audit,
createDispatchAuditEntry(
"request-envelope",
dependencies.context,
correlationId,
ipAddress,
),
"denied",
fail("VALIDATION", "errors.housekeeping.validation", correlationId),
);
}
async function persistReturnedOutcome<T>(
writer: HousekeepingAuditWriter,
entry: AuditEntry,
outcome: "failure" | "denied",
result: HousekeepingResult<T>,
): Promise<HousekeepingResult<T>> {
try {
await writeOutcome(writer, entry, outcome);
} catch {
// The command/preflight failure remains authoritative if evidence is down.
}
return result;
}
async function persistSuccessfulOutcome<T>(
writer: HousekeepingAuditWriter,
entry: AuditEntry,
result: Extract<HousekeepingResult<T>, { ok: true }>,
): Promise<HousekeepingResult<T>> {
try {
await writeOutcome(writer, entry, "success");
return result;
} catch {
let audit: HousekeepingPartialCompletion["audit"] = "persisted";
try {
await writeOutcome(writer, entry, "partial");
} catch {
audit = "unavailable";
}
return withPartialCompletion(result, {
status: "partial",
external: result.completion?.external ?? "not-required",
audit,
});
}
}
async function persistPartialCompletionOutcome<T>(
writer: HousekeepingAuditWriter,
entry: AuditEntry,
result: Extract<HousekeepingResult<T>, { ok: true }>,
): Promise<HousekeepingResult<T>> {
try {
await writeOutcome(writer, entry, "partial");
return result;
} catch {
return withPartialCompletion(result, {
...(result.completion ?? {
status: "partial",
external: "not-required",
audit: "unavailable",
}),
audit: "unavailable",
});
}
}
function withPartialCompletion<T>(
result: Extract<HousekeepingResult<T>, { ok: true }>,
completion: HousekeepingPartialCompletion,
): HousekeepingResult<T> {
return { ...result, completion };
}