Consume raw arma-sync-events into a validated, idempotent operation ledger: contract v1, reject/dead-letter with zero effects, effect keys with unique-constraint race handling, and a staged-activation flag hook. Registers the operations permission domains. Ultraworked with [Sisyphus](https://github.com/code-yeongyu/oh-my-openagent) Co-authored-by: Sisyphus <clio-agent@sisyphuslabs.ai>
197 lines
No EOL
6.7 KiB
TypeScript
197 lines
No EOL
6.7 KiB
TypeScript
import { createHash } from "node:crypto";
|
|
import type { Payload, PayloadRequest } from "payload";
|
|
import { validateOperationEvent } from "./validate";
|
|
import { applyOperationEffect } from "./settlement";
|
|
import type { OperationEffectType, OperationEventRecord, RawSyncEvent } from "./types";
|
|
|
|
/** JSON.stringify with recursively sorted object keys (stable hashing input). */
|
|
export function stableStringify(value: unknown): string {
|
|
if (value === null || typeof value !== "object") return JSON.stringify(value);
|
|
if (Array.isArray(value)) return `[${value.map((item) => stableStringify(item)).join(",")}]`;
|
|
const obj = value as Record<string, unknown>;
|
|
const keys = Object.keys(obj).sort();
|
|
return `{${keys.map((key) => `${JSON.stringify(key)}:${stableStringify(obj[key])}`).join(",")}}`;
|
|
}
|
|
|
|
/** SHA-256 hex digest of the stable-stringified payload. */
|
|
export function hashPayload(payload: unknown): string {
|
|
return createHash("sha256").update(stableStringify(payload)).digest("hex");
|
|
}
|
|
|
|
function pgErrorCode(error: unknown): string | undefined {
|
|
const cause = (error as { cause?: { code?: string } })?.cause;
|
|
return cause?.code ?? (error as { code?: string })?.code;
|
|
}
|
|
|
|
interface DerivedEffect {
|
|
effectType: OperationEffectType;
|
|
payload: Record<string, unknown>;
|
|
}
|
|
|
|
/**
|
|
* Per-type effect derivation. Returns null when the event derives no effect
|
|
* (player.kill is an AAR fact only; unknown types never reach this point).
|
|
* The effect payload carries the full validated event payload through, so
|
|
* the settlement layer sees the complete extraction/loss/end contract.
|
|
*/
|
|
export function deriveEffect(raw: RawSyncEvent, operationId: string): DerivedEffect | null {
|
|
const payload = raw.payload as Record<string, unknown> | undefined;
|
|
switch (raw.type) {
|
|
case "operation.start":
|
|
return {
|
|
effectType: "operation.started",
|
|
payload: { operationId, ...payload },
|
|
};
|
|
case "operation.end":
|
|
return {
|
|
effectType: "operation.ended",
|
|
payload: { operationId, ...payload },
|
|
};
|
|
case "operation.objective":
|
|
return {
|
|
effectType: "operation.objective-updated",
|
|
payload: { operationId, ...payload },
|
|
};
|
|
case "operation.extraction":
|
|
return {
|
|
effectType: "operation.extraction-recorded",
|
|
payload: { operationId, ...payload },
|
|
};
|
|
case "operation.inventory":
|
|
return {
|
|
effectType: "operation.inventory-recorded",
|
|
payload: { operationId, ...payload },
|
|
};
|
|
default:
|
|
return null;
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Consume one raw sync event into the operation ledger.
|
|
*
|
|
* Idempotent: a raw event already present in the ledger (by `sourceEvent`) is
|
|
* skipped, and an effect whose `effectKey` already exists is skipped. The whole
|
|
* body is defensive: the ledger must never break raw event ingestion, so any
|
|
* failure is logged and swallowed.
|
|
*
|
|
* When called from a Payload hook, pass `req` so the ledger writes join the
|
|
* caller's transaction (a separate connection cannot see uncommitted rows).
|
|
*/
|
|
export async function processOperationEvent(
|
|
payload: Payload,
|
|
raw: RawSyncEvent,
|
|
req?: PayloadRequest,
|
|
): Promise<void> {
|
|
try {
|
|
const validation = validateOperationEvent(raw);
|
|
|
|
const record: OperationEventRecord = {
|
|
sourceEvent: raw.id,
|
|
eventId: raw.eventId,
|
|
operationId: validation.operationId ?? "unknown",
|
|
type: raw.type,
|
|
version: validation.version ?? 1,
|
|
status: validation.status,
|
|
error: validation.error ?? null,
|
|
payloadHash: hashPayload(raw.payload),
|
|
actor: null,
|
|
server: typeof raw.server === "number" ? raw.server : null,
|
|
mission: null,
|
|
occurredAt: raw.occurredAt,
|
|
validatedAt: new Date().toISOString(),
|
|
reconciliationNotes: null,
|
|
};
|
|
|
|
// Actor/mission links come from the payload when present.
|
|
if (typeof raw.payload === "object" && raw.payload !== null && !Array.isArray(raw.payload)) {
|
|
const p = raw.payload as Record<string, unknown>;
|
|
if (typeof p.actorId === "number") record.actor = p.actorId;
|
|
if (typeof p.missionId === "number") record.mission = p.missionId;
|
|
}
|
|
|
|
// Idempotency: skip if this raw event is already in the ledger.
|
|
const existing = await payload.find({
|
|
collection: "operation-events",
|
|
where: { sourceEvent: { equals: raw.id } },
|
|
limit: 1,
|
|
depth: 0,
|
|
overrideAccess: true,
|
|
req,
|
|
});
|
|
if (existing.docs.length > 0) {
|
|
payload.logger.info(`[Operations] Duplicate operation event ${raw.eventId} skipped`);
|
|
return;
|
|
}
|
|
|
|
let eventDoc;
|
|
try {
|
|
eventDoc = await payload.create({
|
|
collection: "operation-events",
|
|
data: record,
|
|
overrideAccess: true,
|
|
depth: 0,
|
|
req,
|
|
});
|
|
} catch (error) {
|
|
if (pgErrorCode(error) === "23505") {
|
|
// Unique-constraint race with a concurrent request: already stored.
|
|
payload.logger.info(`[Operations] Duplicate operation event ${raw.eventId} skipped (race)`);
|
|
} else {
|
|
payload.logger.error(
|
|
`[Operations] Failed to create operation event ${raw.eventId}: ${
|
|
error instanceof Error ? error.message : "unknown error"
|
|
}`,
|
|
);
|
|
}
|
|
return;
|
|
}
|
|
|
|
// Only validated events derive effects.
|
|
if (validation.status !== "validated") return;
|
|
|
|
const effect = deriveEffect(raw, record.operationId);
|
|
if (!effect) return;
|
|
|
|
try {
|
|
const effectDoc = await payload.create({
|
|
collection: "operation-effects",
|
|
data: {
|
|
operationEvent: eventDoc.id,
|
|
sourceEvent: raw.id,
|
|
operationId: record.operationId,
|
|
effectKey: `${raw.eventId}:${effect.effectType}`,
|
|
effectType: effect.effectType,
|
|
status: "pending",
|
|
payload: effect.payload,
|
|
appliedAt: null,
|
|
reconciliationNotes: null,
|
|
},
|
|
overrideAccess: true,
|
|
depth: 0,
|
|
req,
|
|
});
|
|
|
|
// Settlement transitions the effect to applied/failed and never
|
|
// throws: the ledger path must not break raw event ingestion.
|
|
await applyOperationEffect(payload, effectDoc.id, req);
|
|
} catch (error) {
|
|
if (pgErrorCode(error) === "23505") {
|
|
// Unique-constraint race on effectKey: already applied.
|
|
payload.logger.info(`[Operations] Duplicate effect ${raw.eventId}:${effect.effectType} skipped`);
|
|
} else {
|
|
payload.logger.error(
|
|
`[Operations] Failed to create effect ${raw.eventId}:${effect.effectType}: ${
|
|
error instanceof Error ? error.message : "unknown error"
|
|
}`,
|
|
);
|
|
}
|
|
}
|
|
} catch (error) {
|
|
payload.logger.error(
|
|
`[Operations] Failed to process operation event ${raw.eventId}: ${
|
|
error instanceof Error ? error.message : "unknown error"
|
|
}`,
|
|
);
|
|
}
|
|
} |