import { execFile } from "node:child_process"; import { randomUUID } from "node:crypto"; import { once } from "node:events"; import { copyFile, mkdir, mkdtemp, rm } from "node:fs/promises"; import { join, resolve } from "node:path"; import { pathToFileURL } from "node:url"; import { promisify } from "node:util"; import { eq, sql } from "drizzle-orm"; import type { RowDataPacket } from "mysql2/promise"; import mysql from "mysql2/promise"; import { GenericContainer, type StartedTestContainer, Wait, } from "testcontainers"; import { afterAll, beforeAll, beforeEach, describe, expect, it, vi, } from "vitest"; import { WebsiteArticles } from "@/db/schema"; import { articleEditToken } from "@/lib/article-edit-token"; // Only request/framework boundaries are replaced; persistence, caches and workers are real. const boundaries = vi.hoisted(() => ({ notify: vi.fn() })); vi.mock("@/lib/admin/guard", () => ({ requirePermission: async () => ({ id: 7, username: "Integration editor" }), })); // Permission constants stay real without loading the session/ACL runtime. vi.mock("@/lib/permissions", () => import("@/lib/permission-slugs")); vi.mock("next-intl/server", () => ({ getTranslations: async () => (key: string) => key, })); // The data cache only exists inside a Next render, so both entry points degrade // to the real work underneath them instead of being stubbed out. vi.mock("next/cache", () => ({ revalidatePath: vi.fn(), revalidateTag: vi.fn(), unstable_cache: (fn: unknown) => fn, })); vi.mock("next/navigation", () => ({ redirect: (url: string) => { throw Error(`Unexpected integration redirect: ${url}`); }, })); vi.mock("@/lib/services/webhook", () => ({ notify: boundaries.notify })); vi.mock("@/lib/auth", () => ({ auth: async () => ({ user: { id: "7" } }) })); vi.mock("@/lib/api-auth", () => ({ bearerUserId: async () => 7 })); const exec = promisify(execFile); const NEWS_REVISION_KEY = "cms:news:revision"; const migrations = [ "0025_article_publication.sql", "0026_article_editor_recovery.sql", "0027_catalog_packages.sql", "0028_history_snapshots.sql", "0029_admin_table_views.sql", "0031_operations_outbox.sql", ]; let maria: StartedTestContainer | undefined; let redisContainer: StartedTestContainer | undefined; let connection: mysql.Connection | undefined; let appDb: typeof import("@/lib/db").db | undefined; let appRedis: typeof import("@/lib/redis").redis; let commands: typeof import("@/features/catalog/server/bulk-offers"); let operations: typeof import("@/features/operations/server"); let cache: typeof import("@/lib/cache"); let articles: typeof import("@/actions/admin-articles"); let publicNews: typeof import("@/lib/services/news-detail"); let commentSubmission: typeof import("@/lib/services/article-comment-submission"); let commentModeration: typeof import("@/lib/services/moderation"); let commentAction: typeof import("@/actions/article-comments"); let commentApi: typeof import("@/app/api/articles/[slug]/comment/route"); let scheduler: typeof import("@/lib/services/news-scheduler"); let worker: typeof import("@/features/operations/worker"); let migrationRoot: string | undefined; let databaseUrl: string; async function rows(query: string) { if (!connection) throw Error("Integration database is not connected"); const [result] = await connection.query(query); return result; } async function migrate(...args: string[]) { if (!migrationRoot) throw Error("Migration fixture is not ready"); // Only explicit test variables: the copied loader cannot see checkout .env files. return exec( process.execPath, [ "--import", pathToFileURL(resolve("node_modules/tsx/dist/loader.mjs")).href, join(migrationRoot, "scripts/apply-migrations.ts"), ...args, ], { cwd: migrationRoot, env: { PATH: process.env.PATH, SystemRoot: process.env.SystemRoot, DATABASE_URL: databaseUrl, NODE_ENV: "test", }, timeout: 30_000, }, ); } beforeAll(async () => { // No fixed names, ports, external URLs, shared volumes or container reuse. const databasePassword = randomUUID(); const redisPassword = randomUUID(); maria = await new GenericContainer("mariadb:11.4.5") .withEnvironment({ MARIADB_ROOT_PASSWORD: randomUUID(), MARIADB_DATABASE: "integration", MARIADB_USER: "integration", MARIADB_PASSWORD: databasePassword, }) .withExposedPorts(3306) .withHealthCheck({ test: ["CMD", "healthcheck.sh", "--connect", "--innodb_initialized"], interval: 1000, timeout: 5000, retries: 60, startPeriod: 1000, }) .withWaitStrategy(Wait.forHealthCheck()) .withStartupTimeout(120_000) .start(); databaseUrl = `mysql://integration:${databasePassword}@${maria.getHost()}:${maria.getMappedPort(3306)}/integration`; // Inspect stored TIMESTAMP values as UTC independently of the host timezone. // The application pool retains its production configuration. connection = await mysql.createConnection({ host: maria.getHost(), port: maria.getMappedPort(3306), user: "integration", password: databasePassword, database: "integration", timezone: "Z", supportBigNumbers: true, bigNumberStrings: true, charset: "utf8mb4", }); redisContainer = await new GenericContainer("redis:7.4.2-alpine") .withCommand(["redis-server", "--requirepass", redisPassword]) .withExposedPorts(6379) .withWaitStrategy(Wait.forLogMessage("Ready to accept connections")) .withStartupTimeout(60_000) .start(); process.env.DATABASE_URL = databaseUrl; process.env.REDIS_URL = `redis://:${redisPassword}@${redisContainer.getHost()}:${redisContainer.getMappedPort(6379)}/0`; delete process.env.SKIP_ENV_VALIDATION; delete process.env.OPENAI_API_KEY; Object.assign(process.env, { NODE_ENV: "test" }); process.env.HOTEL_NAME = "Integration"; await connection.query( "CREATE TABLE catalog_pages (id INT PRIMARY KEY, caption VARCHAR(255) NOT NULL) ENGINE=InnoDB", ); await connection.query( "CREATE TABLE catalog_items (id INT PRIMARY KEY, catalog_name VARCHAR(255) NOT NULL, page_id VARCHAR(25) NOT NULL, cost_credits INT NOT NULL, cost_points INT NOT NULL, points_type INT NOT NULL) ENGINE=InnoDB", ); await connection.query( "CREATE TABLE admin_audit_log (id INT AUTO_INCREMENT PRIMARY KEY, user_id INT NOT NULL, action VARCHAR(191) NOT NULL DEFAULT '', target VARCHAR(191) NOT NULL DEFAULT '', target_id INT NULL, details TEXT NULL, `before` TEXT NULL, `after` TEXT NULL, diff TEXT NULL, ip_address VARCHAR(45) NULL, created_at VARCHAR(64) NOT NULL DEFAULT '', updated_at VARCHAR(64) NULL) ENGINE=InnoDB", ); await connection.query( "CREATE TABLE website_articles (id BIGINT UNSIGNED AUTO_INCREMENT PRIMARY KEY, slug VARCHAR(255) NOT NULL UNIQUE, title VARCHAR(255) NOT NULL, short_story VARCHAR(255) NOT NULL, full_story LONGTEXT NOT NULL, user_id INT NULL, image VARCHAR(255) NOT NULL, created_at TIMESTAMP NULL DEFAULT NULL, updated_at TIMESTAMP NULL DEFAULT NULL) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4", ); await connection.query( "CREATE TABLE website_article_comments (id BIGINT UNSIGNED AUTO_INCREMENT PRIMARY KEY, article_id BIGINT UNSIGNED NOT NULL, user_id INT NOT NULL, comment VARCHAR(255) NOT NULL, created_at TIMESTAMP NULL, updated_at TIMESTAMP NULL) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4", ); await connection.query( "CREATE TABLE website_wordfilter (id BIGINT UNSIGNED AUTO_INCREMENT PRIMARY KEY, word VARCHAR(255) NOT NULL UNIQUE, created_at TIMESTAMP NULL, updated_at TIMESTAMP NULL) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4", ); // The emulator owns core tables. Exercise the real CMS migration CLI over this baseline. migrationRoot = await mkdtemp(join(resolve("integration"), ".migration-")); await mkdir(join(migrationRoot, "scripts")); await mkdir(join(migrationRoot, "drizzle/migrations"), { recursive: true }); for (const file of [ "apply-migrations.ts", "load-env.ts", "db-url.ts", "sql-statements.ts", ]) { await copyFile( resolve("scripts", file), join(migrationRoot, "scripts", file), ); } for (const file of migrations) await copyFile( resolve("drizzle/migrations", file), join(migrationRoot, "drizzle/migrations", file), ); await migrate(); appDb = (await import("@/lib/db")).db; appRedis = (await import("@/lib/redis")).redis; if (!appRedis) throw Error("Redis must be enabled in integration tests"); await appRedis.ping(); commands = await import("@/features/catalog/server/bulk-offers"); cache = await import("@/lib/cache"); operations = await import("@/features/operations/server"); articles = await import("@/actions/admin-articles"); publicNews = await import("@/lib/services/news-detail"); commentSubmission = await import("@/lib/services/article-comment-submission"); commentModeration = await import("@/lib/services/moderation"); commentAction = await import("@/actions/article-comments"); commentApi = await import("@/app/api/articles/[slug]/comment/route"); scheduler = await import("@/lib/services/news-scheduler"); worker = await import("@/features/operations/worker"); }); afterAll(async () => { // Settle every cleanup so a failed setup or close cannot leak the other container. const cleanup = await Promise.allSettled([ appDb?.$client.end(), appRedis?.quit(), connection?.end(), ]); const stops = await Promise.allSettled([ maria?.stop(), redisContainer?.stop(), ]); if (migrationRoot) { const target = resolve(migrationRoot); if ( !target.startsWith(`${resolve("integration")}\\.migration-`) && !target.startsWith(`${resolve("integration")}/.migration-`) ) throw Error("Unsafe migration fixture path"); await rm(target, { recursive: true, force: true }); } for (const result of [...cleanup, ...stops]) if (result.status === "rejected") throw result.reason; }); beforeEach(async () => { if (!connection) throw Error("Integration database is not connected"); vi.clearAllMocks(); await appRedis?.set(NEWS_REVISION_KEY, randomUUID()); await connection.query("DROP TRIGGER IF EXISTS reject_news_effect"); await connection.query( "DROP TRIGGER IF EXISTS reject_second_scheduled_effect", ); await connection.query("DELETE FROM website_article_revisions"); await connection.query("DELETE FROM website_article_drafts"); await connection.query("DELETE FROM website_article_comments"); await connection.query("DELETE FROM website_wordfilter"); commentModeration.reloadWordFilter(); await appRedis?.del( "ratelimit:comment:7", "ratelimit:comment:8", "ratelimit:comment:9", ); await connection.query("DELETE FROM website_articles"); await connection.query("DELETE FROM cms_outbox"); await connection.query("DELETE FROM cms_operations"); await connection.query("DROP TRIGGER IF EXISTS reject_second_history"); await connection.query("DELETE FROM admin_audit_log"); await connection.query("DELETE FROM catalog_items"); await connection.query("DELETE FROM catalog_pages"); await connection.query( "INSERT INTO catalog_pages VALUES (1, 'Original'), (2, 'Destination')", ); await connection.query( "INSERT INTO catalog_items VALUES (1, 'Chair', '1', 10, 0, 0), (2, 'Table', '1', 20, 0, 0)", ); }); const input = { ids: [1, 2], changes: { costCredits: { mode: "add" as const, value: 5 } }, }; describe("MariaDB migrations and catalog transactions", () => { it("runs the actual migration CLI twice without duplicate tracking or lost data", async () => { const before = await rows("SELECT * FROM cms_migrations ORDER BY id"); expect(before.map((row) => row.migration)).toEqual( migrations.map((file) => file.replace(/\.sql$/, "")), ); await connection?.query( "INSERT INTO website_admin_table_views VALUES (1, '/admin/catalog', 'Saved', '{}')", ); expect((await migrate()).stdout).toContain( "All migrations already applied", ); expect(await rows("SELECT * FROM cms_migrations ORDER BY id")).toEqual( before, ); expect(await rows("SELECT name FROM website_admin_table_views")).toEqual([ { name: "Saved" }, ]); expect((await migrate("--status")).stdout).toContain( `${migrations.length}/${migrations.length} applied, 0 pending`, ); const columns = await rows( "SELECT COLUMN_NAME, DATA_TYPE FROM information_schema.COLUMNS WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME='admin_audit_log' AND COLUMN_NAME IN ('before','after') ORDER BY COLUMN_NAME", ); expect(columns.map((row) => row.DATA_TYPE)).toEqual([ "mediumtext", "mediumtext", ]); }); it("commits both offers and history, then restores them through the real undo command", async () => { const preview = await commands.previewBulkOffersCommand(input); const result = await commands.applyBulkOffersCommand( input, preview.fingerprint, 7, ); expect(result.changedCount).toBe(2); expect(result.historyIds).toHaveLength(2); expect( (await rows("SELECT cost_credits FROM catalog_items ORDER BY id")).map( (row) => row.cost_credits, ), ).toEqual([15, 25]); expect( await rows("SELECT target, action FROM admin_audit_log ORDER BY id"), ).toEqual([ { target: "catalog_offer", action: "history_update" }, { target: "catalog_offer", action: "history_update" }, ]); await commands.undoBulkOffersCommand(result.historyIds, 7); expect( (await rows("SELECT cost_credits FROM catalog_items ORDER BY id")).map( (row) => row.cost_credits, ), ).toEqual([10, 20]); expect(await rows("SELECT id FROM admin_audit_log")).toHaveLength(4); }); it("rolls back earlier offer updates and history when the second history insert fails", async () => { await connection?.query( "CREATE TRIGGER reject_second_history BEFORE INSERT ON admin_audit_log FOR EACH ROW BEGIN IF NEW.target_id=2 THEN SIGNAL SQLSTATE '45000' SET MESSAGE_TEXT='forced history failure'; END IF; END", ); const preview = await commands.previewBulkOffersCommand(input); await expect( commands.applyBulkOffersCommand(input, preview.fingerprint, 7), ).rejects.toThrow(); expect( (await rows("SELECT cost_credits FROM catalog_items ORDER BY id")).map( (row) => row.cost_credits, ), ).toEqual([10, 20]); expect(await rows("SELECT id FROM admin_audit_log")).toHaveLength(0); }); it("serializes competing applications of one preview and rejects the stale contender", async () => { const preview = await commands.previewBulkOffersCommand(input); const results = await Promise.allSettled([ commands.applyBulkOffersCommand(input, preview.fingerprint, 7), commands.applyBulkOffersCommand(input, preview.fingerprint, 8), ]); expect( results.filter((result) => result.status === "fulfilled"), ).toHaveLength(1); const rejected = results.find((result) => result.status === "rejected"); expect(rejected?.status === "rejected" && rejected.reason.message).toMatch( /selection.*changed|preview/i, ); expect( (await rows("SELECT cost_credits FROM catalog_items ORDER BY id")).map( (row) => row.cost_credits, ), ).toEqual([15, 25]); expect(await rows("SELECT id FROM admin_audit_log")).toHaveLength(2); }); }); describe("Redis application cache", () => { it("writes to real Redis with expiry and reads it after memory invalidation", async () => { const key = `integration:${randomUUID()}:catalog`; let fetches = 0; const fetch = async () => ({ revision: ++fetches }); expect(await cache.cached(key, 60_000, fetch)).toEqual({ revision: 1 }); expect(await appRedis?.get(key)).toBe('{"revision":1}'); expect(await appRedis?.ttl(key)).toBeGreaterThan(0); cache.invalidateMemory(key); expect(await cache.cached(key, 60_000, fetch)).toEqual({ revision: 1 }); expect(fetches).toBe(1); await appRedis?.del(key); cache.invalidateMemory(key); expect(await cache.cached(key, 60_000, fetch)).toEqual({ revision: 2 }); }); it("keeps two independently named cache entries isolated during invalidation", async () => { const first = `integration:${randomUUID()}:first`; const second = `integration:${randomUUID()}:second`; await cache.cached(first, 60_000, async () => "first"); await cache.cached(second, 60_000, async () => "second"); await appRedis?.del(first); cache.invalidateMemory(first); cache.invalidateMemory(second); expect(await cache.cached(first, 60_000, async () => "updated")).toBe( "updated", ); expect(await cache.cached(second, 60_000, async () => "wrong")).toBe( "second", ); expect(await appRedis?.get(second)).toBe('"second"'); }); }); describe("operation idempotency and transactional outbox", () => { it("replays concurrent identical requests after applying their mutation and effect only once", async () => { const request = { actorId: 7, kind: "integration.test", key: randomUUID(), input: { increment: 3 }, }; let executions = 0; const work = async ( tx: import("@/features/operations/server").OperationTransaction, id: string, ) => { executions++; await tx.execute( sql`UPDATE catalog_items SET cost_credits=cost_credits+3 WHERE id=1`, ); await operations.enqueueEffect(tx, id, "catalog.refresh"); return { operationId: id, changed: 1 }; }; const result = await Promise.all([ operations.runOperation(request, work), operations.runOperation(request, work), ]); expect(result[0]).toEqual(result[1]); expect(executions).toBe(1); expect( await rows("SELECT cost_credits FROM catalog_items WHERE id=1"), ).toEqual([{ cost_credits: 13 }]); expect(await rows("SELECT id FROM cms_operations")).toHaveLength(1); expect(await rows("SELECT topic,status FROM cms_outbox")).toEqual([ { topic: "catalog.refresh", status: "pending" }, ]); await expect( operations.runOperation({ ...request, input: { increment: 9 } }, work), ).rejects.toThrow(); expect(executions).toBe(1); expect( await rows("SELECT cost_credits FROM catalog_items WHERE id=1"), ).toEqual([{ cost_credits: 13 }]); }); it("rolls back the request, mutation and queued effect together, allowing the same request to retry", async () => { const request = { actorId: 7, kind: "integration.test", key: randomUUID(), input: { increment: 3 }, }; await expect( operations.runOperation(request, async (tx, id) => { await tx.execute( sql`UPDATE catalog_items SET cost_credits=99 WHERE id=1`, ); await operations.enqueueEffect(tx, id, "catalog.refresh"); throw Error("failure after effect enqueue"); }), ).rejects.toThrow("failure after effect enqueue"); expect( await rows("SELECT cost_credits FROM catalog_items WHERE id=1"), ).toEqual([{ cost_credits: 10 }]); expect(await rows("SELECT id FROM cms_operations")).toHaveLength(0); expect(await rows("SELECT id FROM cms_outbox")).toHaveLength(0); await expect( operations.runOperation(request, async (tx, id) => { await tx.execute( sql`UPDATE catalog_items SET cost_credits=13 WHERE id=1`, ); await operations.enqueueEffect(tx, id, "catalog.refresh"); return { changed: 1 }; }), ).resolves.toEqual({ changed: 1 }); expect( await rows("SELECT cost_credits FROM catalog_items WHERE id=1"), ).toEqual([{ cost_credits: 13 }]); expect(await rows("SELECT id FROM cms_outbox")).toHaveLength(1); }); it("leases a queued effect to only one consumer and ignores a stale completion token", async () => { await operations.runOperation( { actorId: 7, kind: "integration.test", key: randomUUID(), input: {} }, async (tx, id) => { await operations.enqueueEffect(tx, id, "catalog.refresh"); return { queued: true }; }, ); const claims = await Promise.all([ operations.effectRepository.claim(), operations.effectRepository.claim(), ]); expect(claims.filter(Boolean)).toHaveLength(1); const claim = claims.find((value) => value !== null); if (!claim) throw Error("Expected one effect lease"); await operations.effectRepository.complete({ ...claim, token: randomUUID(), }); expect(await rows("SELECT status FROM cms_outbox")).toEqual([ { status: "running" }, ]); await operations.effectRepository.complete(claim); expect(await rows("SELECT status FROM cms_outbox")).toEqual([ { status: "done" }, ]); }); }); function articleForm(extra: Record = {}) { const form = new FormData(); for (const [key, value] of Object.entries({ requestKey: randomUUID(), title: "Integration news", slug: "integration-news", shortStory: "A real database publication", fullStory: "

Original draft body

", image: "/images/news.png", status: "draft", ...extra, })) form.set(key, value); return form; } async function savedArticle(slug = "integration-news") { if (!appDb) throw Error("Integration database is not connected"); const [article] = await appDb .select() .from(WebsiteArticles) .where(eq(WebsiteArticles.slug, slug)) .limit(1); if (!article) throw Error(`Missing article fixture: ${slug}`); return article; } async function seedScheduledArticles(now: Date) { if (!appDb) throw Error("Integration database is not connected"); const before = new Date(now.getTime() - 60_000); const after = new Date(now.getTime() + 60_000); await appDb.insert(WebsiteArticles).values( [ { id: 1n, slug: "due-earlier", status: "scheduled", publishAt: before, userId: 7, }, { id: 9007199254740993n, slug: "due-now", status: "scheduled", publishAt: now, userId: null, }, { id: 3n, slug: "future", status: "scheduled", publishAt: after, userId: 7, }, { id: 4n, slug: "draft", status: "draft", publishAt: before, userId: 7 }, { id: 5n, slug: "already-public", status: "published", publishAt: before, userId: 7, }, ].map((article) => ({ ...article, title: article.slug, shortStory: "Scheduled integration fixture", fullStory: `

${article.slug}

`, image: "", createdAt: before, updatedAt: before, publishedAt: article.status === "published" ? before : null, })), ); } describe("real news publication, scheduling and cache delivery", () => { it("creates one draft and publishes once under duplicate submits, preserving revision, audit and public content", async () => { const draft = articleForm(); const created = await Promise.all([ articles.createArticle(draft), articles.createArticle(draft), ]); expect(created[0]).toMatchObject({ ok: true }); expect(created[1]).toEqual(created[0]); expect(await rows("SELECT id FROM website_articles")).toHaveLength(1); expect(await rows("SELECT kind FROM cms_operations")).toEqual([ { kind: "news.create" }, ]); expect(await rows("SELECT id FROM website_article_revisions")).toHaveLength( 0, ); expect(boundaries.notify).not.toHaveBeenCalled(); const existing = await savedArticle(); const [creation] = await rows( "SELECT result_json FROM cms_operations WHERE kind='news.create'", ); expect(JSON.parse(creation.result_json)).toMatchObject({ articleId: String(existing.id), articleTitle: existing.title, }); expect(existing.status).toBe("draft"); expect(existing.publishedAt).toBeNull(); expect(await publicNews.getPublishedArticle(existing.slug)).toBeNull(); const negativeRevision = await appRedis?.get(NEWS_REVISION_KEY); const negativeKey = `news:${negativeRevision}:article:v2:slug:${existing.slug}`; expect(await appRedis?.get(negativeKey)).toBe("null"); expect(await appRedis?.ttl(negativeKey)).toBeGreaterThan(0); const body = `

${"Contenuto completo è 📰 ".repeat(4000)}

`; const publish = articleForm({ id: String(existing.id), baseToken: articleEditToken(existing), status: "published", fullStory: body, }); const published = await Promise.all([ articles.updateArticle(publish), articles.updateArticle(publish), ]); expect(published[0]).toMatchObject({ ok: true }); expect(published[1]).toEqual(published[0]); const saved = await savedArticle(); expect(saved).toMatchObject({ status: "published", fullStory: body, publishAt: null, }); expect(saved.publishedAt).toBeInstanceOf(Date); const revisions = await rows( "SELECT payload FROM website_article_revisions", ); expect(revisions).toHaveLength(1); expect(JSON.parse(revisions[0].payload)).toMatchObject({ status: "draft", fullStory: existing.fullStory, }); const history = await rows( "SELECT target,`before`,`after` FROM admin_audit_log", ); expect(history).toHaveLength(1); expect(history[0].target).toBe("news"); expect(JSON.parse(history[0].before)).toMatchObject({ status: "draft", fullStory: existing.fullStory, }); expect(JSON.parse(history[0].after)).toMatchObject({ status: "published", fullStory: body, }); expect(await rows("SELECT kind FROM cms_operations ORDER BY kind")).toEqual( [{ kind: "news.create" }, { kind: "news.update" }], ); expect(await rows("SELECT topic,status FROM cms_outbox")).toEqual([ { topic: "news.refresh", status: "pending" }, { topic: "news.refresh", status: "pending" }, ]); expect(boundaries.notify).toHaveBeenCalledOnce(); await worker.drainOperationEffects(); expect(await rows("SELECT status FROM cms_outbox")).toEqual([ { status: "done" }, { status: "done" }, ]); expect(await appRedis?.get(NEWS_REVISION_KEY)).not.toBe(negativeRevision); expect(await publicNews.getPublishedArticle(existing.slug)).toMatchObject({ id: existing.id, fullStory: body, slug: existing.slug, }); const publicKey = `news:${await appRedis?.get(NEWS_REVISION_KEY)}:article:v2:slug:${existing.slug}`; expect(JSON.parse(String(await appRedis?.get(publicKey)))).toMatchObject({ id: String(existing.id), fullStory: body, }); cache.invalidateMemory(publicKey); const fromRedis = await publicNews.getPublishedArticle(existing.slug); expect(fromRedis?.id).toBe(existing.id); expect(fromRedis?.publishedAt).toEqual(saved.publishedAt); }); it("rolls back publication, revision, audit and operation when queuing its effect fails, then retries the same submit", async () => { expect(await articles.createArticle(articleForm())).toMatchObject({ ok: true, }); const existing = await savedArticle(); const publish = articleForm({ id: String(existing.id), baseToken: articleEditToken(existing), status: "published", }); const operationsBefore = await rows("SELECT * FROM cms_operations"); const effectsBefore = await rows("SELECT * FROM cms_outbox"); await connection?.query( "CREATE TRIGGER reject_news_effect BEFORE INSERT ON cms_outbox FOR EACH ROW BEGIN IF NEW.topic='news.refresh' THEN SIGNAL SQLSTATE '45000' SET MESSAGE_TEXT='forced news effect failure'; END IF; END", ); expect(await articles.updateArticle(publish)).toMatchObject({ ok: false }); expect(await savedArticle()).toEqual(existing); expect(await rows("SELECT id FROM website_article_revisions")).toHaveLength( 0, ); expect(await rows("SELECT id FROM admin_audit_log")).toHaveLength(0); expect(await rows("SELECT * FROM cms_operations")).toEqual( operationsBefore, ); expect(await rows("SELECT * FROM cms_outbox")).toEqual(effectsBefore); expect(boundaries.notify).not.toHaveBeenCalled(); await connection?.query("DROP TRIGGER reject_news_effect"); expect(await articles.updateArticle(publish)).toMatchObject({ ok: true }); expect(await articles.updateArticle(publish)).toMatchObject({ ok: true }); expect((await savedArticle()).status).toBe("published"); expect(await rows("SELECT id FROM website_article_revisions")).toHaveLength( 1, ); expect(await rows("SELECT id FROM admin_audit_log")).toHaveLength(1); expect(await rows("SELECT id FROM cms_operations")).toHaveLength(2); expect(await rows("SELECT id FROM cms_outbox")).toHaveLength(2); expect(boundaries.notify).toHaveBeenCalledOnce(); }); it("rolls back the entire scheduled batch if its second effect cannot be queued", async () => { const now = new Date(Math.floor(Date.now() / 1000) * 1000); await seedScheduledArticles(now); const before = await rows("SELECT * FROM website_articles ORDER BY id"); await connection?.query( "CREATE TRIGGER reject_second_scheduled_effect BEFORE INSERT ON cms_outbox FOR EACH ROW BEGIN IF NEW.topic='news.refresh' AND (SELECT COUNT(*) FROM cms_operations WHERE kind='news.schedule.publish')=2 THEN SIGNAL SQLSTATE '45000' SET MESSAGE_TEXT='forced second scheduled effect failure'; END IF; END", ); await expect(scheduler.publishDueArticles(now)).rejects.toThrow(); expect(await rows("SELECT * FROM website_articles ORDER BY id")).toEqual( before, ); expect(await rows("SELECT id FROM cms_operations")).toHaveLength(0); expect(await rows("SELECT id FROM cms_outbox")).toHaveLength(0); await connection?.query("DROP TRIGGER reject_second_scheduled_effect"); expect(await scheduler.publishDueArticles(now)).toBe(2); expect(await rows("SELECT id FROM cms_operations")).toHaveLength(2); expect(await rows("SELECT topic,status FROM cms_outbox")).toEqual([ { topic: "news.refresh", status: "pending" }, { topic: "news.refresh", status: "pending" }, ]); }); it("serializes competing scheduler ticks without duplicate effects or early publication", async () => { const now = new Date(Math.floor(Date.now() / 1000) * 1000); await seedScheduledArticles(now); const untouched = await rows( "SELECT * FROM website_articles WHERE id IN (3,4,5) ORDER BY id", ); expect(await publicNews.getPublishedArticle("due-earlier")).toBeNull(); expect(await publicNews.getPublishedArticle("due-now")).toBeNull(); const results = await Promise.all([ scheduler.publishDueArticles(now), scheduler.publishDueArticles(now), ]); expect(results.sort()).toEqual([0, 2]); expect(await scheduler.publishDueArticles(now)).toBe(0); expect( await rows( "SELECT * FROM website_articles WHERE id IN (3,4,5) ORDER BY id", ), ).toEqual(untouched); const published = await rows( "SELECT CAST(id AS CHAR) AS id,status,published_at FROM website_articles WHERE id IN (1,9007199254740993) ORDER BY id", ); expect(published).toEqual([ { id: "1", status: "published", published_at: now }, { id: "9007199254740993", status: "published", published_at: now }, ]); const recorded = await rows( "SELECT actor_id,kind,result_json FROM cms_operations ORDER BY actor_id", ); expect(recorded).toHaveLength(2); expect( recorded.map((entry) => ({ actor: entry.actor_id, kind: entry.kind, result: JSON.parse(entry.result_json), })), ).toEqual([ { actor: 0, kind: "news.schedule.publish", result: { articleId: "9007199254740993", articleTitle: "due-now", published: true, }, }, { actor: 7, kind: "news.schedule.publish", result: { articleId: "1", articleTitle: "due-earlier", published: true, }, }, ]); expect(await rows("SELECT topic,status FROM cms_outbox")).toEqual([ { topic: "news.refresh", status: "pending" }, { topic: "news.refresh", status: "pending" }, ]); await worker.drainOperationEffects(); expect(await publicNews.getPublishedArticle("due-earlier")).toMatchObject({ id: 1n, }); expect(await publicNews.getPublishedArticle("due-now")).toMatchObject({ id: 9007199254740993n, }); expect(await publicNews.getPublishedArticle("future")).toBeNull(); expect(await publicNews.getPublishedArticle("draft")).toBeNull(); }); it("keeps a committed publication readable during Redis disconnect and repairs cached absence through a retried delivery", async () => { if (!appRedis) throw Error("Redis must be enabled in integration tests"); const redis = appRedis; expect(await articles.createArticle(articleForm())).toMatchObject({ ok: true, }); await worker.drainOperationEffects(); const existing = await savedArticle(); expect(await publicNews.getPublishedArticle(existing.slug)).toBeNull(); const negativeRevision = await redis.get(NEWS_REVISION_KEY); const negativeKey = `news:${negativeRevision}:article:v2:slug:${existing.slug}`; expect(await redis.get(negativeKey)).toBe("null"); const publish = articleForm({ id: String(existing.id), baseToken: articleEditToken(existing), status: "published", }); const disconnected = once(redis, "end"); redis.disconnect(); await disconnected; try { expect(redis.status).toBe("end"); expect(await articles.updateArticle(publish)).toMatchObject({ ok: true }); expect((await savedArticle()).status).toBe("published"); expect(await publicNews.getPublishedArticle(existing.slug)).toMatchObject( { id: existing.id }, ); await worker.drainOperationEffects(); const [pending] = await rows( "SELECT status,attempts,last_error FROM cms_outbox WHERE status<>'done'", ); expect(pending).toMatchObject({ status: "pending", attempts: 1 }); expect(pending.last_error).toBeTruthy(); await redis.connect(); expect(await redis.ping()).toBe("PONG"); expect(await redis.get(NEWS_REVISION_KEY)).toBe(negativeRevision); expect(await publicNews.getPublishedArticle(existing.slug)).toBeNull(); // Make the real queued retry due without a wall-clock sleep. await connection?.query( "UPDATE cms_outbox SET available_at=UTC_TIMESTAMP(3) WHERE status='pending'", ); await worker.drainOperationEffects(); expect( await rows( "SELECT status,attempts,last_error FROM cms_outbox ORDER BY attempts", ), ).toEqual([ { status: "done", attempts: 1, last_error: null }, { status: "done", attempts: 2, last_error: null }, ]); expect(await redis.get(NEWS_REVISION_KEY)).not.toBe(negativeRevision); expect(await publicNews.getPublishedArticle(existing.slug)).toMatchObject( { id: existing.id, fullStory: existing.fullStory }, ); expect(await articles.updateArticle(publish)).toMatchObject({ ok: true }); expect(await rows("SELECT id FROM website_articles")).toHaveLength(1); expect( await rows("SELECT id FROM website_article_revisions"), ).toHaveLength(1); expect(await rows("SELECT id FROM cms_operations")).toHaveLength(2); expect(await rows("SELECT id FROM cms_outbox")).toHaveLength(2); expect(boundaries.notify).toHaveBeenCalledOnce(); } finally { if (redis.status === "end") await redis.connect(); } }); }); async function seedCommentArticles() { if (!connection) throw Error("Integration database is not connected"); await connection.query( "INSERT INTO website_articles (id,slug,title,short_story,full_story,image,status,publish_at) VALUES (101,'comment-public','Public','','','','published',NULL),(102,'comment-draft','Draft','','','','draft',NULL),(103,'comment-future','Future','','','','published',DATE_ADD(UTC_TIMESTAMP(), INTERVAL 1 DAY)),(104,'comment-due','Due','','','','published',DATE_SUB(UTC_TIMESTAMP(), INTERVAL 1 DAY))", ); } describe("real comment publication, moderation and shared quota", () => { it("checks both target forms with MariaDB and rejects blocked text through the real word filter", async () => { await seedCommentArticles(); for (const userId of [8, 9]) { for (const [id, slug, eligible] of [ [101, "comment-public", true], [102, "comment-draft", false], [103, "comment-future", false], [104, "comment-due", true], ] as const) { const result = await commentSubmission.submitArticleComment({ userId, target: userId === 8 ? { id: String(id) } : { slug }, comment: "Database verified", }); expect(result).toEqual( eligible ? { ok: true, slug } : { ok: false, reason: "not_found" }, ); } } expect( await rows( "SELECT article_id,user_id,comment FROM website_article_comments ORDER BY id", ), ).toEqual([ { article_id: "101", user_id: 8, comment: "Database verified" }, { article_id: "104", user_id: 8, comment: "Database verified" }, { article_id: "101", user_id: 9, comment: "Database verified" }, { article_id: "104", user_id: 9, comment: "Database verified" }, ]); await connection?.query( "INSERT INTO website_wordfilter (word) VALUES ('blocked-content')", ); commentModeration.reloadWordFilter(); expect( await commentSubmission.submitArticleComment({ userId: 8, target: { slug: "comment-public" }, comment: "BLOCKED-CONTENT", }), ).toEqual({ ok: false, reason: "moderated" }); expect(await rows("SELECT id FROM website_article_comments")).toHaveLength( 4, ); }); it("shares the Redis bucket when alternating the real form and API handlers", async () => { if (!appRedis) throw Error("Redis must be enabled in integration tests"); await seedCommentArticles(); const form = new FormData(); form.set("articleId", "101"); form.set("slug", "comment-public"); form.set("comment", "Shared quota"); form.set("userId", "999"); const api = () => commentApi.POST( new Request("https://hotel.test/api/articles/comment-public/comment", { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ comment: "Shared quota", userId: 999 }), }), { params: Promise.resolve({ slug: "comment-public" }) }, ); for (let attempt = 0; attempt < 5; attempt++) { if (attempt % 2 === 0) { await expect(commentAction.postComment(form)).rejects.toThrow( "/news/comment-public?comment=posted", ); } else { const response = await api(); expect(response.status).toBe(200); expect(await response.json()).toEqual({ ok: true }); } } const limited = await api(); expect(limited.status).toBe(429); expect(Number(limited.headers.get("Retry-After"))).toBeGreaterThan(0); await expect(commentAction.postComment(form)).rejects.toThrow( "/news/comment-public?error=ratelimit", ); expect(await appRedis.get("ratelimit:comment:7")).toBe("7"); expect(await appRedis.pttl("ratelimit:comment:7")).toBeGreaterThan(0); expect(await appRedis.pttl("ratelimit:comment:7")).toBeLessThanOrEqual( 30_000, ); expect( await rows("SELECT user_id,comment FROM website_article_comments"), ).toEqual( Array.from({ length: 5 }, () => ({ user_id: 7, comment: "Shared quota", })), ); }); });