diff --git a/.changeset/fix-cron-revision-row-reads.md b/.changeset/fix-cron-revision-row-reads.md new file mode 100644 index 0000000000..42b951bc7c --- /dev/null +++ b/.changeset/fix-cron-revision-row-reads.md @@ -0,0 +1,5 @@ +--- +"emdash": patch +--- + +Fixes scheduled maintenance reading the entire revision history on every cron invocation. diff --git a/packages/core/src/after.ts b/packages/core/src/after.ts index 07d13760b2..97f609ba4c 100644 --- a/packages/core/src/after.ts +++ b/packages/core/src/after.ts @@ -16,6 +16,9 @@ * `after()`, which they don't during type-checking. */ +import { trackDeferredTask } from "./deferred-tasks.js"; +import { getRequestContext } from "./request-context.js"; + export type WaitUntilFn = (promise: Promise) => void; // Resolves to the host's waitUntil if the adapter provided one, or @@ -45,11 +48,13 @@ waitUntilReady.catch(() => {}); * that care about errors should handle them inside `fn`. */ export function after(fn: () => void | Promise): void { - const promise = Promise.resolve() + const task = Promise.resolve() .then(fn) .catch((error) => { console.error("[emdash] deferred task failed:", error); }); + const requestTask = getRequestContext()?.deferredTasks?.track(task) ?? task; + const promise = trackDeferredTask(requestTask); // Defer the lifetime-extender handoff to the microtask that resolves // waitUntilReady. On workerd this is effectively instant (the virtual diff --git a/packages/core/src/api/handlers/revision.ts b/packages/core/src/api/handlers/revision.ts index a31ec0bb8e..0683347380 100644 --- a/packages/core/src/api/handlers/revision.ts +++ b/packages/core/src/api/handlers/revision.ts @@ -4,6 +4,7 @@ import type { Kysely } from "kysely"; +import { after } from "../../after.js"; import { ContentRepository } from "../../database/repositories/content.js"; import { RevisionRepository, type Revision } from "../../database/repositories/revision.js"; import { withTransaction } from "../../database/transaction.js"; @@ -115,7 +116,7 @@ export async function handleRevisionRestore( // Atomically update content and create a new revision to record the restore. // If either operation fails, neither is committed (on engines that support // transactions; on D1, withTransaction falls back to sequential execution). - const item = await withTransaction(db, async (trx) => { + const { item, queuedRevisionId } = await withTransaction(db, async (trx) => { const trxContentRepo = new ContentRepository(trx); const trxRevisionRepo = new RevisionRepository(trx); @@ -124,19 +125,32 @@ export async function handleRevisionRestore( slug: typeof _slug === "string" ? _slug : undefined, }); - await trxRevisionRepo.create({ + const queuedRevision = await trxRevisionRepo.create({ collection: revision.collection, entryId: revision.entryId, data: revision.data, authorId: callerUserId, }); - return updated; + return { item: updated, queuedRevisionId: queuedRevision.id }; }); - // Fire-and-forget: prune old revisions to prevent unbounded growth const pruneRepo = new RevisionRepository(db); - void pruneRepo.pruneOldRevisions(revision.collection, revision.entryId, 50).catch(() => {}); + after(async () => { + try { + await pruneRepo.pruneQueuedEntry( + revision.collection, + revision.entryId, + queuedRevisionId, + 50, + ); + } catch (error) { + console.error( + `[revisions] Failed to prune revisions for ${revision.collection}/${revision.entryId}:`, + error, + ); + } + }); return { success: true, diff --git a/packages/core/src/astro/middleware.ts b/packages/core/src/astro/middleware.ts index 9add416e79..018455ff67 100644 --- a/packages/core/src/astro/middleware.ts +++ b/packages/core/src/astro/middleware.ts @@ -63,7 +63,7 @@ import { createInitLock, type InitLock, initWithLock } from "../utils/init-lock. import type { EmDashConfig } from "./integration/runtime.js"; import { ASTRO_COOKIES_SYMBOL, - deferScopedCloseUntilSettled, + coordinateScopedDbLifecycle, finishScoped, } from "./middleware/scoped-db.js"; import { wrapBodyForStreamMetrics } from "./middleware/stream-end-metrics.js"; @@ -366,26 +366,33 @@ async function runOutsideRequest( // outside a request. Any close-less scope created above is discarded. return fn(runtime); } + const { closed, deferredTasks, lifecycle } = coordinateScopedDbLifecycle(scoped); const parent = getRequestContext(); const ctx = parent - ? { ...parent, db: scoped.db } - : { editMode: false, db: scoped.db, metrics: createRequestMetrics(performance.now()) }; + ? { ...parent, db: scoped.db, deferredTasks } + : { + editMode: false, + db: scoped.db, + metrics: createRequestMetrics(performance.now()), + deferredTasks, + }; try { return await runWithContext(ctx, () => fn(runtime)); } finally { // Guard both so a throw in teardown can't mask the callback result or - // skip close() and leak the connection. Mirrors closeSafely() in scoped-db.ts. + // skip lifecycle settlement and leak the connection. try { - scoped.commit(); + lifecycle.commit(); } catch (error) { console.error("[emdash] event-scoped db commit failed:", error); } try { - scoped.close(); + lifecycle.close?.(); } catch (error) { console.error("[emdash] event-scoped db close failed:", error); } + await closed; } } @@ -687,10 +694,11 @@ export const onRequest = defineMiddleware(async (context, next) => { return wrapBodyForStreamMetrics(finalizeResponse(response, timings)); }; if (anonScoped) { + const { deferredTasks, lifecycle } = coordinateScopedDbLifecycle(anonScoped); const parent = getRequestContext(); const ctx = parent - ? { ...parent, db: anonScoped.db } - : { editMode: false, db: anonScoped.db, metrics }; + ? { ...parent, db: anonScoped.db, deferredTasks } + : { editMode: false, db: anonScoped.db, metrics, deferredTasks }; // Eagerly warm site-global layout data (menus, widget areas, // taxonomy terms, settings) concurrently so the layout's // per-component reads overlap into ~one wall-clock round trip and @@ -710,13 +718,11 @@ export const onRequest = defineMiddleware(async (context, next) => { .trim() .startsWith("text/html"); return runWithContext(ctx, async () => { - const scopedLifecycle = acceptsHtml - ? deferScopedCloseUntilSettled(anonScoped, prefetchLayoutData(), after) - : anonScoped; + if (acceptsHtml) after(() => prefetchLayoutData()); // commit() persists per-request state (e.g. the D1 bookmark cookie) // before the response is returned, even if render throws; close() - // (connection teardown) is deferred to stream-end. See finishScoped. - return finishScoped(scopedLifecycle, runAnon); + // waits for stream-end and request-owned deferred work. See finishScoped. + return finishScoped(lifecycle, runAnon); }); } return runAnon(); @@ -903,15 +909,16 @@ export const onRequest = defineMiddleware(async (context, next) => { }; if (scoped) { + const { deferredTasks, lifecycle } = coordinateScopedDbLifecycle(scoped); const parent = getRequestContext(); const ctx = parent - ? { ...parent, db: scoped.db } - : { editMode: false, db: scoped.db, metrics }; + ? { ...parent, db: scoped.db, deferredTasks } + : { editMode: false, db: scoped.db, metrics, deferredTasks }; return runWithContext(ctx, () => // commit() persists per-request state (e.g. the D1 bookmark cookie) // before the response returns, even if render throws; close() - // (connection teardown) is deferred to stream-end. See finishScoped. - finishScoped(scoped, renderAndFinalize), + // waits for stream-end and request-owned deferred work. See finishScoped. + finishScoped(lifecycle, renderAndFinalize), ); } diff --git a/packages/core/src/astro/middleware/scoped-db.ts b/packages/core/src/astro/middleware/scoped-db.ts index a119419a8d..bb8ec39bbc 100644 --- a/packages/core/src/astro/middleware/scoped-db.ts +++ b/packages/core/src/astro/middleware/scoped-db.ts @@ -6,6 +6,9 @@ * request-scoped db adapter's lifecycle around the response. */ +import { createDeferredTaskTracker } from "../../deferred-tasks.js"; +import type { DeferredTaskTracker } from "../../deferred-tasks.js"; + /** * Astro attaches AstroCookies to outgoing responses via a well-known global * symbol. Cloning a Response (`new Response(body, init)`) drops non-header @@ -21,36 +24,27 @@ interface ScopedDbLifecycle { close?: () => void; } -type DeferredTaskScheduler = (task: () => void | Promise) => void; - -/** - * Keep teardown alive while both the response and request-owned background - * work finish. The returned close callback only signals response completion; - * the adapter's real close runs afterward inside the deferred task. - */ -export function deferScopedCloseUntilSettled( - scoped: ScopedDbLifecycle, - pending: Promise, - defer: DeferredTaskScheduler, -): ScopedDbLifecycle { - if (!scoped.close) { - defer(async () => { - await pending; - }); - return scoped; - } +/** Hold the real adapter close behind both response and deferred-task completion. */ +export function coordinateScopedDbLifecycle(scoped: ScopedDbLifecycle): { + closed?: Promise; + deferredTasks?: DeferredTaskTracker; + lifecycle: ScopedDbLifecycle; +} { + if (!scoped.close) return { lifecycle: scoped }; - let settleResponse!: () => void; - const responseSettled = new Promise((resolve) => { - settleResponse = resolve; - }); const close = scoped.close; - defer(async () => { - await Promise.allSettled([pending, responseSettled]); - close(); + const deferredTasks = createDeferredTaskTracker(() => { + try { + close(); + } catch (error) { + console.error("[emdash] request-scoped db close failed:", error); + } }); - - return { commit: scoped.commit, close: settleResponse }; + return { + closed: deferredTasks.settled, + deferredTasks, + lifecycle: { commit: scoped.commit, close: deferredTasks.settle }, + }; } /** @@ -99,9 +93,9 @@ export function wrapResponseForScopedClose(response: Response, close: () => void * Run the request body under a request-scoped db, then settle its lifecycle: * `commit()` runs before the response is returned (so per-request state like a * D1 bookmark cookie is persisted in the headers, even if render throws), while - * `close()` (if any) is deferred to stream-end so a connection-backed adapter - * isn't torn down mid-render. On error the connection is closed immediately - * before rethrowing so it never leaks. + * `close()` (if any) is deferred to lifecycle settlement so a + * connection-backed adapter isn't torn down mid-render or mid-task. On error + * the lifecycle is settled before rethrowing so it never leaks. * * On the error path both `commit()` and `close()` are defended: a throw from * either is logged and swallowed so it can't replace the propagating render diff --git a/packages/core/src/cleanup.ts b/packages/core/src/cleanup.ts index 36072d8721..0fe84bd7a4 100644 --- a/packages/core/src/cleanup.ts +++ b/packages/core/src/cleanup.ts @@ -10,7 +10,7 @@ */ import { createKyselyAdapter, type AuthTables } from "@emdash-cms/auth/adapters/kysely"; -import { sql, type Kysely } from "kysely"; +import type { Kysely } from "kysely"; import { cleanupExpiredChallenges } from "./auth/challenge-store.js"; import { MediaRepository } from "./database/repositories/media.js"; @@ -36,17 +36,14 @@ export interface CleanupResult { mediaUsageOrphanOccurrences: number; } -/** Max revisions to keep per entry during periodic pruning */ const REVISION_KEEP_COUNT = 50; - -/** Only prune entries that exceed this threshold */ -const REVISION_PRUNE_THRESHOLD = REVISION_KEEP_COUNT; +const REVISION_PRUNE_BATCH_SIZE = 10; /** * Run all system cleanup tasks. * - * Safe to call frequently -- each task is a single DELETE with a WHERE clause, - * so repeated calls with nothing to clean are cheap (no-op queries). + * Safe to call frequently -- each subsystem tolerates repeated calls, and + * repeated calls with nothing to clean are cheap. * * @param db - The database instance * @param storage - Optional storage backend for deleting orphaned files. @@ -136,9 +133,8 @@ export async function runSystemCleanup( console.error("[cleanup] Failed to clean media upload attempts:", error); } - // 5. Revision pruning -- trim entries with excessive revision counts try { - result.revisionsPruned = await pruneExcessiveRevisions(db); + result.revisionsPruned = await pruneQueuedRevisions(db); } catch (error) { console.error("[cleanup] Failed to prune revisions:", error); } @@ -155,31 +151,24 @@ export async function runSystemCleanup( return result; } -/** - * Find entries with more than REVISION_PRUNE_THRESHOLD revisions and prune - * them down to REVISION_KEEP_COUNT. - */ -async function pruneExcessiveRevisions(db: Kysely): Promise { - const entries = await sql<{ collection: string; entry_id: string }>` - SELECT collection, entry_id - FROM revisions - GROUP BY collection, entry_id - HAVING COUNT(*) > ${REVISION_PRUNE_THRESHOLD} - `.execute(db); - - if (entries.rows.length === 0) return 0; - +async function pruneQueuedRevisions(db: Kysely): Promise { + const queued = await db + .selectFrom("_emdash_revision_prune_queue") + .selectAll() + .orderBy("revision_id") + .limit(REVISION_PRUNE_BATCH_SIZE) + .execute(); const revisionRepo = new RevisionRepository(db); let totalPruned = 0; - for (const row of entries.rows) { + for (const row of queued) { try { - const pruned = await revisionRepo.pruneOldRevisions( + totalPruned += await revisionRepo.pruneQueuedEntry( row.collection, row.entry_id, + row.revision_id, REVISION_KEEP_COUNT, ); - totalPruned += pruned; } catch (error) { console.error( `[cleanup] Failed to prune revisions for ${row.collection}/${row.entry_id}:`, diff --git a/packages/core/src/database/migrations/059_revision_prune_queue.ts b/packages/core/src/database/migrations/059_revision_prune_queue.ts new file mode 100644 index 0000000000..65522645c6 --- /dev/null +++ b/packages/core/src/database/migrations/059_revision_prune_queue.ts @@ -0,0 +1,36 @@ +import { sql, type Kysely } from "kysely"; + +const REVISION_KEEP_COUNT = 50; + +export async function up(db: Kysely): Promise { + await db.schema + .createTable("_emdash_revision_prune_queue") + .ifNotExists() + .addColumn("collection", "text", (column) => column.notNull()) + .addColumn("entry_id", "text", (column) => column.notNull()) + .addColumn("revision_id", "text", (column) => column.notNull()) + .addPrimaryKeyConstraint("revision_prune_queue_pk", ["collection", "entry_id"]) + .execute(); + + await db.schema + .createIndex("idx_revision_prune_queue_revision_id") + .ifNotExists() + .on("_emdash_revision_prune_queue") + .column("revision_id") + .execute(); + + await sql` + INSERT INTO _emdash_revision_prune_queue (collection, entry_id, revision_id) + SELECT collection, entry_id, MAX(id) + FROM revisions + WHERE true + GROUP BY collection, entry_id + HAVING COUNT(*) > ${REVISION_KEEP_COUNT} + ON CONFLICT (collection, entry_id) + DO UPDATE SET revision_id = excluded.revision_id + `.execute(db); +} + +export async function down(db: Kysely): Promise { + await db.schema.dropTable("_emdash_revision_prune_queue").ifExists().execute(); +} diff --git a/packages/core/src/database/migrations/runner.ts b/packages/core/src/database/migrations/runner.ts index ceda7dfcbe..6db6207a6d 100644 --- a/packages/core/src/database/migrations/runner.ts +++ b/packages/core/src/database/migrations/runner.ts @@ -61,6 +61,7 @@ import * as m055 from "./055_content_translation_group_locale_index.js"; import * as m056 from "./056_taxonomy_term_sort_order.js"; import * as m057 from "./057_collection_hidden.js"; import * as m058 from "./058_collection_sort_order.js"; +import * as m059 from "./059_revision_prune_queue.js"; const MIGRATIONS: Readonly> = Object.freeze({ "001_initial": m001, @@ -120,6 +121,7 @@ const MIGRATIONS: Readonly> = Object.freeze({ "056_taxonomy_term_sort_order": m056, "057_collection_hidden": m057, "058_collection_sort_order": m058, + "059_revision_prune_queue": m059, }); /** Total number of registered migrations. Exported for use in tests. */ diff --git a/packages/core/src/database/repositories/revision.ts b/packages/core/src/database/repositories/revision.ts index 061896fcaa..aea7e1f4a9 100644 --- a/packages/core/src/database/repositories/revision.ts +++ b/packages/core/src/database/repositories/revision.ts @@ -51,6 +51,26 @@ export class RevisionRepository { if (!revision) { throw new Error("Failed to create revision"); } + + try { + await this.db + .insertInto("_emdash_revision_prune_queue") + .values({ + collection: input.collection, + entry_id: input.entryId, + revision_id: id, + }) + .onConflict((conflict) => + conflict.columns(["collection", "entry_id"]).doUpdateSet({ revision_id: id }), + ) + .execute(); + } catch (error) { + console.error( + `[revisions] Failed to queue revision pruning for ${input.collection}/${input.entryId}:`, + error, + ); + } + return revision; } @@ -133,6 +153,19 @@ export class RevisionRepository { .where("entry_id", "=", entryId) .executeTakeFirst(); + try { + await this.db + .deleteFrom("_emdash_revision_prune_queue") + .where("collection", "=", collection) + .where("entry_id", "=", entryId) + .execute(); + } catch (error) { + console.error( + `[revisions] Failed to clear queued revision pruning for ${collection}/${entryId}:`, + error, + ); + } + return Number(result.numDeletedRows ?? 0); } @@ -172,6 +205,22 @@ export class RevisionRepository { return Number(result.numAffectedRows ?? 0); } + async pruneQueuedEntry( + collection: string, + entryId: string, + queuedRevisionId: string, + keepCount: number, + ): Promise { + const pruned = await this.pruneOldRevisions(collection, entryId, keepCount); + await this.db + .deleteFrom("_emdash_revision_prune_queue") + .where("collection", "=", collection) + .where("entry_id", "=", entryId) + .where("revision_id", "=", queuedRevisionId) + .execute(); + return pruned; + } + async deleteIfUnreferenced( collection: string, entryId: string, diff --git a/packages/core/src/database/types.ts b/packages/core/src/database/types.ts index 0ad701b8de..508eb481f9 100644 --- a/packages/core/src/database/types.ts +++ b/packages/core/src/database/types.ts @@ -13,6 +13,12 @@ export interface RevisionTable { created_at: Generated; } +export interface RevisionPruneQueueTable { + collection: string; + entry_id: string; + revision_id: string; +} + export interface TaxonomyTable { id: string; name: string; @@ -505,6 +511,7 @@ export interface SectionTable { // Note: ec_* content tables are dynamic and not part of this type export interface Database { revisions: RevisionTable; + _emdash_revision_prune_queue: RevisionPruneQueueTable; taxonomies: TaxonomyTable; content_taxonomies: ContentTaxonomyTable; _emdash_taxonomy_defs: TaxonomyDefTable; diff --git a/packages/core/src/deferred-tasks.ts b/packages/core/src/deferred-tasks.ts new file mode 100644 index 0000000000..8626b311e9 --- /dev/null +++ b/packages/core/src/deferred-tasks.ts @@ -0,0 +1,71 @@ +export interface DeferredTaskTracker { + settled: Promise; + track(promise: Promise): Promise; + settle(): void; +} + +// Test database teardown uses this process-local registry to drain tasks that +// run without request middleware. Keep it on globalThis so duplicated SSR +// chunks and their test helpers observe the same task set. +const DEFERRED_TASKS_KEY = Symbol.for("emdash:deferred-tasks"); + +const deferredTasks: Set> = + // eslint-disable-next-line typescript/no-unsafe-type-assertion -- globalThis singleton pattern + ((globalThis as Record)[DEFERRED_TASKS_KEY] as + | Set> + | undefined) ?? + (() => { + const tasks = new Set>(); + (globalThis as Record)[DEFERRED_TASKS_KEY] = tasks; + return tasks; + })(); + +export function trackDeferredTask(promise: Promise): Promise { + let tracked!: Promise; + tracked = promise.finally(() => deferredTasks.delete(tracked)); + deferredTasks.add(tracked); + return tracked; +} + +export async function waitForDeferredTasks(): Promise { + while (deferredTasks.size > 0) { + await Promise.allSettled(deferredTasks); + } +} + +export function createDeferredTaskTracker(onSettled: () => void): DeferredTaskTracker { + let pending = 0; + let responseSettled = false; + let completed = false; + let resolveSettled!: () => void; + const settled = new Promise((resolve) => { + resolveSettled = resolve; + }); + + const completeIfSettled = () => { + if (completed || !responseSettled || pending > 0) return; + completed = true; + try { + onSettled(); + } finally { + resolveSettled(); + } + }; + + return { + settled, + track(promise: Promise): Promise { + pending++; + return promise.finally(() => { + // A task may register another after() call before it settles. The + // counter therefore reaches zero only after the full task chain ends. + pending--; + completeIfSettled(); + }); + }, + settle(): void { + responseSettled = true; + completeIfSettled(); + }, + }; +} diff --git a/packages/core/src/emdash-runtime.ts b/packages/core/src/emdash-runtime.ts index 8b0e7c4acb..895308667c 100644 --- a/packages/core/src/emdash-runtime.ts +++ b/packages/core/src/emdash-runtime.ts @@ -2973,7 +2973,16 @@ export class EmDashRuntime { ); } } else { - void revisionRepo.pruneOldRevisions(collection, resolvedId, 50).catch(() => {}); + after(async () => { + try { + await revisionRepo.pruneQueuedEntry(collection, resolvedId, revision.id, 50); + } catch (error) { + console.error( + `[revisions] Failed to prune revisions for ${collection}/${resolvedId}:`, + error, + ); + } + }); } break; } @@ -3426,10 +3435,21 @@ export class EmDashRuntime { throw error; } - // Fire-and-forget: prune old revisions to prevent unbounded growth - void revisionRepo - .pruneOldRevisions(revision.collection, revision.entryId, 50) - .catch(() => {}); + after(async () => { + try { + await revisionRepo.pruneQueuedEntry( + revision.collection, + revision.entryId, + newDraft.id, + 50, + ); + } catch (error) { + console.error( + `[revisions] Failed to prune revisions for ${revision.collection}/${revision.entryId}:`, + error, + ); + } + }); // Return the freshly-fetched item with the new draft hydrated // onto `data`. Without this the response would echo the live diff --git a/packages/core/src/request-context.ts b/packages/core/src/request-context.ts index 21484d21ba..b50afe1876 100644 --- a/packages/core/src/request-context.ts +++ b/packages/core/src/request-context.ts @@ -20,6 +20,7 @@ import { AsyncLocalStorage } from "node:async_hooks"; import type { QueryRecorder } from "./database/instrumentation.js"; +import type { DeferredTaskTracker } from "./deferred-tasks.js"; /** * Lightweight always-on counters surfaced in Server-Timing. @@ -102,6 +103,8 @@ export interface EmDashRequestContext { * and `requestCached`. */ metrics?: RequestMetrics; + /** Deferred work that owns resources scoped to this request. */ + deferredTasks?: DeferredTaskTracker; } const ALS_KEY = Symbol.for("emdash:request-context"); diff --git a/packages/core/tests/integration/content/revision-retention.test.ts b/packages/core/tests/integration/content/revision-retention.test.ts new file mode 100644 index 0000000000..d3320579a0 --- /dev/null +++ b/packages/core/tests/integration/content/revision-retention.test.ts @@ -0,0 +1,73 @@ +import type { Kysely } from "kysely"; +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; + +const { deferred } = vi.hoisted(() => ({ deferred: [] as Array<() => void | Promise> })); +vi.mock("../../../src/after.js", () => ({ + after: (fn: () => void | Promise) => { + deferred.push(fn); + }, +})); +vi.mock( + "virtual:emdash/object-cache", + () => ({ createObjectCache: undefined, objectCacheConfig: {} }), + { virtual: true }, +); + +import { RevisionRepository } from "../../../src/database/repositories/revision.js"; +import type { Database } from "../../../src/database/types.js"; +import type { EmDashRuntime } from "../../../src/emdash-runtime.js"; +import { SchemaRegistry } from "../../../src/schema/registry.js"; +import { createTestRuntime } from "../../utils/mcp-runtime.js"; +import { setupTestDatabase, teardownTestDatabase } from "../../utils/test-db.js"; + +async function flushDeferred(): Promise { + const tasks = deferred.splice(0); + for (const task of tasks) await task(); +} + +describe("revision retention without scheduled cleanup", () => { + let db: Kysely; + let runtime: EmDashRuntime; + + beforeEach(async () => { + deferred.length = 0; + db = await setupTestDatabase(); + const registry = new SchemaRegistry(db); + await registry.createCollection({ slug: "posts", label: "Posts" }); + await registry.createField("posts", { slug: "title", label: "Title", type: "string" }); + runtime = createTestRuntime(db); + }); + + afterEach(async () => { + deferred.length = 0; + await teardownTestDatabase(db); + }); + + it("bounds history through request-lifetime work without a scheduled tick", async () => { + const created = await runtime.handleContentCreate("posts", { + data: { title: "Initial" }, + slug: "bounded-history", + }); + expect(created.success).toBe(true); + const entryId = created.data!.item.id; + const revisions = new RevisionRepository(db); + + for (let index = 0; index < 50; index++) { + await revisions.create({ + collection: "posts", + entryId, + data: { title: `Version ${index}` }, + }); + } + + const saved = await runtime.handleContentUpdate("posts", entryId, { + data: { title: "Version 50" }, + }); + expect(saved.success).toBe(true); + expect(await revisions.countByEntry("posts", entryId)).toBe(51); + + await flushDeferred(); + + expect(await revisions.countByEntry("posts", entryId)).toBe(50); + }); +}); diff --git a/packages/core/tests/integration/database/migrations.test.ts b/packages/core/tests/integration/database/migrations.test.ts index 1bb3036221..a834bfbdd5 100644 --- a/packages/core/tests/integration/database/migrations.test.ts +++ b/packages/core/tests/integration/database/migrations.test.ts @@ -145,6 +145,7 @@ describe("Database Migrations (Integration)", () => { "056_taxonomy_term_sort_order", "057_collection_hidden", "058_collection_sort_order", + "059_revision_prune_queue", ]; await db.deleteFrom("_emdash_migrations").where("name", "in", trailing).execute(); diff --git a/packages/core/tests/unit/after.test.ts b/packages/core/tests/unit/after.test.ts index c23f47a17b..1cfe068f95 100644 --- a/packages/core/tests/unit/after.test.ts +++ b/packages/core/tests/unit/after.test.ts @@ -11,6 +11,7 @@ vi.mock("virtual:emdash/wait-until", () => ({ waitUntil: undefined }), { virtual // behaviors that don't need a real waitUntil: the callback fires, // errors don't escape, and the caller doesn't block. import { after } from "../../src/after.js"; +import { waitForDeferredTasks } from "../../src/deferred-tasks.js"; describe("after()", () => { it("runs the callback", async () => { @@ -53,4 +54,24 @@ describe("after()", () => { // after() returned already — the callback hasn't completed. expect(ran).toBe(false); }); + + it("exposes deferred work for lifecycle teardown", async () => { + let settle!: () => void; + let completed = false; + after(async () => { + await new Promise((resolve) => { + settle = resolve; + }); + completed = true; + }); + await Promise.resolve(); + + const drained = waitForDeferredTasks(); + await Promise.resolve(); + expect(completed).toBe(false); + + settle(); + await drained; + expect(completed).toBe(true); + }); }); diff --git a/packages/core/tests/unit/astro/with-emdash-runtime.test.ts b/packages/core/tests/unit/astro/with-emdash-runtime.test.ts index 0f5d9f8dd4..1baa69dcc2 100644 --- a/packages/core/tests/unit/astro/with-emdash-runtime.test.ts +++ b/packages/core/tests/unit/astro/with-emdash-runtime.test.ts @@ -68,6 +68,7 @@ vi.mock("../../../src/object-cache/index.js", async (importOriginal) => ({ import { createRequestScopedDb } from "virtual:emdash/dialect"; +import { after } from "../../../src/after.js"; import { withEmDashRuntime } from "../../../src/astro/middleware.js"; import { getRequestContext } from "../../../src/request-context.js"; @@ -150,6 +151,39 @@ describe("withEmDashRuntime (#1887)", () => { expect(close).toHaveBeenCalledTimes(1); }); + it("does not finish an event before deferred work releases its scoped db", async () => { + let release!: () => void; + let returned = false; + const close = vi.fn(); + vi.mocked(createRequestScopedDb).mockReturnValue({ + db: { _marker: "scoped" } as never, + commit: vi.fn(), + close, + }); + + const resultPromise = withEmDashRuntime(async () => { + after( + () => + new Promise((resolve) => { + release = resolve; + }), + ); + return "ok"; + }).then((result) => { + returned = true; + return result; + }); + + await vi.waitFor(() => expect(release).toBeTypeOf("function")); + await Promise.resolve(); + expect(returned).toBe(false); + expect(close).not.toHaveBeenCalled(); + + release(); + await expect(resultPromise).resolves.toBe("ok"); + expect(close).toHaveBeenCalledTimes(1); + }); + it("uses the singleton path for a close-less scope", async () => { const commit = vi.fn(); vi.mocked(createRequestScopedDb).mockReturnValue({ diff --git a/packages/core/tests/unit/cleanup.test.ts b/packages/core/tests/unit/cleanup.test.ts index 2e3bfb4514..6c72a0fde8 100644 --- a/packages/core/tests/unit/cleanup.test.ts +++ b/packages/core/tests/unit/cleanup.test.ts @@ -1,19 +1,12 @@ -/** - * Tests for the cleanup subsystems. - * - * Note: runSystemCleanup() is not tested directly here because it imports - * from @emdash-cms/auth/adapters/kysely, which requires the auth package to - * be built. Instead, we test each subsystem independently: - * - cleanupExpiredChallenges: tested in auth/challenge-store.test.ts - * - deleteExpiredTokens: tested below using direct DB operations - * - cleanupPendingUploads: tested below via MediaRepository - * - pruneOldRevisions: tested below via RevisionRepository - */ - -import type { Kysely } from "kysely"; +/** Tests for cleanup subsystems and scheduled cleanup orchestration. */ + +import BetterSqlite3 from "better-sqlite3"; +import { Kysely, SqliteDialect } from "kysely"; import { ulid } from "ulidx"; import { describe, it, expect, beforeEach, afterEach, vi } from "vitest"; +import { runSystemCleanup } from "../../src/cleanup.js"; +import { runMigrations } from "../../src/database/migrations/runner.js"; import { MediaRepository } from "../../src/database/repositories/media.js"; import { RevisionRepository } from "../../src/database/repositories/revision.js"; import type { Database } from "../../src/database/types.js"; @@ -132,6 +125,96 @@ describe("Revision Pruning", () => { }); }); +describe("Scheduled system cleanup", () => { + it("prunes revision entries queued by revision writes", async () => { + const db = await setupTestDatabaseWithCollections(); + const revisionRepo = new RevisionRepository(db); + const entryId = ulid(); + const { sql } = await import("kysely"); + + try { + await sql` + INSERT INTO ec_post (id, slug, status, created_at, updated_at, version) + VALUES (${entryId}, ${"queued-history"}, ${"draft"}, ${new Date().toISOString()}, ${new Date().toISOString()}, ${1}) + `.execute(db); + for (let i = 0; i < 51; i++) { + await revisionRepo.create({ + collection: "post", + entryId, + data: { title: `Version ${i + 1}` }, + }); + } + + const result = await runSystemCleanup(db); + + expect(result.revisionsPruned).toBe(1); + expect(await revisionRepo.countByEntry("post", entryId)).toBe(50); + expect( + await db.selectFrom("_emdash_revision_prune_queue").select("entry_id").execute(), + ).toEqual([]); + } finally { + await db.destroy(); + } + }); + + it("processes the oldest queued revision entries first", async () => { + const db = await setupTestDatabaseWithCollections(); + const revisionRepo = new RevisionRepository(db); + const { sql } = await import("kysely"); + const entryIds = ["z-old", ...Array.from({ length: 10 }, (_, index) => `a0${index}`)]; + + try { + for (const entryId of entryIds) { + await sql` + INSERT INTO ec_post (id, slug, status, created_at, updated_at, version) + VALUES (${entryId}, ${entryId}, ${"draft"}, ${new Date().toISOString()}, ${new Date().toISOString()}, ${1}) + `.execute(db); + await revisionRepo.create({ + collection: "post", + entryId, + data: { title: entryId }, + }); + } + + await runSystemCleanup(db); + + const remaining = await db + .selectFrom("_emdash_revision_prune_queue") + .select("entry_id") + .execute(); + expect(remaining).toHaveLength(1); + expect(remaining[0]?.entry_id).toMatch(/^a/); + expect(remaining).not.toContainEqual({ entry_id: "z-old" }); + } finally { + await db.destroy(); + } + }); + + it("does not query the full revision history", async () => { + const sqlite = new BetterSqlite3(":memory:"); + const queries: string[] = []; + const db = new Kysely({ + dialect: new SqliteDialect({ database: sqlite }), + log: (event) => { + if (event.level === "query") queries.push(event.query.sql); + }, + }); + + try { + await runMigrations(db); + queries.length = 0; + await runSystemCleanup(db); + + const revisionQueries = queries.filter((query) => + /\b(?:from|delete from)\s+["`]?revisions\b/i.test(query), + ); + expect(revisionQueries).toHaveLength(0); + } finally { + await db.destroy(); + } + }); +}); + describe("MediaRepository.cleanupPendingUploads", () => { let db: Kysely; let mediaRepo: MediaRepository; diff --git a/packages/core/tests/unit/database/migrations/059_revision_prune_queue.test.ts b/packages/core/tests/unit/database/migrations/059_revision_prune_queue.test.ts new file mode 100644 index 0000000000..032f6d598e --- /dev/null +++ b/packages/core/tests/unit/database/migrations/059_revision_prune_queue.test.ts @@ -0,0 +1,43 @@ +import BetterSqlite3 from "better-sqlite3"; +import { Kysely, sql, SqliteDialect } from "kysely"; +import { afterEach, describe, expect, it } from "vitest"; + +import { down, up } from "../../../../src/database/migrations/059_revision_prune_queue.js"; + +describe("059_revision_prune_queue migration", () => { + let db: Kysely> | undefined; + + afterEach(async () => { + await db?.destroy(); + }); + + it("queues existing excessive histories and can be retried", async () => { + const sqlite = new BetterSqlite3(":memory:"); + sqlite.exec(` + CREATE TABLE revisions ( + id TEXT PRIMARY KEY, + collection TEXT NOT NULL, + entry_id TEXT NOT NULL + ) + `); + const insert = sqlite.prepare( + "INSERT INTO revisions (id, collection, entry_id) VALUES (?, ?, ?)", + ); + for (let i = 0; i < 51; i++) insert.run(String(i).padStart(2, "0"), "post", "entry-1"); + + db = new Kysely>({ + dialect: new SqliteDialect({ database: sqlite }), + }); + + await up(db); + await up(db); + + const queued = await sql<{ collection: string; entry_id: string; revision_id: string }>` + SELECT collection, entry_id, revision_id + FROM _emdash_revision_prune_queue + `.execute(db); + expect(queued.rows).toEqual([{ collection: "post", entry_id: "entry-1", revision_id: "50" }]); + + await down(db); + }); +}); diff --git a/packages/core/tests/unit/middleware/scoped-db.test.ts b/packages/core/tests/unit/middleware/scoped-db.test.ts index 6884e4f993..65c7c97cee 100644 --- a/packages/core/tests/unit/middleware/scoped-db.test.ts +++ b/packages/core/tests/unit/middleware/scoped-db.test.ts @@ -1,11 +1,13 @@ import { describe, it, expect, vi } from "vitest"; +import { after } from "../../../src/after.js"; import { ASTRO_COOKIES_SYMBOL, - deferScopedCloseUntilSettled, + coordinateScopedDbLifecycle, finishScoped, wrapResponseForScopedClose, } from "../../../src/astro/middleware/scoped-db.js"; +import { runWithContext } from "../../../src/request-context.js"; /** Build a streaming Response whose body emits the given chunks. */ function streamingResponse(chunks: string[], init?: ResponseInit): Response { @@ -103,65 +105,113 @@ describe("wrapResponseForScopedClose", () => { expect(wrapped.headers.has("content-length")).toBe(false); }); }); - -describe("deferScopedCloseUntilSettled", () => { - it("keeps a bodyless response teardown behind in-flight request work", async () => { +describe("finishScoped", () => { + it("keeps a bodyless response connection open for request deferred work", async () => { let settleWork!: () => void; - const inFlightWork = new Promise((resolve) => { - settleWork = resolve; - }); const close = vi.fn(); - let deferredTask!: Promise; - const scoped = deferScopedCloseUntilSettled( - { commit: vi.fn(), close }, - inFlightWork, - (task) => { - deferredTask = Promise.resolve().then(task); - }, - ); + const { deferredTasks, lifecycle } = coordinateScopedDbLifecycle({ + commit: vi.fn(), + close, + }); - await finishScoped(scoped, async () => new Response(null, { status: 204 })); + await runWithContext({ editMode: false, deferredTasks }, async () => { + after(async () => { + await new Promise((resolve) => { + settleWork = resolve; + }); + }); + await Promise.resolve(); + await finishScoped(lifecycle, async () => new Response(null, { status: 204 })); + }); expect(close).not.toHaveBeenCalled(); settleWork(); - await deferredTask; - expect(close).toHaveBeenCalledTimes(1); + await vi.waitFor(() => expect(close).toHaveBeenCalledTimes(1)); }); - it("keeps teardown behind a streaming response after request work settles", async () => { + it("waits for stream completion after deferred work settles", async () => { + let settleWork!: () => void; const close = vi.fn(); - let deferredTask!: Promise; - const scoped = deferScopedCloseUntilSettled( - { commit: vi.fn(), close }, - Promise.resolve(), - (task) => { - deferredTask = Promise.resolve().then(task); - }, - ); - - const response = await finishScoped(scoped, async () => streamingResponse(["body"])); - await Promise.resolve(); + const { deferredTasks, lifecycle } = coordinateScopedDbLifecycle({ + commit: vi.fn(), + close, + }); + const response = await runWithContext({ editMode: false, deferredTasks }, async () => { + after( + () => + new Promise((resolve) => { + settleWork = resolve; + }), + ); + await Promise.resolve(); + return finishScoped(lifecycle, async () => streamingResponse(["body"])); + }); + settleWork(); + await Promise.resolve(); expect(close).not.toHaveBeenCalled(); + await drain(response); - await deferredTask; expect(close).toHaveBeenCalledTimes(1); }); - it("keeps background work deferred for an adapter without teardown", async () => { - const scoped = { commit: vi.fn() }; - let deferredTask!: () => void | Promise; + it("waits for deferred work after the response stream is cancelled", async () => { + let settleWork!: () => void; + const close = vi.fn(); + const { deferredTasks, lifecycle } = coordinateScopedDbLifecycle({ + commit: vi.fn(), + close, + }); + const response = await runWithContext({ editMode: false, deferredTasks }, async () => { + after( + () => + new Promise((resolve) => { + settleWork = resolve; + }), + ); + await Promise.resolve(); + return finishScoped(lifecycle, async () => streamingResponse(["body"])); + }); + + await response.body!.cancel(); + expect(close).not.toHaveBeenCalled(); + + settleWork(); + await vi.waitFor(() => expect(close).toHaveBeenCalledTimes(1)); + }); - const result = deferScopedCloseUntilSettled(scoped, Promise.resolve(), (task) => { - deferredTask = task; + it("waits for nested deferred work after render failure", async () => { + let settleNested!: () => void; + const close = vi.fn(); + const renderError = new Error("render failed"); + const { deferredTasks, lifecycle } = coordinateScopedDbLifecycle({ + commit: vi.fn(), + close, }); - expect(result).toBe(scoped); - await deferredTask(); + await expect( + runWithContext({ editMode: false, deferredTasks }, async () => { + after(async () => { + after( + () => + new Promise((resolve) => { + settleNested = resolve; + }), + ); + }); + await Promise.resolve(); + return finishScoped(lifecycle, async () => { + throw renderError; + }); + }), + ).rejects.toBe(renderError); + await vi.waitFor(() => expect(settleNested).toBeTypeOf("function")); + expect(close).not.toHaveBeenCalled(); + + settleNested(); + await vi.waitFor(() => expect(close).toHaveBeenCalledTimes(1)); }); -}); -describe("finishScoped", () => { it("commits then defers close to stream-end for a streaming response", async () => { const order: string[] = []; const commit = vi.fn(() => order.push("commit")); diff --git a/packages/core/tests/unit/test-db-deferred-teardown.test.ts b/packages/core/tests/unit/test-db-deferred-teardown.test.ts new file mode 100644 index 0000000000..d746a91732 --- /dev/null +++ b/packages/core/tests/unit/test-db-deferred-teardown.test.ts @@ -0,0 +1,29 @@ +import { sql } from "kysely"; +import { expect, it } from "vitest"; + +import { after } from "../../src/after.js"; +import { setupTestDatabase, teardownTestDatabase } from "../utils/test-db.js"; + +it("drains deferred database work before test database teardown", async () => { + const db = await setupTestDatabase(); + let release!: () => void; + after(async () => { + await new Promise((resolve) => { + release = resolve; + }); + await sql`select 1`.execute(db); + }); + await Promise.resolve(); + + let tornDown = false; + const teardown = teardownTestDatabase(db).then(() => { + tornDown = true; + return null; + }); + await Promise.resolve(); + expect(tornDown).toBe(false); + + release(); + await teardown; + expect(tornDown).toBe(true); +}); diff --git a/packages/core/tests/utils/test-db.ts b/packages/core/tests/utils/test-db.ts index c5dace36a9..3258ddf255 100644 --- a/packages/core/tests/utils/test-db.ts +++ b/packages/core/tests/utils/test-db.ts @@ -10,6 +10,7 @@ import { getMigrationStatus, runMigrations } from "../../src/database/migrations import type { MigrationStatus } from "../../src/database/migrations/runner.js"; import { FailFastPostgresDialect } from "../../src/database/pg-migration-lock.js"; import type { Database as DatabaseSchema } from "../../src/database/types.js"; +import { waitForDeferredTasks } from "../../src/deferred-tasks.js"; import { SchemaRegistry } from "../../src/schema/registry.js"; import { resetTaxonomyDefsCacheForTests } from "../../src/taxonomies/index.js"; @@ -119,6 +120,7 @@ export async function setupTestDatabaseWithCollections(): Promise): Promise { + await waitForDeferredTasks(); await db.destroy(); } @@ -418,6 +420,7 @@ export async function setupTestPostgresDatabaseWithCollections(): Promise { + await waitForDeferredTasks(); // Destroy the test pool first await ctx.db.destroy();