Skip to content

Commit c3d5814

Browse files
authored
Merge pull request #2115 from Automattic/feat/2083-durable-queue-dispatch
2 parents 4a3ae64 + 56ce6ee commit c3d5814

12 files changed

Lines changed: 431 additions & 61 deletions

package.json

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -111,7 +111,8 @@
111111
"test:redaction": "tsx tests/redaction.test.ts",
112112
"test:browser-preview-routing": "tsx --test tests/browser-preview-routing.test.ts",
113113
"test:browser-routed-command-security": "tsx --test tests/browser-routed-command-security.test.ts",
114-
"test:cloudflare-runtime": "node --test tests/cloudflare-d1-provisioner.test.mjs && tsx tests/cloudflare-site-context.test.ts && tsx tests/cloudflare-coordinator-site-partitioning.test.ts && tsx tests/cloudflare-d1-operation-repository.test.ts && tsx tests/cloudflare-allocation-lifecycle.test.ts && tsx tests/cloudflare-provisioning-api.test.ts && tsx tests/cloudflare-runtime.test.ts && tsx tests/cloudflare-public-reader.test.ts && node ./node_modules/typescript/bin/tsc -p packages/runtime-cloudflare --noEmit",
114+
"test:cloudflare-runtime": "node --test tests/cloudflare-d1-provisioner.test.mjs && tsx tests/cloudflare-site-context.test.ts && tsx tests/cloudflare-coordinator-site-partitioning.test.ts && tsx tests/cloudflare-d1-operation-repository.test.ts && tsx tests/cloudflare-allocation-lifecycle.test.ts && tsx tests/cloudflare-provisioning-api.test.ts && tsx tests/cloudflare-queue-batch.test.ts && tsx tests/cloudflare-runtime.test.ts && tsx tests/cloudflare-public-reader.test.ts && node ./node_modules/typescript/bin/tsc -p packages/runtime-cloudflare --noEmit",
115+
"test:cloudflare-queue": "tsx tests/cloudflare-queue-batch.test.ts && tsx tests/cloudflare-d1-operation-repository.test.ts && node ./node_modules/typescript/bin/tsc -p packages/runtime-cloudflare --noEmit",
115116
"test:cloudflare-administrator-claim": "tsx tests/cloudflare-provisioning-api.test.ts",
116117
"test:cloudflare-wordpress-auth": "tsx tests/cloudflare-wordpress-auth.test.ts",
117118
"test:cloudflare-wordpress-archive-corpus": "tsx tests/cloudflare-wordpress-archive-corpus.test.ts",

packages/runtime-cloudflare/src/allocation-lifecycle.ts

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -102,6 +102,7 @@ export class CloudflareAllocationLifecycle {
102102
])
103103
}
104104
if (await tableExists(this.database, "wp_codebox_api_admin_claims")) await this.database.prepare("UPDATE wp_codebox_api_admin_claims SET state = 'expired', updated_at = ? WHERE site_id = ? AND state = 'pending'").bind(now, identity.siteId).run()
105+
if (await tableExists(this.database, "wp_codebox_runtime_dispatches")) await this.database.prepare("UPDATE wp_codebox_runtime_dispatches SET state = 'done', last_error = 'Allocation deletion fenced this dispatch.', completed_at = ?, updated_at = ? WHERE site_id = ? AND state NOT IN ('done','dead-letter')").bind(now, now, identity.siteId).run()
105106
}
106107
private async cleanupWork(siteId: string, limit: number): Promise<{ deleted: number; unresolved: number }> {
107108
let deleted = 0
@@ -114,7 +115,10 @@ export class CloudflareAllocationLifecycle {
114115
if (await tableExists(this.database, "wp_codebox_api_admin_claims")) {
115116
const result = await this.database.prepare("DELETE FROM wp_codebox_api_admin_claims WHERE site_id = ?").bind(siteId).run(); deleted += result.meta.changes
116117
}
117-
const tables = ["wp_codebox_operation_attempts", "wp_codebox_operations", "wp_codebox_api_admin_claims"]
118+
for (const table of ["wp_codebox_runtime_dead_letters", "wp_codebox_runtime_dispatches", "wp_codebox_runtime_fairness"]) if (await tableExists(this.database, table)) {
119+
const result = await this.database.prepare(`DELETE FROM ${table} WHERE site_id = ?`).bind(siteId).run(); deleted += result.meta.changes
120+
}
121+
const tables = ["wp_codebox_operation_attempts", "wp_codebox_operations", "wp_codebox_api_admin_claims", "wp_codebox_runtime_dead_letters", "wp_codebox_runtime_dispatches", "wp_codebox_runtime_fairness"]
118122
let unresolved = 0
119123
for (const table of tables) if (await tableExists(this.database, table)) unresolved += (await this.database.prepare(`SELECT COUNT(*) AS count FROM ${table} WHERE site_id = ?`).bind(siteId).first<{ count: number }>())?.count ?? 0
120124
return { deleted, unresolved }

packages/runtime-cloudflare/src/d1-operation-repository.ts

Lines changed: 105 additions & 0 deletions
Large diffs are not rendered by default.

packages/runtime-cloudflare/src/provisioning-api.ts

Lines changed: 19 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,5 @@
11
import { D1OperationRepository, OperationConflict, type StaticArtifactOperation, type StaticArtifactOperationInput } from "./d1-operation-repository.js"
2+
import { parseRuntimeQueuePolicy, runtimeQueueMessage, type RuntimeQueue } from "./queue-dispatch.js"
23
import { MAX_STATIC_ARTIFACT_BYTES, readBoundedRequestBytes, readStaticArtifactImport, StaticArtifactImportError, validateStaticArtifact } from "./static-artifact-import.js"
34
import { allocatePreviewSiteContext, parseSiteContexts, previewDomain, siteStorageKeys, type SiteContext } from "./site-context.js"
45
import { deriveSiteCredential } from "./wordpress-auth.js"
@@ -17,7 +18,7 @@ export interface ProvisioningAllocation {
1718
artifactSha256: string; artifactSize: number; options: StaticArtifactOperationInput["options"]
1819
}
1920
interface CreateInput { key: string; fingerprint: string; artifactSha256: string; artifactSize: number; options: StaticArtifactOperationInput["options"] }
20-
export interface ProvisioningEnv { WORDPRESS_STATE_DATABASE: D1Database; WORDPRESS_STATE_BUCKET: R2Bucket; WORDPRESS_SITE_CONTEXTS?: string; WORDPRESS_PREVIEW_DOMAIN?: string; WORDPRESS_PREVIEW_HOST_SECRET?: string; WORDPRESS_API_TOKENS?: string; WORDPRESS_ADMIN_CLAIM_SECRET?: string; WORDPRESS_ADMIN_PASSWORD?: string }
21+
export interface ProvisioningEnv { WORDPRESS_STATE_DATABASE: D1Database; WORDPRESS_STATE_BUCKET: R2Bucket; WORDPRESS_RUNTIME_QUEUE?: RuntimeQueue; WORDPRESS_QUEUE_POLICY?: string; WORDPRESS_SITE_CONTEXTS?: string; WORDPRESS_PREVIEW_DOMAIN?: string; WORDPRESS_PREVIEW_HOST_SECRET?: string; WORDPRESS_API_TOKENS?: string; WORDPRESS_ADMIN_CLAIM_SECRET?: string; WORDPRESS_ADMIN_PASSWORD?: string }
2122

2223
export async function routeProvisioningApi(request: Request, env: ProvisioningEnv, operations: D1OperationRepository): Promise<Response> {
2324
const parts = new URL(request.url).pathname.split("/").filter(Boolean)
@@ -38,6 +39,12 @@ export async function routeProvisioningApi(request: Request, env: ProvisioningEn
3839
return notFound()
3940
}
4041

42+
async function dispatchOperation(env: ProvisioningEnv, operations: D1OperationRepository, site: SiteContext, operation: StaticArtifactOperation, principal: string): Promise<"queued" | "backpressured"> {
43+
const message = await operations.stageDispatch(site, "operation", operation.operationId, principal)
44+
if (!await operations.admitDispatch(message, parseRuntimeQueuePolicy(env.WORDPRESS_QUEUE_POLICY))) return "backpressured"
45+
try { if (!env.WORDPRESS_RUNTIME_QUEUE) throw new Error("Runtime queue binding is unavailable."); await env.WORDPRESS_RUNTIME_QUEUE.send(message); await operations.deliveredDispatch(message); return "queued" } catch (error) { await operations.failedDispatch(message, error); console.error("Runtime queue producer is backpressured; reconciliation will retry.", error); return "backpressured" }
46+
}
47+
4148
async function lifecycleStatus(request: Request, env: ProvisioningEnv, siteId: string): Promise<Response> {
4249
const token = await authenticate(request, env, "sites:read"); if (token instanceof Response) return token
4350
const allocation = await new AllocationStore(env.WORDPRESS_STATE_DATABASE).bySite(siteId)
@@ -110,7 +117,9 @@ async function create(request: Request, env: ProvisioningEnv, operations: D1Oper
110117
try {
111118
const operation = await resumeProvisioningAllocation(env, site, operations)
112119
const claim = await new AdministratorClaimStore(env.WORDPRESS_STATE_DATABASE).issue(allocation, env)
113-
return siteResource(site, operation, 202, claim)
120+
const response = siteResource(site, operation, 202, claim)
121+
if (operation) response.headers.set("x-wp-codebox-dispatch", dispatchHeader(await operations.dispatchState(runtimeQueueMessage(site, allocationIdentity(site.id).generation, "operation", operation.operationId))))
122+
return response
114123
} catch (error) { if (error instanceof OperationConflict) return apiError(409, "idempotency_conflict", error.message); if (error instanceof AdministratorClaimError) return apiError(409, "administrator_claim_unavailable", "The administrator claim cannot continue provisioning."); throw error }
115124
}
116125

@@ -153,8 +162,11 @@ async function importSite(request: Request, env: ProvisioningEnv, operations: D1
153162
const input = await readStaticArtifactImport(request, env.WORDPRESS_STATE_BUCKET, site)
154163
if (input.idempotencyKey !== key) return apiError(409, "idempotency_conflict", "Idempotency-Key must match the import request.")
155164
const result = await operations.createOrConverge(site, { ...input, artifact: input.artifactReference })
165+
const dispatch = await dispatchOperation(env, operations, site, result.operation, token.principal)
156166
await store.linkOperation(token.principal, siteId, result.operation.operationId, "import", key)
157-
return operationResource(siteId, result.operation, 202)
167+
const response = operationResource(siteId, await operations.get(siteId, result.operation.operationId) ?? result.operation, 202)
168+
response.headers.set("x-wp-codebox-dispatch", dispatch)
169+
return response
158170
} catch (error) { return error instanceof OperationConflict ? apiError(409, "operation_conflict", error.message) : importError(error) }
159171
}
160172

@@ -187,9 +199,12 @@ export async function resumeProvisioningAllocation(env: ProvisioningEnv, site: S
187199
const result = await operations.createOrConverge(site, input)
188200
await store.bindOperation(allocation, result.operation.operationId)
189201
await store.linkOperation(allocation.principal, site.id, result.operation.operationId, "provision", allocation.key)
190-
return result.operation
202+
await dispatchOperation(env, operations, site, result.operation, allocation.principal)
203+
return await operations.get(site.id, result.operation.operationId) ?? result.operation
191204
}
192205

206+
function dispatchHeader(state: string | null): "queued" | "backpressured" { return state && !["backpressured", "dead-letter"].includes(state) ? "queued" : "backpressured" }
207+
193208
async function readCreate(request: Request, bucket: R2Bucket, key: string): Promise<CreateInput> {
194209
const bytes = await readBoundedRequestBytes(request)
195210
let body: Record<string, unknown>; try { body = JSON.parse(new TextDecoder("utf-8", { fatal: true }).decode(bytes)) } catch { throw new StaticArtifactImportError("Provisioning request must be valid UTF-8 JSON.", 400) }
Lines changed: 33 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,33 @@
1+
import type { RuntimeQueueMessage } from "./queue-dispatch.js"
2+
3+
export interface QueueDelivery { body: unknown; attempts: number; ack(): void; retry(options?: { delaySeconds?: number }): void }
4+
export interface ParsedQueueDelivery { raw: QueueDelivery; value: RuntimeQueueMessage }
5+
export type QueueExecutionResult = "ack" | "retry" | { retryAfterSeconds: number }
6+
7+
/** One heavyweight turn per site keeps a delivery batch fair without cross-site serialization. */
8+
export async function dispatchQueueBatch(
9+
deliveries: QueueDelivery[],
10+
parse: (body: unknown) => RuntimeQueueMessage | null,
11+
select: (siteId: string, deliveries: ParsedQueueDelivery[]) => Promise<ParsedQueueDelivery | null>,
12+
execute: (delivery: ParsedQueueDelivery) => Promise<QueueExecutionResult>,
13+
): Promise<void> {
14+
const lanes = new Map<string, ParsedQueueDelivery[]>()
15+
for (const raw of deliveries) {
16+
const value = parse(raw.body)
17+
if (!value) { raw.ack(); continue }
18+
const lane = lanes.get(value.siteId) ?? []
19+
lane.push({ raw, value })
20+
lanes.set(value.siteId, lane)
21+
}
22+
await Promise.all([...lanes.entries()].map(async ([siteId, lane]) => {
23+
const selected = await select(siteId, lane)
24+
for (const delivery of lane) {
25+
if (delivery !== selected) delivery.raw.retry()
26+
}
27+
if (!selected) return
28+
const result = await execute(selected)
29+
if (result === "ack") selected.raw.ack()
30+
else if (result === "retry") selected.raw.retry()
31+
else selected.raw.retry({ delaySeconds: result.retryAfterSeconds })
32+
}))
33+
}
Lines changed: 42 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,42 @@
1+
import type { SiteContext } from "./site-context.js"
2+
3+
export const RUNTIME_QUEUE_MESSAGE_SCHEMA = "wp-codebox/runtime-dispatch/v1"
4+
export const RUNTIME_QUEUE_MAX_ATTEMPTS = 3
5+
export const DEFAULT_RUNTIME_QUEUE_POLICY: RuntimeQueuePolicy = { maxActive: 100, maxActivePerPrincipal: 10 }
6+
7+
/** Queue messages wake durable work; they never own canonical state. */
8+
export type RuntimeQueueMessage = { schema: typeof RUNTIME_QUEUE_MESSAGE_SCHEMA; siteId: string; generation: number; kind: "operation" | "publication"; identity: string }
9+
export interface RuntimeQueuePolicy { maxActive: number; maxActivePerPrincipal: number }
10+
11+
const operationIdentity = /^[0-9a-f]{8}-(?:[0-9a-f]{4}-){3}[0-9a-f]{12}$/
12+
const publicationIdentity = /^sites\/([a-z0-9][a-z0-9-]{0,127})\/publications\/jobs\/[0-9]{20}-[a-f0-9-]{36}\.json$/
13+
14+
export function validRuntimeQueueIdentity(siteId: string, kind: RuntimeQueueMessage["kind"], identity: string): boolean {
15+
if (kind === "operation") return operationIdentity.test(identity)
16+
return publicationIdentity.test(identity) && publicationIdentity.exec(identity)?.[1] === siteId
17+
}
18+
19+
export function runtimeQueueMessage(site: Pick<SiteContext, "id">, generation: number, kind: RuntimeQueueMessage["kind"], identity: string): RuntimeQueueMessage {
20+
if (!/^[a-z0-9][a-z0-9-]{0,127}$/.test(site.id) || !Number.isSafeInteger(generation) || generation < 1 || !validRuntimeQueueIdentity(site.id, kind, identity)) throw new Error("Runtime queue dispatch identity is invalid.")
21+
return { schema: RUNTIME_QUEUE_MESSAGE_SCHEMA, siteId: site.id, generation, kind, identity }
22+
}
23+
24+
export function parseRuntimeQueueMessage(value: unknown): RuntimeQueueMessage | null {
25+
if (!value || typeof value !== "object") return null
26+
const message = value as Partial<RuntimeQueueMessage>
27+
if (message.schema !== RUNTIME_QUEUE_MESSAGE_SCHEMA || (message.kind !== "operation" && message.kind !== "publication") || typeof message.siteId !== "string" || !/^[a-z0-9][a-z0-9-]{0,127}$/.test(message.siteId) || !Number.isSafeInteger(message.generation) || message.generation! < 1 || typeof message.identity !== "string" || !validRuntimeQueueIdentity(message.siteId, message.kind, message.identity)) return null
28+
return message as RuntimeQueueMessage
29+
}
30+
31+
export interface RuntimeQueue { send(message: RuntimeQueueMessage): Promise<void> }
32+
33+
export function parseRuntimeQueuePolicy(value: string | undefined): RuntimeQueuePolicy {
34+
if (!value) return DEFAULT_RUNTIME_QUEUE_POLICY
35+
let policy: Partial<RuntimeQueuePolicy>
36+
try { policy = JSON.parse(value) as Partial<RuntimeQueuePolicy> } catch { throw new Error("Runtime queue policy is invalid.") }
37+
if (!Number.isSafeInteger(policy.maxActive) || !Number.isSafeInteger(policy.maxActivePerPrincipal)
38+
|| policy.maxActive! < 1 || policy.maxActive! > 10_000 || policy.maxActivePerPrincipal! < 1 || policy.maxActivePerPrincipal! > policy.maxActive!) {
39+
throw new Error("Runtime queue policy is invalid.")
40+
}
41+
return policy as RuntimeQueuePolicy
42+
}

0 commit comments

Comments
 (0)