diff --git a/services/deploy-infra/builder/scripts/deploy-local.ts b/services/deploy-infra/builder/scripts/deploy-local.ts index d869083d1f..e544d4651a 100644 --- a/services/deploy-infra/builder/scripts/deploy-local.ts +++ b/services/deploy-infra/builder/scripts/deploy-local.ts @@ -42,7 +42,6 @@ function loadEnvFile(): Record { const key = trimmed.slice(0, eqIndex); let value = trimmed.slice(eqIndex + 1); - // Remove surrounding quotes if present if ( (value.startsWith('"') && value.endsWith('"')) || (value.startsWith("'") && value.endsWith("'")) @@ -145,13 +144,11 @@ async function main() { const slug = parsed.slug || path.basename(directory); const envVars = parsed.envVars; - // Validate directory exists if (!fs.existsSync(directory)) { console.error(`Error: Directory not found: ${directory}`); process.exit(1); } - // Validate directory is actually a directory if (!fs.statSync(directory).isDirectory()) { console.error(`Error: Not a directory: ${directory}`); process.exit(1); @@ -169,12 +166,10 @@ async function main() { console.log('='.repeat(60)); console.log(''); - // Create tar.gz archive in temp location const archivePath = `/tmp/${slug}-${Date.now()}.tar.gz`; console.log('📦 Creating archive...'); - // Build exclusion list for tar command const excludePatterns = [ '.next', 'node_modules', // Dependencies @@ -212,7 +207,6 @@ async function main() { console.log(''); try { - // Upload to builder console.log('🚀 Uploading to builder...'); const archiveData = fs.readFileSync(archivePath); @@ -246,13 +240,11 @@ async function main() { console.log(` Build ID: ${result.buildId}`); console.log(''); - // Poll for status console.log('⏳ Building...'); let lastStatus = ''; let lastEventCount = 0; while (true) { - // Fetch and display new events const eventsResponse = await fetch(`${BUILDER_URL}/deploy/${result.buildId}/events`, { headers: { Authorization: `Bearer ${BUILDER_AUTH_TOKEN}` }, }); @@ -262,7 +254,6 @@ async function main() { payload: { status?: string }; }[]; - // Show new log events for (let i = lastEventCount; i < events.length; i++) { const event = events[i]; console.log(` ${JSON.stringify(event.payload)}`); @@ -273,7 +264,6 @@ async function main() { } lastEventCount = events.length; - // Check if finished if (['deployed', 'failed', 'cancelled'].includes(lastStatus)) { break; } @@ -297,7 +287,6 @@ async function main() { process.exit(1); } } finally { - // Cleanup archive try { fs.unlinkSync(archivePath); } catch { diff --git a/services/deploy-infra/builder/src/cloudflare-api.ts b/services/deploy-infra/builder/src/cloudflare-api.ts index fb74fe4057..c1a2a18dbc 100644 --- a/services/deploy-infra/builder/src/cloudflare-api.ts +++ b/services/deploy-infra/builder/src/cloudflare-api.ts @@ -2,11 +2,6 @@ import * as z from 'zod'; import type { DeploymentFile, WorkerMetadata } from './types'; import type { PlaintextEnvVar } from '../../../../apps/web/src/lib/user-deployments/env-vars-validation'; -/** - * Cloudflare API Client - Handles interactions with Cloudflare Workers API. - * Provides methods for asset uploads and worker deployments. - */ - /** * Error thrown when a worker is not found in Cloudflare (error code 10007). */ @@ -17,9 +12,6 @@ export class WorkerNotFoundError extends Error { } } -/** - * Standard Cloudflare API response structure. - */ const cloudflareErrorSchema = z.object({ code: z.number(), message: z.string(), @@ -53,31 +45,15 @@ export const parseCloudflareResponse = ( }; }; -/** - * CloudflareAPI class for interacting with Cloudflare Workers API. - * Handles asset upload sessions, batch uploads, and worker deployments. - */ export class CloudflareAPI { private accountId: string; private apiToken: string; - /** - * Create a new CloudflareAPI instance. - * - * @param accountId - Cloudflare account ID - * @param apiToken - Cloudflare API token with Workers deployment permissions - */ constructor(accountId: string, apiToken: string) { this.accountId = accountId; this.apiToken = apiToken; } - /** - * Get headers for API requests. - * - * @param contentType - Optional Content-Type header value - * @returns Headers object with authorization and optional content type - */ private getHeaders(contentType?: string): HeadersInit { const headers: HeadersInit = { Authorization: `Bearer ${this.apiToken}`, @@ -90,9 +66,6 @@ export class CloudflareAPI { return headers; } - /** - * Retry a function with exponential backoff for transient errors - */ private async retryWithBackoff( fn: () => Promise, operation: string, @@ -106,14 +79,12 @@ export class CloudflareAPI { } catch (error) { lastError = error instanceof Error ? error : new Error(String(error)); - // Check if it's a 5xx error (transient) const is5xxError = lastError.message.includes('status: 5'); if (!is5xxError || attempt === maxAttempts) { throw lastError; } - // Exponential backoff: 1s, 2s, 4s, etc., max 30s const delayMs = Math.min(1000 * Math.pow(2, attempt - 1), 30000); console.log( `${operation} failed (attempt ${attempt}/${maxAttempts}), retrying in ${delayMs}ms...` @@ -126,14 +97,6 @@ export class CloudflareAPI { throw lastError ?? new Error('Unknown error in retry logic'); } - /** - * Create an asset upload session for a worker script. - * - * @param scriptName - Name of the worker script - * @param manifest - Manifest mapping file paths to their hash and size - * @param dispatchNamespace - Dispatch namespace for the worker - * @returns Upload session JWT and bucket assignments - */ async createAssetUploadSession( scriptName: string, manifest: Record, @@ -176,14 +139,6 @@ export class CloudflareAPI { }, 'Create asset upload session'); } - /** - * Upload a batch of assets to Cloudflare. - * - * @param uploadToken - JWT token from the upload session - * @param fileHashes - Array of file hashes to upload in this batch - * @param fileContents - Map of hash to file content (Buffer) and MIME type - * @returns Completion JWT if all files uploaded (status 201), null otherwise - */ async uploadAssetBatch( scriptName: string, uploadToken: string, @@ -193,7 +148,6 @@ export class CloudflareAPI { return await this.retryWithBackoff(async () => { const url = `https://api.cloudflare.com/client/v4/accounts/${this.accountId}/workers/assets/upload?base64=true`; - // Build FormData with each file const formData = new FormData(); for (const hash of fileHashes) { @@ -202,10 +156,8 @@ export class CloudflareAPI { throw new Error(`Missing file content for hash ${hash}`); } - // Convert Buffer to base64 string (matches VibeSDK pattern) const base64Content = fileData.buffer.toString('base64'); - // Create blob with base64 content as text const blob = new Blob([base64Content], { type: fileData.mimeType }); // Append to form data with hash as both field name and filename @@ -266,26 +218,17 @@ export class CloudflareAPI { }, 'Upload asset batch'); } - /** - * Deploy a worker script to Cloudflare. - */ async deployWorker(params: { - /** Name of the worker script */ scriptName: string; - /** Worker metadata including main module, compatibility settings, and assets */ metadata: WorkerMetadata; - /** Worker script file */ workerScript: DeploymentFile; - /** Dispatch namespace for the worker */ dispatchNamespace: string; - /** Optional array of artifact files to include in deployment */ artifacts?: DeploymentFile[]; /** Environment variables (only non-secret) */ envVars?: PlaintextEnvVar[]; }): Promise { const { scriptName, metadata, workerScript, dispatchNamespace, artifacts, envVars } = params; - // Assert that envVars don't contain any secret variables if (envVars?.some(v => v.isSecret)) { throw new Error( 'Secret environment variables must be set via setSecrets(), not deployWorker()' @@ -295,7 +238,6 @@ export class CloudflareAPI { return await this.retryWithBackoff(async () => { const url = `https://api.cloudflare.com/client/v4/accounts/${this.accountId}/workers/dispatch/namespaces/${dispatchNamespace}/scripts/${scriptName}`; - // Add plain text environment variables to metadata bindings const plainTextBindings = (envVars || []).map(v => ({ type: 'plain_text', name: v.key, @@ -315,7 +257,6 @@ export class CloudflareAPI { }); formData.append('index.js', workerBlob, 'index.js'); - // Append artifact files if provided if (artifacts && artifacts.length > 0) { for (const artifact of artifacts) { const artifactBlob = new Blob([new Uint8Array(artifact.content)], { @@ -357,7 +298,6 @@ export class CloudflareAPI { const existingClass = messageStr.match(/class "([^"]+)"/)?.[1]; if (existingClass && metadata.migrations) { - // Filter out the existing class from migrations const filteredMigrations = metadata.migrations .map(migration => ({ ...migration, @@ -365,7 +305,6 @@ export class CloudflareAPI { })) .filter(migration => !migration.new_classes || migration.new_classes.length > 0); - // Retry with filtered migrations (preserve envVars for retry) return await this.deployWorker({ scriptName, metadata: { ...metadata, migrations: filteredMigrations }, @@ -386,10 +325,6 @@ export class CloudflareAPI { }, 'Deploy worker'); } - /** - * Sets secrets for a worker in a dispatch namespace using Cloudflare's Secrets API. - * Secrets are set individually via PUT requests, processed in parallel batches. - */ async setSecrets( scriptName: string, dispatchNamespace: string, @@ -437,11 +372,7 @@ export class CloudflareAPI { } /** - * Delete a worker script from a dispatch namespace. * Note: Assets are automatically cleaned up when the script is deleted. - * - * @param scriptName - Name of the worker script to delete - * @param dispatchNamespace - Dispatch namespace containing the worker */ async deleteWorker(scriptName: string, dispatchNamespace: string): Promise { return await this.retryWithBackoff(async () => { diff --git a/services/deploy-infra/builder/src/deployer.ts b/services/deploy-infra/builder/src/deployer.ts index d98912a586..8d1b921a57 100644 --- a/services/deploy-infra/builder/src/deployer.ts +++ b/services/deploy-infra/builder/src/deployer.ts @@ -1,8 +1,3 @@ -/** - * Deployer - Orchestrates the deployment of Cloudflare Workers with assets. - * Handles uploading assets and deploying worker scripts. - */ - import { WorkerNotFoundError, type CloudflareAPI } from './cloudflare-api'; import { DEPLOY_DISPATCH_NAMESPACE } from './dispatch-namespace'; import type { DeploymentArtifacts } from './types'; @@ -40,19 +35,11 @@ export class Deployer { }); } - /** - * Deploy a worker with its assets to Cloudflare - */ async deploy(params: { - /** Worker script and asset files */ artifacts: DeploymentArtifacts; - /** Name of the worker */ workerName: string; - /** Callback for logging progress */ logger: (message: string) => void; - /** Dispatch namespace to deploy to */ dispatchNamespace?: string; - /** Optional environment variables (decrypted) */ envVars?: PlaintextEnvVar[]; }): Promise { const { @@ -71,7 +58,6 @@ export class Deployer { const secretEnvVars = envVars?.filter(v => v.isSecret) ?? []; const plainTextEnvVars = envVars?.filter(v => !v.isSecret) ?? []; - // Set secrets if (secretEnvVars.length > 0) { try { await this.api.setSecrets(workerName, dispatchNamespace, secretEnvVars); @@ -87,7 +73,6 @@ export class Deployer { } } - // Handle case with no assets and no artifacts if (assets.length === 0 && artifactFiles.length === 0) { logger('No assets or artifacts found, deploying worker only'); const metadata = { @@ -112,7 +97,6 @@ export class Deployer { logger(`Found ${assets.length} asset files and ${artifactFiles.length} artifact files`); - // Calculate hashes and build data structures const fileContents = new Map(); const manifest: Record = {}; @@ -140,7 +124,6 @@ export class Deployer { logger(`Built asset manifest with ${assets.length} files, ${totalBytes} total bytes`); - // Create asset upload session const { jwt, buckets } = await this.api.createAssetUploadSession( workerName, manifest, @@ -149,7 +132,6 @@ export class Deployer { logger(`Created asset upload session, ${buckets.length} buckets to upload`); - // Upload assets in batches let completionJwt: string | null = null; if (buckets.length === 0) { @@ -157,7 +139,6 @@ export class Deployer { logger('All assets already exist on Cloudflare'); completionJwt = jwt; } else { - // Upload new/changed assets for (let i = 0; i < buckets.length; i++) { const bucket = buckets[i]; logger(`Uploading batch ${i + 1}/${buckets.length} (${bucket.length} files)`); @@ -176,7 +157,6 @@ export class Deployer { logger('All assets uploaded successfully'); - // Build worker metadata with asset configuration const bindings: Array> = [ { name: 'ASSETS', @@ -195,7 +175,6 @@ export class Deployer { }, }; - // Deploy worker with artifact files await this.api.deployWorker({ scriptName: workerName, metadata, diff --git a/services/deploy-infra/builder/src/deployment-orchestrator.ts b/services/deploy-infra/builder/src/deployment-orchestrator.ts index c14ff9747e..8243c522bb 100644 --- a/services/deploy-infra/builder/src/deployment-orchestrator.ts +++ b/services/deploy-infra/builder/src/deployment-orchestrator.ts @@ -1,8 +1,3 @@ -/** - * DeploymentOrchestrator - Durable Object for managing deployment job lifecycle. - * Handles cloning, building, deploying, and tracking job state and events. - */ - import { DurableObject } from 'cloudflare:workers'; import { type ExecEvent, getSandbox, parseSSEStream } from '@cloudflare/sandbox'; import { stripVTControlCharacters } from 'node:util'; @@ -37,22 +32,13 @@ import { } from './errors'; import { sanitizeGitError } from './sanitize-git-error'; -/** - * DeploymentOrchestrator manages the complete lifecycle of a deployment job. - * Persists job state and a bounded ring buffer of events in storage. - */ export class DeploymentOrchestrator extends DurableObject { - /** In-memory cache of current build state */ private state!: Build; constructor(ctx: DurableObjectState, env: Env) { super(ctx, env); } - /** - * Alarm handler for scheduled tasks. - * Handles job execution when a job is queued and ready to run. - */ async alarm(): Promise { await this.loadState(); @@ -61,9 +47,6 @@ export class DeploymentOrchestrator extends DurableObject { } } - /** - * Load state and events from durable storage. - */ private async loadState(): Promise { const storedState = await this.ctx.storage.get('state'); @@ -72,9 +55,6 @@ export class DeploymentOrchestrator extends DurableObject { } } - /** - * Save current state to durable storage. - */ private async saveState(): Promise { await this.ctx.storage.put('state', this.state); } @@ -84,18 +64,11 @@ export class DeploymentOrchestrator extends DurableObject { return this.env.EventsManager.get(eventsManagerId); } - /** - * Add a log event via EventsManager DO. - * Updates local state timestamp and delegates event storage to EventsManager. - */ private async addLogEvent(message: string): Promise { const eventsManager = this.eventsManager(); await eventsManager.addEvent({ type: 'log', payload: { message } }); } - /** - * Add a status change event via EventsManager DO. - */ private async addStatusChangeEvent(status: BuildStatus): Promise { const eventsManager = this.eventsManager(); @@ -105,9 +78,6 @@ export class DeploymentOrchestrator extends DurableObject { }); } - /** - * Log stdout and stderr from an ExecResult, splitting by line and skipping empty lines. - */ private async logExecResult(result: { stdout: string; stderr: string }): Promise { for (const output of [result.stderr, result.stdout]) { for (const line of output.split('\n')) { @@ -118,9 +88,6 @@ export class DeploymentOrchestrator extends DurableObject { } } - /** - * Update build status - */ private async updateStatus(status: BuildStatus): Promise { if (this.state.status === status) { return; @@ -128,7 +95,6 @@ export class DeploymentOrchestrator extends DurableObject { this.state.status = status; - // Update timestamps based on status if (status === 'building' && !this.state.startedAt) { this.state.startedAt = new Date().toISOString(); } @@ -140,13 +106,9 @@ export class DeploymentOrchestrator extends DurableObject { this.state.updatedAt = new Date().toISOString(); await this.saveState(); - // Emit status change event await this.addStatusChangeEvent(status); } - /** - * RPC method: Start the job. - */ async start(params: { buildId: string; slug: string; @@ -163,22 +125,17 @@ export class DeploymentOrchestrator extends DurableObject { }; await this.saveState(); - // Setup events manager const eventsManager = this.eventsManager(); await eventsManager.initialize(this.state.buildId); await this.addLogEvent('Build created and queued'); - // Schedule alarm to run job asynchronously // Alarms are designed for long-running work that survives context timeouts - await this.ctx.storage.setAlarm(Date.now() + 50); // Run almost immediately + await this.ctx.storage.setAlarm(Date.now() + 50); return { status: this.state.status }; } - /** - * RPC method: Start job from uploaded archive. - */ async startFromArchive(params: ArchiveDeployParams): Promise<{ status: BuildStatus }> { this.state = { buildId: params.buildId, @@ -193,21 +150,16 @@ export class DeploymentOrchestrator extends DurableObject { await this.ctx.storage.put('archiveBuffer', params.archiveBuffer); await this.saveState(); - // Setup events manager const eventsManager = this.eventsManager(); await eventsManager.initialize(this.state.buildId); await this.addLogEvent('Build created from archive'); - // Schedule alarm to run job asynchronously await this.ctx.storage.setAlarm(Date.now() + 50); return { status: this.state.status }; } - /** - * RPC method: Return current state. - */ async status(): Promise { if (!this.state) { await this.loadState(); @@ -220,9 +172,6 @@ export class DeploymentOrchestrator extends DurableObject { return this.state; } - /** - * RPC method: Return events from EventsManager. - */ async events(): Promise { if (!this.state) { await this.loadState(); @@ -232,18 +181,12 @@ export class DeploymentOrchestrator extends DurableObject { throw new Error('Build not found'); } - // Get EventsManager DO stub and fetch events via RPC const eventsManagerId = this.env.EventsManager.idFromName(this.state.buildId); const eventsManager = this.env.EventsManager.get(eventsManagerId); return eventsManager.getEvents(); } - /** - * RPC method: Cancel a running build. - * Destroys the sandbox and updates status to cancelled. - * Returns a detailed result indicating whether the build was cancelled and why. - */ async cancel(reason?: string): Promise { await this.loadState(); @@ -254,7 +197,6 @@ export class DeploymentOrchestrator extends DurableObject { }; } - // Only cancel if build is queued or building (NOT deploying) const cancellableStatuses: BuildStatus[] = ['queued', 'building']; if (!cancellableStatuses.includes(this.state.status)) { return { @@ -288,26 +230,18 @@ export class DeploymentOrchestrator extends DurableObject { }; } - /** - * Setup project from archive source. - * Extracts the archive buffer from storage into the sandbox. - */ private async setupArchiveSource(sandbox: Awaited>): Promise { const archiveBuffer = await this.ctx.storage.get('archiveBuffer'); if (!archiveBuffer) { throw new Error('Archive buffer not found in storage'); } - // Clear archive from storage (no longer needed) await this.ctx.storage.delete('archiveBuffer'); await this.addLogEvent('Extracting archive...'); - // Convert archive to base64 for transfer const base64Archive = Buffer.from(archiveBuffer).toString('base64'); - // Write archive to sandbox using base64 decoding - // We need to write the base64 content to a file, then decode it await sandbox.writeFile('/tmp/project.tar.gz.b64', base64Archive); const decodeResult = await sandbox.exec( 'base64 -d /tmp/project.tar.gz.b64 > /tmp/project.tar.gz && rm /tmp/project.tar.gz.b64' @@ -320,7 +254,6 @@ export class DeploymentOrchestrator extends DurableObject { ); } - // Create project directory and extract await sandbox.exec('mkdir -p /workspace/project'); const extractResult = await sandbox.exec('tar -xzf /tmp/project.tar.gz -C /workspace/project'); @@ -331,16 +264,11 @@ export class DeploymentOrchestrator extends DurableObject { ); } - // Clean up archive await sandbox.exec('rm /tmp/project.tar.gz'); await this.addLogEvent('Archive extracted successfully'); } - /** - * Setup project from git source. - * Clones the repository into the sandbox. - */ private async setupGitSource( sandbox: Awaited>, source: GitSource, @@ -378,7 +306,6 @@ export class DeploymentOrchestrator extends DurableObject { ); } - // Check if Git LFS is needed by looking for .gitattributes with LFS patterns const lfsCheckResult = await sandbox.exec( 'cd /workspace/project && [ -f .gitattributes ] && grep -q "filter=lfs" .gitattributes' ); @@ -407,7 +334,6 @@ export class DeploymentOrchestrator extends DurableObject { } } - // Get the commit hash const commitHashResult = await sandbox.exec('cd /workspace/project && git rev-parse HEAD'); if (!commitHashResult.success) { console.log(commitHashResult.stderr); @@ -417,9 +343,6 @@ export class DeploymentOrchestrator extends DurableObject { await this.addLogEvent(`Repository cloned successfully (commit: ${commitHash})`); } - /** - * Clear sensitive data from state and return the access token if present. - */ private async popAccessTokenAndEnvData(): Promise<{ accessToken?: string; envVars?: EncryptedEnvVar[]; @@ -446,15 +369,10 @@ export class DeploymentOrchestrator extends DurableObject { return { accessToken, envVars }; } - /** - * Clear sensitive and unnecessary data from state and storage on failure. - * This ensures archive buffers, access tokens, and env vars don't persist after failures. - */ private async clearSensitiveDataOnFailure(): Promise { // Clear archive buffer from storage (may not have been extracted yet) await this.ctx.storage.delete('archiveBuffer'); - // Clear sensitive data from state let needsSave = false; if (this.state.source?.type === 'git' && this.state.source.accessToken) { @@ -496,10 +414,6 @@ export class DeploymentOrchestrator extends DurableObject { await this.addLogEvent('Database migrations completed'); } - /** - * Main orchestration method. - * Clones repo, builds project, and deploys to Cloudflare. - */ private async run(): Promise { let sandbox: Awaited> | null = null; @@ -509,23 +423,19 @@ export class DeploymentOrchestrator extends DurableObject { throw new Error('No source configured for build'); } - // Extract and clear sensitive data from state const { accessToken, envVars } = await this.popAccessTokenAndEnvData(); await this.updateStatus('building'); - // Get sandbox instance sandbox = getSandbox(this.env.Sandbox, this.state.buildId); await this.addLogEvent('Build environment ready'); - // Setup source if (source.type === 'archive') { await this.setupArchiveSource(sandbox); } else { await this.setupGitSource(sandbox, source, accessToken); } - // Step 1: Detect project type await this.addLogEvent('Analyzing project...'); const detectResult = await sandbox.exec( 'cd /workspace/project && /workspace/detect-project.sh /workspace/project' @@ -542,7 +452,6 @@ export class DeploymentOrchestrator extends DurableObject { await this.addLogEvent(`Detected: ${detectedType}`); - // Validate detected type against supported project types const parseResult = supportedProjectTypeSchema.safeParse(detectedType); if (!parseResult.success) { @@ -572,7 +481,6 @@ export class DeploymentOrchestrator extends DurableObject { // throw new Error('Failed to set up tool versions'); // } - // Step 3: Build based on project type const buildPipelines: Record< ProjectType, Array<{ message: string; script: string; passEnvVars?: boolean }> @@ -632,7 +540,6 @@ export class DeploymentOrchestrator extends DurableObject { ], }; - // Decrypt env vars const decryptedEnvVars = decryptEnvVars( envVars || [], Buffer.from(this.env.ENV_ENCRYPTION_PRIVATE_KEY, 'base64') @@ -647,20 +554,17 @@ export class DeploymentOrchestrator extends DurableObject { ); } - // Store detected project type this.state.projectType = projectType; await this.saveState(); await this.addLogEvent('Build completed successfully'); - // Run migrations if needed if (await this.needsMigrations(sandbox)) { await this.runMigrations(sandbox, decryptedEnvVars); } await this.updateStatus('deploying'); - // Read artifacts from sandbox based on detected project type const artifactReader = new SandboxArtifactReader(); const artifacts = await artifactReader.readArtifactsByType( sandbox, @@ -668,7 +572,6 @@ export class DeploymentOrchestrator extends DurableObject { (message: string) => this.addLogEvent(message) ); - // Deploy artifacts const api = new CloudflareAPI(this.env.CLOUDFLARE_ACCOUNT_ID, this.env.CLOUDFLARE_API_TOKEN); const deployer = new Deployer(api); @@ -681,10 +584,8 @@ export class DeploymentOrchestrator extends DurableObject { await this.updateStatus('deployed'); } catch (error) { - // Update status await this.updateStatus('failed'); - // Clear sensitive/unnecessary data from storage on failure await this.clearSensitiveDataOnFailure(); Sentry.captureException(error, { @@ -698,7 +599,6 @@ export class DeploymentOrchestrator extends DurableObject { }, }); } finally { - // Always destroy sandbox if (sandbox) { try { await sandbox.destroy(); @@ -714,9 +614,6 @@ export class DeploymentOrchestrator extends DurableObject { } } - /** - * Helper method to run a script in the sandbox and stream its output to logs. - */ private async runScript( sandbox: Awaited>, command: string, diff --git a/services/deploy-infra/builder/src/env-decryptor.ts b/services/deploy-infra/builder/src/env-decryptor.ts index f4ab1f240e..83a2136c20 100644 --- a/services/deploy-infra/builder/src/env-decryptor.ts +++ b/services/deploy-infra/builder/src/env-decryptor.ts @@ -6,14 +6,6 @@ import { } from '../../../../apps/web/src/lib/user-deployments/env-vars-validation'; import { EnvDecryptionError } from './errors'; -/** - * Decrypts secret environment variables using the provided private key. - * Non-secret variables are returned as-is. - * - * @param envVars - Array of encrypted environment variables (secrets have encrypted values) - * @param privateKey - RSA private key in PEM format for decryption - * @returns Array of decrypted plaintext environment variables - */ export default function decryptEnvVars( envVars: EncryptedEnvVar[], privateKey: Buffer @@ -24,15 +16,12 @@ export default function decryptEnvVars( return envVars.map(v => { if (!v.isSecret) { - // Non-secret values are already plaintext return markAsPlaintext({ key: v.key, value: v.value, isSecret: v.isSecret }); } try { - // Parse the encrypted value as JSON to get the envelope const envelope = JSON.parse(v.value) as EncryptedEnvelope; - // Decrypt using the private key const decryptedValue = decryptWithPrivateKey(envelope, privateKey); return markAsPlaintext({ diff --git a/services/deploy-infra/builder/src/event-store.ts b/services/deploy-infra/builder/src/event-store.ts index 8573435fe3..97072eb4df 100644 --- a/services/deploy-infra/builder/src/event-store.ts +++ b/services/deploy-infra/builder/src/event-store.ts @@ -1,33 +1,14 @@ -/** - * EventStore - Manages event ring buffer and persistence. - * Handles event creation, storage, and trimming based on last processed event. - */ - import type { Event } from './types'; const MAX_EVENTS = 5000; -/** - * EventStore manages a bounded ring buffer of events with delivery-aware trimming. - * - * Features: - * - Ring buffer with configurable maximum size (MAX_EVENTS) - * - Delivery-aware trimming to preserve unprocessed events - * - Automatic persistence to durable storage - * - Sequential event ID generation - * - Tracks last processed event ID for trimming decisions - */ export class EventStore { - /** In-memory ring buffer of events */ private eventsList: Event[] = []; /** Last processed event ID (-1 means no events processed yet) */ private lastProcessedId: number = -1; constructor(private storage: DurableObjectStorage) {} - /** - * Load events and lastProcessedId from durable storage into memory. - */ async loadEvents(): Promise { const storedEvents = await this.storage.get('events'); if (storedEvents) { @@ -40,16 +21,7 @@ export class EventStore { } } - /** - * Add an event to the ring buffer. - * Automatically trims oldest events if buffer exceeds size limits. - * Returns the created event so caller can use its timestamp. - * - * @param eventData - The event envelope to add (without id and ts which are auto-generated) - * @returns The created event with id and timestamp - */ async addEvent(eventData: Omit): Promise { - // Calculate next event ID based on last event in list const lastEvent = this.eventsList[this.eventsList.length - 1]; const nextEventId = lastEvent ? lastEvent.id + 1 : 0; @@ -61,37 +33,23 @@ export class EventStore { this.eventsList.push(event); - // Trim ring buffer if it exceeds limits await this.trimEvents(); - // Persist changes await this.storage.put('events', this.eventsList); return event; } - /** - * Get all events in the buffer. - * - * @returns Array of all events - */ getEvents(): Event[] { return this.eventsList; } - /** - * Get unprocessed events (events with id > lastProcessedId). - * - * @param limit - Optional maximum number of events to return - * @returns Array of unprocessed events, optionally limited - */ getUnprocessedEvents(limit?: number): Event[] { const index = this.getFirstUnprocessedEventIndex(); if (index === null) { return []; } - // Slice from calculated index, applying limit if provided return this.eventsList.slice(index, limit !== undefined ? index + limit : undefined); } @@ -109,34 +67,21 @@ export class EventStore { return null; } - // Calculate starting index based on first event ID and lastProcessedId const firstEventId = this.eventsList[0].id; const startIndex = this.lastProcessedId - firstEventId + 1; - // Handle edge cases if (startIndex >= this.eventsList.length) { // All events have been processed return null; } - // If startIndex is negative or 0, start from beginning return Math.max(0, startIndex); } - /** - * Get the last processed event ID. - * - * @returns The last processed event ID (-1 if no events processed yet) - */ getLastProcessedId(): number { return this.lastProcessedId; } - /** - * Set the last processed event ID and persist it to storage. - * - * @param id - The event ID to mark as last processed - */ async setLastProcessedId(id: number): Promise { this.lastProcessedId = id; await this.storage.put('lastProcessedId', this.lastProcessedId); @@ -154,12 +99,10 @@ export class EventStore { * temporarily allowing the buffer to exceed MAX_EVENTS if necessary. */ private async trimEvents(): Promise { - // Only trim events that have been successfully processed while (this.eventsList.length > MAX_EVENTS && this.eventsList[0].id <= this.lastProcessedId) { this.eventsList.shift(); } - // Safety check: warn if we can't trim because all events are unprocessed if (this.eventsList.length > MAX_EVENTS) { console.warn( `Event buffer exceeded MAX_EVENTS (${MAX_EVENTS}) but cannot trim - all events are unprocessed. ` + diff --git a/services/deploy-infra/builder/src/events-manager.ts b/services/deploy-infra/builder/src/events-manager.ts index 59ad699f12..808f9e8d19 100644 --- a/services/deploy-infra/builder/src/events-manager.ts +++ b/services/deploy-infra/builder/src/events-manager.ts @@ -1,31 +1,13 @@ -/** - * EventsManager - Durable Object for managing events and webhook delivery. - * One instance per build, owns EventStore and WebhookDelivery. - */ - import { DurableObject } from 'cloudflare:workers'; import type { Env, Event } from './types'; import { EventStore } from './event-store'; import { WebhookDelivery } from './webhook-delivery'; import * as Sentry from '@sentry/cloudflare'; -/** - * State structure for EventsManager - */ type EventsManagerState = { buildId: string; }; -/** - * EventsManager Durable Object - * - * Responsibilities: - * - Own and manage EventStore instance - * - Handle webhook delivery with WebhookDelivery class - * - Provide RPC methods for adding events - * - Manage alarm for webhook batching and retries - * - Track build state for webhook payloads - */ export class EventsManager extends DurableObject { private state: EventsManagerState = { buildId: '', @@ -37,10 +19,8 @@ export class EventsManager extends DurableObject { constructor(ctx: DurableObjectState, env: Env) { super(ctx, env); - // Initialize EventStore with this DO's storage this.eventStore = new EventStore(this.ctx.storage); - // Initialize WebhookDelivery with storage, env, build state accessor, alarm interface, and event store this.webhookDelivery = new WebhookDelivery( this.ctx.storage, this.env, @@ -61,9 +41,6 @@ export class EventsManager extends DurableObject { } } - /** - * Load state from storage - */ private async loadState(): Promise { if (this.state.buildId !== '') { return; @@ -74,23 +51,15 @@ export class EventsManager extends DurableObject { this.state = stored; } - // Load EventStore events from storage await this.eventStore.loadEvents(); - // Initialize WebhookDelivery await this.webhookDelivery.initialize(); } - /** - * Save state to storage - */ private async saveState(): Promise { await this.ctx.storage.put('state', this.state); } - /** - * Alarm handler for webhook delivery - */ async alarm(): Promise { try { await this.loadState(); @@ -105,22 +74,12 @@ export class EventsManager extends DurableObject { } } - /** - * RPC: Add a new event - * - * @param eventData - Event data without id and ts (will be added by EventStore) - */ async addEvent(eventData: Omit): Promise { await this.loadState(); await this.eventStore.addEvent(eventData); await this.webhookDelivery.scheduleFlush(); } - /** - * RPC: Get all events - * - * @returns Array of all events - */ async getEvents(): Promise { await this.loadState(); return this.eventStore.getEvents(); diff --git a/services/deploy-infra/builder/src/index.ts b/services/deploy-infra/builder/src/index.ts index 240b4e0bbd..df4de3212e 100644 --- a/services/deploy-infra/builder/src/index.ts +++ b/services/deploy-infra/builder/src/index.ts @@ -9,14 +9,12 @@ import { CloudflareAPI } from './cloudflare-api'; import { validateWorkerName } from './utils'; import * as Sentry from '@sentry/cloudflare'; -// Import base Durable Objects import { DeploymentOrchestrator as DeploymentOrchestratorBase } from './deployment-orchestrator'; import { EventsManager as EventsManagerBase } from './events-manager'; import { htmlDeployHandler } from './html-deploy/handler'; import { runEphemeralDeploymentCleanup } from './html-deploy/ephemeral-cleanup'; export { Sandbox } from '@cloudflare/sandbox'; -// Export Sentry-instrumented Durable Objects export const DeploymentOrchestrator = Sentry.instrumentDurableObjectWithSentry( (env: Env) => ({ dsn: env.SENTRY_DSN, @@ -50,13 +48,11 @@ app.post('/deploy-html', htmlDeployHandler); // ── Backend-authenticated routes ─────────────────────────────────────────── -// Authentication middleware app.use( '*', backendAuthMiddleware(c => c.env.BACKEND_AUTH_TOKEN) ); -// Route: POST /deploy app.post('/deploy', async (c: Context) => { let body: DeployRequest; @@ -66,7 +62,6 @@ app.post('/deploy', async (c: Context) => { return c.json({ error: 'Invalid JSON body' }, 400); } - // Handle build cancellations if provided if (body.cancelBuildIds && body.cancelBuildIds.length > 0) { await Promise.allSettled( body.cancelBuildIds.map(async buildId => { @@ -83,26 +78,21 @@ app.post('/deploy', async (c: Context) => { ); } - // Validate required fields if (!body.slug || !body.provider || !body.repoSource) { return c.json({ error: 'Missing required fields: slug, provider, repoSource' }, 400); } - // Validate slug format try { validateWorkerName(body.slug); } catch (_error) { return c.json({ error: 'Invalid worker slug' }, 400); } - // Generate unique job ID const buildId = createDurableObjectBuilderID(); - // Get Durable Object stub const id = c.env.DeploymentOrchestrator.idFromName(buildId); const stub = c.env.DeploymentOrchestrator.get(id); - // Start the job via RPC const result = await stub.start({ buildId, slug: body.slug, @@ -116,7 +106,6 @@ app.post('/deploy', async (c: Context) => { envVars: body.envVars, }); - // Return 202 Accepted with job details const response: DeployResponse = { buildId, slug: body.slug, @@ -126,7 +115,6 @@ app.post('/deploy', async (c: Context) => { return c.json(response, 202); }); -// Route: POST /deploy-archive - Deploy from uploaded tar.gz archive app.post('/deploy-archive', async (c: Context) => { const slug = c.req.header('X-Slug'); @@ -134,14 +122,12 @@ app.post('/deploy-archive', async (c: Context) => { return c.json({ error: 'Missing X-Slug header' }, 400); } - // Validate slug format try { validateWorkerName(slug); } catch (_error) { return c.json({ error: 'Invalid worker slug' }, 400); } - // Parse optional env vars from header const envVarsHeader = c.req.header('X-Env-Vars'); let envVars: DeployRequest['envVars'] | undefined; if (envVarsHeader) { @@ -159,14 +145,11 @@ app.post('/deploy-archive', async (c: Context) => { return c.json({ error: 'Empty archive body' }, 400); } - // Generate unique job ID const buildId = createDurableObjectBuilderID(); - // Get Durable Object stub const id = c.env.DeploymentOrchestrator.idFromName(buildId); const stub = c.env.DeploymentOrchestrator.get(id); - // Start the job via RPC with archive source const result = await stub.startFromArchive({ buildId, slug, @@ -174,7 +157,6 @@ app.post('/deploy-archive', async (c: Context) => { envVars, }); - // Return 202 Accepted with job details const response: DeployResponse = { buildId, slug, @@ -189,12 +171,10 @@ app.get('/deploy/:buildId/status', async (c: Context) => { const buildId = c.req.param('buildId'); if (!buildId) return c.json({ error: 'Missing buildId' }, 400); - // Get Durable Object stub const id = c.env.DeploymentOrchestrator.idFromName(buildId); const stub = c.env.DeploymentOrchestrator.get(id); try { - // Fetch status from Durable Object via RPC const status: StatusResponse = await stub.status(); return c.json(status, 200); } catch (error) { @@ -211,12 +191,10 @@ app.get('/deploy/:buildId/events', async (c: Context) => { const buildId = c.req.param('buildId'); if (!buildId) return c.json({ error: 'Missing buildId' }, 400); - // Get Durable Object stub const id = c.env.DeploymentOrchestrator.idFromName(buildId); const stub = c.env.DeploymentOrchestrator.get(id); try { - // Fetch events from Durable Object via RPC const events = await stub.events(); return c.json(events, 200); } catch (error) { @@ -233,7 +211,6 @@ app.delete('/deploy/:buildId', async (c: Context) => { const buildId = c.req.param('buildId'); if (!buildId) return c.json({ error: 'Missing buildId' }, 400); - // Get Durable Object stub const id = c.env.DeploymentOrchestrator.idFromName(buildId); const stub = c.env.DeploymentOrchestrator.get(id); @@ -251,21 +228,19 @@ app.delete('/deploy/:buildId', async (c: Context) => { }); /** - * Delete a worker from the dispatch namespace * Note: Assets are automatically cleaned up when the script is deleted */ app.delete('/worker/:slug', async (c: Context) => { const slug = c.req.param('slug'); if (!slug) return c.json({ error: 'Missing slug' }, 400); - // Validate slug format try { validateWorkerName(slug); } catch (_error) { return c.json({ error: 'Invalid worker slug' }, 400); } - const dispatchNamespace = 'kilo-deploy'; // Hardcoded for now + const dispatchNamespace = 'kilo-deploy'; try { const cloudflareApi = new CloudflareAPI( @@ -284,7 +259,6 @@ app.delete('/worker/:slug', async (c: Context) => { } }); -// Global error handler const errorHandler = createErrorHandler(console, { includeMessage: false }); app.onError((err, c) => { Sentry.captureException(err, { @@ -314,7 +288,6 @@ export default Sentry.withSentry( { fetch: app.fetch, - // ── Scheduled handler: minute-scale ephemeral cleanup ─────────────────── async scheduled( _controller: ScheduledController, env: Env, diff --git a/services/deploy-infra/builder/src/sandbox-artifact-reader.ts b/services/deploy-infra/builder/src/sandbox-artifact-reader.ts index 7ae4d961cd..51902bf4d0 100644 --- a/services/deploy-infra/builder/src/sandbox-artifact-reader.ts +++ b/services/deploy-infra/builder/src/sandbox-artifact-reader.ts @@ -4,25 +4,13 @@ import { readFolderAsArchive } from './sandbox-file-reader'; import { ArtifactReadError } from './errors'; import staticWorkerContent from './assets/static.worker.js'; -// Type for the sandbox stub returned by getSandbox() type SandboxStub = Awaited>; -/** - * Reads deployment artifacts from a Cloudflare Sandbox - */ export class SandboxArtifactReader { /** * Read worker script and assets from sandbox using tar-based reading for better performance. * This method uses tar archives to efficiently transfer entire directory structures from the sandbox, * which is significantly faster than reading files individually. - * - * @param sandbox - Cloudflare Sandbox instance - * @param bundledPath - Path to directory containing bundled files - * @param entrypointFilename - Name of the entrypoint file (e.g., 'worker.js') - * @param assetsPath - Path to assets directory in sandbox - * @param logger - Optional callback for logging progress - * @param batchSize - Batch size parameter (unused in tar-based implementation, kept for signature compatibility) - * @returns Deployment artifacts ready for deployment */ async readArtifacts( sandbox: SandboxStub, @@ -33,18 +21,15 @@ export class SandboxArtifactReader { ): Promise { const session = await sandbox.createSession(); - // Read bundled files as tar archive, excluding README.md and .map files if (logger) logger('Reading build output...'); const bundledFiles = await readFolderAsArchive(session, bundledPath, ['README.md', '*.map']); let workerScript: DeploymentFile | null = null; const artifacts: DeploymentFile[] = []; - // Process bundled files for (const fileEntry of bundledFiles) { const filename = fileEntry.path.split('/').pop() || ''; - // Check if this is the entrypoint file if (filename === entrypointFilename) { workerScript = fileEntry; } else { @@ -56,7 +41,6 @@ export class SandboxArtifactReader { throw new Error('Build output is incomplete'); } - // Read assets as tar archive if (logger) logger('Reading assets...'); const assetFiles = await readFolderAsArchive(session, assetsPath); @@ -67,12 +51,6 @@ export class SandboxArtifactReader { }; } - /** - * Read artifacts from standard OpenNext output structure (Next.js) - * @param sandbox - Cloudflare Sandbox instance - * @param logger - Optional callback for logging progress - * @returns Deployment artifacts - */ async readOpenNextArtifacts( sandbox: SandboxStub, logger?: (message: string) => void @@ -88,13 +66,6 @@ export class SandboxArtifactReader { } } - /** - * Read assets from static site output structure. - * Returns complete DeploymentArtifacts with the static worker script. - * @param sandbox - Cloudflare Sandbox instance - * @param logger - Optional callback for logging progress - * @returns Deployment artifacts with assets and static worker script - */ async readStaticSiteAssets(sandbox: SandboxStub): Promise { const assetsPath = '/workspace/project/.static-site/assets'; const session = await sandbox.createSession(); @@ -113,13 +84,6 @@ export class SandboxArtifactReader { }; } - /** - * Read artifacts based on detected project type - * @param sandbox - Cloudflare Sandbox instance - * @param projectType - Detected project type - * @param logger - Optional callback for logging progress - * @returns Deployment artifacts - */ async readArtifactsByType( sandbox: SandboxStub, projectType: ProjectType, @@ -128,7 +92,6 @@ export class SandboxArtifactReader { if (projectType === 'nextjs') { return this.readOpenNextArtifacts(sandbox, logger); } - // For static sites, return DeploymentArtifacts with static worker script return this.readStaticSiteAssets(sandbox); } } diff --git a/services/deploy-infra/builder/src/sandbox-file-reader.ts b/services/deploy-infra/builder/src/sandbox-file-reader.ts index 982e438133..3f991f93ac 100644 --- a/services/deploy-infra/builder/src/sandbox-file-reader.ts +++ b/services/deploy-infra/builder/src/sandbox-file-reader.ts @@ -10,18 +10,9 @@ import { Readable, type PassThrough } from 'stream'; import type { DeploymentFile } from './types'; import { getMimeType } from './utils'; -// Type for the sandbox stub returned by getSandbox() type SandboxStub = Awaited>; -/** - * List all files recursively in a directory using the sandbox. - * - * @param sandbox - The Cloudflare Sandbox instance - * @param root - Root directory path to search (e.g., "/workspace/result/assets") - * @returns Array of file paths relative to root - */ export async function listFilesRecursive(sandbox: SandboxStub, root: string): Promise { - // Execute find command to list all files recursively const result = await sandbox.exec(`find ${root} -type f`); if (!result.success) { @@ -33,16 +24,13 @@ export async function listFilesRecursive(sandbox: SandboxStub, root: string): Pr return []; } - // Parse output into file paths const files = result.stdout .split('\n') .map((line: string) => line.trim()) .filter((line: string) => line.length > 0) .map((path: string) => { - // Strip root prefix to get relative paths if (path.startsWith(root)) { const relativePath = path.slice(root.length); - // Remove leading slash if present return relativePath.startsWith('/') ? relativePath.slice(1) : relativePath; } return path; @@ -59,7 +47,6 @@ export async function readFileAsBuffer(session: ExecutionSession, path: string): // Escape path for shell command const escapedPath = path.replace(/'/g, "'\\''"); - // First, get file metadata to determine size and if it's binary const statResult = await session.exec(`stat -c '%s' '${escapedPath}' 2>/dev/null`); if (!statResult.success || statResult.exitCode !== 0) { @@ -71,7 +58,6 @@ export async function readFileAsBuffer(session: ExecutionSession, path: string): throw new Error(`Invalid file size for ${path}`); } - // Read file in chunks const chunkSize = 65535 * 40; let bytesRead = 0; let blockNumber = 0; @@ -110,7 +96,6 @@ export async function readFileAsBuffer(session: ExecutionSession, path: string): break; } - // Convert chunk to Buffer const chunkBuffer = Buffer.from(chunkData, 'base64'); chunks.push(chunkBuffer); @@ -118,19 +103,9 @@ export async function readFileAsBuffer(session: ExecutionSession, path: string): blockNumber++; } - // Concatenate all chunks into a single Buffer return Buffer.concat(chunks); } -/** - * Read a folder from the sandbox by creating a tar archive, reading it, and extracting locally. - * This is useful for efficiently transferring entire directory structures from the sandbox. - * - * @param session - The Cloudflare Sandbox ExecutionSession instance - * @param folderPath - Absolute path to the folder in the sandbox - * @param excludePatterns - Optional array of patterns to exclude (e.g., ['node_modules', '*.log', '.git']) - * @returns Array of DeploymentFile objects with paths relative to the archived folder and their contents as buffers - */ export async function readFolderAsArchive( session: ExecutionSession, folderPath: string, @@ -139,7 +114,6 @@ export async function readFolderAsArchive( // Escape path for shell command const escapedPath = folderPath.replace(/'/g, "'\\''"); - // Create a temporary tar file in the sandbox using a UUID const tmpArchivePath = `/tmp/folder-archive-${crypto.randomUUID()}.tar`; const escapedArchivePath = tmpArchivePath.replace(/'/g, "'\\''"); @@ -158,14 +132,12 @@ export async function readFolderAsArchive( try { const archiveBuffer = await readFileAsBuffer(session, tmpArchivePath); - // Extract the archive const files: DeploymentFile[] = []; const bufferStream = Readable.from(archiveBuffer); const extractStream = extract(); extractStream.on('entry', (header: Headers, stream: PassThrough, next: () => void) => { - // Only process files, not directories if (header.type === 'file') { const chunks: Buffer[] = []; @@ -177,7 +149,6 @@ export async function readFolderAsArchive( const fileBuffer = Buffer.concat(chunks); const mimeType = getMimeType(header.name); - // Normalize path: remove leading "./" or "/" let normalizedPath = header.name; if (normalizedPath.startsWith('./')) { normalizedPath = normalizedPath.slice(2); @@ -197,13 +168,11 @@ export async function readFolderAsArchive( throw new Error(`Error reading file ${header.name} from archive: ${err.message}`); }); } else { - // Skip directories stream.resume(); next(); } }); - // Process the archive await pipeline(bufferStream, extractStream); return files; diff --git a/services/deploy-infra/builder/src/sanitize-git-error.ts b/services/deploy-infra/builder/src/sanitize-git-error.ts index ed78c45203..3fdfd52634 100644 --- a/services/deploy-infra/builder/src/sanitize-git-error.ts +++ b/services/deploy-infra/builder/src/sanitize-git-error.ts @@ -8,7 +8,6 @@ export function sanitizeGitError(error: unknown, accessToken: string | undefined } const errorMessage = error instanceof Error ? error.message : String(error); - // Replace the access token with [REDACTED] in the error message // Using replaceAll instead of regex to avoid issues with special characters in tokens const sanitizedMessage = errorMessage.replaceAll(accessToken, '[REDACTED]'); diff --git a/services/deploy-infra/builder/src/tests/sanitize-git-error.test.ts b/services/deploy-infra/builder/src/tests/sanitize-git-error.test.ts index 68e07ab9f4..ab2eb321b1 100644 --- a/services/deploy-infra/builder/src/tests/sanitize-git-error.test.ts +++ b/services/deploy-infra/builder/src/tests/sanitize-git-error.test.ts @@ -1,10 +1,3 @@ -/** - * Tests for sanitizeGitError function. - * - * Ensures access tokens are properly redacted from error messages, - * including tokens that contain regex special characters. - */ - import { sanitizeGitError } from '../sanitize-git-error'; describe('sanitizeGitError', () => { @@ -61,7 +54,6 @@ describe('sanitizeGitError', () => { }); it('should handle tokens with regex special characters safely', () => { - // Token containing regex special characters: . * + ? ^ $ { } [ ] \ | ( ) const accessToken = 'token.with*special+chars?and^more$chars'; const error = new Error(`Failed with token: ${accessToken}`); diff --git a/services/deploy-infra/builder/src/tests/webhook-delivery.test.ts b/services/deploy-infra/builder/src/tests/webhook-delivery.test.ts index f606e6a130..42f108a1bf 100644 --- a/services/deploy-infra/builder/src/tests/webhook-delivery.test.ts +++ b/services/deploy-infra/builder/src/tests/webhook-delivery.test.ts @@ -1,24 +1,7 @@ -/** - * Comprehensive tests for webhook delivery functionality. - * - * Tests cover the following scenarios from the MVP plan: - * - Happy path single batch delivery - * - Multiple batches for large event streams - * - Retryable failure then success with exponential backoff - * - Stop-after exceeded (delivery permanently stopped) - * - Preserve undelivered events - * - Batch timing (waits up to BATCH_MAX_MS when below threshold) - * - * Uses a mock fetch implementation to simulate backend responses. - */ - import type { DeliveryState, Event, Build, Env, WebhookPayload } from '../types'; import { WebhookDelivery } from '../webhook-delivery'; import { EventStore } from '../event-store'; -/** - * Mock DurableObjectStorage for testing - */ class MockDurableObjectStorage { private storage = new Map(); @@ -45,9 +28,6 @@ class MockDurableObjectStorage { } } -/** - * Mock fetch responses for testing - */ type MockFetchResponse = { ok: boolean; status: number; @@ -60,7 +40,6 @@ let lastFetchPayload: WebhookPayload | null = null; const mockFetch = jest.fn(async (url: string, options?: RequestInit): Promise => { fetchCallCount++; - // Capture the payload for assertions if (options?.body) { lastFetchPayload = JSON.parse(options.body as string) as WebhookPayload; } @@ -73,12 +52,8 @@ const mockFetch = jest.fn(async (url: string, options?: RequestInit): Promise): Env { return { CLOUDFLARE_ACCOUNT_ID: 'test-account', @@ -93,9 +68,6 @@ function createTestEnv(overrides?: Partial): Env { } as Env; } -/** - * Test helper class that wraps WebhookDelivery and EventStore - */ class TestWebhookDeliveryHandler { private storage: MockDurableObjectStorage; private eventStore: EventStore; @@ -181,7 +153,6 @@ class TestWebhookDeliveryHandler { describe('Webhook Delivery', () => { beforeEach(() => { - // Reset mock state mockFetchResponses = []; fetchCallCount = 0; lastFetchPayload = null; @@ -194,18 +165,14 @@ describe('Webhook Delivery', () => { await handler.initialize(); - // Add some events await handler.addEvent('Event 1'); await handler.addEvent('Event 2'); await handler.addEvent('Event 3'); - // Mock successful response mockFetchResponses.push({ ok: true, status: 200 }); - // Flush events await handler.flush(); - // Verify delivery expect(fetchCallCount).toBe(1); expect(lastFetchPayload).toBeTruthy(); expect(lastFetchPayload!.events.length).toBe(3); @@ -213,56 +180,48 @@ describe('Webhook Delivery', () => { expect((lastFetchPayload!.events[0].payload as { message: string }).message).toBe('Event 1'); expect((lastFetchPayload!.events[2].payload as { message: string }).message).toBe('Event 3'); - // Verify delivery state updated const deliveryState = handler.getDeliveryState(); expect(deliveryState).toBeTruthy(); expect(deliveryState!.attempt).toBe(0); expect(deliveryState!.nextAttemptAt).toBe(0); - // Verify last processed ID updated expect(handler.getLastProcessedId()).toBe(2); }); it('should split large event streams into multiple batches', async () => { const env = createTestEnv({ - BACKEND_WEBHOOK_BATCH_MAX_EVENTS: '10', // Small batch size for testing + BACKEND_WEBHOOK_BATCH_MAX_EVENTS: '10', }); const handler = new TestWebhookDeliveryHandler(env); await handler.initialize(); - // Add 25 events for (let i = 0; i < 25; i++) { await handler.addEvent(`Event ${i + 1}`); } - // Mock successful responses for 3 batches mockFetchResponses.push({ ok: true, status: 200 }); mockFetchResponses.push({ ok: true, status: 200 }); mockFetchResponses.push({ ok: true, status: 200 }); - // First batch: events 0-9 await handler.flush(); expect(fetchCallCount).toBe(1); expect(lastFetchPayload!.events.length).toBe(10); expect(lastFetchPayload!.events[0].id).toBe(0); expect(lastFetchPayload!.events[9].id).toBe(9); - // Second batch: events 10-19 await handler.flush(); expect(fetchCallCount).toBe(2); expect(lastFetchPayload!.events.length).toBe(10); expect(lastFetchPayload!.events[0].id).toBe(10); expect(lastFetchPayload!.events[9].id).toBe(19); - // Third batch: events 20-24 await handler.flush(); expect(fetchCallCount).toBe(3); expect(lastFetchPayload!.events.length).toBe(5); expect(lastFetchPayload!.events[0].id).toBe(20); expect(lastFetchPayload!.events[4].id).toBe(24); - // Verify final state expect(handler.getLastProcessedId()).toBe(24); }); @@ -277,7 +236,6 @@ describe('Webhook Delivery', () => { await handler.addEvent('Event 1'); await handler.addEvent('Event 2'); - // First attempt: fail with 503 mockFetchResponses.push({ ok: false, status: 503 }); await handler.flush(); @@ -286,12 +244,10 @@ describe('Webhook Delivery', () => { expect(deliveryState!.attempt).toBe(1); expect(deliveryState!.nextAttemptAt).toBeGreaterThan(Date.now()); - // Calculate expected backoff: 1000 * 2^0 = 1000ms const firstBackoff = deliveryState!.nextAttemptAt - Date.now(); expect(firstBackoff).toBeGreaterThanOrEqual(900); expect(firstBackoff).toBeLessThanOrEqual(1100); - // Second attempt: fail with 500 mockFetchResponses.push({ ok: false, status: 500 }); await handler.flush(); @@ -299,12 +255,10 @@ describe('Webhook Delivery', () => { deliveryState = handler.getDeliveryState(); expect(deliveryState!.attempt).toBe(2); - // Calculate expected backoff: 1000 * 2^1 = 2000ms const secondBackoff = deliveryState!.nextAttemptAt - Date.now(); expect(secondBackoff).toBeGreaterThanOrEqual(1900); expect(secondBackoff).toBeLessThanOrEqual(2100); - // Third attempt: succeed mockFetchResponses.push({ ok: true, status: 200 }); await handler.flush(); @@ -317,7 +271,7 @@ describe('Webhook Delivery', () => { it('should stop retrying after STOP_AFTER_ATTEMPTS is exceeded', async () => { const env = createTestEnv({ - BACKEND_WEBHOOK_STOP_AFTER_ATTEMPTS: '2', // Very low for testing + BACKEND_WEBHOOK_STOP_AFTER_ATTEMPTS: '2', BACKEND_WEBHOOK_BACKOFF_BASE_MS: '10', }); const handler = new TestWebhookDeliveryHandler(env); @@ -326,32 +280,27 @@ describe('Webhook Delivery', () => { await handler.addEvent('Event 1'); - // First failure mockFetchResponses.push({ ok: false, status: 503 }); await handler.flush(); let deliveryState = handler.getDeliveryState(); expect(deliveryState!.attempt).toBe(1); - // Second failure mockFetchResponses.push({ ok: false, status: 503 }); await handler.flush(); deliveryState = handler.getDeliveryState(); expect(deliveryState!.attempt).toBe(2); - // Third failure - should exceed limit mockFetchResponses.push({ ok: false, status: 503 }); await handler.flush(); deliveryState = handler.getDeliveryState(); expect(deliveryState!.attempt).toBe(3); - // Verify no more retries happen (attempt > STOP_AFTER_ATTEMPTS) mockFetchResponses.push({ ok: true, status: 200 }); await handler.flush(); - // fetchCallCount should still be 3 (no fourth attempt) expect(fetchCallCount).toBe(3); }); @@ -363,23 +312,19 @@ describe('Webhook Delivery', () => { await handler.initialize(); - // Add events and deliver some for (let i = 0; i < 10; i++) { await handler.addEvent(`Event ${i + 1}`); } - // Deliver first 5 events mockFetchResponses.push({ ok: true, status: 200 }); await handler.flush(); expect(handler.getLastProcessedId()).toBe(4); - // Add more events for (let i = 10; i < 15; i++) { await handler.addEvent(`Event ${i + 1}`); } - // Verify undelivered events are preserved const unprocessedEvents = handler.getUnprocessedEvents(); expect(unprocessedEvents.length).toBe(10); // Events 5-14 expect(unprocessedEvents[0].id).toBe(5); @@ -395,17 +340,14 @@ describe('Webhook Delivery', () => { await handler.initialize(); - // Add only 3 events (below threshold of 10) await handler.addEvent('Event 1'); await handler.addEvent('Event 2'); await handler.addEvent('Event 3'); - // Verify alarm is scheduled const alarmTime = handler.getAlarmTime(); expect(alarmTime).toBeTruthy(); expect(alarmTime!).toBeGreaterThan(Date.now()); - // Alarm should be scheduled for approximately BATCH_MAX_MS in the future const delay = alarmTime! - Date.now(); expect(delay).toBeGreaterThan(2900); expect(delay).toBeLessThan(3100); @@ -419,18 +361,15 @@ describe('Webhook Delivery', () => { await handler.initialize(); - // Add exactly 5 events (at threshold) for (let i = 0; i < 5; i++) { await handler.addEvent(`Event ${i + 1}`); } - // Should schedule immediate flush (50ms) const alarmTime = handler.getAlarmTime(); expect(alarmTime).toBeTruthy(); const delay = alarmTime! - Date.now(); expect(delay).toBeLessThan(100); - // Flush should send immediately mockFetchResponses.push({ ok: true, status: 200 }); await handler.flush(); @@ -446,17 +385,14 @@ describe('Webhook Delivery', () => { await handler.addEvent('Event 1'); - // Start first flush mockFetchResponses.push({ ok: true, status: 200 }); const firstFlush = handler.flush(); - // Try to start second flush while first is in progress mockFetchResponses.push({ ok: true, status: 200 }); const secondFlush = handler.flush(); await Promise.all([firstFlush, secondFlush]); - // Should only make one fetch call due to reentrancy guard expect(fetchCallCount).toBe(1); }); @@ -481,13 +417,10 @@ describe('Webhook Delivery', () => { await handler.initialize(); - // Try to schedule flush with no events await handler.flush(); - // Should not have made any fetch calls expect(fetchCallCount).toBe(0); - // Should not have scheduled an alarm expect(handler.getAlarmTime()).toBeNull(); }); @@ -501,14 +434,12 @@ describe('Webhook Delivery', () => { await handler.addEvent('Event 1'); - // First attempt: fail mockFetchResponses.push({ ok: false, status: 503 }); await handler.flush(); const deliveryState = handler.getDeliveryState(); expect(deliveryState!.attempt).toBe(1); - // Verify alarm is set to nextAttemptAt const alarmTime = handler.getAlarmTime(); expect(alarmTime).toBe(deliveryState!.nextAttemptAt); }); diff --git a/services/deploy-infra/builder/src/types.ts b/services/deploy-infra/builder/src/types.ts index 0fa6cb75e3..5149f450e1 100644 --- a/services/deploy-infra/builder/src/types.ts +++ b/services/deploy-infra/builder/src/types.ts @@ -1,13 +1,8 @@ -/** - * Type definitions for deployment artifacts and job state. - */ - import { z } from 'zod'; import type { Sandbox } from '@cloudflare/sandbox'; import type { DeploymentOrchestrator } from './deployment-orchestrator'; import type { EventsManager } from './events-manager'; -// Import and re-export shared types from backend import type { BuildStatus, Provider, @@ -32,14 +27,6 @@ export type { CancelBuildResult, }; -/** - * Zod schema for supported project types - * - nextjs: Next.js application (uses OpenNext pipeline) - * - hugo: Hugo static site generator (uses Hugo binary) - * - jekyll: Jekyll static site generator (uses Ruby/Bundler) - * - eleventy: Eleventy/11ty static site generator (uses Node.js) - * - plain-html: Plain HTML site (index.html in root) - */ export const supportedProjectTypeSchema = z.enum([ 'nextjs', 'hugo', @@ -51,51 +38,27 @@ export const supportedProjectTypeSchema = z.enum([ export type ProjectType = z.infer; -/** - * Represents a file to be deployed - */ export type DeploymentFile = { - /** Relative path of the file */ path: string; - /** File content */ content: Buffer; - /** MIME type of the file */ mimeType: string; }; -/** - * Worker metadata for Cloudflare deployment - */ export type WorkerMetadata = { - /** Main module entry point */ main_module: string; - /** Compatibility date for the worker */ compatibility_date: string; - /** Compatibility flags for the worker */ compatibility_flags: string[]; - /** Asset configuration with JWT and config */ assets?: { jwt: string; config: Record }; - /** Worker bindings (e.g., KV, Durable Objects, Assets) */ bindings?: Array>; - /** Durable Object migrations */ migrations?: { tag: string; new_classes?: string[] }[]; }; -/** - * Artifacts needed for worker deployment - */ export type DeploymentArtifacts = { - /** Worker script content */ workerScript: DeploymentFile; - /** Additional artifact files (empty array if no artifacts) */ artifacts: DeploymentFile[]; - /** Asset files (empty array if no assets) */ assets: DeploymentFile[]; }; -/** - * Git repository source for deployment - */ export type GitSource = { type: 'git'; provider: Provider; @@ -105,21 +68,12 @@ export type GitSource = { branch?: string; }; -/** - * Archive source for deployment - */ export type ArchiveSource = { type: 'archive'; }; -/** - * Source for deployment - either a git repository or an archive - */ export type BuildSource = GitSource | ArchiveSource; -/** - * Build - */ export type Build = { buildId: string; slug: string; @@ -129,13 +83,9 @@ export type Build = { updatedAt: string; startedAt?: string; completedAt?: string; - /** Detected project type (set after detection phase) */ projectType?: ProjectType; }; -/** - * Request params for starting deployment from archive - */ export type ArchiveDeployParams = { buildId: string; slug: string; @@ -143,19 +93,12 @@ export type ArchiveDeployParams = { envVars?: EncryptedEnvVar[]; }; -/** - * Webhook delivery state tracking (persisted to durable storage). - */ export type DeliveryState = { /** Epoch milliseconds for the next scheduled delivery attempt (0 means no scheduled attempt) */ nextAttemptAt: number; - /** Number of consecutive delivery failures for exponential backoff calculation */ attempt: number; }; -/** - * Environment bindings for the worker - */ export type Env = { CLOUDFLARE_ACCOUNT_ID: string; CLOUDFLARE_API_TOKEN: string; @@ -166,7 +109,6 @@ export type Env = { /** RSA private key in PEM format for decrypting secret environment variables */ ENV_ENCRYPTION_PRIVATE_KEY: string; - /** Cloudflare version metadata binding */ CF_VERSION_METADATA: { id: string; tag: string; timestamp: string }; BACKEND_AUTH_TOKEN: string; @@ -175,16 +117,12 @@ export type Env = { /** URL endpoint where build events will be sent (REQUIRED) */ BACKEND_EVENTS_URL: string; - /** Maximum number of events to batch before sending (default: 50) */ BACKEND_WEBHOOK_BATCH_MAX_EVENTS?: string; - /** Maximum time in milliseconds to wait before sending a batch (default: 3000) */ BACKEND_WEBHOOK_BATCH_MAX_MS?: string; - /** Base backoff time in milliseconds for retry attempts (default: 2000) */ BACKEND_WEBHOOK_BACKOFF_BASE_MS?: string; - /** Maximum number of attempts before giving up (default: 10) */ BACKEND_WEBHOOK_STOP_AFTER_ATTEMPTS?: string; Sandbox: DurableObjectNamespace; @@ -202,9 +140,6 @@ export type Env = { /** Hono app environment with Worker bindings. */ export type HonoEnv = { Bindings: Env }; -/** - * Request body for POST /deploy - */ export type DeployRequest = { slug: string; provider: Provider; @@ -217,24 +152,17 @@ export type DeployRequest = { envVars?: EncryptedEnvVar[]; }; -/** - * Response for POST /deploy - */ export type DeployResponse = { buildId: string; slug: string; status: BuildStatus; }; -/** - * Response for GET /deploy/:buildId/status - */ export type StatusResponse = { status: BuildStatus; updatedAt: string; startedAt?: string; completedAt?: string; - /** Detected project type (available after build detection phase) */ projectType?: ProjectType; }; diff --git a/services/deploy-infra/builder/src/utils.ts b/services/deploy-infra/builder/src/utils.ts index a1d5645031..c60896a012 100644 --- a/services/deploy-infra/builder/src/utils.ts +++ b/services/deploy-infra/builder/src/utils.ts @@ -1,6 +1,3 @@ -/** - * Test-friendly logging helpers that suppress output during tests - */ const isInTestMode = typeof process !== 'undefined' && process.env?.NODE_ENV === 'test'; const consoleExceptInTest = (kind: 'log' | 'warn' | 'error') => (isInTestMode ? () => {} : console[kind]) satisfies typeof console.log; @@ -10,7 +7,6 @@ export const warnExceptInTest = consoleExceptInTest('warn'); export const errorExceptInTest = consoleExceptInTest('error'); export function validateWorkerName(name: string): void { - // Validate worker name const nameRegex = /^[a-zA-Z0-9_-]{1,64}$/; if (!nameRegex.test(name)) { throw new Error( @@ -19,13 +15,7 @@ export function validateWorkerName(name: string): void { } } -/** - * Calculate SHA-256 hash of a Buffer or string - * @param content - Buffer or string to hash - * @returns Full 64-character hex hash - */ export async function calculateSHA256(content: Buffer | string): Promise { - // Convert to Uint8Array for crypto.subtle.digest const inputBuffer = typeof content === 'string' ? new TextEncoder().encode(content) : new Uint8Array(content); const hashBuffer = await crypto.subtle.digest('SHA-256', inputBuffer); @@ -34,11 +24,6 @@ export async function calculateSHA256(content: Buffer | string): Promise return hashHex; } -/** - * Get the byte size of Buffer or string content - * @param content - Buffer or string - * @returns Size in bytes - */ export function getByteSize(content: Buffer | string): number { if (typeof content === 'string') { return new TextEncoder().encode(content).length; @@ -46,11 +31,6 @@ export function getByteSize(content: Buffer | string): number { return content.length; } -/** - * Get MIME type based on file extension - * @param path - File path - * @returns MIME type string - */ export function getMimeType(path: string): string { // Remove query strings first (e.g., 'file.wasm?module' -> 'file.wasm') const cleanPath = path.split('?')[0]; @@ -58,7 +38,6 @@ export function getMimeType(path: string): string { const ext = cleanPath.split('.').pop()?.toLowerCase(); const mimeTypes: Record = { - // Text html: 'text/html', htm: 'text/html', css: 'text/css', @@ -68,8 +47,6 @@ export function getMimeType(path: string): string { xml: 'application/xml', txt: 'text/plain', md: 'text/markdown', - - // Images png: 'image/png', jpg: 'image/jpeg', jpeg: 'image/jpeg', @@ -78,32 +55,22 @@ export function getMimeType(path: string): string { webp: 'image/webp', ico: 'image/x-icon', bmp: 'image/bmp', - - // Fonts woff: 'font/woff', woff2: 'font/woff2', ttf: 'font/ttf', otf: 'font/otf', eot: 'application/vnd.ms-fontobject', - - // Media mp4: 'video/mp4', webm: 'video/webm', mp3: 'audio/mpeg', wav: 'audio/wav', ogg: 'audio/ogg', - - // Documents pdf: 'application/pdf', zip: 'application/zip', tar: 'application/x-tar', gz: 'application/gzip', - - // Web wasm: 'application/wasm', map: 'application/json', - - // Binary bin: 'application/octet-stream', }; diff --git a/services/deploy-infra/builder/src/webhook-delivery.ts b/services/deploy-infra/builder/src/webhook-delivery.ts index ece5e88bb1..5aeaafe875 100644 --- a/services/deploy-infra/builder/src/webhook-delivery.ts +++ b/services/deploy-infra/builder/src/webhook-delivery.ts @@ -1,8 +1,3 @@ -/** - * WebhookDelivery - Manages webhook delivery with batching and retry logic. - * Handles delivery state, exponential backoff, and batch scheduling. - */ - import type { Env, Event, DeliveryState, WebhookPayload } from './types'; import type { EventStore } from './event-store'; import { logExceptInTest, errorExceptInTest } from './utils'; @@ -11,7 +6,6 @@ import * as Sentry from '@sentry/cloudflare'; const DELIVERY_STATE_KEY = 'deliveryState'; export class WebhookDelivery { - /** In-memory cache of delivery state */ private deliveryState: DeliveryState | null = null; /** Reentrancy guard - only needs to be in-memory since requests are cancelled on sleep */ private isFlushing = false; @@ -27,25 +21,16 @@ export class WebhookDelivery { private eventStore: EventStore ) {} - /** - * Initialize the webhook delivery system by loading delivery state. - */ async initialize(): Promise { await this.loadDeliveryState(); } - /** - * Load delivery state from durable storage. - * Returns a default DeliveryState if not found in storage. - * Caches the result in memory for subsequent access. - */ private async loadDeliveryState(): Promise { const stored = await this.storage.get(DELIVERY_STATE_KEY); if (stored) { this.deliveryState = stored; } else { - // Initialize with default values this.deliveryState = { nextAttemptAt: 0, attempt: 0, @@ -55,22 +40,14 @@ export class WebhookDelivery { return this.deliveryState; } - /** - * Save delivery state to durable storage. - * Persists the current in-memory delivery state. - */ private async saveDeliveryState(): Promise { await this.storage.put(DELIVERY_STATE_KEY, this.deliveryState); } - /** - * Get the current delivery state. - */ getDeliveryState(): DeliveryState | null { return this.deliveryState; } - // Internal: centralized configuration values parsed from env private getConfig() { return { BATCH_MAX_EVENTS: Number(this.env.BACKEND_WEBHOOK_BATCH_MAX_EVENTS) || 100, @@ -80,15 +57,11 @@ export class WebhookDelivery { }; } - // Internal: compute backoff delay private computeBackoffDelay(attempt: number, baseMs: number): number { const pow = attempt > 0 ? attempt - 1 : 0; return baseMs * Math.pow(2, pow); } - /** - * Schedule a flush - */ async scheduleFlush(): Promise { if (!this.deliveryState) { return; @@ -129,35 +102,26 @@ export class WebhookDelivery { await this.alarm.set(nextAlarm); } - /** - * Flush pending events to the backend webhook endpoint. - */ async flush(): Promise { if (!this.deliveryState) { return; } - // Guard against reentrancy if (this.isFlushing) { return; } - // Get webhook configuration const { BATCH_MAX_EVENTS, STOP_AFTER_ATTEMPTS } = this.getConfig(); - // Check if delivery has been permanently stopped if (this.deliveryState.attempt > STOP_AFTER_ATTEMPTS) { return; } try { - // Set reentrancy guard this.isFlushing = true; - // Get batch of events to send const eventsToSend = this.eventStore.getUnprocessedEvents(BATCH_MAX_EVENTS); - // If no events to send, return early if (eventsToSend.length === 0) { return; } @@ -165,10 +129,8 @@ export class WebhookDelivery { const lastDeliveredEventId = await this.sendEvents(eventsToSend); if (lastDeliveredEventId !== null) { - // Events were successfully delivered await this.eventStore.setLastProcessedId(lastDeliveredEventId); - // Reset backoff state this.deliveryState.attempt = 0; this.deliveryState.nextAttemptAt = 0; } else { @@ -181,7 +143,6 @@ export class WebhookDelivery { await this.saveDeliveryState(); await this.scheduleFlush(); } finally { - // Always clear reentrancy guard this.isFlushing = false; } } @@ -192,13 +153,11 @@ export class WebhookDelivery { return events[events.length - 1].id; } - // Build webhook payload const payload: WebhookPayload = { buildId: this.getBuildId(), events: events, }; - // Send webhook to backend try { const response = await fetch(this.env.BACKEND_EVENTS_URL, { method: 'POST',