Skip to content
Merged
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
7 changes: 7 additions & 0 deletions docs/contributing/disaster-recovery.md
Original file line number Diff line number Diff line change
Expand Up @@ -188,6 +188,13 @@ Creates `kody-dr-drill-{day}-{hex}` in `DRILL_ACCOUNT_ID` (must differ from
success (keeps it on failure for inspection). This path does not restore
StorageRunner/R2/artifacts into production.

D1 remote import enforces foreign keys while applying `CREATE TABLE`, but
Cloudflare exports are not topologically ordered. Before upload, the control
plane verifies the unmodified SQL MD5 against the signed manifest R2 ETag, then
prefixes `PRAGMA foreign_keys=OFF;` and uses the prepared body's MD5 for the D1
import init/ingest etag. Without that prelude, drills fail with errors like
`no such table: main.users`.

### Graduated production restore

1. **Prepare** — sealed full manifest + D1 manifest signatures must verify; UI
Expand Down
128 changes: 120 additions & 8 deletions packages/backup-control-plane/d1-import-api.node.test.ts
Original file line number Diff line number Diff line change
@@ -1,32 +1,72 @@
import assert from 'node:assert/strict'
import { createHash } from 'node:crypto'

import { test, vi } from 'vitest'

import { importSqlIntoD1 } from './d1-import-api.ts'
import {
d1ImportForeignKeysOffPrefix,
importSqlIntoD1,
prepareD1ImportUpload,
} from './d1-import-api.ts'
import { BackupError } from './backup-policy.ts'

const ACCOUNT = '11111111-1111-4111-8111-111111111111'
const DATABASE = '22222222-2222-4222-8222-222222222222'
const MD5 = 'a'.repeat(32)
const SOURCE_SQL = 'CREATE TABLE t(id INTEGER);\n'
const SOURCE_MD5 = createHash('md5').update(SOURCE_SQL).digest('hex')
const UPLOAD_MD5 = createHash('md5')
.update(d1ImportForeignKeysOffPrefix)
.update(SOURCE_SQL)
.digest('hex')

function importPollSequence(pollResponses: Array<unknown>) {
let phase: 'init' | 'upload' | 'ingest' | 'poll' = 'init'
let pollIndex = 0
let uploadedText: string | null = null
let initEtag: string | null = null
const fetcher: typeof fetch = async (input, init) => {
const url = String(input)
if (url.includes('upload.example')) {
phase = 'ingest'
const body = init?.body
if (typeof body === 'string') {
uploadedText = body
} else if (body instanceof Uint8Array) {
uploadedText = new TextDecoder().decode(body)
} else if (body instanceof ArrayBuffer) {
uploadedText = new TextDecoder().decode(body)
} else if (body instanceof ReadableStream) {
const reader = body.getReader()
const chunks: Array<Uint8Array> = []
for (;;) {
const { done, value } = await reader.read()
if (done) break
if (value) chunks.push(value)
}
const total = chunks.reduce((sum, chunk) => sum + chunk.byteLength, 0)
const merged = new Uint8Array(total)
let offset = 0
for (const chunk of chunks) {
merged.set(chunk, offset)
offset += chunk.byteLength
}
uploadedText = new TextDecoder().decode(merged)
} else {
throw new Error(`unexpected upload body type: ${typeof body}`)
}
return new Response(null, {
status: 200,
headers: { etag: `"${MD5}"` },
headers: { etag: `"${UPLOAD_MD5}"` },
})
}
const body = JSON.parse(String(init?.body ?? '{}')) as {
action?: string
etag?: string
}
switch (body.action) {
case 'init':
phase = 'upload'
initEtag = typeof body.etag === 'string' ? body.etag : null
return Response.json({
success: true,
result: {
Expand Down Expand Up @@ -55,10 +95,43 @@ function importPollSequence(pollResponses: Array<unknown>) {
return {
fetcher,
getPollCount: () => pollIndex,
getUploadedText: () => uploadedText,
getInitEtag: () => initEtag,
}
}

test('importSqlIntoD1 completes on terminal status and final bookmark shapes', async () => {
test('prepareD1ImportUpload verifies source MD5 and prefixes foreign_keys=OFF', async () => {
let loads = 0
const prepared = await prepareD1ImportUpload({
sourceMd5Etag: SOURCE_MD5,
loadSqlBody: async () => {
loads += 1
return SOURCE_SQL
},
})
assert.equal(loads, 2)
assert.equal(prepared.uploadMd5Hex, UPLOAD_MD5)
assert.equal(prepared.sourceBytes, SOURCE_SQL.length)
assert.ok(prepared.uploadBody instanceof Uint8Array)
assert.equal(
new TextDecoder().decode(prepared.uploadBody),
`${d1ImportForeignKeysOffPrefix}${SOURCE_SQL}`,
)
})

test('prepareD1ImportUpload rejects source MD5 mismatches', async () => {
await assert.rejects(
prepareD1ImportUpload({
sourceMd5Etag: 'b'.repeat(32),
loadSqlBody: async () => SOURCE_SQL,
}),
(error: unknown) =>
error instanceof BackupError &&
error.code === 'import-source-etag-mismatch',
)
})

test('importSqlIntoD1 uploads FK-off-prefixed SQL and uses its MD5', async () => {
for (const pollResponses of [
[
{ type: 'import', success: true, status: 'active' },
Expand All @@ -78,8 +151,8 @@ test('importSqlIntoD1 completes on terminal status and final bookmark shapes', a
accountId: ACCOUNT,
databaseId: DATABASE,
token: 'token',
sqlBody: 'CREATE TABLE t(id INTEGER);\n',
md5Etag: MD5,
sourceMd5Etag: SOURCE_MD5,
loadSqlBody: async () => SOURCE_SQL,
options: {
fetcher: sequence.fetcher,
sleep: async () => undefined,
Expand All @@ -88,9 +161,48 @@ test('importSqlIntoD1 completes on terminal status and final bookmark shapes', a
},
})
assert.equal(sequence.getPollCount(), 2)
assert.equal(sequence.getInitEtag(), UPLOAD_MD5)
assert.equal(
sequence.getUploadedText(),
`${d1ImportForeignKeysOffPrefix}${SOURCE_SQL}`,
)
}
})

test('importSqlIntoD1 streams reopen through loadSqlBody without buffering the source twice in one handle', async () => {
const sequence = importPollSequence([
{ type: 'import', success: true, status: 'complete' },
])
let loads = 0
await importSqlIntoD1({
accountId: ACCOUNT,
databaseId: DATABASE,
token: 'token',
sourceMd5Etag: SOURCE_MD5,
loadSqlBody: async () => {
loads += 1
return new ReadableStream<Uint8Array>({
start(controller) {
controller.enqueue(new TextEncoder().encode(SOURCE_SQL))
controller.close()
},
})
},
options: {
fetcher: sequence.fetcher,
sleep: async () => undefined,
maxPollAttempts: 2,
pollDelayMs: 1,
},
})
assert.equal(loads, 2)
assert.equal(sequence.getInitEtag(), UPLOAD_MD5)
assert.equal(
sequence.getUploadedText(),
`${d1ImportForeignKeysOffPrefix}${SOURCE_SQL}`,
)
})

test('importSqlIntoD1 fails closed on non-terminal, expired, and error polls', async () => {
const consoleError = vi.spyOn(console, 'error')
consoleError.mockImplementation(() => undefined)
Expand Down Expand Up @@ -130,8 +242,8 @@ test('importSqlIntoD1 fails closed on non-terminal, expired, and error polls', a
accountId: ACCOUNT,
databaseId: DATABASE,
token: 'token',
sqlBody: 'CREATE TABLE t(id INTEGER);\n',
md5Etag: MD5,
sourceMd5Etag: SOURCE_MD5,
loadSqlBody: async () => SOURCE_SQL,
options: {
fetcher: sequence.fetcher,
sleep: async () => undefined,
Expand Down
Loading
Loading