Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -243,6 +243,8 @@ curl --location 'https://starbasedb.YOUR-ID-HERE.workers.dev/export/dump' \
</code>
</pre>

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.

<h3>JSON Data Export</h3>
<pre>
<code>
Expand Down
112 changes: 112 additions & 0 deletions src/do.ts
Original file line number Diff line number Diff line change
@@ -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
Expand Down Expand Up @@ -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),
}
}

Expand Down Expand Up @@ -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<DumpState>(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;',
Expand Down Expand Up @@ -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<string, string>
): Promise<Response> {
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<Response> {
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[]
Expand Down
140 changes: 140 additions & 0 deletions src/export/chunkedDump.sqlite.test.ts
Original file line number Diff line number Diff line change
@@ -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<string, unknown>()
return {
get: vi.fn(async <T>(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<typeof makeStorage>
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<string, unknown>[]
})
})

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()
})
})
Loading