From 3384d163b06d9839137fadb4fae87adb51474ad5 Mon Sep 17 00:00:00 2001 From: Saurabh Verma Date: Sun, 6 Sep 2026 12:48:14 +0530 Subject: [PATCH] fix(core): reclaim deleted database pages incrementally Signed-off-by: Saurabh Verma --- packages/core/src/database/database.ts | 8 +- packages/core/src/database/storage.ts | 29 +++ packages/core/src/event.ts | 6 +- packages/core/test/database-storage.test.ts | 195 ++++++++++++++++++++ 4 files changed, 235 insertions(+), 3 deletions(-) create mode 100644 packages/core/src/database/storage.ts create mode 100644 packages/core/test/database-storage.test.ts diff --git a/packages/core/src/database/database.ts b/packages/core/src/database/database.ts index d61adf047eac..915b81585684 100644 --- a/packages/core/src/database/database.ts +++ b/packages/core/src/database/database.ts @@ -7,6 +7,7 @@ import { Global } from "../global" import { Flag } from "../flag/flag" import { isAbsolute, join } from "path" import { DatabaseMigration } from "./migration" +import { DatabaseStorage } from "./storage" import { InstallationChannel } from "../installation/version" import { makeGlobalNode } from "../effect/app-node" @@ -24,20 +25,23 @@ const layer = Layer.effect( Effect.gen(function* () { const db = yield* makeDatabase + yield* db.run("PRAGMA busy_timeout = 5000") + yield* DatabaseStorage.configure(db) yield* db.run("PRAGMA journal_mode = WAL") yield* db.run("PRAGMA synchronous = NORMAL") - yield* db.run("PRAGMA busy_timeout = 5000") yield* db.run("PRAGMA cache_size = -64000") yield* db.run("PRAGMA foreign_keys = ON") yield* db.run("PRAGMA wal_checkpoint(PASSIVE)") yield* DatabaseMigration.apply(db) + yield* DatabaseStorage.reclaim(db) return { db } }).pipe(Effect.orDie), ) export function layerFromPath(filename: string) { - return layer.pipe(Layer.provide(sqliteLayer({ filename }))) + // Configure new files before enabling WAL, which initializes the database header. + return layer.pipe(Layer.provide(sqliteLayer({ filename, disableWAL: true }))) } export function path() { diff --git a/packages/core/src/database/storage.ts b/packages/core/src/database/storage.ts new file mode 100644 index 000000000000..6530a6a40f29 --- /dev/null +++ b/packages/core/src/database/storage.ts @@ -0,0 +1,29 @@ +export * as DatabaseStorage from "./storage" + +import type { EffectDrizzleSqlite } from "@opencode-ai/effect-drizzle-sqlite" +import { sql } from "drizzle-orm" +import { Effect } from "effect" + +export function configure(db: EffectDrizzleSqlite.EffectSQLiteDatabase) { + return Effect.gen(function* () { + const mode = yield* db.get<{ auto_vacuum: number }>(sql`PRAGMA auto_vacuum`) + if (mode?.auto_vacuum !== 0) return + if (yield* db.get(sql`SELECT 1 FROM sqlite_schema LIMIT 1`)) return + // Existing databases require an offline VACUUM to change modes. Never + // rebuild them on startup while other processes may be using them. + yield* db.run("PRAGMA auto_vacuum = INCREMENTAL") + }) +} + +export function reclaim(db: EffectDrizzleSqlite.EffectSQLiteDatabase) { + return Effect.gen(function* () { + const mode = yield* db.get<{ auto_vacuum: number }>(sql`PRAGMA auto_vacuum`) + if (mode?.auto_vacuum !== 2) return + const free = yield* db.get<{ freelist_count: number }>(sql`PRAGMA freelist_count`) + if (!free?.freelist_count) return + // Bound each pass rather than draining an arbitrarily large freelist. + // Checkpointing still depends on readers releasing their WAL snapshots. + yield* db.run("PRAGMA incremental_vacuum(256)") + yield* db.run("PRAGMA wal_checkpoint(PASSIVE)") + }).pipe(Effect.catch((error) => Effect.logWarning("Failed to reclaim database pages", error))) +} diff --git a/packages/core/src/event.ts b/packages/core/src/event.ts index c92ac0ac2ce3..8b750379aa50 100644 --- a/packages/core/src/event.ts +++ b/packages/core/src/event.ts @@ -5,6 +5,7 @@ import { Event } from "@opencode-ai/schema/event" import type { Data, Definition, Payload } from "@opencode-ai/schema/event" import { and, asc, eq, gt, inArray } from "drizzle-orm" import { Database } from "./database/database" +import { DatabaseStorage } from "./database/storage" import { EventSequenceTable, EventTable } from "./event/sql" import { Location } from "./location" import { makeGlobalNode } from "./effect/app-node" @@ -519,7 +520,10 @@ export const layerWith = (options?: LayerOptions) => yield* db.delete(EventTable).where(eq(EventTable.aggregate_id, aggregateID)).run() }), ) - .pipe(Effect.orDie) + .pipe( + Effect.orDie, + Effect.tap(() => DatabaseStorage.reclaim(db)), + ) } function claim(aggregateID: string, ownerID: string) { diff --git a/packages/core/test/database-storage.test.ts b/packages/core/test/database-storage.test.ts new file mode 100644 index 000000000000..8b48a19da202 --- /dev/null +++ b/packages/core/test/database-storage.test.ts @@ -0,0 +1,195 @@ +import { describe, expect, test } from "bun:test" +import path from "node:path" +import { sql } from "drizzle-orm" +import { Effect, Schema, type Scope } from "effect" +import { Database } from "@opencode-ai/core/database/database" +import { DatabaseStorage } from "@opencode-ai/core/database/storage" +import { AppNodeBuilder } from "@opencode-ai/core/effect/app-node-builder" +import { LayerNode } from "@opencode-ai/core/effect/layer-node" +import { EventV2 } from "@opencode-ai/core/event" +import { tmpdir } from "./fixture/tmpdir" + +const sqlite = await import("bun:sqlite") + +const run = (filename: string, effect: Effect.Effect) => + Effect.runPromise( + effect.pipe( + Effect.provide( + AppNodeBuilder.build(LayerNode.group([Database.node, EventV2.node]), [ + [Database.node, Database.layerFromPath(filename)], + ]), + ), + Effect.scoped, + ), + ) + +async function fixture(filename: string, mode: "NONE" | "FULL" | "INCREMENTAL") { + await run(filename, Effect.void) + const db = new sqlite.Database(filename) + db.run(`PRAGMA auto_vacuum = ${mode}`) + db.run("VACUUM") + db.run("CREATE TABLE storage_fixture (id TEXT PRIMARY KEY, data BLOB)") + db.run("INSERT INTO storage_fixture VALUES ('keep', 'retained')") + db.run("INSERT INTO storage_fixture VALUES ('remove', zeroblob(4194304))") + db.run("DELETE FROM storage_fixture WHERE id = 'remove'") + db.run("PRAGMA wal_checkpoint(TRUNCATE)") + db.close() +} + +describe("database storage", () => { + test("new databases can reclaim deleted pages incrementally", async () => { + await using tmp = await tmpdir() + await run( + path.join(tmp.path, "new.db"), + Effect.gen(function* () { + const database = yield* Database.Service + expect(yield* database.db.get(sql`PRAGMA auto_vacuum`)).toEqual({ auto_vacuum: 2 }) + }), + ) + }) + + test("reopening an incremental database reclaims a bounded number of free pages", async () => { + await using tmp = await tmpdir() + const filename = path.join(tmp.path, "incremental.db") + await fixture(filename, "INCREMENTAL") + const db = new sqlite.Database(filename) + const before = db.query<{ freelist_count: number }, []>("PRAGMA freelist_count").get()!.freelist_count + db.close() + const size = Bun.file(filename).size + expect(before).toBeGreaterThan(256) + + await run( + filename, + Effect.gen(function* () { + const database = yield* Database.Service + const after = yield* database.db.get<{ freelist_count: number }>(sql`PRAGMA freelist_count`) + expect(before - after!.freelist_count).toBeGreaterThan(0) + expect(before - after!.freelist_count).toBeLessThanOrEqual(256) + expect(after!.freelist_count).toBeGreaterThan(0) + expect(yield* database.db.all(sql`SELECT id, data FROM storage_fixture`)).toEqual([ + { id: "keep", data: "retained" }, + ]) + expect(yield* database.db.get(sql`PRAGMA integrity_check`)).toEqual({ integrity_check: "ok" }) + }), + ) + expect(Bun.file(filename).size).toBeLessThan(size) + }) + + test.each(["NONE", "FULL"] as const)("preserves an existing %s database without rebuilding it", async (mode) => { + await using tmp = await tmpdir() + const filename = path.join(tmp.path, "existing.db") + await fixture(filename, mode) + const db = new sqlite.Database(filename) + const before = db.query("PRAGMA freelist_count").get() + db.close() + const size = Bun.file(filename).size + + await run( + filename, + Effect.gen(function* () { + const database = yield* Database.Service + expect(yield* database.db.get(sql`PRAGMA auto_vacuum`)).toEqual({ auto_vacuum: mode === "NONE" ? 0 : 1 }) + expect(yield* database.db.get(sql`PRAGMA freelist_count`)).toEqual(before) + expect(yield* database.db.all(sql`SELECT id, data FROM storage_fixture`)).toEqual([ + { id: "keep", data: "retained" }, + ]) + }), + ) + expect(Bun.file(filename).size).toBe(size) + }) + + test("removing an aggregate reclaims pages while preserving other durable events", async () => { + await using tmp = await tmpdir() + const filename = path.join(tmp.path, "events.db") + await fixture(filename, "INCREMENTAL") + const event = EventV2.define({ + type: "test.storage", + durable: { version: 1, aggregate: "id" }, + schema: { id: Schema.String, text: Schema.String }, + }) + await run( + filename, + Effect.gen(function* () { + const database = yield* Database.Service + const events = yield* EventV2.Service + yield* events.publish(event, { id: "remove", text: "x".repeat(4 * 1024 * 1024) }) + yield* events.publish(event, { id: "keep", text: "retained" }) + const before = yield* database.db.get<{ page_count: number }>(sql`PRAGMA page_count`) + yield* events.remove("remove") + const after = yield* database.db.get<{ page_count: number }>(sql`PRAGMA page_count`) + expect(after!.page_count).toBeLessThan(before!.page_count) + const rows = yield* database.db.all<{ aggregate_id: string; data: string }>( + sql`SELECT aggregate_id, data FROM event`, + ) + expect(rows.map((row) => ({ id: row.aggregate_id, data: JSON.parse(row.data) }))).toEqual([ + { id: "keep", data: { id: "keep", text: "retained" } }, + ]) + expect(yield* database.db.all(sql`SELECT aggregate_id FROM event_sequence`)).toEqual([{ aggregate_id: "keep" }]) + expect(yield* database.db.get(sql`PRAGMA integrity_check`)).toEqual({ integrity_check: "ok" }) + }), + ) + }) + + test("a competing writer defers maintenance without failing the caller", async () => { + await using tmp = await tmpdir() + const filename = path.join(tmp.path, "busy.db") + await fixture(filename, "INCREMENTAL") + await run( + filename, + Effect.gen(function* () { + const database = yield* Database.Service + yield* database.db.run("PRAGMA busy_timeout = 0") + const before = yield* database.db.get<{ freelist_count: number }>(sql`PRAGMA freelist_count`) + expect(before!.freelist_count).toBeGreaterThan(0) + const writer = new sqlite.Database(filename) + yield* Effect.addFinalizer(() => + Effect.sync(() => { + writer.run("ROLLBACK") + writer.close() + }), + ) + writer.run("BEGIN IMMEDIATE") + yield* DatabaseStorage.reclaim(database.db) + expect(yield* database.db.get(sql`PRAGMA freelist_count`)).toEqual(before) + expect(yield* database.db.all(sql`SELECT id, data FROM storage_fixture`)).toEqual([ + { id: "keep", data: "retained" }, + ]) + }), + ) + }) + + test("a WAL reader keeps its snapshot while reclamation makes bounded progress", async () => { + await using tmp = await tmpdir() + const filename = path.join(tmp.path, "reader.db") + await fixture(filename, "INCREMENTAL") + const reader = new sqlite.Database(filename) + reader.run("BEGIN") + expect(reader.query("SELECT id, data FROM storage_fixture").all()).toEqual([{ id: "keep", data: "retained" }]) + try { + await run( + filename, + Effect.gen(function* () { + const database = yield* Database.Service + const before = yield* database.db.get<{ freelist_count: number }>(sql`PRAGMA freelist_count`) + yield* DatabaseStorage.reclaim(database.db) + const after = yield* database.db.get<{ freelist_count: number }>(sql`PRAGMA freelist_count`) + expect(before!.freelist_count - after!.freelist_count).toBe(256) + expect(reader.query("SELECT id, data FROM storage_fixture").all()).toEqual([{ id: "keep", data: "retained" }]) + }), + ) + } finally { + reader.run("ROLLBACK") + reader.close() + } + await run( + filename, + Effect.gen(function* () { + const database = yield* Database.Service + expect(yield* database.db.get(sql`PRAGMA integrity_check`)).toEqual({ integrity_check: "ok" }) + expect(yield* database.db.all(sql`SELECT id, data FROM storage_fixture`)).toEqual([ + { id: "keep", data: "retained" }, + ]) + }), + ) + }) +})