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" },
+ ])
+ }),
+ )
+ })
+})