fix(housekeeping): converge rank synchronization

This commit is contained in:
Simo committed 2026-09-04 20:38:10 +02:00
1 parent ce4a3ba32c
commit 3d53321575
2 files changed
+165 -34

No files matched your search

@@ -293,6 +293,9 @@ function synchronizationData(
) { ) {
return { return {
...extra, ...extra,
...(intent.kind === "set-rank"
? { userId: intent.userId, rank: intent.rank }
: {}),
operation: intent.operation, operation: intent.operation,
recoveryId, recoveryId,
synchronization, synchronization,
@@ -306,10 +309,45 @@ function synchronizationData(
async function executeExternalSynchronization( async function executeExternalSynchronization(
intent: SystemExternalSyncIntent, intent: SystemExternalSyncIntent,
): Promise<boolean> { ): Promise<{
return intent.kind === "set-rank" readonly intent: SystemExternalSyncIntent;
? rcon.setRank(intent.userId, intent.rank) readonly synchronized: boolean;
: rcon.send("updatepermissions"); }> {
if (intent.kind !== "set-rank") {
return {
intent,
synchronized: await rcon.send("updatepermissions"),
};
}
return db.transaction(async (tx) => {
const [target] = await tx
.select({ rank: User.rank })
.from(User)
.where(eq(User.id, intent.userId))
.for("update");
if (!target) {
throw new SystemMutationFailure(
"NOT_FOUND",
"errors.housekeeping.system.userNotFound",
);
}
const currentIntent = {
...intent,
rank: positiveInteger(target.rank),
} as const satisfies SystemExternalSyncIntent;
let synchronized = false;
try {
synchronized = await rcon.setRank(
currentIntent.userId,
currentIntent.rank,
);
} catch {
// Preserve the resolved current intent for a later convergent retry.
}
return { intent: currentIntent, synchronized };
});
} }
async function completeExternalSynchronization( async function completeExternalSynchronization(
@@ -318,32 +356,39 @@ async function completeExternalSynchronization(
recoveryId = context.correlationId, recoveryId = context.correlationId,
extra: Readonly<Record<string, unknown>> = {}, extra: Readonly<Record<string, unknown>> = {},
): Promise<unknown> { ): Promise<unknown> {
let currentIntent = intent;
let synchronized = false; let synchronized = false;
try { try {
synchronized = await executeExternalSynchronization(intent); const execution = await executeExternalSynchronization(intent);
} catch { currentIntent = execution.intent;
// The committed intent below remains the retry source of truth. synchronized = execution.synchronized;
} catch (error) {
if (error instanceof SystemMutationFailure) throw error;
} }
if (!synchronized) { if (!synchronized) {
try { try {
await logAudit(syncAuditEntry(intent, context, recoveryId, "partial")); await logAudit(
syncAuditEntry(currentIntent, context, recoveryId, "partial"),
);
} catch { } catch {
// The transactionally persisted intent is sufficient for recovery. // The transactionally persisted intent is sufficient for recovery.
} }
throw new SystemCommittedExternalFailure( throw new SystemCommittedExternalFailure(
synchronizationData(intent, recoveryId, "pending", extra), synchronizationData(currentIntent, recoveryId, "pending", extra),
"failed", "failed",
"persisted", "persisted",
); );
} }
try { try {
await logAudit(syncAuditEntry(intent, context, recoveryId, "success")); await logAudit(
syncAuditEntry(currentIntent, context, recoveryId, "success"),
);
} catch { } catch {
throw new SystemCommittedExternalFailure( throw new SystemCommittedExternalFailure(
{ {
...synchronizationData(intent, recoveryId, "completed", extra), ...synchronizationData(currentIntent, recoveryId, "completed", extra),
pending: ["completion-audit"], pending: ["completion-audit"],
}, },
"completed", "completed",
@@ -351,7 +396,7 @@ async function completeExternalSynchronization(
); );
} }
return synchronizationData(intent, recoveryId, "completed", extra); return synchronizationData(currentIntent, recoveryId, "completed", extra);
} }
function parseExternalSyncIntent( function parseExternalSyncIntent(
@@ -525,6 +570,22 @@ async function executeAccessMutation(
rankId: id, rankId: id,
} as const satisfies SystemExternalSyncIntent; } as const satisfies SystemExternalSyncIntent;
await db.transaction(async (tx) => { await db.transaction(async (tx) => {
let rankRows: { id: number }[] = [];
try {
const [rows] = await tx.execute(
sql`SELECT id FROM permission_ranks WHERE id = ${id} LIMIT 1 FOR UPDATE`,
);
rankRows = rows as unknown as { id: number }[];
} catch {
rankRows = [];
}
if (rankRows.length === 0) {
throw new SystemMutationFailure(
"NOT_FOUND",
"errors.housekeeping.system.rankNotFound",
);
}
const [userCount] = await tx const [userCount] = await tx
.select({ total: count() }) .select({ total: count() })
.from(User) .from(User)
@@ -918,17 +979,6 @@ async function executeRconMutation(
rank, rank,
} as const satisfies SystemExternalSyncIntent; } as const satisfies SystemExternalSyncIntent;
await db.transaction(async (tx) => { await db.transaction(async (tx) => {
const [target] = await tx
.select({ rank: User.rank })
.from(User)
.where(eq(User.id, userId))
.for("update");
if (!target) {
throw new SystemMutationFailure(
"NOT_FOUND",
"errors.housekeeping.system.userNotFound",
);
}
let rankRows: { id: number }[] = []; let rankRows: { id: number }[] = [];
try { try {
const [rows] = await tx.execute( const [rows] = await tx.execute(
@@ -944,6 +994,17 @@ async function executeRconMutation(
"errors.housekeeping.system.rankNotFound", "errors.housekeeping.system.rankNotFound",
); );
} }
const [target] = await tx
.select({ rank: User.rank })
.from(User)
.where(eq(User.id, userId))
.for("update");
if (!target) {
throw new SystemMutationFailure(
"NOT_FOUND",
"errors.housekeeping.system.userNotFound",
);
}
if (!context.capability.isSuperAdmin) { if (!context.capability.isSuperAdmin) {
const actorRank = context.capability.actor.rank; const actorRank = context.capability.actor.rank;
if (target.rank >= actorRank) { if (target.rank >= actorRank) {
@@ -257,6 +257,7 @@ describe("durable rank synchronization", () => {
"access.rank.delete", "access.rank.delete",
{ id: 7 }, { id: 7 },
() => { () => {
doubles.dbExecute.mockResolvedValueOnce([[{ id: 7 }]]);
doubles.dbSelect doubles.dbSelect
.mockImplementationOnce(() => plainSelection([{ total: 0 }])) .mockImplementationOnce(() => plainSelection([{ total: 0 }]))
.mockImplementationOnce(() => limitedSelection([])); .mockImplementationOnce(() => limitedSelection([]));
@@ -312,11 +313,24 @@ describe("durable rank synchronization", () => {
expect(doubles.rconSetRank).not.toHaveBeenCalled(); expect(doubles.rconSetRank).not.toHaveBeenCalled();
}); });
it("locks the user row before committing a set-rank intent", async () => { it("locks the target rank before the user when committing a set-rank intent", async () => {
doubles.dbSelect.mockImplementationOnce(() => const lockOrder: string[] = [];
limitedSelection([{ rank: 3 }]), doubles.dbSelect
); .mockImplementationOnce(() => ({
doubles.dbExecute.mockResolvedValue([[{ id: 4 }]]); from: () => ({
where: () => ({
for: async () => {
lockOrder.push("user");
return [{ rank: 3 }];
},
}),
}),
}))
.mockImplementationOnce(() => limitedSelection([{ rank: 4 }]));
doubles.dbExecute.mockImplementationOnce(async () => {
lockOrder.push("rank");
return [[{ id: 4 }]];
});
doubles.rconSetRank.mockResolvedValue(true); doubles.rconSetRank.mockResolvedValue(true);
const result = await systemMutationService.execute( const result = await systemMutationService.execute(
@@ -326,13 +340,37 @@ describe("durable rank synchronization", () => {
); );
expect(result).toMatchObject({ ok: true }); expect(result).toMatchObject({ ok: true });
expect(doubles.dbForUpdate).toHaveBeenCalledWith("update"); expect(lockOrder.slice(0, 2)).toEqual(["rank", "user"]);
});
it("locks a rank before checking whether users still reference it", async () => {
const lockOrder: string[] = [];
doubles.dbExecute.mockImplementationOnce(async () => {
lockOrder.push("rank");
return [[{ id: 7 }]];
});
doubles.dbSelect
.mockImplementationOnce(() => {
lockOrder.push("users");
return plainSelection([{ total: 0 }]);
})
.mockImplementationOnce(() => limitedSelection([]));
doubles.rconSend.mockResolvedValue(true);
const result = await systemMutationService.execute(
serviceContext(PERMS.PERMISSIONS_MANAGE),
"access.rank.delete",
{ id: 7 },
);
expect(result).toMatchObject({ ok: true });
expect(lockOrder.slice(0, 2)).toEqual(["rank", "users"]);
}); });
it("returns the durable recovery key when set-rank RCON fails after commit", async () => { it("returns the durable recovery key when set-rank RCON fails after commit", async () => {
doubles.dbSelect.mockImplementationOnce(() => doubles.dbSelect
limitedSelection([{ rank: 3 }]), .mockImplementationOnce(() => limitedSelection([{ rank: 3 }]))
); .mockImplementationOnce(() => limitedSelection([{ rank: 4 }]));
doubles.dbExecute.mockResolvedValue([[{ id: 4 }]]); doubles.dbExecute.mockResolvedValue([[{ id: 4 }]]);
doubles.rconSetRank.mockResolvedValue(false); doubles.rconSetRank.mockResolvedValue(false);
@@ -382,7 +420,7 @@ describe("durable rank synchronization", () => {
expect(doubles.updateEmulatorRank).not.toHaveBeenCalled(); expect(doubles.updateEmulatorRank).not.toHaveBeenCalled();
}); });
it("retries set-rank from its committed desired-state intent without another database mutation", async () => { it("retries set-rank using the current committed database rank", async () => {
doubles.dbExecute.mockResolvedValueOnce([ doubles.dbExecute.mockResolvedValueOnce([
[ [
{ {
@@ -396,6 +434,9 @@ describe("durable rank synchronization", () => {
}, },
], ],
]); ]);
doubles.dbSelect.mockImplementationOnce(() =>
limitedSelection([{ rank: 5 }]),
);
doubles.rconSetRank.mockResolvedValue(true); doubles.rconSetRank.mockResolvedValue(true);
const result = await systemMutationService.execute( const result = await systemMutationService.execute(
@@ -411,10 +452,39 @@ describe("durable rank synchronization", () => {
synchronization: "completed", synchronization: "completed",
}, },
}); });
expect(doubles.rconSetRank).toHaveBeenCalledWith(8, 4); expect(doubles.rconSetRank).toHaveBeenCalledWith(8, 5);
expect(doubles.dbUpdate).not.toHaveBeenCalled(); expect(doubles.dbUpdate).not.toHaveBeenCalled();
}); });
it("holds the user-row transaction lock while delivering set-rank to RCON", async () => {
let activeTransactions = 0;
let activeTransactionsDuringRcon = 0;
doubles.dbTransaction.mockImplementation(async (run) => {
activeTransactions += 1;
try {
return await run(transactionToken);
} finally {
activeTransactions -= 1;
}
});
doubles.dbSelect
.mockImplementationOnce(() => limitedSelection([{ rank: 3 }]))
.mockImplementationOnce(() => limitedSelection([{ rank: 4 }]));
doubles.dbExecute.mockResolvedValue([[{ id: 4 }]]);
doubles.rconSetRank.mockImplementation(async () => {
activeTransactionsDuringRcon = activeTransactions;
return true;
});
const result = await systemMutationService.execute(
serviceContext(PERMS.RCON_EXECUTE),
"rcon.set-rank",
{ userId: 8, rank: 4 },
);
expect(result).toMatchObject({ ok: true });
expect(activeTransactionsDuringRcon).toBeGreaterThan(0);
});
it("does not repeat an already completed external synchronization", async () => { it("does not repeat an already completed external synchronization", async () => {
doubles.dbExecute.mockResolvedValueOnce([ doubles.dbExecute.mockResolvedValueOnce([
[ [