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
11 changes: 9 additions & 2 deletions docs/contributing/disaster-recovery.md
Original file line number Diff line number Diff line change
Expand Up @@ -370,8 +370,15 @@ config, and expect users to reauthorize OAuth and remote connectors.
| Hourly `:45` | Control plane | D1 freshness (identity, size ceiling, manifest age ≤26h, R2 size/ETag) + seal recent complete days |

Hourly freshness does not SHA-256 the SQL bytes; drills do. Page yourself on
`freshness-stale`, size-ceiling hits, missing manifests, seal failures,
`backup-unrestorable-statements`, or unexpected disablement of the enable gates.
`freshness-unrestorable` (the SQL contains statements D1 cannot import),
`freshness-stale`, size-ceiling hits, missing manifests or required SQL stats,
seal failures, `backup-unrestorable-statements`, or unexpected disablement of
the enable gates. Post-cutover unrestorable exports never receive a canonical
day manifest; catch-up retries can re-export the day after the oversized row or
write path is bounded. Historical bad days may already have canonical and full
manifests; immutable media is intentionally left unchanged, so freshness,
dashboard, drill, and production-restore gates remain essential until it ages
out.

## Offline CLI fallback

Expand Down
88 changes: 85 additions & 3 deletions packages/backup-control-plane/backup-control-plane-test-support.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,11 @@ import {
canonicalBackupManifestPayload,
type BackupManifestPayload,
} from '@kody-internal/shared/backup-manifest.ts'
import {
backupSqlStatsKey,
backupSqlStatsSchemaVersion,
type BackupSqlStats,
} from '@kody-internal/shared/backup-sql-stats.ts'
import { type BackupRuntimeStep } from './backup-runtime.ts'
import { type DurableExportStep } from './durable-export.ts'
import { BackupError, objectKeyForBookmark } from './backup-policy.ts'
Expand Down Expand Up @@ -81,6 +86,7 @@ export class MemoryBucket {
readonly puts: Array<{ key: string; options: R2PutOptions }> = []
private readonly objects = new Map<string, Uint8Array>()
private readonly reportedSizes = new Map<string, number>()
private readonly failedGets = new Set<string>()
private nextPutRace: { key: string; bytes: Uint8Array } | undefined
private lockPolicyEnabled = false

Expand Down Expand Up @@ -126,7 +132,12 @@ export class MemoryBucket {
return this.objects.has(key) ? this.metadata(key) : null
}

async delete(key: string): Promise<void> {
this.objects.delete(key)
}

async get(key: string): Promise<R2ObjectBody | null> {
if (this.failedGets.has(key)) throw new Error('simulated R2 get failure')
const bytes = this.objects.get(key)
if (!bytes) return null
const metadata = this.metadata(key)
Expand All @@ -151,6 +162,10 @@ export class MemoryBucket {
this.reportedSizes.set(key, size)
}

failGetFor(key: string): void {
this.failedGets.add(key)
}

raceOnNextPut(key: string, value: string): void {
this.nextPutRace = { key, bytes: new TextEncoder().encode(value) }
}
Expand Down Expand Up @@ -179,6 +194,36 @@ export const DRILL_ACCOUNT_ID = '33333333-3333-4333-8333-333333333333'
export const BASELINE_SHA256 = 'b'.repeat(64)
const manifestSigningKeys = generateKeyPairSync('ed25519')

export function badSqlStatsFixture(
day: string,
objectKey: string,
): BackupSqlStats {
return {
schemaVersion: backupSqlStatsSchemaVersion,
day,
objectKey,
maxStatementBytes: 1_495_663,
oversizedStatementCount: 1,
importStatementLimitBytes: 100_000,
}
}

export async function putSqlStatsFixture(
bucket: MemoryBucket,
day: string,
objectKey: string,
options: { oversized?: boolean } = {},
): Promise<void> {
const stats = options.oversized
? badSqlStatsFixture(day, objectKey)
: {
...badSqlStatsFixture(day, objectKey),
maxStatementBytes: 50_000,
oversizedStatementCount: 0,
}
await bucket.put(backupSqlStatsKey(objectKey), JSON.stringify(stats))
}

export function environment(bucket = new MemoryBucket()): BackupEnvironment {
return {
BACKUP_BUCKET: bucket as unknown as R2Bucket,
Expand Down Expand Up @@ -307,11 +352,11 @@ export class RetryAfterCommitStep implements BackupRuntimeStep {

export class CachedUploadStep implements BackupRuntimeStep {
private readonly cache = new Map<string, unknown>()
private readonly afterUpload: () => void
private readonly afterUpload: () => void | Promise<void>
private readonly afterFinalization?: () => Promise<void>

constructor(
afterUpload: () => void,
afterUpload: () => void | Promise<void>,
afterFinalization?: () => Promise<void>,
) {
this.afterUpload = afterUpload
Expand All @@ -336,7 +381,7 @@ export class CachedUploadStep implements BackupRuntimeStep {
: callback!
const value = await execute()
this.cache.set(name, value)
if (name === 'stream-export-to-immutable-r2') this.afterUpload()
if (name === 'stream-export-to-immutable-r2') await this.afterUpload()
if (
name === 'verify-stored-object-and-write-immutable-manifest' &&
this.afterFinalization
Expand All @@ -349,6 +394,43 @@ export class CachedUploadStep implements BackupRuntimeStep {
async sleep(): Promise<void> {}
}

export class PreStatsUploadStep implements BackupRuntimeStep {
private readonly cache = new Map<string, unknown>()

async do<T>(
name: string,
config: unknown,
callback: () => Promise<T>,
): Promise<T>
async do<T>(name: string, callback: () => Promise<T>): Promise<T>
async do<T>(
name: string,
configOrCallback: unknown,
callback?: () => Promise<T>,
): Promise<T> {
if (this.cache.has(name)) return this.cache.get(name) as T
const execute =
typeof configOrCallback === 'function'
? (configOrCallback as () => Promise<T>)
: callback!
const value = await execute()
const persisted =
name === 'stream-export-to-immutable-r2' &&
value !== null &&
typeof value === 'object'
? Object.fromEntries(
Object.entries(value).filter(
([key]) => key !== 'sqlStatementStats',
),
)
: value
this.cache.set(name, persisted)
return persisted as T
}

async sleep(): Promise<void> {}
}

export class RetryUploadStep implements BackupRuntimeStep {
readonly uploadAttempts: number[] = []
private readonly cache = new Map<string, unknown>()
Expand Down
163 changes: 163 additions & 0 deletions packages/backup-control-plane/backup-runtime.node.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,8 +15,10 @@ import {
CachedUploadStep,
DATABASE_ID,
MemoryBucket,
PreStatsUploadStep,
RetryAfterCommitStep,
RetryUploadStep,
badSqlStatsFixture,
environment,
exportEnvelope,
identityEnvelope,
Expand Down Expand Up @@ -287,6 +289,167 @@ test('workflow retry reuses an upload committed before step persistence and writ
assert.ok(stats.maxStatementBytes > 0)
})

test('oversized SQL writes stats then fails retryably without a day manifest', async () => {
const consoleError = vi.spyOn(console, 'error')
consoleError.mockImplementation(() => undefined)
const bucket = new MemoryBucket()
const env = environment(bucket)
const payload = backupPayload(env, new Date('2026-07-31T02:15:00Z'))
const objectKey = objectKeyForBookmark(payload.objectPrefix, 'bookmark-1')
const sql = `INSERT INTO t VALUES ('${'x'.repeat(100_001)}');`

await assert.rejects(
runBackupRuntime(
env,
{
instanceId: workflowInstanceId(DATABASE_ID, payload.day),
payload,
timestamp: new Date('2026-07-31T02:15:01Z'),
},
new CachedUploadStep(() => undefined),
{
api: {
fetcher: async (input) =>
String(input).endsWith('/export')
? exportEnvelope('complete')
: identityEnvelope(1_000),
sleep: async () => undefined,
},
downloadFetcher: async () =>
new Response(sql, {
headers: { 'content-length': String(sql.length) },
}),
},
),
(error: unknown) =>
error instanceof BackupError &&
error.code === 'backup-unrestorable-statements' &&
error.retryable,
)

const statsObject = await bucket.get(`${objectKey}.stats.json`)
assert.notEqual(statsObject, null)
const stats = (await statsObject!.json()) as {
oversizedStatementCount: number
}
assert.equal(stats.oversizedStatementCount, 1)
assert.equal(await bucket.head(payload.manifestKey), null)
const events = consoleError.mock.calls.map(([record]) =>
JSON.parse(String(record)),
) as Array<{ event: string }>
assert.ok(
events.some(({ event }) => event === 'backup-unrestorable-statements'),
)
assert.ok(events.some(({ event }) => event === 'backup-failure'))
})

test('cached pre-stats uploads are allowed only for legacy backup days', async () => {
const consoleError = vi.spyOn(console, 'error')
consoleError.mockImplementation(() => undefined)
const consoleLog = vi.spyOn(console, 'log')
consoleLog.mockImplementation(() => undefined)
const options = {
api: {
fetcher: async (input: RequestInfo | URL) =>
String(input).endsWith('/export')
? exportEnvelope('complete')
: identityEnvelope(1_000),
sleep: async () => undefined,
},
downloadFetcher: async () =>
new Response('valid', { headers: { 'content-length': '5' } }),
}

const legacyBucket = new MemoryBucket()
const legacyEnv = environment(legacyBucket)
const legacyPayload = backupPayload(
legacyEnv,
new Date('2026-07-27T02:15:00Z'),
)
await runBackupRuntime(
legacyEnv,
{
instanceId: workflowInstanceId(DATABASE_ID, legacyPayload.day),
payload: legacyPayload,
timestamp: new Date('2026-07-27T02:15:01Z'),
},
new PreStatsUploadStep(),
options,
)
assert.notEqual(await legacyBucket.head(legacyPayload.manifestKey), null)
const legacyEvents = consoleLog.mock.calls.map(([record]) =>
JSON.parse(String(record)),
) as Array<{ event: string }>
assert.ok(
legacyEvents.some(({ event }) => event === 'backup-stats-legacy-missing'),
)

const requiredBucket = new MemoryBucket()
const requiredEnv = environment(requiredBucket)
const requiredPayload = backupPayload(
requiredEnv,
new Date('2026-07-28T02:15:00Z'),
)
await assert.rejects(
runBackupRuntime(
requiredEnv,
{
instanceId: workflowInstanceId(DATABASE_ID, requiredPayload.day),
payload: requiredPayload,
timestamp: new Date('2026-07-28T02:15:01Z'),
},
new PreStatsUploadStep(),
options,
),
(error: unknown) =>
error instanceof BackupError &&
error.code === 'backup-sql-stats-missing' &&
error.retryable,
)
assert.equal(await requiredBucket.head(requiredPayload.manifestKey), null)
})

test('conflicting immutable SQL stats prevent manifest publication', async () => {
const consoleError = vi.spyOn(console, 'error')
consoleError.mockImplementation(() => undefined)
const bucket = new MemoryBucket()
const env = environment(bucket)
const payload = backupPayload(env, new Date('2026-07-31T02:15:00Z'))
const objectKey = objectKeyForBookmark(payload.objectPrefix, 'bookmark-1')

await assert.rejects(
runBackupRuntime(
env,
{
instanceId: workflowInstanceId(DATABASE_ID, payload.day),
payload,
timestamp: new Date('2026-07-31T02:15:01Z'),
},
new CachedUploadStep(async () => {
await bucket.put(
`${objectKey}.stats.json`,
JSON.stringify(badSqlStatsFixture(payload.day, objectKey)),
)
}),
{
api: {
fetcher: async (input) =>
String(input).endsWith('/export')
? exportEnvelope('complete')
: identityEnvelope(1_000),
sleep: async () => undefined,
},
downloadFetcher: async () =>
new Response('valid', { headers: { 'content-length': '5' } }),
},
),
(error: unknown) =>
error instanceof BackupError &&
error.code === 'backup-sql-stats-conflict',
)
assert.equal(await bucket.head(payload.manifestKey), null)
})

test('initial upload ignores a stale cached signed URL and refreshes it in the callback', async () => {
const bucket = new MemoryBucket()
const env = environment(bucket)
Expand Down
Loading
Loading