From 37d28893de9dffa9ab399cc6fff1e818c3ad3946 Mon Sep 17 00:00:00 2001 From: Z8MB1E Date: Sun, 20 Sep 2026 02:45:00 -0400 Subject: [PATCH] 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 --- docs/operations/runbook.md | 75 +++++ src/app/api/arma/v1/commands/ack/route.ts | 35 +-- src/app/api/arma/v1/commands/pull/route.ts | 4 + src/lib/arma-bridge/commands.ts | 104 +++++++ src/lib/operations/reconcile.ts | 126 ++++++++ src/lib/operations/rollout.ts | 75 +++++ tests/int/operation-rollout.int.spec.ts | 338 +++++++++++++++++++++ 7 files changed, 737 insertions(+), 20 deletions(-) create mode 100644 docs/operations/runbook.md create mode 100644 src/lib/arma-bridge/commands.ts create mode 100644 src/lib/operations/reconcile.ts create mode 100644 src/lib/operations/rollout.ts create mode 100644 tests/int/operation-rollout.int.spec.ts diff --git a/docs/operations/runbook.md b/docs/operations/runbook.md new file mode 100644 index 0000000..1043d8a --- /dev/null +++ b/docs/operations/runbook.md @@ -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` (`:`, 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. diff --git a/src/app/api/arma/v1/commands/ack/route.ts b/src/app/api/arma/v1/commands/ack/route.ts index bd3aff5..b31c1d7 100644 --- a/src/app/api/arma/v1/commands/ack/route.ts +++ b/src/app/api/arma/v1/commands/ack/route.ts @@ -1,6 +1,7 @@ import { NextRequest, NextResponse } from "next/server"; import { asJson } from "@/lib/arma-bridge/json"; import { requireArmaServer } from "@/lib/arma-bridge/server"; +import { acknowledgeCommand } from "@/lib/arma-bridge/commands"; export const runtime = "nodejs"; 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 }); } - const commands = await payload.find({ - collection: "arma-commands", - where: { and: [{ commandId: { equals: body.id } }, { server: { equals: server.id } }] }, - limit: 1, - overrideAccess: true, - }); - 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, + const outcome = await acknowledgeCommand(payload, server.id, { + id: body.id, + success: body.success, + result: asJson(body.result, null), + error: typeof body.error === "string" ? body.error : undefined, }); - 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 }); } diff --git a/src/app/api/arma/v1/commands/pull/route.ts b/src/app/api/arma/v1/commands/pull/route.ts index 4a372c5..d3a7dfa 100644 --- a/src/app/api/arma/v1/commands/pull/route.ts +++ b/src/app/api/arma/v1/commands/pull/route.ts @@ -1,5 +1,6 @@ import { NextRequest, NextResponse } from "next/server"; import { requireArmaServer } from "@/lib/arma-bridge/server"; +import { recoverExpiredCommands } from "@/lib/arma-bridge/commands"; export const runtime = "nodejs"; export const dynamic = "force-dynamic"; @@ -9,6 +10,9 @@ export async function POST(req: NextRequest) { if (auth instanceof NextResponse) return auth; const { payload, server } = auth; + // Expired leases (delivered, never acked) rejoin the queued pool. + await recoverExpiredCommands(payload, server.id); + const commands = await payload.find({ collection: "arma-commands", where: { and: [{ server: { equals: server.id } }, { status: { equals: "queued" } }] }, diff --git a/src/lib/arma-bridge/commands.ts b/src/lib/arma-bridge/commands.ts new file mode 100644 index 0000000..91caf49 --- /dev/null +++ b/src/lib/arma-bridge/commands.ts @@ -0,0 +1,104 @@ +import type { Payload } from "payload"; +import { asJson } from "./json"; + +type PayloadType = Awaited>; + +/** 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 { + 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 { + 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 }; +} diff --git a/src/lib/operations/reconcile.ts b/src/lib/operations/reconcile.ts new file mode 100644 index 0000000..4da3b05 --- /dev/null +++ b/src/lib/operations/reconcile.ts @@ -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>; + +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): 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 { + 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); + 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" }; + } +} diff --git a/src/lib/operations/rollout.ts b/src/lib/operations/rollout.ts new file mode 100644 index 0000000..da8db85 --- /dev/null +++ b/src/lib/operations/rollout.ts @@ -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>; + +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 { + 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), options?.req); + processed += 1; + } + payload.logger.info(`[Operations] Backfill processed ${processed} events (${skipped} already ledgered)`); + return { processed, skipped, forbidden: false }; +} diff --git a/tests/int/operation-rollout.int.spec.ts b/tests/int/operation-rollout.int.spec.ts new file mode 100644 index 0000000..12e4e98 --- /dev/null +++ b/tests/int/operation-rollout.int.spec.ts @@ -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 { + 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 { + 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"); + }); +});