From 275a574203632cec27de8af3a33573d18bebabfc Mon Sep 17 00:00:00 2001 From: simoleo89 Date: Sun, 13 Sep 2026 20:23:29 +0200 Subject: [PATCH] fix(news): retry rolled-back scheduled publication deadlocks --- src/lib/services/news-scheduler.test.ts | 97 ++++++++++++++++++++++++- src/lib/services/news-scheduler.ts | 87 +++++++++++++++------- 2 files changed, 155 insertions(+), 29 deletions(-) diff --git a/src/lib/services/news-scheduler.test.ts b/src/lib/services/news-scheduler.test.ts index 685e19c5..b4854efa 100644 --- a/src/lib/services/news-scheduler.test.ts +++ b/src/lib/services/news-scheduler.test.ts @@ -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) => { + 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); +}); diff --git a/src/lib/services/news-scheduler.ts b/src/lib/services/news-scheduler.ts index 554e2ac8..74bb74fd 100644 --- a/src/lib/services/news-scheduler.ts +++ b/src/lib/services/news-scheduler.ts @@ -7,40 +7,71 @@ 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 { // 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) => { - 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`, - ); - let published = 0; - for (const article of rows as unknown as Array<{ - id: string; - userId: number | null; - title: string | null; - publishAt: Date | string; - }>) { - 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) - continue; - const operationId = randomUUID(); - const hash = operationHash({ - articleId: article.id, - publishAt: new Date(article.publishAt).toISOString(), + 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`, + ); + let published = 0; + for (const article of rows as unknown as Array<{ + id: string; + userId: number | null; + title: string | null; + publishAt: Date | string; + }>) { + 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 + ) + continue; + const operationId = randomUUID(); + const hash = operationHash({ + articleId: article.id, + publishAt: new Date(article.publishAt).toISOString(), + }); + // Legacy articles without an author remain attributed to system (0). + await tx.execute( + sql`INSERT INTO cms_operations (id,actor_id,kind,request_key,request_hash,result_json) VALUES (${operationId},${article.userId ?? 0},${"news.schedule.publish"},${randomUUID()},${hash},${JSON.stringify({ articleId: article.id, articleTitle: article.title, published: true })})`, + ); + await enqueueEffect(tx, operationId, "news.refresh"); + published++; + } + return published; }); - // Legacy articles without an author remain attributed to system (0). - await tx.execute( - sql`INSERT INTO cms_operations (id,actor_id,kind,request_key,request_hash,result_json) VALUES (${operationId},${article.userId ?? 0},${"news.schedule.publish"},${randomUUID()},${hash},${JSON.stringify({ articleId: article.id, articleTitle: article.title, published: true })})`, - ); - await enqueueEffect(tx, operationId, "news.refresh"); - published++; + break; + } catch (error) { + // Retry only known rollbacks, never timeouts or an unknown commit outcome. + if (attempt >= 2 || !isDeadlock(error)) throw error; } - return published; - }); + } if (published > 0) { try { await invalidateNewsCache();