diff --git a/README.md b/README.md index 5931b1c..43e17bc 100644 --- a/README.md +++ b/README.md @@ -243,6 +243,8 @@ curl --location 'https://starbasedb.YOUR-ID-HERE.workers.dev/export/dump' \ +Internal database dumps use one temporary-chunk retention policy. Chunks remain in Durable Object storage and, when configured, in R2 while a dump is active. After the R2 multipart object is finalized successfully, both copies of the temporary chunks are deleted. Without R2, the Durable Object chunk records remain the completed job's downloadable artifact. Large internal requests automatically use the resumable job flow; `?job=1` remains supported for explicit job requests. +
diff --git a/src/do.ts b/src/do.ts
index b6bb2b6..816de25 100644
--- a/src/do.ts
+++ b/src/do.ts
@@ -1,4 +1,7 @@
import { DurableObject } from 'cloudflare:workers'
+import { DUMP_STATE_KEY, type DumpState } from './export/chunkedDump'
+import type { DataSource } from './types'
+import type { StarbaseDBConfiguration } from './handler'
export class StarbaseDBDurableObject extends DurableObject {
// Durable storage for the SQL database
@@ -72,6 +75,8 @@ export class StarbaseDBDurableObject extends DurableObject {
deleteAlarm: this.deleteAlarm.bind(this),
getStatistics: this.getStatistics.bind(this),
executeQuery: this.executeQuery.bind(this),
+ startDumpJob: this.startDumpJob.bind(this),
+ dumpJobStatus: this.dumpJobStatus.bind(this),
}
}
@@ -106,6 +111,55 @@ export class StarbaseDBDurableObject extends DurableObject {
async alarm() {
try {
+ // A chunked dump job in flight takes priority: resume its next
+ // bounded cycle (breathing interval already elapsed). The dump job
+ // always targets the internal source, which is this DO itself.
+ const dumpState = await this.storage.get(DUMP_STATE_KEY)
+ if (dumpState && !dumpState.completedAt) {
+ const { runDumpJob } = await import('./export/dump')
+ await runDumpJob(
+ {
+ storage: this.storage,
+ env: {
+ R2_DUMP_BUCKET: (
+ this.env as Env & { R2_DUMP_BUCKET?: R2Bucket }
+ ).R2_DUMP_BUCKET,
+ },
+ dataSource: this.dumpJobDataSource(),
+ config: { role: 'admin' } as StarbaseDBConfiguration,
+ setAlarm: (time, options) =>
+ this.setAlarm(time, options),
+ },
+ new URLSearchParams()
+ )
+ return
+ }
+
+ // A finished dump that has not been consolidated yet: merge its
+ // chunk records into the single R2 object (multipart, resumable)
+ // so presigned download URLs become available.
+ if (
+ dumpState &&
+ dumpState.completedAt &&
+ (this.env as Env & { R2_DUMP_BUCKET?: R2Bucket })
+ .R2_DUMP_BUCKET &&
+ (!dumpState.finalizedAt || !dumpState.temporaryChunksCleanedAt)
+ ) {
+ const { runDumpFinalize } = await import('./export/dump')
+ await runDumpFinalize({
+ storage: this.storage,
+ env: {
+ R2_DUMP_BUCKET: (
+ this.env as Env & { R2_DUMP_BUCKET?: R2Bucket }
+ ).R2_DUMP_BUCKET,
+ },
+ dataSource: this.dumpJobDataSource(),
+ config: { role: 'admin' } as StarbaseDBConfiguration,
+ setAlarm: (time, options) => this.setAlarm(time, options),
+ })
+ return
+ }
+
// Fetch all the tasks that are marked to emit an event for this cycle.
const task = (await this.executeQuery({
sql: 'SELECT * FROM tmp_cron_tasks WHERE is_active = 1;',
@@ -284,6 +338,64 @@ export class StarbaseDBDurableObject extends DurableObject {
}
}
+ /**
+ * Internal data source for dump jobs. Built in-process (no RPC hop): the
+ * engine's queries execute directly against this DO's SQLite storage.
+ */
+ private dumpJobDataSource(): DataSource {
+ return {
+ source: 'internal',
+ rpc: this.init(),
+ } as unknown as DataSource
+ }
+
+ /**
+ * Chunked dump job entry point, executed inside the Durable Object so it
+ * can drive bounded work cycles, persist progress, mirror chunks to R2
+ * (when bound) and resume itself through the DO alarm. Takes only plain
+ * config/params so the RPC surface avoids circular DataSource typings.
+ */
+ public async startDumpJob(
+ config: StarbaseDBConfiguration,
+ searchParams: Record
+ ): Promise {
+ const { runDumpJob } = await import('./export/dump')
+ return runDumpJob(
+ {
+ storage: this.storage,
+ env: {
+ // Optional binding; deployers add it to wrangler.toml when
+ // they want R2-backed dumps. Absent = storage-only mode.
+ R2_DUMP_BUCKET: (
+ this.env as Env & { R2_DUMP_BUCKET?: R2Bucket }
+ ).R2_DUMP_BUCKET,
+ },
+ dataSource: this.dumpJobDataSource(),
+ config,
+ setAlarm: (time, options) => this.setAlarm(time, options),
+ },
+ new URLSearchParams(searchParams)
+ )
+ }
+
+ /** Status/fetch endpoint for a chunked dump job. */
+ public async dumpJobStatus(
+ config: StarbaseDBConfiguration
+ ): Promise {
+ const { dumpJobStatus } = await import('./export/dump')
+ return dumpJobStatus({
+ storage: this.storage,
+ env: {
+ R2_DUMP_BUCKET: (
+ this.env as Env & { R2_DUMP_BUCKET?: R2Bucket }
+ ).R2_DUMP_BUCKET,
+ },
+ dataSource: this.dumpJobDataSource(),
+ config,
+ setAlarm: (time, options) => this.setAlarm(time, options),
+ })
+ }
+
public async executeQuery(opts: {
sql: string
params?: unknown[]
diff --git a/src/export/chunkedDump.sqlite.test.ts b/src/export/chunkedDump.sqlite.test.ts
new file mode 100644
index 0000000..ec5bcd6
--- /dev/null
+++ b/src/export/chunkedDump.sqlite.test.ts
@@ -0,0 +1,140 @@
+import { createClient, type Client } from '@libsql/client'
+import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
+import { executeOperation } from './index'
+import {
+ ChunkedDumpEngine,
+ DEFAULT_DUMP_OPTIONS,
+ DUMP_STATE_KEY,
+ type DumpState,
+} from './chunkedDump'
+import type { DataSource } from '../types'
+import type { StarbaseDBConfiguration } from '../handler'
+
+vi.mock('./index', () => ({
+ executeOperation: vi.fn(),
+}))
+
+const makeStorage = () => {
+ const values = new Map()
+ return {
+ get: vi.fn(async (key: string) => values.get(key) as T),
+ put: vi.fn(async (key: string, value: unknown) => {
+ values.set(key, value)
+ }),
+ delete: vi.fn(async (key: string) => {
+ values.delete(key)
+ }),
+ }
+}
+
+const makeDataSource = (): DataSource =>
+ ({ source: 'internal', rpc: {} }) as unknown as DataSource
+
+const makeConfig = (): StarbaseDBConfiguration => ({ role: 'admin' })
+
+describe('ChunkedDumpEngine with SQLite', () => {
+ let client: Client
+ let storage: ReturnType
+ let queries: { sql: string; params?: unknown[] }[]
+
+ beforeEach(async () => {
+ client = createClient({ url: 'file::memory:' })
+ storage = makeStorage()
+ queries = []
+ vi.mocked(executeOperation).mockImplementation(async (batch: any[]) => {
+ const query = batch[0]
+ queries.push(query)
+ const result = await client.execute({
+ sql: query.sql,
+ args: query.params ?? [],
+ })
+ return result.rows as Record[]
+ })
+ })
+
+ afterEach(() => {
+ client.close()
+ vi.clearAllMocks()
+ })
+
+ it('keeps negative, zero, and positive rowids and handles composite WITHOUT ROWID keys', async () => {
+ await client.execute(
+ 'CREATE TABLE "odd names" (id INTEGER, value TEXT)'
+ )
+ await client.execute({
+ sql: 'INSERT INTO "odd names" (rowid, id, value) VALUES (?, ?, ?), (?, ?, ?), (?, ?, ?)',
+ args: [-7, 10, 'negative', 0, 20, 'zero', 9, 30, 'positive'],
+ })
+ await client.execute(
+ 'CREATE TABLE "without rowid" ("key one" TEXT NOT NULL, part INTEGER NOT NULL, value TEXT, PRIMARY KEY ("key one", part)) WITHOUT ROWID'
+ )
+ await client.execute({
+ sql: 'INSERT INTO "without rowid" ("key one", part, value) VALUES (?, ?, ?), (?, ?, ?)',
+ args: ['a', 2, 'a2', 'a', 1, 'a1'],
+ })
+ await client.execute('CREATE TABLE "shadow" ("rowid" TEXT, value TEXT)')
+ await client.execute({
+ sql: 'INSERT INTO "shadow" ("rowid", value) VALUES (?, ?), (?, ?)',
+ args: ['first', 1, 'second', 2],
+ })
+
+ const engine = new ChunkedDumpEngine(
+ storage as unknown as DurableObjectStorage,
+ undefined,
+ makeDataSource(),
+ makeConfig(),
+ { ...DEFAULT_DUMP_OPTIONS, rowsPerBatch: 2, chunkTargetBytes: 80 }
+ )
+ let state = await engine.startDump()
+ for (let i = 0; i < 20 && !state.completedAt; i++) {
+ state = await engine.runCycle()
+ }
+ expect(state.completedAt).toBeDefined()
+ expect(state.totalRows).toBe(7)
+
+ const firstDataQuery = queries.find(
+ (query) =>
+ query.sql.includes('FROM "odd names"') &&
+ query.sql.includes('ORDER BY')
+ )
+ expect(firstDataQuery?.sql).not.toContain('WHERE "rowid" >')
+ expect(
+ queries.some(
+ (query) =>
+ query.sql.includes('FROM "odd names"') &&
+ query.sql.includes('WHERE "rowid" > ?') &&
+ query.params?.[0] === 0
+ )
+ ).toBe(true)
+ expect(
+ queries.some(
+ (query) =>
+ query.sql.includes('FROM "without rowid"') &&
+ query.sql.includes('("key one", "part") > (?, ?)')
+ )
+ ).toBe(true)
+ expect(
+ queries.some(
+ (query) =>
+ query.sql.includes('FROM "shadow"') &&
+ query.sql.includes('ORDER BY "_rowid_"')
+ )
+ ).toBe(true)
+
+ const stream = await engine.assembleDump(state)
+ expect(stream).not.toBeNull()
+ const dump = await new Response(stream as ReadableStream).text()
+ expect(dump).toContain('INSERT INTO "odd names"')
+ expect(dump).toContain('negative')
+ expect(dump).toContain('zero')
+ expect(dump).toContain('positive')
+ expect(dump).toContain('INSERT INTO "without rowid"')
+ expect(dump).toContain('a1')
+ expect(dump).toContain('a2')
+ expect(dump).toContain('INSERT INTO "shadow"')
+ const persisted = (await storage.get(DUMP_STATE_KEY)) as
+ | DumpState
+ | undefined
+ expect(persisted?.completedAt).toBeDefined()
+ })
+})
diff --git a/src/export/chunkedDump.test.ts b/src/export/chunkedDump.test.ts
new file mode 100644
index 0000000..f410705
--- /dev/null
+++ b/src/export/chunkedDump.test.ts
@@ -0,0 +1,866 @@
+import { describe, it, expect, vi, beforeEach } from 'vitest'
+import {
+ ChunkedDumpEngine,
+ DEFAULT_DUMP_OPTIONS,
+ MIN_R2_PART_SIZE_BYTES,
+ normalizePartSizeBytes,
+ isSafeIdentifier,
+ quoteIdentifier,
+ makeDumpFileName,
+ serializeRows,
+ DUMP_STATE_KEY,
+ DUMP_CHUNK_KEY,
+ type ChunkRecord,
+ type DumpOptions,
+ type DumpState,
+} from './chunkedDump'
+import { executeOperation } from './index'
+import { parseDumpOptions } from './dump'
+import type { DataSource } from '../types'
+import type { StarbaseDBConfiguration } from '../handler'
+
+vi.mock('./index', () => ({
+ executeOperation: vi.fn(),
+}))
+
+type StorageMap = Map
+
+const makeStorage = () => {
+ const map: StorageMap = new Map()
+ return {
+ get: vi.fn(async (key: string) => map.get(key) as T),
+ put: vi.fn(async (key: string, value: unknown) => {
+ map.set(key, value)
+ }),
+ delete: vi.fn(async (key: string) => {
+ map.delete(key)
+ }),
+ _map: map,
+ }
+}
+
+const makeR2 = () => ({
+ put: vi.fn(async () => undefined),
+ get: vi.fn(async () => null),
+})
+
+const makeDataSource = (): DataSource =>
+ ({ source: 'internal', rpc: {} }) as unknown as DataSource
+
+const makeConfig = (): StarbaseDBConfiguration => ({ role: 'admin' })
+
+const makeEngine = (overrides: Partial = {}) => {
+ const storage = makeStorage()
+ const r2 = makeR2()
+ const engine = new ChunkedDumpEngine(
+ storage as unknown as DurableObjectStorage,
+ r2 as unknown as R2Bucket,
+ makeDataSource(),
+ makeConfig(),
+ { ...DEFAULT_DUMP_OPTIONS, ...overrides }
+ )
+ return { engine, storage, r2 }
+}
+
+beforeEach(() => {
+ vi.clearAllMocks()
+})
+
+describe('identifier safety', () => {
+ it('accepts ordinary table names', () => {
+ expect(isSafeIdentifier('users')).toBe(true)
+ expect(isSafeIdentifier('order_items2')).toBe(true)
+ })
+
+ it('rejects injection attempts', () => {
+ expect(isSafeIdentifier('users; DROP TABLE users')).toBe(false)
+ expect(isSafeIdentifier('users"')).toBe(false)
+ expect(isSafeIdentifier('tmp_cache')).toBe(true) // syntactically safe
+ })
+})
+
+describe('row serialization', () => {
+ it('escapes single quotes in strings', () => {
+ const { content } = serializeRows('users', [{ name: "O'Brien" }])
+ expect(content).toContain("'O''Brien'")
+ expect(content).toContain('INSERT INTO "users"')
+ })
+
+ it('renders NULLs, numbers and booleans unquoted', () => {
+ const { content } = serializeRows('t', [{ a: null, b: 1.5, c: true }])
+ expect(content).toContain('(NULL, 1.5, true)')
+ })
+
+ it('returns empty content for zero rows', () => {
+ const { content, rowCount } = serializeRows('t', [])
+ expect(content).toBe('')
+ expect(rowCount).toBe(0)
+ })
+
+ it('quotes identifiers and serializes null and binary values', () => {
+ const { content } = serializeRows('order "items"', [
+ { 'value "x"': null, payload: new Uint8Array([0, 15, 255]) },
+ ])
+ expect(content).toContain('INSERT INTO "order ""items"""')
+ expect(content).toContain('("value ""x""", "payload")')
+ expect(content).toContain("(NULL, X'000fff')")
+ })
+})
+
+describe('dump file naming', () => {
+ it('follows dump_YYYYMMDD-HHMMSS.sql in UTC', () => {
+ const name = makeDumpFileName(new Date('2024-01-01T17:00:00Z'))
+ expect(name).toBe('dump_20240101-170000.sql')
+ })
+})
+
+describe('ChunkedDumpEngine cycles', () => {
+ const setupTables = (tables: string[], schemas: Record) => {
+ vi.mocked(executeOperation).mockImplementation(async (queries: any) => {
+ const sql: string = queries[0].sql
+ const params: unknown[] = queries[0].params ?? []
+ if (sql.includes("type='table' AND name NOT LIKE")) {
+ return tables.map((name) => ({ name }))
+ }
+ if (sql.includes('SELECT sql FROM sqlite_master')) {
+ const name = String(params[0] ?? '')
+ return schemas[name] ? [{ sql: schemas[name] }] : []
+ }
+ if (sql.includes('PRAGMA table_info')) {
+ return [{ name: 'id', pk: 1 }]
+ }
+ if (sql.includes('LIMIT 0')) {
+ return []
+ }
+ if (sql.includes('ORDER BY')) {
+ const since = Number(params[0] ?? -1)
+ if (since < 2) {
+ return since < 0
+ ? [
+ {
+ __starbase_dump_cursor_0: 0,
+ id: 0,
+ },
+ {
+ __starbase_dump_cursor_0: 1,
+ id: 1,
+ },
+ ]
+ : [
+ {
+ __starbase_dump_cursor_0: since + 1,
+ id: since + 1,
+ },
+ ]
+ }
+ return []
+ }
+ return []
+ })
+ }
+
+ it('completes a small dump in one cycle and writes chunks to R2', async () => {
+ setupTables(['users'], { users: 'CREATE TABLE users (id INTEGER)' })
+ vi.mocked(executeOperation).mockImplementation(async (queries: any) => {
+ const sql: string = queries[0].sql
+ const params: unknown[] = queries[0].params ?? []
+ if (sql.includes("name NOT LIKE 'tmp_%'")) {
+ return [{ name: 'users' }]
+ }
+ if (sql.includes('SELECT sql FROM sqlite_master')) {
+ return [{ sql: 'CREATE TABLE users (id INTEGER)' }]
+ }
+ if (sql.includes('PRAGMA table_info')) {
+ return [{ name: 'id', pk: 1 }]
+ }
+ if (sql.includes('LIMIT 0')) {
+ return []
+ }
+ if (sql.includes('ORDER BY')) {
+ const since = Number(params[0] ?? -1)
+ if (since < 0) {
+ return [
+ { __starbase_dump_cursor_0: -2, id: -2 },
+ { __starbase_dump_cursor_0: 0, id: 0 },
+ ]
+ }
+ return since < 1 ? [{ __starbase_dump_cursor_0: 1, id: 1 }] : []
+ }
+ return []
+ })
+
+ const { engine, storage, r2 } = makeEngine()
+ const state = await engine.startDump()
+
+ expect(state.completedAt).toBeDefined()
+ expect(state.phase).toBe('complete')
+ expect(state.totalRows).toBe(3)
+ expect(state.chunkIndex).toBeGreaterThan(0)
+ expect(r2.put).toHaveBeenCalled()
+ // Progress persisted in DO storage for resumability.
+ const persisted = (await storage.get(DUMP_STATE_KEY)) as
+ | DumpState
+ | undefined
+ expect(persisted?.completedAt).toBeDefined()
+ })
+
+ it('yields mid-dump when the cycle time budget is exhausted, then resumes', async () => {
+ setupTables(['users'], { users: 'CREATE TABLE users (id INTEGER)' })
+ let call = 0
+ vi.mocked(executeOperation).mockImplementation(async (queries: any) => {
+ const sql: string = queries[0].sql
+ const params: unknown[] = queries[0].params ?? []
+ if (sql.includes("name NOT LIKE 'tmp_%'"))
+ return [{ name: 'users' }]
+ if (sql.includes('SELECT sql FROM sqlite_master')) {
+ return [{ sql: 'CREATE TABLE users (id INTEGER)' }]
+ }
+ if (sql.includes('PRAGMA table_info')) {
+ return [{ name: 'id', pk: 1 }]
+ }
+ if (sql.includes('LIMIT 0')) {
+ return []
+ }
+ if (sql.includes('ORDER BY')) {
+ const since = Number(params[0] ?? -1)
+ call++
+ return since < 3
+ ? [{ __starbase_dump_cursor_0: since + 1, id: since + 1 }]
+ : []
+ }
+ return []
+ })
+
+ // 0ms budget: every data batch closes the cycle → resumable steps.
+ const { engine } = makeEngine({ cycleTimeBudgetMs: 0 })
+ const first = await engine.startDump()
+ expect(first.completedAt).toBeUndefined()
+ expect(['schema', 'table-data']).toContain(first.phase)
+
+ // Resume cycles until complete.
+ let state = first
+ for (let i = 0; i < 50 && !state.completedAt; i++) {
+ state = await engine.runCycle()
+ }
+ expect(state.completedAt).toBeDefined()
+ expect(state.totalRows).toBeGreaterThan(0)
+ })
+
+ it('quotes valid and unusual table names instead of dropping them', async () => {
+ const dataSqls: string[] = []
+ vi.mocked(executeOperation).mockImplementation(async (queries: any) => {
+ const sql: string = queries[0].sql
+ const params: unknown[] = queries[0].params ?? []
+ if (sql.includes("name NOT LIKE 'tmp_%'")) {
+ return [{ name: 'good_table' }, { name: 'order items "2026"' }]
+ }
+ if (sql.includes('SELECT sql FROM sqlite_master')) {
+ return [
+ { sql: 'CREATE TABLE "order items ""2026""" (id INTEGER)' },
+ ]
+ }
+ if (sql.includes('PRAGMA table_info')) {
+ return [{ name: 'id', pk: 1 }]
+ }
+ if (sql.includes('LIMIT 0')) return []
+ if (sql.includes('ORDER BY')) {
+ dataSqls.push(sql)
+ return Number(params[0] ?? -1) < 0
+ ? [{ __starbase_dump_cursor_0: 1, id: 1 }]
+ : []
+ }
+ return []
+ })
+
+ const { engine, storage } = makeEngine()
+ const state = await engine.startDump()
+ const persisted = (await storage.get(DUMP_STATE_KEY)) as
+ | DumpState
+ | undefined
+ expect(persisted?.tables).toEqual(['good_table', 'order items "2026"'])
+ expect(
+ dataSqls.some((sql) =>
+ sql.includes(quoteIdentifier('order items "2026"'))
+ )
+ ).toBe(true)
+ expect(state.totalRows).toBe(2)
+ })
+
+ it('startDump resumes an in-progress dump instead of restarting', async () => {
+ setupTables(['users'], {})
+ const { engine, storage } = makeEngine({ cycleTimeBudgetMs: 0 })
+ // Seed an in-progress state.
+ const seed: DumpState = {
+ dumpId: 'dump_seed',
+ fileName: 'dump_seed.sql',
+ phase: 'table-data',
+ tables: ['users'],
+ tableIndex: 0,
+ lastFetchedRowId: 3,
+ chunkRowOffset: 0,
+ bytesWritten: 0,
+ chunkIndex: 0,
+ startedAt: Date.now() - 1000,
+ updatedAt: Date.now() - 500,
+ totalRows: 0,
+ }
+ await storage.put(DUMP_STATE_KEY, seed)
+
+ let dataSql = ''
+ let dataParams: unknown[] = []
+ vi.mocked(executeOperation).mockImplementation(async (queries: any) => {
+ const sql: string = queries[0].sql
+ if (sql.includes('ORDER BY')) {
+ dataSql = sql
+ dataParams = queries[0].params ?? []
+ return []
+ }
+ return []
+ })
+
+ const state = await engine.startDump()
+ expect(state.dumpId).toBe('dump_seed')
+ expect(state.tables).toEqual(['users'])
+ expect(dataSql).toContain('WHERE "rowid" > ?')
+ expect(dataParams).toEqual([3])
+ })
+
+ it('assembleDump concatenates persisted chunks in order', async () => {
+ const { engine, storage } = makeEngine()
+ const state: DumpState = {
+ dumpId: 'dump_x',
+ fileName: 'dump_x.sql',
+ phase: 'complete',
+ tables: ['t'],
+ tableIndex: 1,
+ lastFetchedRowId: 3,
+ chunkRowOffset: 0,
+ bytesWritten: 10,
+ chunkIndex: 2,
+ startedAt: 1,
+ updatedAt: 2,
+ completedAt: 3,
+ totalRows: 3,
+ }
+ const c0: ChunkRecord = {
+ dumpId: 'dump_x',
+ chunkIndex: 0,
+ content: 'CREATE TABLE t (id INTEGER);\n',
+ bytes: 30,
+ createdAt: 1,
+ }
+ const c1: ChunkRecord = {
+ dumpId: 'dump_x',
+ chunkIndex: 1,
+ content: 'INSERT INTO "t" ("id") VALUES (1);\n',
+ bytes: 35,
+ createdAt: 2,
+ }
+ await storage.put(`${DUMP_CHUNK_KEY}:0`, c0)
+ await storage.put(`${DUMP_CHUNK_KEY}:1`, c1)
+
+ const stream = await engine.assembleDump(state)
+ expect(stream).not.toBeNull()
+ const text = await new Response(stream as ReadableStream).text()
+ expect(text).toContain('CREATE TABLE t')
+ expect(text).toContain('VALUES (1)')
+ })
+
+ it('respects custom rowsPerBatch from options', async () => {
+ let captured = ''
+ vi.mocked(executeOperation).mockImplementation(async (queries: any) => {
+ const sql: string = queries[0].sql
+ if (sql.includes('ORDER BY')) {
+ captured = sql
+ return []
+ }
+ if (sql.includes("name NOT LIKE 'tmp_%'"))
+ return [{ name: 'users' }]
+ if (sql.includes('SELECT sql FROM sqlite_master')) {
+ return [{ sql: 'CREATE TABLE users (id INTEGER)' }]
+ }
+ return []
+ })
+ const { engine } = makeEngine({ rowsPerBatch: 42 })
+ await engine.startDump()
+ expect(captured).toContain('LIMIT 42')
+ })
+})
+
+// ---------------------------------------------------------------------------
+// R2 multipart finalize + presigned URL
+// ---------------------------------------------------------------------------
+
+type UploadedPart = { partNumber: number; etag: string }
+
+const makeMultipartR2 = () => {
+ const uploaded: { partNumber: number; size: number }[] = []
+ const uploadedContents: { partNumber: number; content: string }[] = []
+ let completed: UploadedPart[] | null = null
+ let resumedWith: string | null = null
+ let createdFor: string | null = null
+ let aborted = 0
+ let signedUrlResult: { url?: string } | null | undefined = undefined
+ const deletedKeys: string[] = []
+ const puts: { key: string; size: number; content: string }[] = []
+ const objects = new Map }>()
+
+ const mpu = {
+ key: '',
+ uploadId: 'mpu-1',
+ uploadPart: async (partNumber: number, value: Uint8Array) => {
+ uploaded.push({ partNumber, size: value.length })
+ uploadedContents.push({
+ partNumber,
+ content: new TextDecoder().decode(value),
+ })
+ return { partNumber, etag: `etag-${partNumber}` }
+ },
+ abort: async () => {
+ aborted++
+ },
+ complete: async (parts: UploadedPart[]) => {
+ completed = parts
+ const total = uploaded.reduce((sum, p) => sum + p.size, 0)
+ return { key: mpu.key, size: total }
+ },
+ }
+
+ const r2 = {
+ createMultipartUpload: async (key: string) => {
+ createdFor = key
+ mpu.key = key
+ return mpu
+ },
+ resumeMultipartUpload: (key: string, uploadId: string) => {
+ resumedWith = uploadId
+ mpu.key = key
+ return mpu
+ },
+ put: async (key: string, value: string | Uint8Array) => {
+ const bytes =
+ typeof value === 'string'
+ ? new TextEncoder().encode(value)
+ : value
+ const content = new TextDecoder().decode(bytes)
+ puts.push({ key, size: bytes.length, content })
+ objects.set(key, {
+ body: new ReadableStream({
+ start(controller) {
+ if (bytes.length > 0) controller.enqueue(bytes)
+ controller.close()
+ },
+ }),
+ })
+ return { key, size: bytes.length }
+ },
+ get: async (key: string) => objects.get(key) ?? null,
+ delete: async (key: string) => {
+ deletedKeys.push(key)
+ },
+ }
+ return {
+ r2,
+ uploaded,
+ uploadedContents: () => uploadedContents,
+ puts: () => puts,
+ aborted: () => aborted,
+ deleted: () => deletedKeys,
+ completed: () => completed,
+ resumedWith: () => resumedWith,
+ createdFor: () => createdFor,
+ objects,
+ signedUrlResult: () => signedUrlResult,
+ setSigned: (v: { url?: string } | null | undefined) => {
+ signedUrlResult = v
+ },
+ }
+}
+
+const seedCompletedState = async (
+ storage: ReturnType,
+ chunks: string[],
+ extra: Partial = {}
+) => {
+ const state: DumpState = {
+ dumpId: 'dump_fin',
+ fileName: 'dump_fin.sql',
+ phase: 'complete',
+ tables: ['t'],
+ tableIndex: 1,
+ lastFetchedRowId: null,
+ chunkRowOffset: 0,
+ bytesWritten: chunks.join('').length,
+ chunkIndex: chunks.length,
+ startedAt: 1,
+ updatedAt: 2,
+ completedAt: 3,
+ totalRows: 0,
+ finalizePartSizeBytes: MIN_R2_PART_SIZE_BYTES,
+ ...extra,
+ }
+ for (let i = 0; i < chunks.length; i++) {
+ const record: ChunkRecord = {
+ dumpId: 'dump_fin',
+ chunkIndex: i,
+ content: chunks[i],
+ bytes: chunks[i].length,
+ createdAt: i,
+ }
+ await storage.put(`${DUMP_CHUNK_KEY}:${i}`, record)
+ }
+ await storage.put(DUMP_STATE_KEY, state)
+ return state
+}
+
+describe('finalizeDump (R2 multipart upload + presigned URL)', () => {
+ it('aggregates uneven chunks into fixed-size parts with a small tail', async () => {
+ const storage = makeStorage()
+ const m = makeMultipartR2()
+ const partSize = MIN_R2_PART_SIZE_BYTES
+ const chunks = ['A'.repeat(partSize + 1), 'B'.repeat(partSize - 1), 'C']
+ await seedCompletedState(storage, chunks)
+ const engine = new ChunkedDumpEngine(
+ storage as unknown as DurableObjectStorage,
+ m.r2 as unknown as R2Bucket,
+ makeDataSource(),
+ makeConfig(),
+ DEFAULT_DUMP_OPTIONS
+ )
+ const { done, state } = await engine.finalizeDump({ partSizeBytes: 1 })
+
+ expect(done).toBe(true)
+ expect(m.createdFor()).toBe('dumps/dump_fin/dump_fin.sql')
+ expect(m.uploaded.map((p) => p.size)).toEqual([partSize, partSize, 1])
+ expect(m.uploaded.map((p) => p.partNumber)).toEqual([1, 2, 3])
+ expect(
+ m
+ .uploadedContents()
+ .map((p) => p.content)
+ .join('')
+ ).toBe(chunks.join(''))
+ expect(m.completed()).toEqual([
+ { partNumber: 1, etag: 'etag-1' },
+ { partNumber: 2, etag: 'etag-2' },
+ { partNumber: 3, etag: 'etag-3' },
+ ])
+ expect(state.finalObjectKey).toBe('dumps/dump_fin/dump_fin.sql')
+ expect(state.finalObjectSize).toBe(partSize * 2 + 1)
+ expect(state.finalizePartSizeBytes).toBe(MIN_R2_PART_SIZE_BYTES)
+ expect(state.finalizedAt).toBeDefined()
+ expect(state.temporaryChunksCleanedAt).toBeDefined()
+ expect(m.deleted().length).toBe(3)
+ expect(await storage.get(`${DUMP_CHUNK_KEY}:0`)).toBeUndefined()
+ })
+
+ it('retries temporary chunk cleanup after a finalized upload', async () => {
+ const storage = makeStorage()
+ const m = makeMultipartR2()
+ await seedCompletedState(storage, ['data'])
+ const originalDelete = m.r2.delete
+ let attempts = 0
+ m.r2.delete = vi.fn(async (key: string) => {
+ attempts++
+ if (attempts === 1) throw new Error('temporary delete failure')
+ return originalDelete(key)
+ })
+ const engine = new ChunkedDumpEngine(
+ storage as unknown as DurableObjectStorage,
+ m.r2 as unknown as R2Bucket,
+ makeDataSource(),
+ makeConfig(),
+ DEFAULT_DUMP_OPTIONS
+ )
+
+ const first = await engine.finalizeDump()
+ expect(first.done).toBe(false)
+ expect(first.state.temporaryChunksCleanedAt).toBeUndefined()
+ expect(await storage.get(`${DUMP_CHUNK_KEY}:0`)).toBeDefined()
+ const second = await engine.finalizeDump()
+ expect(second.done).toBe(true)
+ expect(second.state.temporaryChunksCleanedAt).toBeDefined()
+ })
+
+ it('resumes after multiple compliant parts without re-uploading bytes', async () => {
+ const storage = makeStorage()
+ const m = makeMultipartR2()
+ const partSize = MIN_R2_PART_SIZE_BYTES
+ const chunks = [
+ 'A'.repeat(partSize + 1),
+ 'B'.repeat(partSize - 1),
+ 'C'.repeat(partSize + 1),
+ 'D'.repeat(partSize),
+ ]
+ await seedCompletedState(storage, chunks, {
+ finalObjectKey: 'dumps/dump_fin/dump_fin.sql',
+ finalizeUploadId: 'mpu-9',
+ finalizeParts: [
+ { partNumber: 1, etag: 'etag-1' },
+ { partNumber: 2, etag: 'etag-2' },
+ ],
+ finalizeBytes: partSize * 2,
+ finalizePartSizeBytes: partSize,
+ })
+ const engine = new ChunkedDumpEngine(
+ storage as unknown as DurableObjectStorage,
+ m.r2 as unknown as R2Bucket,
+ makeDataSource(),
+ makeConfig(),
+ DEFAULT_DUMP_OPTIONS
+ )
+ const { done } = await engine.finalizeDump({ partSizeBytes: 1 })
+
+ expect(done).toBe(true)
+ expect(m.resumedWith()).toBe('mpu-9')
+ expect(m.uploaded).toEqual([
+ { partNumber: 3, size: partSize },
+ { partNumber: 4, size: partSize },
+ { partNumber: 5, size: 1 },
+ ])
+ expect(m.completed()).toEqual([
+ { partNumber: 1, etag: 'etag-1' },
+ { partNumber: 2, etag: 'etag-2' },
+ { partNumber: 3, etag: 'etag-3' },
+ { partNumber: 4, etag: 'etag-4' },
+ { partNumber: 5, etag: 'etag-5' },
+ ])
+ })
+
+ it('normalizes under-minimum API and persisted part sizes', async () => {
+ expect(normalizePartSizeBytes(1)).toBe(MIN_R2_PART_SIZE_BYTES)
+ expect(
+ parseDumpOptions(new URLSearchParams('partBytes=1'))
+ .finalizePartSizeBytes
+ ).toBe(MIN_R2_PART_SIZE_BYTES)
+ expect(
+ parseDumpOptions(new URLSearchParams('partBytes=0'))
+ .finalizePartSizeBytes
+ ).toBe(MIN_R2_PART_SIZE_BYTES)
+ expect(
+ parseDumpOptions(new URLSearchParams('partBytes=invalid'))
+ .finalizePartSizeBytes
+ ).toBe(MIN_R2_PART_SIZE_BYTES)
+
+ const storage = makeStorage()
+ const m = makeMultipartR2()
+ await seedCompletedState(storage, ['x'], {
+ finalizePartSizeBytes: 1,
+ })
+ const engine = new ChunkedDumpEngine(
+ storage as unknown as DurableObjectStorage,
+ m.r2 as unknown as R2Bucket,
+ makeDataSource(),
+ makeConfig(),
+ DEFAULT_DUMP_OPTIONS
+ )
+ const result = await engine.finalizeDump({ partSizeBytes: 1 })
+
+ expect(result.done).toBe(true)
+ expect(result.state.finalizePartSizeBytes).toBe(MIN_R2_PART_SIZE_BYTES)
+ expect(m.uploaded).toEqual([{ partNumber: 1, size: 1 }])
+ })
+
+ it('aborts an invalid persisted multipart before restarting it', async () => {
+ const storage = makeStorage()
+ const m = makeMultipartR2()
+ await seedCompletedState(storage, ['x'], {
+ finalizeUploadId: 'old-upload',
+ finalizeParts: [{ partNumber: 1, etag: 'old-etag' }],
+ finalizeBytes: 1,
+ finalizePartSizeBytes: 1,
+ })
+ const engine = new ChunkedDumpEngine(
+ storage as unknown as DurableObjectStorage,
+ m.r2 as unknown as R2Bucket,
+ makeDataSource(),
+ makeConfig(),
+ DEFAULT_DUMP_OPTIONS
+ )
+ const result = await engine.finalizeDump({ partSizeBytes: 1 })
+
+ expect(result.done).toBe(true)
+ expect(m.aborted()).toBe(1)
+ expect(m.uploaded).toEqual([{ partNumber: 1, size: 1 }])
+ expect(m.completed()?.length).toBe(1)
+ expect(result.state.finalizePartSizeBytes).toBe(MIN_R2_PART_SIZE_BYTES)
+ })
+
+ it('writes an empty downloadable object without creating a multipart upload', async () => {
+ const storage = makeStorage()
+ const m = makeMultipartR2()
+ await seedCompletedState(storage, [])
+ const engine = new ChunkedDumpEngine(
+ storage as unknown as DurableObjectStorage,
+ m.r2 as unknown as R2Bucket,
+ makeDataSource(),
+ makeConfig(),
+ DEFAULT_DUMP_OPTIONS
+ )
+ const result = await engine.finalizeDump({ partSizeBytes: 1 })
+
+ expect(result.done).toBe(true)
+ expect(m.createdFor()).toBeNull()
+ expect(m.uploaded).toEqual([])
+ expect(m.completed()).toBeNull()
+ expect(m.aborted()).toBe(0)
+ expect(m.puts()).toEqual([
+ {
+ key: 'dumps/dump_fin/dump_fin.sql',
+ size: 0,
+ content: '',
+ },
+ ])
+ expect(result.state.finalObjectKey).toBe('dumps/dump_fin/dump_fin.sql')
+ expect(result.state.finalObjectSize).toBe(0)
+ expect(result.state.finalizeUploadId).toBeUndefined()
+ expect(result.state.finalizeParts).toBeUndefined()
+ expect(result.state.finalizedAt).toBeDefined()
+
+ const stream = await engine.assembleDump(result.state)
+ expect(stream).not.toBeNull()
+ expect(await new Response(stream as ReadableStream).text()).toBe('')
+ })
+
+ it('returns done:false when the time budget runs out, then finishes compliant parts', async () => {
+ const storage = makeStorage()
+ const m = makeMultipartR2()
+ const partSize = MIN_R2_PART_SIZE_BYTES
+ const chunks = ['A'.repeat(partSize + 1), 'B'.repeat(partSize - 1), 'C']
+ await seedCompletedState(storage, chunks)
+ const engine = new ChunkedDumpEngine(
+ storage as unknown as DurableObjectStorage,
+ m.r2 as unknown as R2Bucket,
+ makeDataSource(),
+ makeConfig(),
+ DEFAULT_DUMP_OPTIONS
+ )
+ const first = await engine.finalizeDump({ timeBudgetMs: -1 })
+ expect(first.done).toBe(false)
+ expect(first.state.finalizeUploadId).toBe('mpu-1')
+
+ const second = await engine.finalizeDump()
+ expect(second.done).toBe(true)
+ expect(m.uploaded.map((p) => p.size)).toEqual([partSize, partSize, 1])
+ expect(second.state.finalizedAt).toBeDefined()
+ })
+
+ it('finalizes without R2 binding (streaming-only environments)', async () => {
+ const storage = makeStorage()
+ await seedCompletedState(storage, ['data'])
+ const engine = new ChunkedDumpEngine(
+ storage as unknown as DurableObjectStorage,
+ undefined as unknown as R2Bucket,
+ makeDataSource(),
+ makeConfig(),
+ { ...DEFAULT_DUMP_OPTIONS }
+ )
+ const { done, state } = await engine.finalizeDump()
+ expect(done).toBe(true)
+ expect(state.finalObjectKey).toBeUndefined()
+ expect(state.finalizedAt).toBeDefined()
+ })
+
+ it('getPresignedUrl returns the signed URL when the runtime supports createSignedUrl', async () => {
+ const storage = makeStorage()
+ const m = makeMultipartR2()
+ await seedCompletedState(storage, ['data'], {
+ finalObjectKey: 'dumps/dump_fin/dump_fin.sql',
+ finalizedAt: 9,
+ })
+ m.setSigned({ url: 'https://signed.example/dump?sig=abc' })
+ const r2 = Object.assign(m.r2, {
+ createSignedUrl: async () =>
+ m.signedUrlResult() as { url?: string },
+ })
+ const engine = new ChunkedDumpEngine(
+ storage as unknown as DurableObjectStorage,
+ r2 as unknown as R2Bucket,
+ makeDataSource(),
+ makeConfig(),
+ { ...DEFAULT_DUMP_OPTIONS }
+ )
+ const url = await engine.getPresignedUrl(600)
+ expect(url).toBe('https://signed.example/dump?sig=abc')
+ })
+
+ it('getPresignedUrl returns null when the binding lacks createSignedUrl or the URL shape is unexpected', async () => {
+ const storage = makeStorage()
+ const m = makeMultipartR2()
+ await seedCompletedState(storage, ['data'], {
+ finalObjectKey: 'dumps/dump_fin/dump_fin.sql',
+ finalizedAt: 9,
+ })
+ // Older binding: no createSignedUrl at all.
+ const engineA = new ChunkedDumpEngine(
+ storage as unknown as DurableObjectStorage,
+ m.r2 as unknown as R2Bucket,
+ makeDataSource(),
+ makeConfig(),
+ { ...DEFAULT_DUMP_OPTIONS }
+ )
+ await expect(engineA.getPresignedUrl()).resolves.toBeNull()
+
+ // Unexpected result shape (missing url) must not throw.
+ m.setSigned({})
+ const r2 = Object.assign(m.r2, {
+ createSignedUrl: async () =>
+ m.signedUrlResult() as { url?: string },
+ })
+ const engineB = new ChunkedDumpEngine(
+ storage as unknown as DurableObjectStorage,
+ r2 as unknown as R2Bucket,
+ makeDataSource(),
+ makeConfig(),
+ { ...DEFAULT_DUMP_OPTIONS }
+ )
+ await expect(engineB.getPresignedUrl()).resolves.toBeNull()
+ })
+
+ it('does not return an empty fallback after temporary chunks are cleaned', async () => {
+ const storage = makeStorage()
+ const m = makeMultipartR2()
+ const state = await seedCompletedState(storage, ['chunk-content'], {
+ finalObjectKey: 'dumps/dump_fin/dump_fin.sql',
+ finalizedAt: 9,
+ temporaryChunksCleanedAt: 10,
+ })
+ const engine = new ChunkedDumpEngine(
+ storage as unknown as DurableObjectStorage,
+ m.r2 as unknown as R2Bucket,
+ makeDataSource(),
+ makeConfig(),
+ DEFAULT_DUMP_OPTIONS
+ )
+
+ await expect(engine.assembleDump(state)).resolves.toBeNull()
+ })
+
+ it('assembleDump prefers the consolidated final object', async () => {
+ const storage = makeStorage()
+ const m = makeMultipartR2()
+ await seedCompletedState(storage, ['chunk-content'], {
+ finalObjectKey: 'dumps/dump_fin/dump_fin.sql',
+ finalizedAt: 9,
+ })
+ const bytes = new TextEncoder().encode('CONSOLIDATED')
+ m.objects.set('dumps/dump_fin/dump_fin.sql', {
+ body: new ReadableStream({
+ start(c) {
+ c.enqueue(bytes)
+ c.close()
+ },
+ }),
+ })
+ const engine = new ChunkedDumpEngine(
+ storage as unknown as DurableObjectStorage,
+ m.r2 as unknown as R2Bucket,
+ makeDataSource(),
+ makeConfig(),
+ { ...DEFAULT_DUMP_OPTIONS }
+ )
+ const state = (await storage.get(DUMP_STATE_KEY)) as DumpState
+ const stream = await engine.assembleDump(state)
+ const text = await new Response(stream as ReadableStream).text()
+ expect(text).toBe('CONSOLIDATED')
+ })
+})
diff --git a/src/export/chunkedDump.ts b/src/export/chunkedDump.ts
new file mode 100644
index 0000000..5c81e72
--- /dev/null
+++ b/src/export/chunkedDump.ts
@@ -0,0 +1,885 @@
+import { DataSource } from '../types'
+import { StarbaseDBConfiguration } from '../handler'
+import { executeOperation } from './index'
+
+export const DUMP_STATE_KEY = 'tmp_dump_state'
+export const DUMP_CHUNK_KEY = 'tmp_dump_chunk'
+export const MIN_R2_PART_SIZE_BYTES = 5 * 1024 * 1024
+
+export const DEFAULT_DUMP_OPTIONS = {
+ cycleTimeBudgetMs: 5_000,
+ breathingIntervalMs: 5_000,
+ rowsPerBatch: 500,
+ chunkTargetBytes: 512 * 1024,
+ finalizePartSizeBytes: MIN_R2_PART_SIZE_BYTES,
+ finalizeTimeBudgetMs: 20_000,
+} as const
+
+export interface DumpOptions {
+ cycleTimeBudgetMs?: number
+ breathingIntervalMs?: number
+ rowsPerBatch?: number
+ chunkTargetBytes?: number
+ finalizePartSizeBytes?: number
+ finalizeTimeBudgetMs?: number
+}
+
+export type DumpPhase = 'schema' | 'table-data' | 'complete'
+export type DumpCursorMode = 'rowid' | 'primary-key'
+export type DumpCursorValue = string | number | bigint | null
+
+export interface DumpState {
+ dumpId: string
+ fileName: string
+ phase: DumpPhase
+ tables: string[]
+ tableIndex: number
+ lastFetchedRowId: number | null
+ chunkRowOffset: number
+ bytesWritten: number
+ chunkIndex: number
+ startedAt: number
+ updatedAt: number
+ completedAt?: number
+ totalRows: number
+ finalObjectKey?: string
+ finalObjectSize?: number
+ finalizedAt?: number
+ finalizeUploadId?: string
+ finalizeParts?: R2UploadedPart[]
+ finalizeBytes?: number
+ finalizePartSizeBytes?: number
+ currentTable?: string
+ cursorMode?: DumpCursorMode
+ cursorColumns?: string[]
+ cursorAliases?: string[]
+ cursorValues?: DumpCursorValue[] | null
+ temporaryChunksCleanedAt?: number
+}
+
+export interface ChunkRecord {
+ dumpId: string
+ chunkIndex: number
+ content: string
+ bytes: number
+ createdAt: number
+}
+
+const IDENTIFIER_PATTERN = /^[A-Za-z_][A-Za-z0-9_]*$/
+const CURSOR_ALIAS_PREFIX = '__starbase_dump_cursor_'
+
+export function isSafeIdentifier(name: string): boolean {
+ return IDENTIFIER_PATTERN.test(name)
+}
+
+export function quoteIdentifier(name: string): string {
+ return `"${name.replace(/"/g, '""')}"`
+}
+
+export function sqlCommentLabel(value: string): string {
+ return value.replace(/[\r\n]/g, ' ')
+}
+
+function sqlQuote(value: unknown): string {
+ if (value === null || value === undefined) {
+ return 'NULL'
+ }
+ if (typeof value === 'number') {
+ return Number.isFinite(value) ? String(value) : 'NULL'
+ }
+ if (typeof value === 'bigint' || typeof value === 'boolean') {
+ return String(value)
+ }
+ if (value instanceof ArrayBuffer) {
+ return `X'${Array.from(new Uint8Array(value), (b) =>
+ b.toString(16).padStart(2, '0')
+ ).join('')}'`
+ }
+ if (ArrayBuffer.isView(value)) {
+ return `X'${Array.from(
+ new Uint8Array(value.buffer, value.byteOffset, value.byteLength),
+ (b) => b.toString(16).padStart(2, '0')
+ ).join('')}'`
+ }
+ if (typeof value === 'string') {
+ return `'${value.replace(/'/g, "''")}'`
+ }
+ const serialized = JSON.stringify(value)
+ return `'${(serialized ?? String(value)).replace(/'/g, "''")}'`
+}
+
+export function serializeRows(
+ table: string,
+ rows: Record[],
+ columns?: string[]
+): { content: string; rowCount: number } {
+ if (rows.length === 0) {
+ return { content: '', rowCount: 0 }
+ }
+ const selectedColumns = columns ?? Object.keys(rows[0])
+ const columnList = selectedColumns.map((c) => quoteIdentifier(c)).join(', ')
+ const lines = rows.map((row) => {
+ const values = selectedColumns.map((c) => sqlQuote(row[c]))
+ return `INSERT INTO ${quoteIdentifier(table)} (${columnList}) VALUES (${values.join(', ')});`
+ })
+ return { content: `${lines.join('\n')}\n`, rowCount: rows.length }
+}
+
+export function makeDumpFileName(now = new Date()): string {
+ const pad = (n: number) => String(n).padStart(2, '0')
+ const stamp =
+ `${now.getUTCFullYear()}${pad(now.getUTCMonth() + 1)}${pad(now.getUTCDate())}` +
+ `-${pad(now.getUTCHours())}${pad(now.getUTCMinutes())}${pad(now.getUTCSeconds())}`
+ return `dump_${stamp}.sql`
+}
+
+type TableCursorPlan = {
+ mode: DumpCursorMode
+ columns: string[]
+ aliases: string[]
+}
+
+type CursorRow = Record
+
+function dumpChunkKey(dumpId: string, chunkIndex: number): string {
+ return `${dumpId}/${String(chunkIndex).padStart(8, '0')}.sql`
+}
+
+function normalizeCursorValue(value: unknown): DumpCursorValue {
+ if (value === null || value === undefined) {
+ return null
+ }
+ if (
+ typeof value === 'string' ||
+ typeof value === 'number' ||
+ typeof value === 'bigint'
+ ) {
+ return value
+ }
+ return String(value)
+}
+
+function stripCursorColumns(
+ row: CursorRow,
+ aliases: string[]
+): Record {
+ const result = { ...row }
+ for (const alias of aliases) {
+ delete result[alias]
+ }
+ return result
+}
+
+function rowCursorValues(row: CursorRow, aliases: string[]): DumpCursorValue[] {
+ return aliases.map((alias) => normalizeCursorValue(row[alias]))
+}
+
+function makeCursorAliases(columns: string[]): string[] {
+ const names = new Set(columns.map((column) => column.toLowerCase()))
+ return columns.map((_, index) => {
+ let alias = `${CURSOR_ALIAS_PREFIX}${index}`
+ while (names.has(alias.toLowerCase())) {
+ alias += '_'
+ }
+ return alias
+ })
+}
+
+export function normalizePartSizeBytes(value: unknown): number {
+ const numeric = typeof value === 'number' ? value : Number(value)
+ if (!Number.isFinite(numeric)) {
+ return MIN_R2_PART_SIZE_BYTES
+ }
+ return Math.max(MIN_R2_PART_SIZE_BYTES, Math.floor(numeric))
+}
+
+export class ChunkedDumpEngine {
+ constructor(
+ private readonly storage: DurableObjectStorage,
+ private readonly r2: R2Bucket | undefined,
+ private readonly dataSource: DataSource,
+ private readonly config: StarbaseDBConfiguration,
+ options: Required
+ ) {
+ this.options = {
+ ...options,
+ finalizePartSizeBytes: normalizePartSizeBytes(
+ options.finalizePartSizeBytes
+ ),
+ }
+ }
+
+ private readonly options: Required
+
+ async startDump(): Promise {
+ const existing = await this.storage.get(DUMP_STATE_KEY)
+ if (existing) {
+ return this.runCycle()
+ }
+
+ const tablesResult = await executeOperation(
+ [
+ {
+ sql: "SELECT name FROM sqlite_master WHERE type='table' AND name NOT LIKE 'tmp_%';",
+ },
+ ],
+ this.dataSource,
+ this.config
+ )
+ const tables = tablesResult
+ .map((row: Record) => String(row.name))
+ .filter((name) => name.length > 0)
+
+ const now = Date.now()
+ const state: DumpState = {
+ dumpId: `dump_${now.toString(36)}`,
+ fileName: makeDumpFileName(new Date(now)),
+ phase: 'schema',
+ tables,
+ tableIndex: 0,
+ lastFetchedRowId: null,
+ chunkRowOffset: 0,
+ bytesWritten: 0,
+ chunkIndex: 0,
+ startedAt: now,
+ updatedAt: now,
+ totalRows: 0,
+ finalizePartSizeBytes: this.options.finalizePartSizeBytes,
+ }
+ await this.storage.put(DUMP_STATE_KEY, state)
+ return this.runCycle()
+ }
+
+ async getState(): Promise {
+ return this.storage.get(DUMP_STATE_KEY)
+ }
+
+ async runCycle(): Promise {
+ const state = (await this.storage.get(
+ DUMP_STATE_KEY
+ )) as DumpState
+ if (!state) {
+ throw new Error('No dump in progress')
+ }
+ if (state.completedAt) {
+ return state
+ }
+
+ const cycleStart = Date.now()
+ let content = ''
+ let contentBytes = 0
+ let chunkDirty = false
+
+ const flushChunk = async () => {
+ if (!chunkDirty) {
+ return
+ }
+ const record: ChunkRecord = {
+ dumpId: state.dumpId,
+ chunkIndex: state.chunkIndex,
+ content,
+ bytes: contentBytes,
+ createdAt: Date.now(),
+ }
+ if (this.r2) {
+ await this.r2.put(
+ dumpChunkKey(state.dumpId, record.chunkIndex),
+ record.content
+ )
+ }
+ await this.storage.put(
+ `${DUMP_CHUNK_KEY}:${record.chunkIndex}`,
+ record
+ )
+ state.chunkIndex += 1
+ state.bytesWritten += contentBytes
+ content = ''
+ contentBytes = 0
+ chunkDirty = false
+ }
+
+ const persist = async () => {
+ state.updatedAt = Date.now()
+ await this.storage.put(DUMP_STATE_KEY, state)
+ }
+
+ if (state.phase === 'schema') {
+ while (state.tableIndex < state.tables.length) {
+ const table = state.tables[state.tableIndex]
+ state.tableIndex++
+ const schemaResult = await executeOperation(
+ [
+ {
+ sql: "SELECT sql FROM sqlite_master WHERE type='table' AND name=?;",
+ params: [table],
+ },
+ ],
+ this.dataSource,
+ this.config
+ )
+ if (schemaResult.length && schemaResult[0]?.sql) {
+ const ddl = `\n-- Table: ${sqlCommentLabel(table)}\n${String(schemaResult[0].sql)};\n\n`
+ content += ddl
+ contentBytes += ddl.length
+ chunkDirty = true
+ }
+ if (
+ contentBytes >= this.options.chunkTargetBytes ||
+ Date.now() - cycleStart >= this.options.cycleTimeBudgetMs
+ ) {
+ await flushChunk()
+ await persist()
+ return state
+ }
+ }
+ state.phase = 'table-data'
+ state.tableIndex = 0
+ }
+
+ while (state.tableIndex < state.tables.length) {
+ const table = state.tables[state.tableIndex]
+ const hasLegacyRowIdCursor =
+ state.currentTable === undefined &&
+ state.cursorMode === undefined &&
+ state.lastFetchedRowId !== null
+ if (state.currentTable !== table) {
+ state.currentTable = table
+ state.cursorMode = undefined
+ state.cursorColumns = undefined
+ state.cursorAliases = undefined
+ state.cursorValues = null
+ if (!hasLegacyRowIdCursor) {
+ state.lastFetchedRowId = null
+ }
+ }
+
+ const plan = await this.resolveTableCursor(table, state)
+ const aliases = state.cursorAliases ?? plan.aliases
+ const cursorColumns = state.cursorColumns ?? plan.columns
+ const cursorValues =
+ state.cursorValues ??
+ (state.cursorMode === 'rowid' && state.lastFetchedRowId !== null
+ ? [state.lastFetchedRowId]
+ : null)
+ if (cursorValues && !state.cursorValues) {
+ state.cursorValues = cursorValues
+ }
+ const limit = Math.max(1, Math.floor(this.options.rowsPerBatch))
+ const selectedCursorColumns = cursorColumns
+ .map(
+ (column, index) =>
+ `${quoteIdentifier(column)} AS ${quoteIdentifier(aliases[index])}`
+ )
+ .join(', ')
+ const orderColumns = cursorColumns
+ .map((column) => quoteIdentifier(column))
+ .join(', ')
+ const where =
+ cursorValues && cursorValues.length === cursorColumns.length
+ ? state.cursorMode === 'rowid'
+ ? ` WHERE ${quoteIdentifier(cursorColumns[0])} > ?`
+ : ` WHERE (${cursorColumns.map((column) => quoteIdentifier(column)).join(', ')}) > (${cursorValues.map(() => '?').join(', ')})`
+ : ''
+ const params = cursorValues ? [...cursorValues] : []
+ const rowsResult = ((await executeOperation(
+ [
+ {
+ sql: `SELECT ${selectedCursorColumns}, * FROM ${quoteIdentifier(table)}${where} ORDER BY ${orderColumns} LIMIT ${limit};`,
+ params,
+ },
+ ],
+ this.dataSource,
+ this.config
+ )) ?? []) as CursorRow[]
+
+ if (rowsResult.length === 0) {
+ state.tableIndex++
+ state.currentTable = undefined
+ state.cursorMode = undefined
+ state.cursorColumns = undefined
+ state.cursorAliases = undefined
+ state.cursorValues = null
+ state.lastFetchedRowId = null
+ continue
+ }
+
+ const nextCursor = rowCursorValues(
+ rowsResult[rowsResult.length - 1],
+ aliases
+ )
+ state.cursorValues = nextCursor
+ if (state.cursorMode === 'rowid') {
+ const numeric = Number(nextCursor[0])
+ state.lastFetchedRowId = Number.isFinite(numeric)
+ ? numeric
+ : null
+ }
+
+ const rows = rowsResult.map((row) =>
+ stripCursorColumns(row, aliases)
+ )
+ const { content: batchSql, rowCount } = serializeRows(table, rows)
+ content += batchSql
+ contentBytes += batchSql.length
+ state.totalRows += rowCount
+ state.chunkRowOffset += rowCount
+ chunkDirty = true
+
+ if (
+ contentBytes >= this.options.chunkTargetBytes ||
+ Date.now() - cycleStart >= this.options.cycleTimeBudgetMs
+ ) {
+ await flushChunk()
+ await persist()
+ return state
+ }
+ }
+
+ const done = `\n-- Dump complete: ${state.totalRows} rows, ${state.chunkIndex + (chunkDirty ? 1 : 0)} chunks.\n`
+ content += done
+ contentBytes += done.length
+ chunkDirty = true
+ await flushChunk()
+
+ state.phase = 'complete'
+ state.completedAt = Date.now()
+ await persist()
+ return state
+ }
+
+ shouldContinue(state: DumpState, requestStart: number): boolean {
+ if (state.completedAt) return false
+ return Date.now() - requestStart < this.options.cycleTimeBudgetMs
+ }
+
+ breathingDelayMs(): number {
+ return this.options.breathingIntervalMs
+ }
+
+ async finalizeDump(
+ options: { partSizeBytes?: number; timeBudgetMs?: number } = {}
+ ): Promise<{ done: boolean; state: DumpState }> {
+ const state = (await this.storage.get(
+ DUMP_STATE_KEY
+ )) as DumpState
+ if (!state) {
+ throw new Error('No dump in progress')
+ }
+
+ const persistedPartSize = state.finalizePartSizeBytes
+ const requestedPartSize =
+ persistedPartSize ??
+ options.partSizeBytes ??
+ this.options.finalizePartSizeBytes
+ const partSize = normalizePartSizeBytes(requestedPartSize)
+ const partSizeChanged = state.finalizePartSizeBytes !== partSize
+ const resetMultipart =
+ Boolean(state.finalizeUploadId) &&
+ (persistedPartSize === undefined || persistedPartSize !== partSize)
+ if (partSizeChanged) {
+ state.finalizePartSizeBytes = partSize
+ }
+ const persistPartSize = async () => {
+ if (!partSizeChanged) return
+ state.updatedAt = Date.now()
+ await this.storage.put(DUMP_STATE_KEY, state)
+ }
+
+ if (!state.completedAt) {
+ await persistPartSize()
+ return { done: false, state }
+ }
+ if (state.finalizedAt) {
+ await persistPartSize()
+ if (this.r2 && !state.temporaryChunksCleanedAt) {
+ const cleaned = await this.cleanupTemporaryChunks(state)
+ return { done: cleaned, state }
+ }
+ return { done: true, state }
+ }
+ if (!this.r2) {
+ await persistPartSize()
+ state.finalizedAt = Date.now()
+ state.updatedAt = state.finalizedAt
+ await this.storage.put(DUMP_STATE_KEY, state)
+ return { done: true, state }
+ }
+
+ const budget = options.timeBudgetMs ?? this.options.finalizeTimeBudgetMs
+ const key =
+ state.finalObjectKey ?? `dumps/${state.dumpId}/${state.fileName}`
+ let multipartReset = false
+ if (resetMultipart && state.finalizeUploadId) {
+ const upload = this.r2.resumeMultipartUpload(
+ key,
+ state.finalizeUploadId
+ )
+ try {
+ await upload.abort()
+ } catch {
+ return { done: false, state }
+ }
+ state.finalizeUploadId = undefined
+ state.finalizeParts = undefined
+ state.finalizeBytes = undefined
+ multipartReset = true
+ }
+ if (partSizeChanged || multipartReset) {
+ state.updatedAt = Date.now()
+ await this.storage.put(DUMP_STATE_KEY, state)
+ }
+ const hasContent = await this.hasChunkContent(state)
+ if (!hasContent) {
+ if (state.bytesWritten > 0) {
+ throw new Error('Dump chunk data is missing')
+ }
+ return this.finalizeEmptyDump(state, key)
+ }
+
+ const cycleStart = Date.now()
+ let mpu: R2MultipartUpload
+ if (state.finalizeUploadId) {
+ mpu = this.r2.resumeMultipartUpload(key, state.finalizeUploadId)
+ } else {
+ mpu = await this.r2.createMultipartUpload(key)
+ state.finalObjectKey = key
+ state.finalizeUploadId = mpu.uploadId
+ state.finalizeParts = []
+ state.finalizeBytes = 0
+ await this.storage.put(DUMP_STATE_KEY, state)
+ }
+
+ const parts = state.finalizeParts ?? []
+ const encoder = new TextEncoder()
+ const persist = async () => {
+ state.finalizeParts = parts
+ state.updatedAt = Date.now()
+ await this.storage.put(DUMP_STATE_KEY, state)
+ }
+
+ let skipped = Math.max(0, state.finalizeBytes ?? 0)
+ let partNumber = parts.length + 1
+ const buffer = new Uint8Array(partSize)
+ let bufferLength = 0
+
+ for (let i = 0; i < state.chunkIndex; i++) {
+ const record = await this.storage.get(
+ `${DUMP_CHUNK_KEY}:${i}`
+ )
+ if (!record?.content) continue
+ const encoded = encoder.encode(record.content)
+ let offset = 0
+ if (skipped >= encoded.length) {
+ skipped -= encoded.length
+ continue
+ }
+ if (skipped > 0) {
+ offset = skipped
+ skipped = 0
+ }
+
+ while (offset < encoded.length) {
+ const length = Math.min(
+ partSize - bufferLength,
+ encoded.length - offset
+ )
+ buffer.set(
+ encoded.subarray(offset, offset + length),
+ bufferLength
+ )
+ bufferLength += length
+ offset += length
+
+ if (bufferLength === partSize) {
+ if (Date.now() - cycleStart >= budget) {
+ await persist()
+ return { done: false, state }
+ }
+ const part = await mpu.uploadPart(
+ partNumber,
+ buffer.slice(0, partSize)
+ )
+ parts.push({ partNumber, etag: part.etag })
+ state.finalizeBytes = (state.finalizeBytes ?? 0) + partSize
+ bufferLength = 0
+ await persist()
+ partNumber++
+ }
+ }
+ }
+
+ if (bufferLength > 0) {
+ if (Date.now() - cycleStart >= budget) {
+ await persist()
+ return { done: false, state }
+ }
+ const part = await mpu.uploadPart(
+ partNumber,
+ buffer.slice(0, bufferLength)
+ )
+ parts.push({ partNumber, etag: part.etag })
+ state.finalizeBytes = (state.finalizeBytes ?? 0) + bufferLength
+ bufferLength = 0
+ await persist()
+ }
+
+ if (parts.length === 0) {
+ return { done: false, state }
+ }
+
+ const object = await mpu.complete(parts)
+ state.finalObjectKey = key
+ state.finalObjectSize = object?.size ?? state.finalizeBytes
+ state.finalizedAt = Date.now()
+ state.updatedAt = state.finalizedAt
+ state.finalizeUploadId = undefined
+ state.finalizeParts = undefined
+ state.finalizeBytes = undefined
+ await this.storage.put(DUMP_STATE_KEY, state)
+
+ const cleaned = await this.cleanupTemporaryChunks(state)
+ return { done: cleaned, state }
+ }
+
+ async getPresignedUrl(expiresInSeconds = 3600): Promise {
+ const state = (await this.storage.get(
+ DUMP_STATE_KEY
+ )) as DumpState
+ if (!state?.finalObjectKey || !this.r2) return null
+ const creator = (
+ this.r2 as R2Bucket & { createSignedUrl?: R2SignedUrlCreator }
+ ).createSignedUrl
+ if (typeof creator !== 'function') return null
+ try {
+ const signed = await creator.call(
+ this.r2,
+ state.finalObjectKey,
+ expiresInSeconds
+ )
+ return signed?.url ?? null
+ } catch {
+ return null
+ }
+ }
+
+ async assembleDump(
+ state: DumpState
+ ): Promise | null> {
+ if (!state.completedAt) return null
+
+ if (state.finalObjectKey && this.r2) {
+ const final = await this.r2.get(state.finalObjectKey)
+ if (final?.body) return final.body
+ }
+
+ if (state.temporaryChunksCleanedAt) {
+ return null
+ }
+
+ if (this.r2) {
+ const stream = await this.concatenateR2(state)
+ if (stream) return stream
+ }
+
+ const storage = this.storage
+ return new ReadableStream({
+ async start(controller) {
+ const encoder = new TextEncoder()
+ for (let i = 0; i < state.chunkIndex; i++) {
+ const record = await storage.get(
+ `${DUMP_CHUNK_KEY}:${i}`
+ )
+ if (record?.content) {
+ controller.enqueue(encoder.encode(record.content))
+ }
+ }
+ controller.close()
+ },
+ })
+ }
+
+ private async resolveTableCursor(
+ table: string,
+ state: DumpState
+ ): Promise {
+ if (state.cursorMode && state.cursorColumns && state.cursorAliases) {
+ return {
+ mode: state.cursorMode,
+ columns: state.cursorColumns,
+ aliases: state.cursorAliases,
+ }
+ }
+
+ const tableInfo = ((await executeOperation(
+ [{ sql: `PRAGMA table_info(${quoteIdentifier(table)});` }],
+ this.dataSource,
+ this.config
+ )) ?? []) as Record[]
+ const columns = tableInfo.map((row) => String(row.name))
+ const primaryKey = tableInfo
+ .filter((row) => Number(row.pk) > 0)
+ .sort((left, right) => Number(left.pk) - Number(right.pk))
+ .map((row) => String(row.name))
+
+ const rowidAliases = ['rowid', '_rowid_', 'oid'].filter(
+ (alias) => !columns.some((column) => column.toLowerCase() === alias)
+ )
+ for (const alias of rowidAliases) {
+ try {
+ await executeOperation(
+ [
+ {
+ sql: `SELECT ${quoteIdentifier(alias)} AS ${quoteIdentifier(`${CURSOR_ALIAS_PREFIX}probe`)} FROM ${quoteIdentifier(table)} LIMIT 0;`,
+ },
+ ],
+ this.dataSource,
+ this.config
+ )
+ const plan = {
+ mode: 'rowid' as const,
+ columns: [alias],
+ aliases: makeCursorAliases([alias]),
+ }
+ state.cursorMode = plan.mode
+ state.cursorColumns = plan.columns
+ state.cursorAliases = plan.aliases
+ state.cursorValues = null
+ return plan
+ } catch {}
+ }
+
+ if (primaryKey.length === 0) {
+ throw new Error(
+ `Cannot determine a stable cursor for table ${table}`
+ )
+ }
+ const plan = {
+ mode: 'primary-key' as const,
+ columns: primaryKey,
+ aliases: makeCursorAliases(primaryKey),
+ }
+ state.cursorMode = plan.mode
+ state.cursorColumns = plan.columns
+ state.cursorAliases = plan.aliases
+ state.cursorValues = null
+ return plan
+ }
+
+ private async hasChunkContent(state: DumpState): Promise {
+ for (let i = 0; i < state.chunkIndex; i++) {
+ const record = await this.storage.get(
+ `${DUMP_CHUNK_KEY}:${i}`
+ )
+ if (record?.content) {
+ return true
+ }
+ }
+ return false
+ }
+
+ private async finalizeEmptyDump(
+ state: DumpState,
+ key: string
+ ): Promise<{ done: boolean; state: DumpState }> {
+ const r2 = this.r2
+ if (!r2) {
+ throw new Error(
+ 'R2 binding is required for empty dump finalization'
+ )
+ }
+ if (state.finalizeUploadId) {
+ const upload = r2.resumeMultipartUpload(key, state.finalizeUploadId)
+ try {
+ await upload.abort()
+ } catch {
+ return { done: false, state }
+ }
+ state.finalizeUploadId = undefined
+ state.finalizeParts = undefined
+ state.finalizeBytes = undefined
+ state.updatedAt = Date.now()
+ await this.storage.put(DUMP_STATE_KEY, state)
+ }
+
+ await r2.put(key, new Uint8Array(0))
+ state.finalObjectKey = key
+ state.finalObjectSize = 0
+ state.finalizedAt = Date.now()
+ state.updatedAt = state.finalizedAt
+ state.finalizeUploadId = undefined
+ state.finalizeParts = undefined
+ state.finalizeBytes = undefined
+ await this.storage.put(DUMP_STATE_KEY, state)
+ const cleaned = await this.cleanupTemporaryChunks(state)
+ return { done: cleaned, state }
+ }
+
+ private async cleanupTemporaryChunks(state: DumpState): Promise {
+ let complete = true
+ for (let i = 0; i < state.chunkIndex; i++) {
+ let r2Deleted = true
+ if (this.r2) {
+ try {
+ await this.r2.delete(dumpChunkKey(state.dumpId, i))
+ } catch {
+ r2Deleted = false
+ complete = false
+ }
+ }
+ if (!r2Deleted) {
+ continue
+ }
+ try {
+ await this.storage.delete(`${DUMP_CHUNK_KEY}:${i}`)
+ } catch {
+ complete = false
+ }
+ }
+ if (complete) {
+ state.temporaryChunksCleanedAt = Date.now()
+ }
+ state.updatedAt = Date.now()
+ await this.storage.put(DUMP_STATE_KEY, state)
+ return complete
+ }
+
+ private async concatenateR2(
+ state: DumpState
+ ): Promise | null> {
+ if (!this.r2) return null
+ const head = await this.r2.get(dumpChunkKey(state.dumpId, 0))
+ if (!head) return null
+ const parts: ReadableStream[] = []
+ for (let i = 0; i < state.chunkIndex; i++) {
+ const obj = await this.r2.get(dumpChunkKey(state.dumpId, i))
+ if (obj?.body) {
+ parts.push(obj.body)
+ }
+ }
+ return concatStreams(parts)
+ }
+}
+
+type R2SignedUrlCreator = (
+ key: string,
+ expiresInSeconds: number
+) => Promise<{ url?: string } | null>
+
+function concatStreams(
+ streams: ReadableStream[]
+): ReadableStream {
+ return new ReadableStream({
+ async start(controller) {
+ for (const stream of streams) {
+ const reader = stream.getReader()
+ for (;;) {
+ const { done, value } = await reader.read()
+ if (done) break
+ controller.enqueue(value)
+ }
+ reader.releaseLock()
+ }
+ controller.close()
+ },
+ })
+}
diff --git a/src/export/dump.job.test.ts b/src/export/dump.job.test.ts
new file mode 100644
index 0000000..bc35f7e
--- /dev/null
+++ b/src/export/dump.job.test.ts
@@ -0,0 +1,151 @@
+import { beforeEach, describe, expect, it, vi } from 'vitest'
+import { runDumpJob, dumpJobStatus, type DumpEngineHost } from './dump'
+import { executeOperation } from './index'
+import {
+ DUMP_STATE_KEY,
+ MIN_R2_PART_SIZE_BYTES,
+ type DumpState,
+} from './chunkedDump'
+import type { DataSource } from '../types'
+import type { StarbaseDBConfiguration } from '../handler'
+
+vi.mock('./index', () => ({
+ executeOperation: vi.fn(),
+}))
+
+const makeStorage = () => {
+ const values = new Map()
+ return {
+ get: vi.fn(async (key: string) => values.get(key) as T),
+ put: vi.fn(async (key: string, value: unknown) => {
+ values.set(key, value)
+ }),
+ delete: vi.fn(async (key: string) => {
+ values.delete(key)
+ }),
+ }
+}
+
+const makeDataSource = (): DataSource =>
+ ({ source: 'internal', rpc: {} }) as unknown as DataSource
+
+const makeConfig = (): StarbaseDBConfiguration => ({ role: 'admin' })
+
+const makeHost = () => {
+ const alarms: number[] = []
+ const host: DumpEngineHost = {
+ storage: makeStorage() as unknown as DurableObjectStorage,
+ env: {},
+ dataSource: makeDataSource(),
+ config: makeConfig(),
+ setAlarm: vi.fn(async (time: number) => {
+ alarms.push(time)
+ }),
+ }
+ return { host, alarms }
+}
+
+const setupCompleteQueries = () => {
+ vi.mocked(executeOperation).mockImplementation(async (queries: any[]) => {
+ const query = queries[0]
+ const sql: string = query.sql
+ const params: unknown[] = query.params ?? []
+ if (sql.includes('sqlite_master') && sql.includes('name NOT LIKE')) {
+ return [{ name: 't' }]
+ }
+ if (sql.includes('SELECT sql FROM sqlite_master')) {
+ return [{ sql: 'CREATE TABLE t (id INTEGER)' }]
+ }
+ if (sql.includes('PRAGMA table_info')) {
+ return [{ name: 'id', pk: 1 }]
+ }
+ if (sql.includes('LIMIT 0')) {
+ return []
+ }
+ if (sql.includes('ORDER BY')) {
+ return params.length === 0
+ ? [{ __starbase_dump_cursor_0: 1, id: 1 }]
+ : []
+ }
+ return []
+ })
+}
+
+describe('dump job lifecycle', () => {
+ beforeEach(() => {
+ vi.clearAllMocks()
+ })
+
+ it('reuses a completed job instead of planning a second dump', async () => {
+ setupCompleteQueries()
+ const { host } = makeHost()
+ const fetchSpy = vi.spyOn(globalThis, 'fetch')
+
+ const first = await runDumpJob(
+ host,
+ new URLSearchParams({
+ callbackUrl: 'https://example.invalid/callback',
+ })
+ )
+ const firstText = await first.text()
+ const firstState = (await host.storage.get(DUMP_STATE_KEY)) as DumpState
+ const planCallsAfterFirst = vi
+ .mocked(executeOperation)
+ .mock.calls.filter(([queries]) =>
+ String((queries as any[])[0].sql).includes('name NOT LIKE')
+ ).length
+
+ const second = await runDumpJob(host, new URLSearchParams())
+ const secondText = await second.text()
+ const secondState = (await host.storage.get(
+ DUMP_STATE_KEY
+ )) as DumpState
+ const planCallsAfterSecond = vi
+ .mocked(executeOperation)
+ .mock.calls.filter(([queries]) =>
+ String((queries as any[])[0].sql).includes('name NOT LIKE')
+ ).length
+ const status = await dumpJobStatus(host)
+
+ expect(first.status).toBe(200)
+ expect(second.status).toBe(200)
+ expect(secondText).toBe(firstText)
+ expect(secondState.dumpId).toBe(firstState.dumpId)
+ expect(planCallsAfterSecond).toBe(planCallsAfterFirst)
+ expect(status.status).toBe(200)
+ expect(await status.text()).toBe(firstText)
+ expect(fetchSpy).not.toHaveBeenCalled()
+ expect('callbackUrl' in firstState).toBe(false)
+ fetchSpy.mockRestore()
+ })
+
+ it('normalizes the part size in persisted job state', async () => {
+ setupCompleteQueries()
+ const { host } = makeHost()
+ const response = await runDumpJob(
+ host,
+ new URLSearchParams({ partBytes: '1' })
+ )
+ await response.arrayBuffer()
+ const state = (await host.storage.get(DUMP_STATE_KEY)) as DumpState
+
+ expect(state.finalizePartSizeBytes).toBe(MIN_R2_PART_SIZE_BYTES)
+ })
+
+ it('returns a resumable response and schedules the next cycle when work remains', async () => {
+ setupCompleteQueries()
+ const { host, alarms } = makeHost()
+ const response = await runDumpJob(
+ host,
+ new URLSearchParams({ chunkBytes: '1' }),
+ Date.now() - 10_000
+ )
+ const body = (await response.json()) as {
+ result: { status: string }
+ }
+
+ expect(response.status).toBe(202)
+ expect(body.result.status).toBe('in-progress')
+ expect(alarms.length).toBeGreaterThan(0)
+ })
+})
diff --git a/src/export/dump.ts b/src/export/dump.ts
index 91a2e89..c14fa3b 100644
--- a/src/export/dump.ts
+++ b/src/export/dump.ts
@@ -1,71 +1,319 @@
-import { executeOperation } from '.'
+import { executeOperation } from './index'
import { StarbaseDBConfiguration } from '../handler'
import { DataSource } from '../types'
import { createResponse } from '../utils'
+import {
+ ChunkedDumpEngine,
+ DEFAULT_DUMP_OPTIONS,
+ normalizePartSizeBytes,
+ quoteIdentifier,
+ sqlCommentLabel,
+ isSafeIdentifier,
+ type DumpOptions,
+} from './chunkedDump'
-export async function dumpDatabaseRoute(
+export const AUTO_JOB_THRESHOLD_BYTES = 8 * 1024 * 1024
+
+export interface DumpJobEnv {
+ R2_DUMP_BUCKET?: R2Bucket
+}
+
+export interface DumpEngineHost {
+ storage: DurableObjectStorage
+ env: DumpJobEnv
+ dataSource: DataSource
+ config: StarbaseDBConfiguration
+ setAlarm: (
+ time: number,
+ options?: DurableObjectSetAlarmOptions
+ ) => Promise
+}
+
+export function parseDumpOptions(searchParams: URLSearchParams): DumpOptions {
+ const options: DumpOptions = {}
+ const readNumber = (key: string, target: keyof DumpOptions) => {
+ const raw = searchParams.get(key)
+ if (raw === null) return
+ const value = Number(raw)
+ if (Number.isFinite(value) && value > 0) {
+ options[target] = value as never
+ }
+ }
+ readNumber('cycleMs', 'cycleTimeBudgetMs')
+ readNumber('breathMs', 'breathingIntervalMs')
+ readNumber('rows', 'rowsPerBatch')
+ readNumber('chunkBytes', 'chunkTargetBytes')
+ const rawPartBytes = searchParams.get('partBytes')
+ if (rawPartBytes !== null) {
+ options.finalizePartSizeBytes = normalizePartSizeBytes(
+ Number(rawPartBytes)
+ )
+ }
+ readNumber('finalizeMs', 'finalizeTimeBudgetMs')
+ return options
+}
+
+function legacyIdentifier(name: string): string {
+ return isSafeIdentifier(name) ? name : quoteIdentifier(name)
+}
+
+function legacyValue(value: unknown): string {
+ if (value === null || value === undefined) return 'NULL'
+ if (typeof value === 'number')
+ return Number.isFinite(value) ? String(value) : 'NULL'
+ if (typeof value === 'bigint' || typeof value === 'boolean')
+ return String(value)
+ if (value instanceof ArrayBuffer) {
+ return `X'${Array.from(new Uint8Array(value), (b) => b.toString(16).padStart(2, '0')).join('')}'`
+ }
+ if (ArrayBuffer.isView(value)) {
+ return `X'${Array.from(new Uint8Array(value.buffer, value.byteOffset, value.byteLength), (b) => b.toString(16).padStart(2, '0')).join('')}'`
+ }
+ if (typeof value === 'string') return `'${value.replace(/'/g, "''")}'`
+ return `'${(JSON.stringify(value) ?? String(value)).replace(/'/g, "''")}'`
+}
+
+export async function legacyDump(
dataSource: DataSource,
config: StarbaseDBConfiguration
): Promise {
try {
- // Get all table names
const tablesResult = await executeOperation(
[{ sql: "SELECT name FROM sqlite_master WHERE type='table';" }],
dataSource,
config
)
- const tables = tablesResult.map((row: any) => row.name)
- let dumpContent = 'SQLite format 3\0' // SQLite file header
+ const tables = tablesResult
+ .map((row: Record) => String(row.name))
+ .filter((name: string) => name.length > 0)
+ let dumpContent = 'SQLite format 3\0'
- // Iterate through all tables
for (const table of tables) {
- // Get table schema
const schemaResult = await executeOperation(
[
{
- sql: `SELECT sql FROM sqlite_master WHERE type='table' AND name='${table}';`,
+ sql: "SELECT sql FROM sqlite_master WHERE type='table' AND name=?;",
+ params: [table],
},
],
dataSource,
config
)
- if (schemaResult.length) {
- const schema = schemaResult[0].sql
- dumpContent += `\n-- Table: ${table}\n${schema};\n\n`
+ if (schemaResult.length && schemaResult[0]?.sql) {
+ dumpContent += `\n-- Table: ${sqlCommentLabel(table)}\n${String(schemaResult[0].sql)};\n\n`
}
- // Get table data
const dataResult = await executeOperation(
- [{ sql: `SELECT * FROM ${table};` }],
+ [{ sql: `SELECT * FROM ${quoteIdentifier(table)};` }],
dataSource,
config
)
for (const row of dataResult) {
const values = Object.values(row).map((value) =>
- typeof value === 'string'
- ? `'${value.replace(/'/g, "''")}'`
- : value
+ legacyValue(value)
)
- dumpContent += `INSERT INTO ${table} VALUES (${values.join(', ')});\n`
+ dumpContent += `INSERT INTO ${legacyIdentifier(table)} VALUES (${values.join(', ')});\n`
}
dumpContent += '\n'
}
- // Create a Blob from the dump content
const blob = new Blob([dumpContent], { type: 'application/x-sqlite3' })
-
const headers = new Headers({
'Content-Type': 'application/x-sqlite3',
'Content-Disposition': 'attachment; filename="database_dump.sql"',
})
-
return new Response(blob, { headers })
} catch (error: any) {
console.error('Database Dump Error:', error)
return createResponse(undefined, 'Failed to create database dump', 500)
}
}
+
+export async function shouldUseResumableDump(
+ dataSource: DataSource,
+ config: StarbaseDBConfiguration
+): Promise {
+ try {
+ const readPragma = async (sql: string): Promise => {
+ const result = await executeOperation([{ sql }], dataSource, config)
+ const row = result[0] as Record | undefined
+ const value = row ? Object.values(row)[0] : undefined
+ const number = Number(value)
+ return Number.isFinite(number) ? number : 0
+ }
+ const pageCount = await readPragma('PRAGMA page_count;')
+ const pageSize = await readPragma('PRAGMA page_size;')
+ return pageCount * pageSize >= AUTO_JOB_THRESHOLD_BYTES
+ } catch {
+ return true
+ }
+}
+
+export async function runDumpJob(
+ host: DumpEngineHost,
+ searchParams: URLSearchParams,
+ requestStart: number = Date.now()
+): Promise {
+ try {
+ const parsedOptions = parseDumpOptions(searchParams)
+ const options = { ...DEFAULT_DUMP_OPTIONS, ...parsedOptions }
+ const engine = new ChunkedDumpEngine(
+ host.storage,
+ host.env.R2_DUMP_BUCKET,
+ host.dataSource,
+ host.config,
+ options
+ )
+
+ let state = await engine.getState()
+ if (!state) {
+ state = await engine.startDump()
+ } else if (!state.completedAt) {
+ state = await engine.runCycle()
+ }
+
+ while (
+ !state.completedAt &&
+ Date.now() - requestStart < options.cycleTimeBudgetMs
+ ) {
+ state = await engine.runCycle()
+ }
+
+ if (state.completedAt) {
+ if (
+ host.env.R2_DUMP_BUCKET &&
+ (!state.finalizedAt || !state.temporaryChunksCleanedAt)
+ ) {
+ try {
+ await host.setAlarm(Date.now() + 1_000)
+ } catch (alarmError) {
+ console.error(
+ 'Failed to schedule dump finalize:',
+ alarmError
+ )
+ }
+ }
+ const stream = await engine.assembleDump(state)
+ if (stream) {
+ return new Response(stream, {
+ headers: {
+ 'Content-Type': 'application/x-sqlite3',
+ 'Content-Disposition': `attachment; filename="${state.fileName}"`,
+ },
+ })
+ }
+ }
+
+ const resumeAt = Date.now() + options.breathingIntervalMs
+ await host.setAlarm(resumeAt)
+ return createResponse(
+ {
+ dumpId: state.dumpId,
+ status: 'in-progress',
+ phase: state.phase,
+ progress: {
+ tablesTotal: state.tables.length,
+ tableIndex: state.tableIndex,
+ totalRows: state.totalRows,
+ bytesWritten: state.bytesWritten,
+ chunkIndex: state.chunkIndex,
+ },
+ resumeAt,
+ fileName: state.fileName,
+ },
+ undefined,
+ 202
+ )
+ } catch (error: any) {
+ console.error('Database Dump Error:', error)
+ return createResponse(undefined, 'Failed to create database dump', 500)
+ }
+}
+
+export async function dumpJobStatus(host: DumpEngineHost): Promise {
+ const engine = new ChunkedDumpEngine(
+ host.storage,
+ host.env.R2_DUMP_BUCKET,
+ host.dataSource,
+ host.config,
+ DEFAULT_DUMP_OPTIONS
+ )
+ const state = await engine.getState()
+ if (!state) {
+ return createResponse(undefined, 'No dump job found', 404)
+ }
+ if (!state.completedAt) {
+ return createResponse(
+ {
+ dumpId: state.dumpId,
+ status: 'in-progress',
+ phase: state.phase,
+ progress: {
+ tablesTotal: state.tables.length,
+ tableIndex: state.tableIndex,
+ totalRows: state.totalRows,
+ bytesWritten: state.bytesWritten,
+ chunkIndex: state.chunkIndex,
+ },
+ fileName: state.fileName,
+ },
+ undefined,
+ 202
+ )
+ }
+ if (state.finalObjectKey) {
+ const downloadUrl = await engine.getPresignedUrl()
+ if (downloadUrl) {
+ return createResponse(
+ {
+ dumpId: state.dumpId,
+ status: 'complete',
+ downloadUrl,
+ downloadUrlExpiresInSeconds: 3600,
+ downloadType: 'presigned-url',
+ finalObjectKey: state.finalObjectKey,
+ size: state.finalObjectSize ?? state.bytesWritten,
+ totalRows: state.totalRows,
+ fileName: state.fileName,
+ },
+ undefined,
+ 200
+ )
+ }
+ }
+ const stream = await engine.assembleDump(state)
+ if (!stream) {
+ return createResponse(undefined, 'Dump data unavailable', 410)
+ }
+ return new Response(stream, {
+ headers: {
+ 'Content-Type': 'application/x-sqlite3',
+ 'Content-Disposition': `attachment; filename="${state.fileName}"`,
+ },
+ })
+}
+
+export async function runDumpFinalize(host: DumpEngineHost): Promise {
+ const engine = new ChunkedDumpEngine(
+ host.storage,
+ host.env.R2_DUMP_BUCKET,
+ host.dataSource,
+ host.config,
+ DEFAULT_DUMP_OPTIONS
+ )
+ const { done } = await engine.finalizeDump()
+ if (!done) {
+ await host.setAlarm(Date.now() + 1_000)
+ }
+}
+
+export async function dumpDatabaseRoute(
+ dataSource: DataSource,
+ config: StarbaseDBConfiguration
+): Promise {
+ return legacyDump(dataSource, config)
+}
diff --git a/src/handler.dump.test.ts b/src/handler.dump.test.ts
new file mode 100644
index 0000000..5ebc7db
--- /dev/null
+++ b/src/handler.dump.test.ts
@@ -0,0 +1,87 @@
+import { beforeEach, describe, expect, it, vi } from 'vitest'
+import { StarbaseDB } from './handler'
+import type { DataSource } from './types'
+
+const makeContext = () =>
+ ({
+ waitUntil: vi.fn(),
+ }) as unknown as ExecutionContext
+
+const makeSource = (large: boolean) => {
+ const startDumpJob = vi.fn(async () => new Response('job', { status: 202 }))
+ const executeQuery = vi.fn(async ({ sql }: { sql: string }) => {
+ if (sql.includes('PRAGMA page_count')) {
+ return [{ page_count: large ? 4096 : 1 }]
+ }
+ if (sql.includes('PRAGMA page_size')) {
+ return [{ page_size: 4096 }]
+ }
+ if (sql.includes('SELECT name FROM sqlite_master')) {
+ return [{ name: 't' }]
+ }
+ if (sql.includes('SELECT sql FROM sqlite_master')) {
+ return [{ sql: 'CREATE TABLE t (id INTEGER)' }]
+ }
+ return []
+ })
+ const dataSource = {
+ source: 'internal',
+ rpc: {
+ executeQuery,
+ startDumpJob,
+ dumpJobStatus: vi.fn(async () => new Response('status')),
+ },
+ } as unknown as DataSource
+ return { dataSource, startDumpJob, executeQuery }
+}
+
+describe('StarbaseDB dump routing', () => {
+ beforeEach(() => {
+ vi.clearAllMocks()
+ })
+
+ it('automatically routes a large internal dump to a resumable job', async () => {
+ const { dataSource, startDumpJob } = makeSource(true)
+ const app = new StarbaseDB({
+ dataSource,
+ config: {
+ role: 'admin',
+ features: { export: true, rls: false, allowlist: false },
+ },
+ })
+ const response = await app.handle(
+ new Request('https://example.test/export/dump'),
+ makeContext()
+ )
+
+ expect(response.status).toBe(202)
+ expect(await response.text()).toBe('job')
+ expect(startDumpJob).toHaveBeenCalledOnce()
+ })
+
+ it('keeps small internal dumps synchronous and preserves explicit job requests', async () => {
+ const { dataSource, startDumpJob } = makeSource(false)
+ const app = new StarbaseDB({
+ dataSource,
+ config: {
+ role: 'admin',
+ features: { export: true, rls: false, allowlist: false },
+ },
+ })
+
+ const small = await app.handle(
+ new Request('https://example.test/export/dump'),
+ makeContext()
+ )
+ expect(small.status).toBe(200)
+ expect(await small.text()).toContain('CREATE TABLE t')
+ expect(startDumpJob).not.toHaveBeenCalled()
+
+ const explicit = await app.handle(
+ new Request('https://example.test/export/dump?job=1'),
+ makeContext()
+ )
+ expect(explicit.status).toBe(202)
+ expect(startDumpJob).toHaveBeenCalledOnce()
+ })
+})
diff --git a/src/handler.ts b/src/handler.ts
index 3fa0085..5dc471e 100644
--- a/src/handler.ts
+++ b/src/handler.ts
@@ -6,7 +6,7 @@ import { DataSource } from './types'
import { LiteREST } from './literest'
import { executeQuery, executeTransaction } from './operation'
import { createResponse, QueryRequest, QueryTransactionRequest } from './utils'
-import { dumpDatabaseRoute } from './export/dump'
+import { dumpDatabaseRoute, shouldUseResumableDump } from './export/dump'
import { exportTableToJsonRoute } from './export/json'
import { exportTableToCsvRoute } from './export/csv'
import { importDumpRoute } from './import/dump'
@@ -120,10 +120,55 @@ export class StarbaseDB {
}
if (this.getFeature('export')) {
- this.app.get('/export/dump', this.isInternalSource, async () => {
+ this.app.get('/export/dump', this.isInternalSource, async (c) => {
+ const url = new URL(c.req.raw.url)
+ let wantsJob = url.searchParams.get('job') === '1'
+
+ if (
+ !wantsJob &&
+ this.dataSource.source === 'internal' &&
+ (await shouldUseResumableDump(this.dataSource, this.config))
+ ) {
+ wantsJob = true
+ }
+
+ if (wantsJob && this.dataSource.source === 'internal') {
+ const searchParams: Record = {}
+ url.searchParams.forEach((value, key) => {
+ searchParams[key] = value
+ })
+ const rpc = this.dataSource.rpc as unknown as {
+ startDumpJob: (
+ config: StarbaseDBConfiguration,
+ searchParams: Record
+ ) => Promise
+ }
+ return await rpc.startDumpJob(this.config, searchParams)
+ }
+
return dumpDatabaseRoute(this.dataSource, this.config)
})
+ this.app.get(
+ '/export/dump/status',
+ this.isInternalSource,
+ async () => {
+ if (this.dataSource.source === 'internal') {
+ const rpc = this.dataSource.rpc as unknown as {
+ dumpJobStatus: (
+ config: StarbaseDBConfiguration
+ ) => Promise
+ }
+ return await rpc.dumpJobStatus(this.config)
+ }
+ return createResponse(
+ undefined,
+ 'Chunked dump status requires the internal data source',
+ 400
+ )
+ }
+ )
+
this.app.get(
'/export/json/:tableName',
this.isInternalSource,