chore(operations): rollout controls, reconciliation, and bridge recovery
Env kill switch with raw-evidence retention, idempotent backfill, in-place dead-letter replay with reconciliation notes, command lease recovery on pull, ack terminal-state guards, and an operations runbook. Ultraworked with [Sisyphus](https://github.com/code-yeongyu/oh-my-openagent) Co-authored-by: Sisyphus <clio-agent@sisyphuslabs.ai>
This commit is contained in:
parent
c19d0d05ee
commit
37d28893de
7 changed files with 737 additions and 20 deletions
75
docs/operations/runbook.md
Normal file
75
docs/operations/runbook.md
Normal file
|
|
@ -0,0 +1,75 @@
|
||||||
|
# Operations runbook
|
||||||
|
|
||||||
|
Operational procedures for the operation ledger, its rollout controls, and the
|
||||||
|
Arma command bridge. Everything here assumes the roadmap in
|
||||||
|
`.omo/plans/ptf-4x-tarkov-milsim-roadmap.md` has shipped.
|
||||||
|
|
||||||
|
## Migration
|
||||||
|
|
||||||
|
- All operation schema lives in named Payload migrations under `src/migrations/`:
|
||||||
|
`20260918_180851_add_operation_ledger`, `20260918_194842_add_operation_allocation`,
|
||||||
|
`20260918_222549_add_operation_readiness`, `20260920_044743_add_zone_pressure`.
|
||||||
|
- Apply with `bun run payload migrate` (dev uses `push: false`; never `bun run db push`,
|
||||||
|
which fails on `spatial_ref_sys` ownership).
|
||||||
|
- Verify status with `bun run payload migrate --status` (or `bun run migrate --status`).
|
||||||
|
- Regenerated artifacts after schema or collection changes: `bun run generate:types`
|
||||||
|
and `bun run generate:importmap`. Never hand-edit `src/payload-types.ts`,
|
||||||
|
`src/payload-generated-schema.ts`, or `src/app/(payload)/admin/importMap.js`.
|
||||||
|
|
||||||
|
## Rollout and rollback-by-disable
|
||||||
|
|
||||||
|
- Kill switch: set `OPERATION_LEDGER_ENABLED=false` in the environment. Raw
|
||||||
|
`arma-sync-events` evidence continues to be stored (that collection's writes
|
||||||
|
never depend on the flag); only the derived ledger (operation events and
|
||||||
|
effects) is skipped. This is the staged-activation control.
|
||||||
|
- Re-enable by removing the flag or setting it to anything other than `"false"`.
|
||||||
|
- Catch-up after a disabled window or a deploy gap: run the backfill action
|
||||||
|
(`reconcileBackfill` in `src/app/(frontend)/operations/actions.ts`, requires
|
||||||
|
`operation-events:update`). It replays raw events with no ledger row;
|
||||||
|
`processOperationEvent` is idempotent by `sourceEvent`, so it is safe to run
|
||||||
|
repeatedly.
|
||||||
|
|
||||||
|
## Duplicate and rejected event recovery
|
||||||
|
|
||||||
|
- Inbound dedup is by `eventId` (`<serverId>:<event.id>`, unique). A duplicate
|
||||||
|
batch is a no-op.
|
||||||
|
- Rejected (`invalid-payload`, `unsupported-version`) and dead-lettered
|
||||||
|
(`unsupported-type`) rows are retained as evidence and derive zero effects.
|
||||||
|
- Replay a single rejected/dead-letter row from the AAR page (Reconciliation
|
||||||
|
panel) or the `reconcileOperationEvent` action. Replay re-validates the
|
||||||
|
immutable raw event and updates the row in place with a reconciliation note;
|
||||||
|
if still invalid it stays rejected and no effects are created.
|
||||||
|
- Raw sync events are never deleted or rewritten, by code or by hand.
|
||||||
|
|
||||||
|
## Command lease expiry and recovery
|
||||||
|
|
||||||
|
- `arma-commands` lifecycle: `queued` -> `delivered` (lease starts, `deliveredAt`
|
||||||
|
set) -> `succeeded`/`failed` (terminal, `completedAt` set).
|
||||||
|
- Lease timeout: 5 minutes (`COMMAND_LEASE_TIMEOUT_MS` in
|
||||||
|
`src/lib/arma-bridge/commands.ts`). Every pull request first re-queues
|
||||||
|
delivered commands whose lease expired without an ack, so a lost command is
|
||||||
|
re-delivered on the next pull.
|
||||||
|
- Acks are guarded: a duplicate ack is an accepted no-op (first result wins);
|
||||||
|
acking a command that was never delivered is rejected with HTTP 409.
|
||||||
|
|
||||||
|
## Inventory and economy repair
|
||||||
|
|
||||||
|
- Extraction deposits and locker credits are preflighted before any write; a
|
||||||
|
full HQ or locker produces a `failed` effect (never a silent loss, never a
|
||||||
|
partial deposit). Repair a failed effect by freeing capacity, then replaying
|
||||||
|
the failed effect through reconciliation after fixing the root cause.
|
||||||
|
- Economy deltas are bounded (faucet/drain caps of 1000) and provenance-marked
|
||||||
|
(`ECONOMY_NOTE_PREFIX`). `resetEconomyState` rebuilds market-state docs and
|
||||||
|
preserves provenance; operation provenance lives in the operation ledger and
|
||||||
|
is never consumed by the reset.
|
||||||
|
- Payments never precede the locker preflight in `buyListing`; a credit failure
|
||||||
|
after payment is compensated by an adjustment refund. Investigate any
|
||||||
|
`adjustment` transaction referencing an operation effect as evidence of that
|
||||||
|
compensation path.
|
||||||
|
|
||||||
|
## Metrics and logging
|
||||||
|
|
||||||
|
- All ledger activity logs under the `[Operations]` prefix via `payload.logger`.
|
||||||
|
- Bridge lease recovery logs under `[ArmaBridge]`.
|
||||||
|
- Ledger rows carry `reconciliationNotes`; every manual replay appends a note
|
||||||
|
with the acting user id and outcome.
|
||||||
|
|
@ -1,6 +1,7 @@
|
||||||
import { NextRequest, NextResponse } from "next/server";
|
import { NextRequest, NextResponse } from "next/server";
|
||||||
import { asJson } from "@/lib/arma-bridge/json";
|
import { asJson } from "@/lib/arma-bridge/json";
|
||||||
import { requireArmaServer } from "@/lib/arma-bridge/server";
|
import { requireArmaServer } from "@/lib/arma-bridge/server";
|
||||||
|
import { acknowledgeCommand } from "@/lib/arma-bridge/commands";
|
||||||
|
|
||||||
export const runtime = "nodejs";
|
export const runtime = "nodejs";
|
||||||
export const dynamic = "force-dynamic";
|
export const dynamic = "force-dynamic";
|
||||||
|
|
@ -17,26 +18,20 @@ export async function POST(req: NextRequest) {
|
||||||
return NextResponse.json({ error: "id and success are required." }, { status: 400 });
|
return NextResponse.json({ error: "id and success are required." }, { status: 400 });
|
||||||
}
|
}
|
||||||
|
|
||||||
const commands = await payload.find({
|
const outcome = await acknowledgeCommand(payload, server.id, {
|
||||||
collection: "arma-commands",
|
id: body.id,
|
||||||
where: { and: [{ commandId: { equals: body.id } }, { server: { equals: server.id } }] },
|
success: body.success,
|
||||||
limit: 1,
|
result: asJson(body.result, null),
|
||||||
overrideAccess: true,
|
error: typeof body.error === "string" ? body.error : undefined,
|
||||||
});
|
|
||||||
const command = commands.docs[0];
|
|
||||||
if (!command) return NextResponse.json({ error: "Command not found." }, { status: 404 });
|
|
||||||
|
|
||||||
await payload.update({
|
|
||||||
collection: "arma-commands",
|
|
||||||
id: command.id,
|
|
||||||
data: {
|
|
||||||
status: body.success ? "succeeded" : "failed",
|
|
||||||
result: asJson(body.result),
|
|
||||||
error: typeof body.error === "string" ? body.error : undefined,
|
|
||||||
completedAt: new Date().toISOString(),
|
|
||||||
},
|
|
||||||
overrideAccess: true,
|
|
||||||
});
|
});
|
||||||
|
|
||||||
return NextResponse.json({ ok: true });
|
if (!outcome.ok) {
|
||||||
|
const status = outcome.error === "not-found" ? 404 : 409;
|
||||||
|
const message =
|
||||||
|
outcome.error === "not-found"
|
||||||
|
? "Command not found."
|
||||||
|
: "Command was not delivered to this server; ack rejected.";
|
||||||
|
return NextResponse.json({ error: message }, { status });
|
||||||
|
}
|
||||||
|
return NextResponse.json({ ok: true, duplicate: outcome.duplicate === true });
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -1,5 +1,6 @@
|
||||||
import { NextRequest, NextResponse } from "next/server";
|
import { NextRequest, NextResponse } from "next/server";
|
||||||
import { requireArmaServer } from "@/lib/arma-bridge/server";
|
import { requireArmaServer } from "@/lib/arma-bridge/server";
|
||||||
|
import { recoverExpiredCommands } from "@/lib/arma-bridge/commands";
|
||||||
|
|
||||||
export const runtime = "nodejs";
|
export const runtime = "nodejs";
|
||||||
export const dynamic = "force-dynamic";
|
export const dynamic = "force-dynamic";
|
||||||
|
|
@ -9,6 +10,9 @@ export async function POST(req: NextRequest) {
|
||||||
if (auth instanceof NextResponse) return auth;
|
if (auth instanceof NextResponse) return auth;
|
||||||
const { payload, server } = auth;
|
const { payload, server } = auth;
|
||||||
|
|
||||||
|
// Expired leases (delivered, never acked) rejoin the queued pool.
|
||||||
|
await recoverExpiredCommands(payload, server.id);
|
||||||
|
|
||||||
const commands = await payload.find({
|
const commands = await payload.find({
|
||||||
collection: "arma-commands",
|
collection: "arma-commands",
|
||||||
where: { and: [{ server: { equals: server.id } }, { status: { equals: "queued" } }] },
|
where: { and: [{ server: { equals: server.id } }, { status: { equals: "queued" } }] },
|
||||||
|
|
|
||||||
104
src/lib/arma-bridge/commands.ts
Normal file
104
src/lib/arma-bridge/commands.ts
Normal file
|
|
@ -0,0 +1,104 @@
|
||||||
|
import type { Payload } from "payload";
|
||||||
|
import { asJson } from "./json";
|
||||||
|
|
||||||
|
type PayloadType = Awaited<ReturnType<typeof import("payload").getPayload>>;
|
||||||
|
|
||||||
|
/** A delivered command older than this is treated as lost and re-queued. */
|
||||||
|
export const COMMAND_LEASE_TIMEOUT_MS = 5 * 60_000;
|
||||||
|
|
||||||
|
export interface AckInput {
|
||||||
|
id: string;
|
||||||
|
success: boolean;
|
||||||
|
result?: unknown;
|
||||||
|
error?: string;
|
||||||
|
}
|
||||||
|
|
||||||
|
export interface AckOutcome {
|
||||||
|
ok: boolean;
|
||||||
|
duplicate?: boolean;
|
||||||
|
error?: "not-found" | "not-delivered";
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Re-queue delivered commands whose lease expired without an ack. Returns the
|
||||||
|
* number of commands recovered.
|
||||||
|
*/
|
||||||
|
export async function recoverExpiredCommands(
|
||||||
|
payload: PayloadType,
|
||||||
|
serverId: number,
|
||||||
|
): Promise<number> {
|
||||||
|
const cutoff = new Date(Date.now() - COMMAND_LEASE_TIMEOUT_MS).toISOString();
|
||||||
|
const expired = await payload.find({
|
||||||
|
collection: "arma-commands",
|
||||||
|
where: {
|
||||||
|
and: [
|
||||||
|
{ server: { equals: serverId } },
|
||||||
|
{ status: { equals: "delivered" } },
|
||||||
|
{ deliveredAt: { less_than: cutoff } },
|
||||||
|
],
|
||||||
|
},
|
||||||
|
limit: 100,
|
||||||
|
depth: 0,
|
||||||
|
overrideAccess: true,
|
||||||
|
});
|
||||||
|
|
||||||
|
for (const command of expired.docs) {
|
||||||
|
await payload.update({
|
||||||
|
collection: "arma-commands",
|
||||||
|
id: command.id,
|
||||||
|
data: { status: "queued", deliveredAt: null },
|
||||||
|
overrideAccess: true,
|
||||||
|
depth: 0,
|
||||||
|
});
|
||||||
|
}
|
||||||
|
if (expired.docs.length > 0) {
|
||||||
|
payload.logger.warn(
|
||||||
|
`[ArmaBridge] Recovered ${expired.docs.length} expired command lease(s) for server ${serverId}`,
|
||||||
|
);
|
||||||
|
}
|
||||||
|
return expired.docs.length;
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Apply a command acknowledgement with terminal-state and delivery guards.
|
||||||
|
* Duplicate acks are accepted as no-ops (the first result wins); acking a
|
||||||
|
* queued (never delivered) command is rejected.
|
||||||
|
*/
|
||||||
|
export async function acknowledgeCommand(
|
||||||
|
payload: PayloadType,
|
||||||
|
serverId: number,
|
||||||
|
input: AckInput,
|
||||||
|
): Promise<AckOutcome> {
|
||||||
|
const found = await payload.find({
|
||||||
|
collection: "arma-commands",
|
||||||
|
where: {
|
||||||
|
and: [{ commandId: { equals: input.id } }, { server: { equals: serverId } }],
|
||||||
|
},
|
||||||
|
limit: 1,
|
||||||
|
depth: 0,
|
||||||
|
overrideAccess: true,
|
||||||
|
});
|
||||||
|
const command = found.docs[0] as { id: number; status: string } | undefined;
|
||||||
|
if (!command) return { ok: false, error: "not-found" };
|
||||||
|
|
||||||
|
if (command.status === "succeeded" || command.status === "failed") {
|
||||||
|
return { ok: true, duplicate: true };
|
||||||
|
}
|
||||||
|
if (command.status !== "delivered") {
|
||||||
|
return { ok: false, error: "not-delivered" };
|
||||||
|
}
|
||||||
|
|
||||||
|
await payload.update({
|
||||||
|
collection: "arma-commands",
|
||||||
|
id: command.id,
|
||||||
|
data: {
|
||||||
|
status: input.success ? "succeeded" : "failed",
|
||||||
|
result: asJson(input.result, null),
|
||||||
|
error: input.error,
|
||||||
|
completedAt: new Date().toISOString(),
|
||||||
|
},
|
||||||
|
overrideAccess: true,
|
||||||
|
depth: 0,
|
||||||
|
});
|
||||||
|
return { ok: true };
|
||||||
|
}
|
||||||
126
src/lib/operations/reconcile.ts
Normal file
126
src/lib/operations/reconcile.ts
Normal file
|
|
@ -0,0 +1,126 @@
|
||||||
|
import type { Payload, PayloadRequest } from "payload";
|
||||||
|
import { hasPermission } from "@/utils/access-control/hasPermission";
|
||||||
|
import { findByIdOrNull } from "./settlementResolve";
|
||||||
|
import { validateOperationEvent } from "./validate";
|
||||||
|
import { deriveEffect } from "./process";
|
||||||
|
import { applyOperationEffect } from "./settlement";
|
||||||
|
import type { RawSyncEvent } from "./types";
|
||||||
|
|
||||||
|
type PayloadType = Awaited<ReturnType<typeof import("payload").getPayload>>;
|
||||||
|
|
||||||
|
export interface ReplayResult {
|
||||||
|
ok: boolean;
|
||||||
|
status?: string;
|
||||||
|
effectsCreated?: number;
|
||||||
|
error?: string;
|
||||||
|
}
|
||||||
|
|
||||||
|
function relationId(value: unknown): number | null {
|
||||||
|
if (typeof value === "number") return value;
|
||||||
|
if (value !== null && typeof value === "object" && "id" in (value as { id: unknown })) {
|
||||||
|
const id = (value as { id: unknown }).id;
|
||||||
|
return typeof id === "number" ? id : null;
|
||||||
|
}
|
||||||
|
return null;
|
||||||
|
}
|
||||||
|
|
||||||
|
export function toRawSyncEvent(doc: Record<string, unknown>): RawSyncEvent {
|
||||||
|
return {
|
||||||
|
id: doc.id as number,
|
||||||
|
eventId: doc.eventId as string,
|
||||||
|
server: relationId(doc.server),
|
||||||
|
type: doc.type as string,
|
||||||
|
occurredAt: doc.occurredAt as string,
|
||||||
|
payload: doc.payload,
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Re-run validation and effect derivation for one rejected or dead-lettered
|
||||||
|
* ledger row. The raw evidence is immutable and never modified; the row is
|
||||||
|
* updated in place with a reconciliation note. Effect creation stays idempotent
|
||||||
|
* by effectKey, so replays can never double-apply.
|
||||||
|
*/
|
||||||
|
export async function replayOperationEvent(
|
||||||
|
payload: PayloadType,
|
||||||
|
user: { id: number } | null | undefined,
|
||||||
|
operationEventId: number,
|
||||||
|
req?: PayloadRequest,
|
||||||
|
): Promise<ReplayResult> {
|
||||||
|
const allowed = await hasPermission(payload, user, "operation-events:update");
|
||||||
|
if (!allowed) return { ok: false, error: "forbidden" };
|
||||||
|
|
||||||
|
const record = await findByIdOrNull(payload, "operation-events", operationEventId);
|
||||||
|
if (!record) return { ok: false, error: "not-found" };
|
||||||
|
const row = record as unknown as {
|
||||||
|
status: string;
|
||||||
|
sourceEvent: unknown;
|
||||||
|
reconciliationNotes: string | null;
|
||||||
|
};
|
||||||
|
if (row.status === "validated") return { ok: false, error: "already-validated" };
|
||||||
|
|
||||||
|
const source = await findByIdOrNull(payload, "arma-sync-events", relationId(row.sourceEvent) ?? -1);
|
||||||
|
if (!source) return { ok: false, error: "source-missing" };
|
||||||
|
|
||||||
|
const raw = toRawSyncEvent(source as unknown as Record<string, unknown>);
|
||||||
|
const validation = validateOperationEvent(raw);
|
||||||
|
const note = `Replayed by user ${user?.id ?? "unknown"}: ${validation.status}${
|
||||||
|
validation.error ? ` (${validation.error.code})` : ""
|
||||||
|
}`;
|
||||||
|
const notes = [row.reconciliationNotes, note].filter(Boolean).join(" | ");
|
||||||
|
|
||||||
|
await payload.update({
|
||||||
|
collection: "operation-events",
|
||||||
|
id: operationEventId,
|
||||||
|
data: {
|
||||||
|
status: validation.status,
|
||||||
|
error: validation.error ?? null,
|
||||||
|
validatedAt: new Date().toISOString(),
|
||||||
|
reconciliationNotes: notes,
|
||||||
|
},
|
||||||
|
overrideAccess: true,
|
||||||
|
depth: 0,
|
||||||
|
req,
|
||||||
|
});
|
||||||
|
|
||||||
|
if (validation.status !== "validated") {
|
||||||
|
payload.logger.warn(`[Operations] Replay of ${raw.eventId} remains ${validation.status}`);
|
||||||
|
return { ok: true, status: validation.status, effectsCreated: 0 };
|
||||||
|
}
|
||||||
|
|
||||||
|
const effect = deriveEffect(raw, validation.operationId ?? "unknown");
|
||||||
|
if (!effect) return { ok: true, status: "validated", effectsCreated: 0 };
|
||||||
|
|
||||||
|
try {
|
||||||
|
const effectDoc = await payload.create({
|
||||||
|
collection: "operation-effects",
|
||||||
|
data: {
|
||||||
|
operationEvent: operationEventId,
|
||||||
|
sourceEvent: raw.id,
|
||||||
|
operationId: validation.operationId ?? "unknown",
|
||||||
|
effectKey: `${raw.eventId}:${effect.effectType}`,
|
||||||
|
effectType: effect.effectType,
|
||||||
|
status: "pending",
|
||||||
|
payload: effect.payload,
|
||||||
|
appliedAt: null,
|
||||||
|
reconciliationNotes: `Created by replay: ${note}`,
|
||||||
|
},
|
||||||
|
overrideAccess: true,
|
||||||
|
depth: 0,
|
||||||
|
req,
|
||||||
|
});
|
||||||
|
await applyOperationEffect(payload, effectDoc.id, req);
|
||||||
|
return { ok: true, status: "validated", effectsCreated: 1 };
|
||||||
|
} catch (error) {
|
||||||
|
const code = (error as { cause?: { code?: string } })?.cause?.code;
|
||||||
|
if (code === "23505") {
|
||||||
|
return { ok: true, status: "validated", effectsCreated: 0 };
|
||||||
|
}
|
||||||
|
payload.logger.error(
|
||||||
|
`[Operations] Replay effect failed for ${raw.eventId}: ${
|
||||||
|
error instanceof Error ? error.message : "unknown error"
|
||||||
|
}`,
|
||||||
|
);
|
||||||
|
return { ok: false, error: "effect-failed" };
|
||||||
|
}
|
||||||
|
}
|
||||||
75
src/lib/operations/rollout.ts
Normal file
75
src/lib/operations/rollout.ts
Normal file
|
|
@ -0,0 +1,75 @@
|
||||||
|
import type { Payload, PayloadRequest } from "payload";
|
||||||
|
import { hasPermission } from "@/utils/access-control/hasPermission";
|
||||||
|
import { processOperationEvent } from "./process";
|
||||||
|
import { toRawSyncEvent } from "./reconcile";
|
||||||
|
|
||||||
|
type PayloadType = Awaited<ReturnType<typeof import("payload").getPayload>>;
|
||||||
|
|
||||||
|
export const OPERATION_LEDGER_FLAG = "OPERATION_LEDGER_ENABLED";
|
||||||
|
|
||||||
|
/** Staged-activation kill switch: env `OPERATION_LEDGER_ENABLED=false` disables
|
||||||
|
* derived-effect processing while raw evidence ingestion continues. */
|
||||||
|
export function isOperationLedgerEnabled(): boolean {
|
||||||
|
return process.env[OPERATION_LEDGER_FLAG] !== "false";
|
||||||
|
}
|
||||||
|
|
||||||
|
export interface BackfillResult {
|
||||||
|
processed: number;
|
||||||
|
skipped: number;
|
||||||
|
forbidden: boolean;
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Replay raw sync events that have no ledger row yet (for example events
|
||||||
|
* ingested while the ledger flag was disabled). Idempotent: processOperationEvent
|
||||||
|
* skips source events that already have a row.
|
||||||
|
*/
|
||||||
|
export async function backfillOperationEvents(
|
||||||
|
payload: PayloadType,
|
||||||
|
user: { id: number } | null | undefined,
|
||||||
|
options?: { limit?: number; req?: PayloadRequest },
|
||||||
|
): Promise<BackfillResult> {
|
||||||
|
const allowed = await hasPermission(payload, user, "operation-events:update");
|
||||||
|
if (!allowed) return { processed: 0, skipped: 0, forbidden: true };
|
||||||
|
|
||||||
|
const limit = options?.limit ?? 200;
|
||||||
|
const raw = await payload.find({
|
||||||
|
collection: "arma-sync-events",
|
||||||
|
sort: "-createdAt",
|
||||||
|
limit,
|
||||||
|
depth: 0,
|
||||||
|
overrideAccess: true,
|
||||||
|
req: options?.req,
|
||||||
|
});
|
||||||
|
|
||||||
|
const rawIds = raw.docs.map((doc) => doc.id as number);
|
||||||
|
const existing = rawIds.length
|
||||||
|
? await payload.find({
|
||||||
|
collection: "operation-events",
|
||||||
|
where: { sourceEvent: { in: rawIds } },
|
||||||
|
limit: rawIds.length,
|
||||||
|
depth: 0,
|
||||||
|
overrideAccess: true,
|
||||||
|
req: options?.req,
|
||||||
|
})
|
||||||
|
: { docs: [] as { sourceEvent: unknown }[] };
|
||||||
|
const alreadyLedgered = new Set(
|
||||||
|
existing.docs.map((doc) => {
|
||||||
|
const source = (doc as { sourceEvent: unknown }).sourceEvent;
|
||||||
|
return typeof source === "number" ? source : null;
|
||||||
|
}),
|
||||||
|
);
|
||||||
|
|
||||||
|
let processed = 0;
|
||||||
|
let skipped = 0;
|
||||||
|
for (const doc of raw.docs) {
|
||||||
|
if (alreadyLedgered.has(doc.id as number)) {
|
||||||
|
skipped += 1;
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
await processOperationEvent(payload, toRawSyncEvent(doc as unknown as Record<string, unknown>), options?.req);
|
||||||
|
processed += 1;
|
||||||
|
}
|
||||||
|
payload.logger.info(`[Operations] Backfill processed ${processed} events (${skipped} already ledgered)`);
|
||||||
|
return { processed, skipped, forbidden: false };
|
||||||
|
}
|
||||||
338
tests/int/operation-rollout.int.spec.ts
Normal file
338
tests/int/operation-rollout.int.spec.ts
Normal file
|
|
@ -0,0 +1,338 @@
|
||||||
|
import { getPayload, Payload } from "payload";
|
||||||
|
import config from "@/payload.config";
|
||||||
|
import { afterAll, beforeAll, describe, expect, it } from "vitest";
|
||||||
|
import { acknowledgeCommand, recoverExpiredCommands } from "@/lib/arma-bridge/commands";
|
||||||
|
import { replayOperationEvent } from "@/lib/operations/reconcile";
|
||||||
|
import { backfillOperationEvents, isOperationLedgerEnabled } from "@/lib/operations/rollout";
|
||||||
|
|
||||||
|
let payload: Payload;
|
||||||
|
const RUN = `ops-rollout-${Date.now().toString(36)}`;
|
||||||
|
|
||||||
|
const syncEventIds: number[] = [];
|
||||||
|
const opEventIds: number[] = [];
|
||||||
|
const opEffectIds: number[] = [];
|
||||||
|
const commandIds: number[] = [];
|
||||||
|
const userIds: number[] = [];
|
||||||
|
const roleIds: number[] = [];
|
||||||
|
let serverId: number;
|
||||||
|
let adminUserId: number;
|
||||||
|
let plainUserId: number;
|
||||||
|
let ledgerUserId: number;
|
||||||
|
|
||||||
|
async function makeUser(label: string, roleDocIds?: number[]): Promise<number> {
|
||||||
|
const user = await payload.create({
|
||||||
|
collection: "users",
|
||||||
|
data: {
|
||||||
|
username: `${RUN}-${label}`,
|
||||||
|
discordUsername: `${RUN}-${label}-discord`,
|
||||||
|
displayName: `${RUN} ${label}`,
|
||||||
|
steamId: `${RUN}-steam-${label}`,
|
||||||
|
password: "Test1234",
|
||||||
|
...(roleDocIds && roleDocIds.length > 0 ? { roleDocs: roleDocIds } : {}),
|
||||||
|
},
|
||||||
|
overrideAccess: true,
|
||||||
|
depth: 0,
|
||||||
|
});
|
||||||
|
userIds.push(user.id);
|
||||||
|
return user.id;
|
||||||
|
}
|
||||||
|
|
||||||
|
async function makeCommand(
|
||||||
|
label: string,
|
||||||
|
status: "queued" | "delivered" | "succeeded" | "failed",
|
||||||
|
deliveredAt?: string,
|
||||||
|
): Promise<number> {
|
||||||
|
const command = await payload.create({
|
||||||
|
collection: "arma-commands",
|
||||||
|
data: {
|
||||||
|
commandId: `${RUN}-cmd-${label}`,
|
||||||
|
server: serverId,
|
||||||
|
type: "test.command",
|
||||||
|
payload: { action: "test" },
|
||||||
|
status,
|
||||||
|
...(deliveredAt ? { deliveredAt } : {}),
|
||||||
|
},
|
||||||
|
overrideAccess: true,
|
||||||
|
depth: 0,
|
||||||
|
});
|
||||||
|
commandIds.push(command.id);
|
||||||
|
return command.id;
|
||||||
|
}
|
||||||
|
|
||||||
|
async function opEventFor(rawId: number): Promise<{ id: number; status: string } | null> {
|
||||||
|
const found = await payload.find({
|
||||||
|
collection: "operation-events",
|
||||||
|
where: { sourceEvent: { equals: rawId } },
|
||||||
|
limit: 1,
|
||||||
|
depth: 0,
|
||||||
|
overrideAccess: true,
|
||||||
|
});
|
||||||
|
return (found.docs[0] as { id: number; status: string } | undefined) ?? null;
|
||||||
|
}
|
||||||
|
|
||||||
|
beforeAll(async () => {
|
||||||
|
const payloadConfig = await config;
|
||||||
|
payload = await getPayload({ config: payloadConfig });
|
||||||
|
|
||||||
|
const server = await payload.create({
|
||||||
|
collection: "game-servers",
|
||||||
|
data: { serverId: `${RUN}-srv`, name: `${RUN} Server`, status: "online" },
|
||||||
|
overrideAccess: true,
|
||||||
|
depth: 0,
|
||||||
|
});
|
||||||
|
serverId = server.id;
|
||||||
|
|
||||||
|
const ledgerRole = await payload.create({
|
||||||
|
collection: "roles",
|
||||||
|
data: {
|
||||||
|
name: `${RUN} Ledger Op`,
|
||||||
|
slug: `${RUN}-ledger-op`,
|
||||||
|
permissions: ["operation-events:update"],
|
||||||
|
},
|
||||||
|
overrideAccess: true,
|
||||||
|
depth: 0,
|
||||||
|
});
|
||||||
|
roleIds.push(ledgerRole.id);
|
||||||
|
|
||||||
|
const superRole = await payload.create({
|
||||||
|
collection: "roles",
|
||||||
|
data: { name: `${RUN} Rollout Admin`, slug: `${RUN}-rollout-admin`, isSuperuser: true },
|
||||||
|
overrideAccess: true,
|
||||||
|
depth: 0,
|
||||||
|
});
|
||||||
|
roleIds.push(superRole.id);
|
||||||
|
|
||||||
|
plainUserId = await makeUser("plain");
|
||||||
|
ledgerUserId = await makeUser("ledger", [ledgerRole.id]);
|
||||||
|
adminUserId = await makeUser("admin", [superRole.id]);
|
||||||
|
});
|
||||||
|
|
||||||
|
afterAll(async () => {
|
||||||
|
if (!payload) return;
|
||||||
|
process.env.OPERATION_LEDGER_ENABLED = "true";
|
||||||
|
for (const id of opEffectIds) {
|
||||||
|
await payload.delete({ collection: "operation-effects", id, overrideAccess: true }).catch(() => {});
|
||||||
|
}
|
||||||
|
for (const id of opEventIds) {
|
||||||
|
await payload.delete({ collection: "operation-events", id, overrideAccess: true }).catch(() => {});
|
||||||
|
}
|
||||||
|
for (const id of syncEventIds) {
|
||||||
|
await payload.delete({ collection: "arma-sync-events", id, overrideAccess: true }).catch(() => {});
|
||||||
|
}
|
||||||
|
for (const id of commandIds) {
|
||||||
|
await payload.delete({ collection: "arma-commands", id, overrideAccess: true }).catch(() => {});
|
||||||
|
}
|
||||||
|
for (const id of userIds) {
|
||||||
|
await payload.delete({ collection: "users", id, overrideAccess: true }).catch(() => {});
|
||||||
|
}
|
||||||
|
for (const id of roleIds) {
|
||||||
|
await payload.delete({ collection: "roles", id, overrideAccess: true }).catch(() => {});
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
describe("rollout: staged activation", () => {
|
||||||
|
it("defaults the ledger flag to enabled", () => {
|
||||||
|
expect(isOperationLedgerEnabled()).toBe(true);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("stores raw evidence but no derived ledger while disabled, then backfills", async () => {
|
||||||
|
const previous = process.env.OPERATION_LEDGER_ENABLED;
|
||||||
|
process.env.OPERATION_LEDGER_ENABLED = "false";
|
||||||
|
try {
|
||||||
|
const raw = await payload.create({
|
||||||
|
collection: "arma-sync-events",
|
||||||
|
data: {
|
||||||
|
eventId: `${RUN}-evt-disabled`,
|
||||||
|
server: serverId,
|
||||||
|
type: "operation.start",
|
||||||
|
occurredAt: new Date().toISOString(),
|
||||||
|
payload: { operationId: `${RUN}-op-disabled`, version: 1 },
|
||||||
|
},
|
||||||
|
overrideAccess: true,
|
||||||
|
depth: 0,
|
||||||
|
});
|
||||||
|
syncEventIds.push(raw.id);
|
||||||
|
|
||||||
|
expect(await opEventFor(raw.id)).toBeNull();
|
||||||
|
|
||||||
|
const admin = await payload.findByID({ collection: "users", id: adminUserId, overrideAccess: true, depth: 0 });
|
||||||
|
expect(admin.id).toBe(adminUserId);
|
||||||
|
const backfill = await backfillOperationEvents(payload, { id: adminUserId });
|
||||||
|
expect(backfill.forbidden).toBe(false);
|
||||||
|
expect(backfill.processed).toBeGreaterThanOrEqual(1);
|
||||||
|
|
||||||
|
const event = await opEventFor(raw.id);
|
||||||
|
expect(event).not.toBeNull();
|
||||||
|
expect(event?.status).toBe("validated");
|
||||||
|
if (event) opEventIds.push(event.id);
|
||||||
|
|
||||||
|
const effects = await payload.find({
|
||||||
|
collection: "operation-effects",
|
||||||
|
where: { operationEvent: { equals: event?.id ?? -1 } },
|
||||||
|
limit: 10,
|
||||||
|
depth: 0,
|
||||||
|
overrideAccess: true,
|
||||||
|
});
|
||||||
|
expect(effects.docs).toHaveLength(1);
|
||||||
|
opEffectIds.push(effects.docs[0].id);
|
||||||
|
|
||||||
|
const again = await backfillOperationEvents(payload, { id: adminUserId });
|
||||||
|
expect(again.processed).toBe(0);
|
||||||
|
} finally {
|
||||||
|
if (previous === undefined) delete process.env.OPERATION_LEDGER_ENABLED;
|
||||||
|
else process.env.OPERATION_LEDGER_ENABLED = previous;
|
||||||
|
}
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
|
describe("rollout: command lease and ack", () => {
|
||||||
|
it("re-queues only commands whose lease expired", async () => {
|
||||||
|
const staleId = await makeCommand(
|
||||||
|
"stale",
|
||||||
|
"delivered",
|
||||||
|
new Date(Date.now() - 10 * 60_000).toISOString(),
|
||||||
|
);
|
||||||
|
const freshId = await makeCommand("fresh", "delivered", new Date().toISOString());
|
||||||
|
|
||||||
|
const recovered = await recoverExpiredCommands(payload, serverId);
|
||||||
|
expect(recovered).toBeGreaterThanOrEqual(1);
|
||||||
|
|
||||||
|
const stale = await payload.findByID({ collection: "arma-commands", id: staleId, overrideAccess: true, depth: 0 });
|
||||||
|
expect((stale as unknown as { status: string }).status).toBe("queued");
|
||||||
|
expect((stale as unknown as { deliveredAt: unknown }).deliveredAt ?? null).toBeNull();
|
||||||
|
|
||||||
|
const fresh = await payload.findByID({ collection: "arma-commands", id: freshId, overrideAccess: true, depth: 0 });
|
||||||
|
expect((fresh as unknown as { status: string }).status).toBe("delivered");
|
||||||
|
});
|
||||||
|
|
||||||
|
it("accepts a first ack, treats duplicates as no-ops, and rejects acking undelivered commands", async () => {
|
||||||
|
const deliveredId = await makeCommand("acked", "delivered", new Date().toISOString());
|
||||||
|
const queuedId = await makeCommand("never-delivered", "queued");
|
||||||
|
|
||||||
|
const first = await acknowledgeCommand(payload, serverId, {
|
||||||
|
id: `${RUN}-cmd-acked`,
|
||||||
|
success: true,
|
||||||
|
result: { done: true },
|
||||||
|
});
|
||||||
|
expect(first.ok).toBe(true);
|
||||||
|
expect(first.duplicate).toBeUndefined();
|
||||||
|
|
||||||
|
const afterFirst = await payload.findByID({
|
||||||
|
collection: "arma-commands",
|
||||||
|
id: deliveredId,
|
||||||
|
overrideAccess: true,
|
||||||
|
depth: 0,
|
||||||
|
});
|
||||||
|
const completedAt = (afterFirst as unknown as { completedAt: string }).completedAt;
|
||||||
|
|
||||||
|
const duplicate = await acknowledgeCommand(payload, serverId, {
|
||||||
|
id: `${RUN}-cmd-acked`,
|
||||||
|
success: false,
|
||||||
|
error: "late duplicate",
|
||||||
|
});
|
||||||
|
expect(duplicate.ok).toBe(true);
|
||||||
|
expect(duplicate.duplicate).toBe(true);
|
||||||
|
|
||||||
|
const afterDuplicate = await payload.findByID({
|
||||||
|
collection: "arma-commands",
|
||||||
|
id: deliveredId,
|
||||||
|
overrideAccess: true,
|
||||||
|
depth: 0,
|
||||||
|
});
|
||||||
|
expect((afterDuplicate as unknown as { status: string }).status).toBe("succeeded");
|
||||||
|
expect((afterDuplicate as unknown as { completedAt: string }).completedAt).toBe(completedAt);
|
||||||
|
|
||||||
|
const undelivered = await acknowledgeCommand(payload, serverId, {
|
||||||
|
id: `${RUN}-cmd-never-delivered`,
|
||||||
|
success: true,
|
||||||
|
});
|
||||||
|
expect(undelivered.ok).toBe(false);
|
||||||
|
expect(undelivered.error).toBe("not-delivered");
|
||||||
|
|
||||||
|
const missing = await acknowledgeCommand(payload, serverId, {
|
||||||
|
id: `${RUN}-cmd-does-not-exist`,
|
||||||
|
success: true,
|
||||||
|
});
|
||||||
|
expect(missing.ok).toBe(false);
|
||||||
|
expect(missing.error).toBe("not-found");
|
||||||
|
|
||||||
|
void queuedId;
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
|
describe("rollout: reconciliation authorization and replay", () => {
|
||||||
|
it("forbids replay without the operation-events:update permission", async () => {
|
||||||
|
const raw = await payload.create({
|
||||||
|
collection: "arma-sync-events",
|
||||||
|
data: {
|
||||||
|
eventId: `${RUN}-evt-dead`,
|
||||||
|
server: serverId,
|
||||||
|
type: "operation.unsupported.thing",
|
||||||
|
occurredAt: new Date().toISOString(),
|
||||||
|
payload: { anything: true },
|
||||||
|
},
|
||||||
|
overrideAccess: true,
|
||||||
|
depth: 0,
|
||||||
|
});
|
||||||
|
syncEventIds.push(raw.id);
|
||||||
|
const event = await opEventFor(raw.id);
|
||||||
|
expect(event?.status).toBe("dead-letter");
|
||||||
|
if (event) opEventIds.push(event.id);
|
||||||
|
|
||||||
|
const forbidden = await replayOperationEvent(payload, { id: plainUserId }, event?.id ?? -1);
|
||||||
|
expect(forbidden.ok).toBe(false);
|
||||||
|
expect(forbidden.error).toBe("forbidden");
|
||||||
|
});
|
||||||
|
|
||||||
|
it("replays a dead-letter row in place, keeps invalid evidence invalid, and creates no effects", async () => {
|
||||||
|
const raw = await payload.create({
|
||||||
|
collection: "arma-sync-events",
|
||||||
|
data: {
|
||||||
|
eventId: `${RUN}-evt-dead-2`,
|
||||||
|
server: serverId,
|
||||||
|
type: "operation.still-unsupported",
|
||||||
|
occurredAt: new Date().toISOString(),
|
||||||
|
payload: { anything: true },
|
||||||
|
},
|
||||||
|
overrideAccess: true,
|
||||||
|
depth: 0,
|
||||||
|
});
|
||||||
|
syncEventIds.push(raw.id);
|
||||||
|
const event = await opEventFor(raw.id);
|
||||||
|
expect(event?.status).toBe("dead-letter");
|
||||||
|
if (event) opEventIds.push(event.id);
|
||||||
|
|
||||||
|
const result = await replayOperationEvent(payload, { id: ledgerUserId }, event?.id ?? -1);
|
||||||
|
expect(result.ok).toBe(true);
|
||||||
|
expect(result.status).toBe("dead-letter");
|
||||||
|
expect(result.effectsCreated).toBe(0);
|
||||||
|
|
||||||
|
const refreshed = await opEventFor(raw.id);
|
||||||
|
expect((refreshed as unknown as { reconciliationNotes: string | null } | null)?.reconciliationNotes ?? "").toContain(
|
||||||
|
"Replayed by user",
|
||||||
|
);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("rejects replay of an already-validated row", async () => {
|
||||||
|
const raw = await payload.create({
|
||||||
|
collection: "arma-sync-events",
|
||||||
|
data: {
|
||||||
|
eventId: `${RUN}-evt-valid`,
|
||||||
|
server: serverId,
|
||||||
|
type: "operation.start",
|
||||||
|
occurredAt: new Date().toISOString(),
|
||||||
|
payload: { operationId: `${RUN}-op-valid`, version: 1 },
|
||||||
|
},
|
||||||
|
overrideAccess: true,
|
||||||
|
depth: 0,
|
||||||
|
});
|
||||||
|
syncEventIds.push(raw.id);
|
||||||
|
const event = await opEventFor(raw.id);
|
||||||
|
expect(event?.status).toBe("validated");
|
||||||
|
if (event) opEventIds.push(event.id);
|
||||||
|
|
||||||
|
const result = await replayOperationEvent(payload, { id: ledgerUserId }, event?.id ?? -1);
|
||||||
|
expect(result.ok).toBe(false);
|
||||||
|
expect(result.error).toBe("already-validated");
|
||||||
|
});
|
||||||
|
});
|
||||||
Loading…
Reference in a new issue