fix(news): persist publication requests and retry cache delivery
CI / check (push) Successful in 3m12s
CI / deploy (push) Successful in 1m13s
CI / publish-container (push) Successful in 42s

This commit is contained in:
Simo committed 2026-09-13 18:33:45 +02:00
1 parent c977fe95ba
commit 52f6d1491f
40 files changed
+670 -71

No files matched your search

+48
View File
@@ -0,0 +1,48 @@
import { expect, test } from "@playwright/test";
test("news editor retries a lost response with the same request across reload and clears it after success", async ({
page,
}) => {
const requests: Array<Record<string, string>> = [];
let succeed = false;
await page.route("**/fixture/actions/article", async (route) => {
requests.push(route.request().postDataJSON());
if (!succeed) await route.abort("failed");
else
await route.fulfill({
contentType: "application/json",
body: JSON.stringify({ ok: true, data: {} }),
});
});
await page.goto("/admin/articles/new");
await expect(page.locator("iframe.tox-edit-area__iframe")).toBeVisible();
const save = page.getByRole("button", { name: "Save article", exact: true });
await save.click();
await expect(page.getByRole("alert")).toBeVisible();
await save.click();
await expect.poll(() => requests.length).toBe(2);
expect(requests[1].requestKey).toBe(requests[0].requestKey);
await page.reload();
await expect(page.locator("iframe.tox-edit-area__iframe")).toBeVisible();
await save.click();
await expect.poll(() => requests.length).toBe(3);
expect(requests[2].requestKey).toBe(requests[0].requestKey);
const stored = await page.evaluate(() =>
sessionStorage.getItem("cms:news:request:current:new"),
);
expect(stored).not.toContain("Community update");
succeed = true;
await save.click();
await expect
.poll(() =>
page.evaluate(() =>
sessionStorage.getItem("cms:news:request:current:new"),
),
)
.toBeNull();
expect(requests[3].requestKey).toBe(requests[0].requestKey);
await page.locator('input[name="title"]').fill("Another article");
await save.click();
await expect.poll(() => requests.length).toBe(5);
expect(requests[4].requestKey).not.toBe(requests[0].requestKey);
});
+4 -25
View File
@@ -2,14 +2,9 @@ import { drainOperationEffects } from "../src/features/operations/worker";
import { drainFurnitureImports } from "../src/lib/services/furni-job-worker";
import "./load-env";
import { Cron } from "croner";
import { and, eq, lt, lte, sql } from "drizzle-orm";
import { lt, sql } from "drizzle-orm";
import { env } from "../src/env";
import {
db,
PasswordReset,
WebsiteArticles,
WebsiteLoginLogs,
} from "../src/lib/db";
import { db, PasswordReset, WebsiteLoginLogs } from "../src/lib/db";
import { logger } from "../src/lib/logger";
import { redis } from "../src/lib/redis";
import {
@@ -23,7 +18,7 @@ import {
diskLevel,
parseDfOutput,
} from "../src/lib/services/disk-usage";
import { invalidateNewsCache } from "../src/lib/services/news-cache";
import { publishDueArticles } from "../src/lib/services/news-scheduler";
import { rcon } from "../src/lib/services/rcon";
function captureWorkerError(err: unknown, context: string): void {
@@ -347,24 +342,8 @@ async function pruneDockerCache(force = false): Promise<void> {
async function publishScheduledArticles(): Promise<void> {
try {
const now = new Date();
const result = await db
.update(WebsiteArticles)
.set({
status: "published",
publishedAt: now,
updatedAt: now,
})
.where(
and(
eq(WebsiteArticles.status, "scheduled"),
lte(WebsiteArticles.publishAt, now),
),
);
const [info] = result;
const published = info.affectedRows ?? 0;
const published = await publishDueArticles();
if (published > 0) {
await invalidateNewsCache();
logger.info(`Published ${published} scheduled article(s)`, {
module: "jobs",
});
+105
View File
@@ -2,6 +2,11 @@ import { beforeEach, describe, expect, it, vi } from "vitest";
const state = vi.hoisted(() => ({
rows: [] as unknown[][],
operations: new Map<
string,
{ id: string; requestHash: string; resultJson: string | null }
>(),
effects: [] as unknown[],
insert: vi.fn(),
update: vi.fn(),
invalidate: vi.fn(),
@@ -32,7 +37,32 @@ vi.mock("next/navigation", () => ({
},
}));
vi.mock("@/lib/db", async () => {
const { MySqlDialect } = await import("drizzle-orm/mysql-core");
const dialect = new MySqlDialect();
const database = {
execute: async (statement: import("drizzle-orm").SQL) => {
const { sql, params } = dialect.sqlToQuery(statement);
if (sql.startsWith("INSERT INTO cms_operations")) {
const [id, actor, kind, key, hash] = params;
const identity = `${actor}:${kind}:${key}`;
if (!state.operations.has(identity))
state.operations.set(identity, {
id: String(id),
requestHash: String(hash),
resultJson: null,
});
} else if (sql.startsWith("SELECT id,request_hash")) {
return [[state.operations.get(params.join(":"))]];
} else if (sql.startsWith("UPDATE cms_operations")) {
const row = [...state.operations.values()].find(
(row) => row.id === params[1],
);
if (row) row.resultJson = String(params[0]);
} else if (sql.startsWith("INSERT INTO cms_outbox"))
state.effects.push(params);
else throw new Error(`Unexpected SQL: ${sql}`);
return [[]];
},
select: () => ({
from: () => ({
where: () => ({
@@ -62,6 +92,7 @@ import { createArticle, updateArticle } from "./admin-articles";
function form(extra: Record<string, string> = {}) {
const f = new FormData();
for (const [k, v] of Object.entries({
requestKey: "11111111-1111-4111-8111-111111111111",
title: "Test news",
slug: "test-news",
shortStory: "Summary",
@@ -76,8 +107,30 @@ describe("news saves", () => {
beforeEach(() => {
vi.clearAllMocks();
state.rows = [];
state.operations.clear();
state.effects = [];
state.invalidate.mockReset();
state.insert.mockResolvedValue([{ insertId: 1 }]);
});
it("replays a committed creation without another article or notification", async () => {
const first = await createArticle(form());
const retry = await createArticle(form());
expect(first.ok).toBe(true);
expect(retry).toEqual(first);
expect(state.insert).toHaveBeenCalledOnce();
expect(state.notify).toHaveBeenCalledOnce();
expect(state.effects).toHaveLength(1);
});
it("returns success when cache refresh fails after the article commits", async () => {
state.invalidate.mockRejectedValue(Error("cache unavailable"));
await expect(createArticle(form())).resolves.toMatchObject({ ok: true });
expect(state.effects).toHaveLength(1);
});
it("rejects reuse of a committed request key with different content", async () => {
expect((await createArticle(form())).ok).toBe(true);
expect((await createArticle(form({ title: "Different" }))).ok).toBe(false);
expect(state.insert).toHaveBeenCalledOnce();
});
it("preserves a recoverable error and logs its reference", async () => {
state.insert.mockRejectedValueOnce(Error("Unknown column status"));
const result = await createArticle(form());
@@ -132,6 +185,58 @@ describe("news saves", () => {
);
expect(state.invalidate).toHaveBeenCalledOnce();
});
it("replays a committed update without stale-token failure or another revision", async () => {
const existing = {
title: "Old",
slug: "test-news",
image: "",
shortStory: "",
fullStory: "",
status: "draft",
publishAt: null,
publishedAt: null,
};
state.rows = [[existing]];
const data = form({ id: "1", baseToken: articleEditToken(existing) });
const first = await updateArticle(data);
expect(first.ok).toBe(true);
expect(await updateArticle(data)).toEqual(first);
expect(state.update).toHaveBeenCalledOnce();
expect(state.insert).toHaveBeenCalledOnce();
expect(state.notify).toHaveBeenCalledOnce();
expect(state.effects).toHaveLength(1);
});
it("accepts long article bodies while detecting body changes on replay", async () => {
const body = "a".repeat(100_000);
expect((await createArticle(form({ fullStory: body }))).ok).toBe(true);
expect((await createArticle(form({ fullStory: `${body}b` }))).ok).toBe(
false,
);
expect(state.insert).toHaveBeenCalledOnce();
});
it("preserves the scheduled instant without publishing early", async () => {
expect(
(
await createArticle(
form({ status: "scheduled", publishAt: "2099-01-01T00:00:00Z" }),
)
).ok,
).toBe(true);
expect(state.insert).toHaveBeenCalledWith(
expect.objectContaining({
status: "scheduled",
publishAt: new Date("2099-01-01T00:00:00Z"),
publishedAt: null,
}),
);
expect(state.notify).not.toHaveBeenCalled();
});
it("does not report a committed article as failed when notification dispatch throws", async () => {
state.notify.mockImplementationOnce(() => {
throw Error("notification failure");
});
await expect(createArticle(form())).resolves.toMatchObject({ ok: true });
});
it("rejects an outdated editor without inserting a revision or overwriting the article", async () => {
state.rows = [
[
+73 -17
View File
@@ -1,11 +1,13 @@
"use server";
import { createHash, randomUUID } from "node:crypto";
import { and, eq, ne } from "drizzle-orm";
import { revalidatePath } from "next/cache";
import { redirect } from "next/navigation";
import { getTranslations } from "next-intl/server";
import { ArticleDrafts, ArticleRevisions } from "@/db/article-editor";
import { historySnapshot, recordHistory } from "@/features/history/server";
import { enqueueEffect, runOperation } from "@/features/operations/server";
import { requirePermission } from "@/lib/admin/guard";
import { articleEditToken } from "@/lib/article-edit-token";
import {
@@ -77,18 +79,48 @@ async function refreshNews() {
revalidatePath("/me");
revalidatePath("/");
}
function operationInput(
form: FormData,
input: ReturnType<typeof readArticleInput>,
) {
return {
...input,
fullStory: createHash("sha256").update(input.fullStory).digest("hex"),
publishAt: input.publishAt?.toISOString() ?? null,
id: String(form.get("id") ?? ""),
baseToken: String(form.get("baseToken") ?? ""),
};
}
async function refreshSavedNews(id?: bigint) {
try {
await refreshNews();
if (id) revalidatePath(`/admin/articles/${id}`);
} catch (error) {
logger.error("News saved; cache refresh will be retried", {
module: "news",
error,
});
}
}
export async function createArticle(
formData: FormData,
): Promise<ArticleSaveResult> {
const staff = await requirePermission(PERMS.NEWS_EDIT);
let input: ReturnType<typeof readArticleInput>;
let slug: string;
let notification: Parameters<typeof notify>[0] | undefined;
try {
input = readArticleInput(formData);
const input = readArticleInput(formData);
await runOperation(
{
actorId: staff.id,
kind: "news.create",
key: String(formData.get("requestKey") ?? randomUUID()),
input: operationInput(formData, input),
},
async (tx, operationId) => {
const { rawSlug, ...fields } = input;
slug = await uniqueSlug(rawSlug || fields.title);
const slug = await uniqueSlug(rawSlug || fields.title, undefined, tx);
const now = new Date();
await db.insert(WebsiteArticles).values({
await tx.insert(WebsiteArticles).values({
...fields,
slug,
userId: staff.id,
@@ -96,18 +128,28 @@ export async function createArticle(
updatedAt: now,
publishedAt: fields.status === "published" ? now : null,
});
} catch (error) {
return saveFailure(error, "create");
}
// Auxiliary services must not turn a committed insert into an apparent failure.
await enqueueEffect(tx, operationId, "news.refresh");
if (input.status === "published")
notify({
notification = {
action: "news_publish",
actor: staff.username,
target: input.title,
details: slug,
});
await refreshNews();
};
return { ok: true };
},
);
} catch (error) {
return saveFailure(error, "create");
}
if (notification) {
try {
notify(notification);
} catch (error) {
logger.error("News notification failed", { module: "news", error });
}
}
await refreshSavedNews();
return { ok: true, data: { redirectTo: "/admin/articles" } };
}
export async function updateArticle(
@@ -122,7 +164,14 @@ export async function updateArticle(
let slug = "";
try {
input = readArticleInput(formData);
await db.transaction(async (tx) => {
await runOperation(
{
actorId: staff.id,
kind: "news.update",
key: String(formData.get("requestKey") ?? randomUUID()),
input: operationInput(formData, input),
},
async (tx, operationId) => {
const [existing] = await tx
.select()
.from(WebsiteArticles)
@@ -169,19 +218,26 @@ export async function updateArticle(
})
.where(eq(WebsiteArticles.id, id));
await recordHistory(tx, "news", historyId, staff.id, historyBefore);
});
await enqueueEffect(tx, operationId, "news.refresh");
return { ok: true };
},
);
} catch (error) {
return saveFailure(error, "update");
}
if (becamePublished)
if (becamePublished) {
try {
notify({
action: "news_publish",
actor: staff.username,
target: input.title,
details: slug,
});
await refreshNews();
revalidatePath(`/admin/articles/${id}`);
} catch (error) {
logger.error("News notification failed", { module: "news", error });
}
}
await refreshSavedNews(id);
return { ok: true, data: { redirectTo: "/admin/articles" } };
}
+7 -1
View File
@@ -26,7 +26,13 @@ export default async function DeliveryPage() {
<li key={item.id} className="admin-card min-w-0 space-y-2">
<div className="flex flex-wrap justify-between gap-3">
<strong>
{t(item.topic === "catalog.refresh" ? "hotel" : "export")}
{t(
item.topic === "catalog.refresh"
? "hotel"
: item.topic === "news.refresh"
? "news"
: "export",
)}
</strong>
<span>{t(`states.${item.status}`)}</span>
</div>
+25 -1
View File
@@ -7,6 +7,7 @@ import { useFormRecovery } from "@/hooks/use-form-recovery";
import { useUnsavedChanges } from "@/hooks/use-unsaved-changes";
import type { ArticleDraft } from "@/lib/article-draft";
import type { ArticleSaveResult } from "@/lib/article-input";
import { type ArticleRequest, articleRequest } from "@/lib/article-request";
import { slugify } from "@/lib/format";
import { publicationIssues } from "@/lib/publication-preflight";
import { ArticlePreview, type ArticlePreviewData } from "./article-preview";
@@ -46,6 +47,8 @@ export function ArticleForm({
const [saveError, setSaveError] = useState<string | null>(null);
const router = useRouter();
const saving = useRef(false);
const request = useRef<ArticleRequest | undefined>(undefined);
const requestScope = `${userId ?? "current"}:${articleKey}`;
const editVersion = useRef(0);
function markDirty() {
localRecovery.changed();
@@ -156,6 +159,19 @@ export function ArticleForm({
startTransition(async () => {
try {
await recovery.wait();
let storage: Storage | undefined;
try {
storage = window.sessionStorage;
} catch {
/* Storage may be disabled. */
}
request.current = await articleRequest(
data,
requestScope,
request.current,
storage,
);
data.set("requestKey", request.current.key);
const result = await action(data);
if (result && !result.ok) {
setSaveError(result.error);
@@ -166,7 +182,15 @@ export function ArticleForm({
setDirty(false);
localRecovery.clear();
}
if (result?.ok) await recovery.clear().catch(() => {});
if (result?.ok) {
request.current = undefined;
try {
storage?.removeItem(`cms:news:request:${requestScope}`);
} catch {
/* Storage may be disabled. */
}
await recovery.clear().catch(() => {});
}
if (result?.ok && result.data?.redirectTo)
router.push(result.data.redirectTo);
} catch (error) {
+4 -1
View File
@@ -53,7 +53,10 @@ export function operationHash(value: unknown) {
export function retryDelay(attempt: number) {
return Math.min(3600, 60 * 2 ** Math.max(0, attempt - 1));
}
export type EffectTopic = "catalog.refresh" | "catalog.export.request";
export type EffectTopic =
| "catalog.refresh"
| "catalog.export.request"
| "news.refresh";
export interface EffectClaim {
id: string;
token: string;
+56
View File
@@ -0,0 +1,56 @@
import { beforeEach, expect, it, vi } from "vitest";
const mocks = vi.hoisted(() => ({
claim: vi.fn(),
complete: vi.fn(),
fail: vi.fn(),
news: vi.fn(),
catalog: vi.fn(),
exportEnabled: vi.fn(),
request: vi.fn(),
error: vi.fn(),
}));
vi.mock("./server", () => ({
effectRepository: {
claim: mocks.claim,
complete: mocks.complete,
fail: mocks.fail,
},
}));
vi.mock("@/lib/services/news-cache", () => ({
refreshNewsCacheForDelivery: mocks.news,
}));
vi.mock("@/features/catalog/server/sync-status", () => ({
sendCatalogUpdate: mocks.catalog,
}));
vi.mock("@/lib/services/catalog-git-queue", () => ({
catalogExportEnabled: mocks.exportEnabled,
catalogExportQueue: () => ({ request: mocks.request }),
}));
vi.mock("@/lib/logger", () => ({ logger: { error: mocks.error } }));
import { drainOperationEffects } from "./worker";
const effect = {
id: "one",
token: "lease",
attempts: 1,
topic: "news.refresh",
};
beforeEach(() => {
vi.resetAllMocks();
mocks.claim.mockResolvedValueOnce(effect).mockResolvedValue(null);
mocks.news.mockResolvedValue(undefined);
});
it("invalidates the shared news cache before acknowledging delivery", async () => {
await drainOperationEffects();
expect(mocks.news).toHaveBeenCalledTimes(1);
expect(mocks.complete).toHaveBeenCalledWith(effect);
expect(mocks.fail).not.toHaveBeenCalled();
});
it("keeps a failed news invalidation retryable", async () => {
mocks.news.mockRejectedValue(new Error("Redis unavailable"));
await drainOperationEffects();
expect(mocks.fail).toHaveBeenCalledWith(effect);
expect(mocks.complete).not.toHaveBeenCalled();
});
+3
View File
@@ -5,6 +5,7 @@ import {
catalogExportEnabled,
catalogExportQueue,
} from "@/lib/services/catalog-git-queue";
import { refreshNewsCacheForDelivery } from "@/lib/services/news-cache";
import { dispatchEffects } from "./dispatcher";
import { effectRepository } from "./server";
@@ -19,6 +20,8 @@ export async function drainOperationEffects() {
throw new Error("Hotel update not delivered");
} else if (claim.topic === "catalog.export.request") {
if (catalogExportEnabled()) await catalogExportQueue().request();
} else if (claim.topic === "news.refresh") {
await refreshNewsCacheForDelivery();
} else throw new Error("Unknown delivery topic");
});
} catch (error) {
+46
View File
@@ -0,0 +1,46 @@
import { beforeEach, expect, it } from "vitest";
import { articleRequest } from "./article-request";
const values = new Map<string, string>();
const storage = {
getItem: (key: string) => values.get(key) ?? null,
setItem: (key: string, value: string) => {
values.set(key, value);
},
};
function form(body = "Private article body") {
const form = new FormData();
form.set("fullStory", body);
form.set("status", "published");
return form;
}
beforeEach(() => values.clear());
it("reuses the request after a lost response and an editor reload", async () => {
const first = await articleRequest(form(), "staff:1:new", undefined, storage);
const retry = await articleRequest(form(), "staff:1:new", undefined, storage);
expect(retry.key).toBe(first.key);
expect([...values.values()].join()).not.toContain("Private article body");
});
it("uses a new request for changed content or a different editor", async () => {
const first = await articleRequest(form(), "staff:1:new", undefined, storage);
expect(
(await articleRequest(form("Changed"), "staff:1:new", first, storage)).key,
).not.toBe(first.key);
expect(
(await articleRequest(form(), "staff:2:new", undefined, storage)).key,
).not.toBe(first.key);
});
it("retains the in-memory request when browser storage is unavailable", async () => {
const blocked = {
getItem: () => {
throw Error("blocked");
},
setItem: () => {
throw Error("blocked");
},
};
const first = await articleRequest(form(), "staff:1:new", undefined, blocked);
expect(
(await articleRequest(form(), "staff:1:new", first, blocked)).key,
).toBe(first.key);
});
+45
View File
@@ -0,0 +1,45 @@
export interface ArticleRequest {
signature: string;
key: string;
}
/** Retain only a fingerprint and request ID, never editorial content. */
export async function articleRequest(
form: FormData,
scope: string,
previous?: ArticleRequest,
storage?: Pick<Storage, "getItem" | "setItem">,
): Promise<ArticleRequest> {
const entries = [...form.entries()]
.filter(([key]) => key !== "requestKey")
.sort(([a], [b]) => a.localeCompare(b));
const digest = await crypto.subtle.digest(
"SHA-256",
new TextEncoder().encode(JSON.stringify([scope, entries])),
);
const signature = Array.from(new Uint8Array(digest), (byte) =>
byte.toString(16).padStart(2, "0"),
).join("");
const storageKey = `cms:news:request:${scope}`;
let saved = previous;
if (!saved) {
try {
saved = JSON.parse(storage?.getItem(storageKey) ?? "null") ?? undefined;
} catch {
/* Storage may be disabled. */
}
}
if (
saved?.signature === signature &&
/^[a-f0-9]{8}-[a-f0-9]{4}-[a-f0-9]{4}-[a-f0-9]{4}-[a-f0-9]{12}$/i.test(
saved.key,
)
)
return saved;
const next = { signature, key: crypto.randomUUID() };
try {
storage?.setItem(storageKey, JSON.stringify(next));
} catch {
/* The editor retains the key in memory. */
}
return next;
}
+14 -1
View File
@@ -18,7 +18,11 @@ vi.mock("@/lib/cache", () => ({
},
}));
import { cacheNews, invalidateNewsCache } from "./news-cache";
import {
cacheNews,
invalidateNewsCache,
refreshNewsCacheForDelivery,
} from "./news-cache";
beforeEach(() => {
state.values.clear();
@@ -65,3 +69,12 @@ it("reads fresh data when Redis is unavailable", async () => {
"fresh",
]);
});
it("keeps failed durable cache invalidations retryable", async () => {
state.set.mockRejectedValue(Error("offline"));
await expect(refreshNewsCacheForDelivery()).rejects.toThrow("offline");
});
it("does not acknowledge an ended Redis connection as refreshed", async () => {
state.status = "end";
await expect(refreshNewsCacheForDelivery()).rejects.toThrow();
});
+8
View File
@@ -30,3 +30,11 @@ export async function invalidateNewsCache(): Promise<void> {
});
}
}
/** Worker delivery must fail visibly so the outbox can retry the invalidation. */
export async function refreshNewsCacheForDelivery(): Promise<void> {
if (!redis) return;
if (redis.status === "end")
throw new Error("News cache connection is closed");
await redis.set(REVISION_KEY, randomUUID());
}
+130
View File
@@ -0,0 +1,130 @@
import type { SQL } from "drizzle-orm";
import { MySqlDialect } from "drizzle-orm/mysql-core";
import { beforeEach, expect, it, vi } from "vitest";
const state = vi.hoisted(() => ({
articles: [] as Array<{
id: string;
userId: number | null;
status: string;
publishAt: Date;
}>,
operations: [] as unknown[][],
effects: [] as unknown[][],
failOutbox: false,
loseRace: false,
refresh: vi.fn(),
}));
vi.mock("@/lib/logger", () => ({ logger: { error: vi.fn() } }));
vi.mock("./news-cache", () => ({ invalidateNewsCache: state.refresh }));
vi.mock("@/lib/db", () => {
const dialect = new MySqlDialect();
const tx = {
execute: async (query: SQL) => {
const { sql, params } = dialect.sqlToQuery(query);
if (sql.startsWith("SELECT"))
return [
state.articles
.filter(
(a) =>
a.status === "scheduled" && a.publishAt <= (params[0] as Date),
)
.slice(0, 100),
];
if (sql.startsWith("UPDATE website_articles")) {
const article = state.articles.find((a) => a.id === params[2]);
if (
state.loseRace ||
!article ||
article.status !== "scheduled" ||
article.publishAt > (params[3] as Date)
)
return [{ affectedRows: 0 }];
article.status = "published";
return [{ affectedRows: 1 }];
}
if (sql.startsWith("INSERT INTO cms_operations"))
state.operations.push(params);
else if (sql.startsWith("INSERT INTO cms_outbox")) {
if (state.failOutbox) throw Error("outbox unavailable");
state.effects.push(params);
} else throw Error(`Unexpected SQL ${sql}`);
return [[]];
},
};
return {
db: {
transaction: async (work: (tx: unknown) => Promise<unknown>) => {
const before = structuredClone({
articles: state.articles,
operations: state.operations,
effects: state.effects,
});
try {
return await work(tx);
} catch (error) {
Object.assign(state, before);
throw error;
}
},
},
};
});
import { publishDueArticles } from "./news-scheduler";
const now = new Date("2030-01-02T00:00:00Z");
beforeEach(() => {
state.articles = [
{
id: "1",
userId: 0,
status: "scheduled",
publishAt: new Date("2030-01-01T00:00:00Z"),
},
];
state.operations = [];
state.effects = [];
state.failOutbox = false;
state.loseRace = false;
state.refresh.mockReset();
});
it("publishes once with the original legacy author and a durable cache refresh", async () => {
expect(await publishDueArticles(now)).toBe(1);
expect(await publishDueArticles(now)).toBe(0);
expect(state.articles[0].status).toBe("published");
expect(state.operations).toHaveLength(1);
expect(state.operations[0][1]).toBe(0);
expect(state.effects).toHaveLength(1);
expect(state.effects[0][2]).toBe("news.refresh");
});
it("rolls publication back when durable cache enqueue fails", async () => {
state.failOutbox = true;
await expect(publishDueArticles(now)).rejects.toThrow("outbox unavailable");
expect(state.articles[0].status).toBe("scheduled");
expect(state.operations).toEqual([]);
expect(state.effects).toEqual([]);
});
it("leaves future articles and concurrent changed articles alone", async () => {
state.articles[0].publishAt = new Date("2031-01-01T00:00:00Z");
expect(await publishDueArticles(now)).toBe(0);
state.articles[0].publishAt = new Date("2030-01-01T00:00:00Z");
state.loseRace = true;
expect(await publishDueArticles(now)).toBe(0);
expect(state.operations).toEqual([]);
expect(state.effects).toEqual([]);
});
it("records real authors and represents missing legacy authors as system", async () => {
state.articles[0].userId = null;
state.articles.push({ ...state.articles[0], id: "2", userId: 42 });
expect(await publishDueArticles(now)).toBe(2);
expect(state.operations.map((row) => row[1])).toEqual([0, 42]);
});
it("reports committed publication even when immediate cache refresh fails", async () => {
state.refresh.mockRejectedValue(Error("cache offline"));
expect(await publishDueArticles(now)).toBe(1);
expect(state.articles[0].status).toBe("published");
expect(state.effects).toHaveLength(1);
expect(state.refresh).toHaveBeenCalledOnce();
});
+52
View File
@@ -0,0 +1,52 @@
import "server-only";
import { randomUUID } from "node:crypto";
import { sql } from "drizzle-orm";
import { operationHash } from "@/features/operations/model";
import { enqueueEffect } from "@/features/operations/server";
import { db } from "@/lib/db";
import { logger } from "@/lib/logger";
import { invalidateNewsCache } from "./news-cache";
/** Publish a bounded batch and persist its refresh work before committing. */
export async function publishDueArticles(now = new Date()): Promise<number> {
const published = await db.transaction(async (tx) => {
const [rows] = await tx.execute(
sql`SELECT CAST(id AS CHAR) AS id,user_id AS userId,publish_at AS publishAt FROM website_articles WHERE status='scheduled' AND publish_at<=${now} 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;
publishAt: Date | string;
}>) {
const [result] = await tx.execute(
sql`UPDATE website_articles SET status='published',published_at=${now},updated_at=${now} WHERE id=${article.id} AND status='scheduled' AND publish_at<=${now}`,
);
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, published: true })})`,
);
await enqueueEffect(tx, operationId, "news.refresh");
published++;
}
return published;
});
if (published > 0) {
try {
await invalidateNewsCache();
} catch (error) {
logger.error("Scheduled news published; cache refresh remains queued", {
module: "news",
error,
});
}
}
return published;
}
+2 -1
View File
@@ -3552,7 +3552,8 @@
"running": "Delivering",
"done": "Delivered",
"failed": "Needs attention"
}
},
"news": "News cache refresh"
}
},
"myDashboard": {
+2 -1
View File
@@ -4375,7 +4375,8 @@
"running": "Delivering",
"done": "Delivered",
"failed": "Needs attention"
}
},
"news": "News cache refresh"
}
},
"error": {
+2 -1
View File
@@ -4375,7 +4375,8 @@
"running": "Delivering",
"done": "Delivered",
"failed": "Needs attention"
}
},
"news": "News cache refresh"
}
},
"error": {
+2 -1
View File
@@ -4375,7 +4375,8 @@
"running": "Delivering",
"done": "Delivered",
"failed": "Needs attention"
}
},
"news": "News cache refresh"
}
},
"error": {
+2 -1
View File
@@ -4345,7 +4345,8 @@
"running": "Delivering",
"done": "Delivered",
"failed": "Needs attention"
}
},
"news": "News cache refresh"
}
},
"radioRequests": {
+2 -1
View File
@@ -4375,7 +4375,8 @@
"running": "Delivering",
"done": "Delivered",
"failed": "Needs attention"
}
},
"news": "News cache refresh"
}
},
"error": {
+2 -1
View File
@@ -4548,7 +4548,8 @@
"running": "Delivering",
"done": "Delivered",
"failed": "Needs attention"
}
},
"news": "News cache refresh"
}
},
"mod": {
+2 -1
View File
@@ -4345,7 +4345,8 @@
"running": "Delivering",
"done": "Delivered",
"failed": "Needs attention"
}
},
"news": "News cache refresh"
}
},
"radioRequests": {
+2 -1
View File
@@ -3552,7 +3552,8 @@
"running": "Delivering",
"done": "Delivered",
"failed": "Needs attention"
}
},
"news": "News cache refresh"
}
},
"myDashboard": {
+2 -1
View File
@@ -4345,7 +4345,8 @@
"running": "Delivering",
"done": "Delivered",
"failed": "Needs attention"
}
},
"news": "News cache refresh"
}
},
"radioRequests": {
+2 -1
View File
@@ -4375,7 +4375,8 @@
"running": "Delivering",
"done": "Delivered",
"failed": "Needs attention"
}
},
"news": "News cache refresh"
}
},
"error": {
+2 -1
View File
@@ -4375,7 +4375,8 @@
"running": "Delivering",
"done": "Delivered",
"failed": "Needs attention"
}
},
"news": "News cache refresh"
}
},
"error": {
+2 -1
View File
@@ -4517,7 +4517,8 @@
"running": "In consegna",
"done": "Consegnata",
"failed": "Da verificare"
}
},
"news": "Aggiornamento news"
}
},
"radioRequests": {
+2 -1
View File
@@ -3552,7 +3552,8 @@
"running": "Delivering",
"done": "Delivered",
"failed": "Needs attention"
}
},
"news": "News cache refresh"
}
},
"myDashboard": {
+2 -1
View File
@@ -4548,7 +4548,8 @@
"running": "Wordt afgeleverd",
"done": "Afgeleverd",
"failed": "Aandacht nodig"
}
},
"news": "Nieuws bijwerken"
}
},
"mod": {
+2 -1
View File
@@ -4375,7 +4375,8 @@
"running": "Delivering",
"done": "Delivered",
"failed": "Needs attention"
}
},
"news": "News cache refresh"
}
},
"error": {
+2 -1
View File
@@ -4375,7 +4375,8 @@
"running": "Delivering",
"done": "Delivered",
"failed": "Needs attention"
}
},
"news": "News cache refresh"
}
},
"error": {
+2 -1
View File
@@ -4375,7 +4375,8 @@
"running": "Delivering",
"done": "Delivered",
"failed": "Needs attention"
}
},
"news": "News cache refresh"
}
},
"error": {
+2 -1
View File
@@ -4375,7 +4375,8 @@
"running": "Delivering",
"done": "Delivered",
"failed": "Needs attention"
}
},
"news": "News cache refresh"
}
},
"error": {
+2 -1
View File
@@ -4373,7 +4373,8 @@
"running": "Delivering",
"done": "Delivered",
"failed": "Needs attention"
}
},
"news": "News cache refresh"
}
},
"error": {
+2 -1
View File
@@ -4375,7 +4375,8 @@
"running": "Delivering",
"done": "Delivered",
"failed": "Needs attention"
}
},
"news": "News cache refresh"
}
},
"error": {
+2 -1
View File
@@ -4375,7 +4375,8 @@
"running": "Delivering",
"done": "Delivered",
"failed": "Needs attention"
}
},
"news": "News cache refresh"
}
},
"error": {
+2 -1
View File
@@ -4375,7 +4375,8 @@
"running": "Delivering",
"done": "Delivered",
"failed": "Needs attention"
}
},
"news": "News cache refresh"
}
},
"error": {
+2 -1
View File
@@ -4375,7 +4375,8 @@
"running": "Delivering",
"done": "Delivered",
"failed": "Needs attention"
}
},
"news": "News cache refresh"
}
},
"error": {
+2 -1
View File
@@ -4375,7 +4375,8 @@
"running": "Delivering",
"done": "Delivered",
"failed": "Needs attention"
}
},
"news": "News cache refresh"
}
},
"error": {