Skip to content

Commit ec65be5

Browse files
committed
Preserve lossless D1 rollback state
1 parent 46487a4 commit ec65be5

10 files changed

Lines changed: 120 additions & 28 deletions

packages/runtime-cloudflare/README.md

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -36,7 +36,9 @@ Canonical Cloudflare boots patch only the assembled PHP MEMFS copy of `/wordpres
3636

3737
The authenticated mutation-fence endpoints provide a bounded cutover window without interrupting anonymous published reads. `POST ?phase=operator-fence-acquire` accepts `{"ttlSeconds":30}` through `{"ttlSeconds":600}` and returns an opaque token plus expiry. While active, both coordinators reject new canonical leases, reset, adoption, manual publication, scheduled publication, and cron work with a conflict. `POST ?phase=operator-fence-renew` accepts the token and a fresh bounded TTL; `POST ?phase=operator-fence-release` accepts the token. Expiry automatically reopens mutations, and status responses never expose the token.
3838

39-
Use authenticated `GET ?phase=operator-fence-status` after acquisition to export a coherent cutover envelope. It includes the coordinator store, pointer, version, matching commit receipt, validated R2 manifest identity, fence expiry, and a `coherent` verdict. Coherence also reads every referenced canonical object and verifies its declared size and SHA-256, so an incomplete target cannot be promoted. Adopt that exact pointer/version into the target coordinator, require the target status envelope to match, then promote the target Worker before the source fence expires. The first target mutation must commit version `N+1`. Rollback uses the same sequence in reverse after fencing the active target; relying on an unfenced status read is not a lossless cutover procedure.
39+
Use authenticated `GET ?phase=operator-fence-status` after acquisition to export a coherent cutover envelope. It includes the selected coordinator store, pointer, version, matching commit receipt, validated R2 manifest identity, fence expiry, and a `coherent` verdict. Coherence also reads every referenced canonical object and verifies its declared size and SHA-256, so an incomplete target cannot be promoted. Adopt that exact pointer/version into the target coordinator, require the target status envelope to match, then promote the target Worker before the source fence expires. The first target mutation must commit version `N+1`.
40+
41+
For lossless D1-to-Durable-Object rollback, keep D1 deployed as the active Worker, acquire its fence, and capture its coherent status. Address the retained Durable Object only with `coordinator=durable-object` on authenticated `operator-fence-*` and `operator-adopt` requests. Fence that selected coordinator first, then adopt the D1 pointer and version with its matching `fenceToken`; selected-coordinator adoption fails unless that fence is active and matches. Require its selected status envelope to be coherent and identical, then promote the Durable Object Worker before the D1 fence expires and release its fence. Its first mutation must commit at `N+1`. The Durable Object binding and historical `WordPressStateCoordinator` export remain in the D1 profile for this procedure. Non-operator requests always retain the active coordinator, and selectors on unsupported or unknown operator routes fail closed.
4042

4143
## Static Artifact Import
4244

packages/runtime-cloudflare/src/d1-revision-coordinator.ts

Lines changed: 16 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -66,8 +66,8 @@ export class D1RevisionCoordinator implements RevisionCoordinator {
6666
await releaseMutationFence(this.database, this.siteId, token)
6767
}
6868

69-
adopt(pointer: MarkdownPointer, version: number): Promise<{ pointer: MarkdownPointer; version: number }> {
70-
return adoptWordPressState(this.database, this.siteId, pointer, version)
69+
adopt(pointer: MarkdownPointer, version: number, fenceToken?: string, requireFence = false): Promise<{ pointer: MarkdownPointer; version: number }> {
70+
return adoptWordPressState(this.database, this.siteId, pointer, version, fenceToken, requireFence)
7171
}
7272

7373
async reset(): Promise<void> {
@@ -202,34 +202,41 @@ async function releaseMutationFence(database: D1Database, siteId: string, token:
202202
return { released: true }
203203
}
204204

205-
async function adoptWordPressState(database: D1Database, siteId: string, pointer: MarkdownPointer, version: number): Promise<{ pointer: MarkdownPointer; version: number }> {
205+
async function adoptWordPressState(database: D1Database, siteId: string, pointer: MarkdownPointer, version: number, fenceToken?: string, requireFence = false): Promise<{ pointer: MarkdownPointer; version: number }> {
206206
validatePointer(pointer)
207207
if (!Number.isSafeInteger(version) || version < 1) throw new RevisionConflict("A positive canonical version is required for D1 adoption.")
208208
await ensureSchema(database, siteId)
209209
const now = Date.now()
210+
const fence = await readMutationFence(database, siteId)
211+
if (requireFence && !fence.active) throw new RevisionConflict("Selected coordinator adoption requires an active cutover fence.")
212+
if (fence.active && !fenceToken) throw new RevisionConflict("D1 coordinator adoption is blocked by an unmatched fence.", fence.expiresAt)
210213
await database.batch([
211214
database.prepare(`UPDATE wp_codebox_state
212215
SET revision = ?, manifest_key = ?, persisted_at = ?, version = ?,
213216
lease_token = NULL, lease_base_revision = NULL, lease_version = NULL, lease_expires_at = NULL
214217
WHERE site_id = ? AND (lease_token IS NULL OR lease_expires_at <= ?)
215-
AND NOT EXISTS (SELECT 1 FROM wp_codebox_fences WHERE site_id = ? AND expires_at > ?)
218+
AND (NOT EXISTS (SELECT 1 FROM wp_codebox_fences WHERE site_id = ? AND expires_at > ?)
219+
OR EXISTS (SELECT 1 FROM wp_codebox_fences WHERE site_id = ? AND token = ? AND expires_at > ?))
216220
AND ((revision IS NULL AND manifest_key IS NULL AND persisted_at IS NULL)
217-
OR (version = ? AND revision = ? AND manifest_key = ? AND persisted_at = ?))
221+
OR (version = ? AND revision = ? AND manifest_key = ? AND persisted_at = ?)
222+
OR version < ?)
218223
AND (NOT EXISTS (SELECT 1 FROM wp_codebox_commits WHERE site_id = ? AND version = ?)
219224
OR EXISTS (SELECT 1 FROM wp_codebox_commits WHERE site_id = ? AND version = ?
220225
AND revision = ? AND manifest_key = ? AND persisted_at = ?))`)
221226
.bind(pointer.revision, pointer.manifestKey, pointer.persistedAt, version, siteId, now, siteId, now,
222-
version, pointer.revision, pointer.manifestKey, pointer.persistedAt,
227+
siteId, fenceToken ?? null, now,
228+
version, pointer.revision, pointer.manifestKey, pointer.persistedAt, version,
223229
siteId, version, siteId, version, pointer.revision, pointer.manifestKey, pointer.persistedAt),
224230
database.prepare(`INSERT OR IGNORE INTO wp_codebox_commits (site_id, version, revision, manifest_key, persisted_at)
225231
SELECT site_id, version, revision, manifest_key, persisted_at FROM wp_codebox_state
226232
WHERE site_id = ? AND version = ? AND revision = ? AND manifest_key = ? AND persisted_at = ? AND lease_token IS NULL
227-
AND NOT EXISTS (SELECT 1 FROM wp_codebox_fences WHERE site_id = ? AND expires_at > ?)`)
228-
.bind(siteId, version, pointer.revision, pointer.manifestKey, pointer.persistedAt, siteId, now),
233+
AND (NOT EXISTS (SELECT 1 FROM wp_codebox_fences WHERE site_id = ? AND expires_at > ?)
234+
OR EXISTS (SELECT 1 FROM wp_codebox_fences WHERE site_id = ? AND token = ? AND expires_at > ?))`)
235+
.bind(siteId, version, pointer.revision, pointer.manifestKey, pointer.persistedAt, siteId, now, siteId, fenceToken ?? null, now),
229236
])
230237
const [state, committed] = await Promise.all([readRow(database, siteId), readCommittedPointer(database, siteId, version)])
231238
if (state.version !== version || !samePointer(pointerFromRow(state), pointer) || !samePointer(committed, pointer)) {
232-
throw new RevisionConflict("D1 coordinator adoption requires empty or exactly matching state without an active lease or fence.")
239+
throw new RevisionConflict("D1 coordinator adoption requires empty, exactly matching, or monotonic forward state without an active lease or an unmatched fence.")
233240
}
234241
return { pointer, version }
235242
}
Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,14 @@
1+
import type { RevisionCoordinator } from "./revision-coordinator.js"
2+
import type { WorkerRequestRoute } from "./request-routing.js"
3+
4+
export function selectOperatorCoordinator(
5+
active: RevisionCoordinator,
6+
route: WorkerRequestRoute,
7+
selector: string | null,
8+
resolve?: (selector: string) => RevisionCoordinator | undefined,
9+
): { coordinator: RevisionCoordinator; selected: boolean } | null {
10+
if (!selector) return { coordinator: active, selected: false }
11+
if (route.kind !== "operator-adopt" && route.kind !== "operator-fence") return route.kind.startsWith("operator-") ? null : { coordinator: active, selected: false }
12+
const coordinator = resolve?.(selector)
13+
return coordinator ? { coordinator, selected: true } : null
14+
}

packages/runtime-cloudflare/src/revision-coordinator.ts

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -56,7 +56,7 @@ export interface RevisionCoordinator {
5656
acquireFence(ttlMs: number): Promise<MutationFence>
5757
renewFence(token: string, ttlMs: number): Promise<MutationFence>
5858
releaseFence(token: string): Promise<void>
59-
adopt(pointer: MarkdownPointer, version: number): Promise<{ pointer: MarkdownPointer; version: number }>
59+
adopt(pointer: MarkdownPointer, version: number, fenceToken?: string, requireFence?: boolean): Promise<{ pointer: MarkdownPointer; version: number }>
6060
reset(): Promise<void>
6161
}
6262

packages/runtime-cloudflare/src/state-coordinator.ts

Lines changed: 8 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -71,8 +71,8 @@ export class DurableObjectRevisionCoordinator implements RevisionCoordinator {
7171
await this.call("fence-release", { token })
7272
}
7373

74-
adopt(pointer: MarkdownPointer, version: number): Promise<{ pointer: MarkdownPointer; version: number }> {
75-
return this.call("adopt", { pointer, version })
74+
adopt(pointer: MarkdownPointer, version: number, fenceToken?: string, requireFence = false): Promise<{ pointer: MarkdownPointer; version: number }> {
75+
return this.call("adopt", { pointer, version, fenceToken, requireFence })
7676
}
7777

7878
async reset(): Promise<void> {
@@ -252,7 +252,8 @@ export class WordPressStateCoordinator implements DurableObject {
252252
private async adopt(siteId: string, body: Record<string, unknown>): Promise<{ pointer: MarkdownPointer; version: number }> {
253253
const record = await this.record(siteId)
254254
const fence = await this.activeFence(record)
255-
if (fence) throw new RevisionConflict("Coordinator adoption is blocked by an active cutover fence.", fence.expiresAt)
255+
if (body.requireFence === true && !fence) throw new RevisionConflict("Selected coordinator adoption requires an active cutover fence.")
256+
if (fence && body.fenceToken !== fence.token) throw new RevisionConflict("Coordinator adoption is blocked by an active cutover fence.", fence.expiresAt)
256257
const pointer = body.pointer as MarkdownPointer
257258
const version = body.version
258259
if (!pointer || typeof pointer.revision !== "string" || typeof pointer.manifestKey !== "string" || typeof pointer.persistedAt !== "string"
@@ -261,7 +262,10 @@ export class WordPressStateCoordinator implements DurableObject {
261262
if (record.lease) delete record.lease
262263
const empty = record.version === 0 && record.pointer === null
263264
const exact = record.version === version && samePointer(record.pointer, pointer)
264-
if (!empty && !exact) throw new RevisionConflict("Coordinator adoption requires empty or exactly matching state.")
265+
const forward = (version as number) > record.version
266+
if (!empty && !exact && !forward) throw new RevisionConflict("Coordinator adoption requires empty, exactly matching, or monotonic forward state.")
267+
const committed = await this.state.storage.get<MarkdownPointer>(`wordpress-state-commit/${version}`)
268+
if (committed && !samePointer(committed, pointer)) throw new RevisionConflict("Coordinator adoption cannot replace an existing commit receipt.")
265269
record.pointer = pointer
266270
record.version = version as number
267271
await this.state.storage.transaction(async (transaction) => {
Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,12 +1,16 @@
11
import { D1RevisionCoordinator } from "./d1-revision-coordinator.js"
2-
export { WordPressStateCoordinator } from "./state-coordinator.js"
2+
import { DurableObjectRevisionCoordinator, WordPressStateCoordinator } from "./state-coordinator.js"
3+
export { WordPressStateCoordinator }
34
import type { SiteContext } from "./site-context.js"
45
import { createCloudflareRuntime, type RuntimeEnv } from "./worker.js"
56

67
interface D1RuntimeEnv extends RuntimeEnv {
78
WORDPRESS_STATE_DATABASE: D1Database
9+
WORDPRESS_STATE: DurableObjectNamespace
810
}
911

1012
export default createCloudflareRuntime<D1RuntimeEnv>((env, site: SiteContext) => (
1113
new D1RevisionCoordinator(env.WORDPRESS_STATE_DATABASE, site.id)
14+
), (env, site: SiteContext, selector) => (
15+
selector === "durable-object" ? new DurableObjectRevisionCoordinator(env.WORDPRESS_STATE.getByName(site.id), site.id) : undefined
1216
))

packages/runtime-cloudflare/src/worker.ts

Lines changed: 17 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,7 @@ import phpWasmModule from "../../../node_modules/@php-wasm/web-8-5/asyncify/8_5_
88
import { CLOUDFLARE_RUNTIME_HEALTH_MARKER, CLOUDFLARE_RUNTIME_HEALTH_SCHEMA, cloudflareRuntimeHealthResponse } from "./health-envelope.js"
99
import { leaseRetryDelayMs } from "./lease-retry.js"
1010
import { logMutationPhase, MutationRetainedBytes } from "./mutation-memory.js"
11+
import { selectOperatorCoordinator } from "./operator-coordinator.js"
1112
import { canonicalPublicRoute, MAX_PUBLISHED_PAGE_BYTES, MAX_PUBLISHED_REVISION_BYTES, MAX_PUBLISHED_ROUTES, normalizePublishedRoutes, PUBLISHED_PAGE_SCHEMA, PUBLISHED_REVISION_SCHEMA, publishedPageObjectKey, publishedRevisionObjectKey, validatePublishedRevision, type PublishedRevision } from "./published-reader.js"
1213
import { RevisionConflict, type MarkdownPointer, type MutationFence, type RevisionCoordinator, type RevisionLease } from "./revision-coordinator.js"
1314
import { routeWorkerRequest } from "./request-routing.js"
@@ -210,7 +211,10 @@ export interface RuntimeEnv {
210211
WORDPRESS_OPERATOR_TOKEN?: string
211212
}
212213

213-
export function createCloudflareRuntime<Env extends RuntimeEnv>(resolveCoordinator: (env: Env, site: SiteContext) => RevisionCoordinator) {
214+
export function createCloudflareRuntime<Env extends RuntimeEnv>(
215+
resolveCoordinator: (env: Env, site: SiteContext) => RevisionCoordinator,
216+
resolveOperatorCoordinator?: (env: Env, site: SiteContext, selector: string) => RevisionCoordinator | undefined,
217+
) {
214218
return {
215219
async fetch(request: Request, env: Env): Promise<Response> {
216220
let site: SiteContext
@@ -223,17 +227,20 @@ export function createCloudflareRuntime<Env extends RuntimeEnv>(resolveCoordinat
223227
if (new URL(request.url).pathname === "/wp-cron.php") return new Response("WordPress cron is managed by the Cloudflare scheduled handler.", { status: 404 })
224228
const publishedResponse = await servePublishedWordPressPage(request, env.WORDPRESS_STATE_BUCKET, site)
225229
if (publishedResponse) return publishedResponse
226-
const coordinator = resolveCoordinator(env, site)
230+
const route = routeWorkerRequest(request)
231+
const selector = new URL(request.url).searchParams.get("coordinator")
232+
const selection = selectOperatorCoordinator(resolveCoordinator(env, site), route, selector, (selected) => resolveOperatorCoordinator?.(env, site, selected))
233+
if (!selection) return new Response("Unsupported operator coordinator selector.", { status: 400 })
234+
const coordinator = selection.coordinator
227235
const wpContentResponse = await serveWordPressWpContent(request, env.WORDPRESS_STATE_BUCKET, coordinator, site)
228236
if (wpContentResponse) return wpContentResponse
229237
const staticResponse = await serveWordPressStaticAsset(request, env.WORDPRESS_STATE_BUCKET)
230238
if (staticResponse) return staticResponse
231-
const route = routeWorkerRequest(request)
232239
const uploadResponse = await serveWordPressUpload(request, env.WORDPRESS_STATE_BUCKET, coordinator, site)
233240
if (uploadResponse) return uploadResponse
234241
if (route.kind === "operator-reset") return resetCanonicalWordPress(request, env, coordinator, site)
235242
if (route.kind === "operator-restore") return restoreCanonicalWordPress(request, env, coordinator, site)
236-
if (route.kind === "operator-adopt") return adoptCanonicalWordPress(request, env, coordinator, site)
243+
if (route.kind === "operator-adopt") return adoptCanonicalWordPress(request, env, coordinator, site, selection.selected)
237244
if (route.kind === "operator-fence") return operateCanonicalMutationFence(request, env, coordinator, site, route.action)
238245
if (route.kind === "operator-static-artifact-import") return importCanonicalStaticArtifact(request, env, coordinator, site)
239246
if (route.kind === "operator-publish") return publishCanonicalWordPressPages(request, env, coordinator, site)
@@ -418,28 +425,31 @@ async function restoreCanonicalWordPress(request: Request, env: RuntimeEnv, coor
418425
}
419426
}
420427

421-
async function adoptCanonicalWordPress(request: Request, env: RuntimeEnv, coordinator: RevisionCoordinator, site: SiteContext): Promise<Response> {
428+
async function adoptCanonicalWordPress(request: Request, env: RuntimeEnv, coordinator: RevisionCoordinator, site: SiteContext, requireFence = false): Promise<Response> {
422429
if (request.method !== "POST") return new Response("Canonical adoption requires POST.", { status: 405 })
423430
if (!await isAuthorizedOperator(request, env, site)) {
424431
return new Response("Canonical adoption authorization failed.", { status: 401 })
425432
}
426433
let pointer: MarkdownPointer
427434
let version: number
435+
let fenceToken: string | undefined
428436
try {
429-
const body = await request.json<{ pointer?: unknown; version?: unknown }>()
437+
const body = await request.json<{ pointer?: unknown; version?: unknown; fenceToken?: unknown }>()
430438
pointer = body.pointer as MarkdownPointer
431439
version = body.version as number
440+
fenceToken = typeof body.fenceToken === "string" ? body.fenceToken : undefined
432441
} catch {
433442
return new Response("Canonical adoption requires JSON state.", { status: 400 })
434443
}
444+
if (requireFence && !fenceToken) return new Response("Selected coordinator adoption requires its active fence token.", { status: 400 })
435445
if (!isCanonicalRestorePointer(pointer, site) || !Number.isSafeInteger(version) || version < 1) return new Response("Canonical adoption state is invalid.", { status: 400 })
436446
const manifest = await readMarkdownManifest(env.WORDPRESS_STATE_BUCKET, pointer, site)
437447
if (!manifest || manifest.revision !== pointer.revision || manifest.manifestKey !== pointer.manifestKey || manifest.persistedAt !== pointer.persistedAt || !Array.isArray(manifest.files)) {
438448
return new Response("Canonical adoption manifest is unavailable or inconsistent.", { status: 409 })
439449
}
440450
let adopted: { pointer: MarkdownPointer; version: number }
441451
try {
442-
adopted = await coordinator.adopt(pointer, version)
452+
adopted = await coordinator.adopt(pointer, version, fenceToken, requireFence)
443453
} catch (error) {
444454
if (!(error instanceof RevisionConflict)) throw error
445455
return coordinatorConflictResponse(error)

packages/runtime-cloudflare/wrangler.d1.jsonc

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,14 @@
1616
"database_id": "00000000-0000-0000-0000-000000000000"
1717
}
1818
],
19+
"durable_objects": {
20+
"bindings": [
21+
{ "name": "WORDPRESS_STATE", "class_name": "WordPressStateCoordinator" }
22+
]
23+
},
24+
"migrations": [
25+
{ "tag": "v1", "new_sqlite_classes": ["WordPressStateCoordinator"] }
26+
],
1927
"rules": [
2028
{ "type": "CompiledWasm", "globs": ["**/*.wasm"], "fallthrough": false },
2129
{ "type": "Data", "globs": ["**/*.sqlite", "**/*-runtime.zip", "**/*-canonical-seed.zip"], "fallthrough": false }

0 commit comments

Comments
 (0)