fix(news): retry rolled-back scheduled publication deadlocks
CI / check (push) Successful in 3m35s
CI / deploy (push) Successful in 24s
CI / publish-container (push) Successful in 1m23s

This commit is contained in:
Simo committed 2026-09-13 20:23:29 +02:00
1 parent 9a2d73a6f6
commit 275a574203
2 files changed
+129 -3

No files matched your search

+96 -1
View File
@@ -15,6 +15,8 @@ const state = vi.hoisted(() => ({
operations: [] as unknown[][],
effects: [] as unknown[][],
failOutbox: false,
transactionCalls: 0,
failures: [] as Array<{ at: "select" | "outbox" | "commit"; error: unknown }>,
loseRace: false,
refresh: vi.fn(),
}));
@@ -26,10 +28,14 @@ vi.mock("@/lib/db", () => {
value instanceof Date
? value
: new Date(`${String(value).replace(" ", "T")}Z`);
const failAt = (at: "select" | "outbox" | "commit") => {
if (state.failures[0]?.at === at) throw state.failures.shift()?.error;
};
const tx = {
execute: async (query: SQL) => {
const { sql, params } = dialect.sqlToQuery(query);
state.queries.push({ sql, params });
if (sql.startsWith("SELECT")) failAt("select");
if (sql.startsWith("SELECT"))
return [
state.articles
@@ -54,6 +60,7 @@ vi.mock("@/lib/db", () => {
if (sql.startsWith("INSERT INTO cms_operations"))
state.operations.push(params);
else if (sql.startsWith("INSERT INTO cms_outbox")) {
failAt("outbox");
if (state.failOutbox) throw Error("outbox unavailable");
state.effects.push(params);
} else throw Error(`Unexpected SQL ${sql}`);
@@ -63,17 +70,22 @@ vi.mock("@/lib/db", () => {
return {
db: {
transaction: async (work: (tx: unknown) => Promise<unknown>) => {
state.transactionCalls++;
const before = structuredClone({
articles: state.articles,
operations: state.operations,
effects: state.effects,
});
let result: unknown;
try {
return await work(tx);
result = await work(tx);
} catch (error) {
Object.assign(state, before);
throw error;
}
// A lost commit response may follow a durable commit: preserve its writes.
failAt("commit");
return result;
},
},
};
@@ -95,6 +107,8 @@ beforeEach(() => {
state.operations = [];
state.effects = [];
state.failOutbox = false;
state.transactionCalls = 0;
state.failures = [];
state.loseRace = false;
state.refresh.mockReset();
});
@@ -174,3 +188,84 @@ it("records the article title and exact identifier for the operation screen", as
published: true,
});
});
const deadlock = () =>
Object.assign(new Error("Deadlock found"), {
code: "ER_LOCK_DEADLOCK",
errno: 1213,
});
it.each(["select", "outbox"] as const)(
"retries a rolled-back deadlock at %s with the same cutoff and one durable effect",
async (at) => {
state.failures = [
{ at, error: new Error("Drizzle query failed", { cause: deadlock() }) },
];
expect(await publishDueArticles(now)).toBe(1);
expect(state.transactionCalls).toBe(2);
expect(state.articles[0].status).toBe("published");
expect(state.operations).toHaveLength(1);
expect(state.effects).toHaveLength(1);
expect(state.refresh).toHaveBeenCalledOnce();
const cutoffs = state.queries
.filter(({ sql }) => sql.startsWith("SELECT"))
.map(({ params }) => params[0]);
expect(cutoffs).toEqual([
"2030-01-02 00:00:00.000",
"2030-01-02 00:00:00.000",
]);
},
);
it("bounds deadlock retries to three whole transactions and preserves the final error", async () => {
const failure = deadlock();
state.failures = Array.from({ length: 3 }, () => ({
at: "outbox" as const,
error: failure,
}));
await expect(publishDueArticles(now)).rejects.toBe(failure);
expect(state.transactionCalls).toBe(3);
expect(state.articles[0].status).toBe("scheduled");
expect(state.operations).toHaveLength(0);
expect(state.effects).toHaveLength(0);
expect(state.refresh).not.toHaveBeenCalled();
});
it("recognizes the numeric MariaDB deadlock code through nested driver causes", async () => {
state.failures = [
{ at: "select", error: { cause: { cause: { errno: 1213 } } } },
];
expect(await publishDueArticles(now)).toBe(1);
expect(state.transactionCalls).toBe(2);
});
it.each([
Object.assign(new Error("lock wait timeout"), {
code: "ER_LOCK_WAIT_TIMEOUT",
errno: 1205,
}),
new Error("ER_LOCK_DEADLOCK in an untrusted message"),
Object.assign(new Error("connection lost"), {
code: "PROTOCOL_CONNECTION_LOST",
}),
])("does not retry unrelated errors: %s", async (failure) => {
state.failures = [{ at: "outbox", error: failure }];
await expect(publishDueArticles(now)).rejects.toBe(failure);
expect(state.transactionCalls).toBe(1);
expect(state.operations).toHaveLength(0);
expect(state.effects).toHaveLength(0);
});
it("never retries an unknown commit outcome", async () => {
const failure = Object.assign(new Error("Commit response lost"), {
code: "PROTOCOL_CONNECTION_LOST",
});
state.failures = [{ at: "commit", error: failure }];
await expect(publishDueArticles(now)).rejects.toBe(failure);
expect(state.transactionCalls).toBe(1);
expect(state.operations).toHaveLength(1);
expect(state.effects).toHaveLength(1);
expect(state.refresh).not.toHaveBeenCalled();
});
it("bounds traversal of cyclic error causes", async () => {
const failure: { cause?: unknown } = {};
failure.cause = failure;
state.failures = [{ at: "select", error: failure }];
await expect(publishDueArticles(now)).rejects.toBe(failure);
expect(state.transactionCalls).toBe(1);
});
+33 -2
View File
@@ -7,11 +7,34 @@ import { db } from "@/lib/db";
import { logger } from "@/lib/logger";
import { invalidateNewsCache } from "./news-cache";
/** MySQL rolls the entire victim transaction back on this specific error. */
function isDeadlock(error: unknown): boolean {
let cause = error;
for (
let depth = 0;
cause && typeof cause === "object" && depth < 5;
depth++
) {
const record = cause as {
code?: unknown;
errno?: unknown;
cause?: unknown;
};
if (record.code === "ER_LOCK_DEADLOCK" || record.errno === 1213)
return true;
cause = record.cause;
}
return false;
}
/** Publish a bounded batch and persist its refresh work before committing. */
export async function publishDueArticles(now = new Date()): Promise<number> {
// Match Drizzle's TIMESTAMP writes; raw Date parameters use the host timezone.
const utcNow = now.toISOString().slice(0, -1).replace("T", " ");
const published = await db.transaction(async (tx) => {
let published: number;
for (let attempt = 0; ; attempt++) {
try {
published = await db.transaction(async (tx) => {
const [rows] = await tx.execute(
sql`SELECT CAST(id AS CHAR) AS id,user_id AS userId,title,publish_at AS publishAt FROM website_articles WHERE status='scheduled' AND publish_at<=${utcNow} ORDER BY publish_at,id LIMIT 100 FOR UPDATE`,
);
@@ -25,7 +48,9 @@ export async function publishDueArticles(now = new Date()): Promise<number> {
const [result] = await tx.execute(
sql`UPDATE website_articles SET status='published',published_at=${utcNow},updated_at=${utcNow} WHERE id=${article.id} AND status='scheduled' AND publish_at<=${utcNow}`,
);
if ((result as unknown as { affectedRows: number }).affectedRows !== 1)
if (
(result as unknown as { affectedRows: number }).affectedRows !== 1
)
continue;
const operationId = randomUUID();
const hash = operationHash({
@@ -41,6 +66,12 @@ export async function publishDueArticles(now = new Date()): Promise<number> {
}
return published;
});
break;
} catch (error) {
// Retry only known rollbacks, never timeouts or an unknown commit outcome.
if (attempt >= 2 || !isDeadlock(error)) throw error;
}
}
if (published > 0) {
try {
await invalidateNewsCache();