use crate::cli::commands::{
CloudflareCmd, CloudflareContainersCmd, CloudflareContainersSubCommand, CloudflareDoctorCmd,
CloudflareEnvCmd, CloudflareEnvSubCommand, CloudflareInitCmd, CloudflareJobWorkflow,
CloudflareJobsCmd, CloudflareJobsSubCommand, CloudflareRollout, CloudflareSubCommand,
CloudflareTunnelCmd, CloudflareTunnelCreateCmd, CloudflareTunnelDeleteCmd,
CloudflareTunnelInfoCmd, CloudflareTunnelQuickStartCmd, CloudflareTunnelRunCmd,
CloudflareTunnelSubCommand, CloudflareWorkflowsCmd, CloudflareWorkflowsSubCommand,
};
use crate::commands::version::check_domain_for_cloudflare_release;
use crate::commands::workers::project::{
ensure_selected_workers_configured, parse_multiline_env_file, resolve_worker_targets,
write_worker_configs_to_project,
};
use crate::commands::workers::wrangler::{self, WranglerOutput};
use crate::strategies::{WorkerConfig, WorkerContainerConfig, XbpConfig};
use crate::utils::{find_xbp_config_upwards, parse_config_with_auto_heal};
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use serde_json::Value;
use std::collections::{HashMap, HashSet};
use std::env;
use std::fs;
use std::path::{Path, PathBuf};
use std::process::{Command, Stdio};
use uuid::Uuid;
const CLOUDFLARE_JOB_DIR: &str = ".xbp/jobs/cloudflare";
#[derive(Debug, Clone)]
struct ResolvedCloudflareApp {
project_root: PathBuf,
config_path: PathBuf,
worker_root: PathBuf,
app_name: String,
worker: WorkerConfig,
container: Option<WorkerContainerConfig>,
wrangler_config: PathBuf,
}
impl ResolvedCloudflareApp {
fn require_container(&self) -> Result<&WorkerContainerConfig, String> {
self.container.as_ref().ok_or_else(|| {
format!(
"Worker app `{}` does not define `workers[].container` in {}. Run `xbp cloudflare init --app {} --worker-root {} --container-port <port> --write` first (only required for Container-backed Workers).",
self.app_name,
self.config_path.display(),
self.app_name,
self.worker_root.display()
)
})
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
struct CloudflareWorkflowPayload {
app: Option<String>,
token: Option<String>,
account_id: Option<String>,
workflow: CloudflareJobWorkflow,
rollout: CloudflareRollout,
version: Option<String>,
domain: Option<String>,
dry_run: bool,
skip_deploy: bool,
allow_unchanged_container_image: bool,
prune_old_images: bool,
keep_image_tag_count: Option<usize>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[serde(rename_all = "kebab-case")]
enum CloudflareJobStatus {
Queued,
Running,
Succeeded,
Failed,
Cancelled,
}
impl CloudflareJobStatus {
fn as_str(&self) -> &'static str {
match self {
Self::Queued => "queued",
Self::Running => "running",
Self::Succeeded => "succeeded",
Self::Failed => "failed",
Self::Cancelled => "cancelled",
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
struct CloudflareJobRecord {
id: String,
status: CloudflareJobStatus,
payload: CloudflareWorkflowPayload,
created_at: DateTime<Utc>,
updated_at: DateTime<Utc>,
started_at: Option<DateTime<Utc>>,
finished_at: Option<DateTime<Utc>>,
error: Option<String>,
logs: Vec<String>,
}
pub async fn run_cloudflare(cmd: CloudflareCmd, _debug: bool) -> Result<(), String> {
match cmd.command {
CloudflareSubCommand::Doctor(doctor_cmd) => {
let app = resolve_cloudflare_app(cmd.root.as_deref(), cmd.app.as_deref())?;
run_doctor(&app, &doctor_cmd).await
}
CloudflareSubCommand::Init(init_cmd) => {
run_init_contract(cmd.root.as_deref(), cmd.app.as_deref(), init_cmd)
}
CloudflareSubCommand::Deploy(deploy_cmd) => {
let app = resolve_cloudflare_app(cmd.root.as_deref(), cmd.app.as_deref())?;
let payload = CloudflareWorkflowPayload {
app: Some(app.app_name.clone()),
token: cmd.token,
account_id: cmd.account_id,
workflow: CloudflareJobWorkflow::Deploy,
rollout: deploy_cmd.rollout,
version: None,
domain: None,
dry_run: deploy_cmd.dry_run,
skip_deploy: deploy_cmd.skip_deploy,
allow_unchanged_container_image: deploy_cmd.allow_unchanged_container_image,
prune_old_images: deploy_cmd.prune_old_images,
keep_image_tag_count: deploy_cmd.keep_image_tag_count,
};
run_deploy_workflow(&app, &payload).await
}
CloudflareSubCommand::Release(release_cmd) => {
let mut app = resolve_cloudflare_app(cmd.root.as_deref(), cmd.app.as_deref())?;
let payload = CloudflareWorkflowPayload {
app: Some(app.app_name.clone()),
token: cmd.token,
account_id: cmd.account_id,
workflow: CloudflareJobWorkflow::Release,
rollout: release_cmd.rollout,
version: Some(release_cmd.version),
domain: release_cmd.domain,
dry_run: false,
skip_deploy: release_cmd.skip_deploy,
allow_unchanged_container_image: release_cmd.allow_unchanged_container_image,
prune_old_images: release_cmd.prune_old_images,
keep_image_tag_count: release_cmd.keep_image_tag_count,
};
run_release_workflow(&mut app, &payload).await
}
CloudflareSubCommand::Containers(containers_cmd) => {
let app = resolve_cloudflare_app(cmd.root.as_deref(), cmd.app.as_deref())?;
run_containers(&app, containers_cmd)
}
CloudflareSubCommand::Tunnel(tunnel_cmd) => {
let project_root = resolve_project_root(cmd.root.as_deref())?;
run_tunnel(&project_root, tunnel_cmd)
}
CloudflareSubCommand::Workflows(workflows_cmd) => {
let project_root = resolve_project_root(cmd.root.as_deref())?;
run_workflows(&project_root, workflows_cmd)
}
CloudflareSubCommand::Env(env_cmd) => {
let worker_root = resolve_worker_root(cmd.root.as_deref(), cmd.app.as_deref())?;
run_cloudflare_env(&worker_root, env_cmd)
}
CloudflareSubCommand::Jobs(jobs_cmd) => {
run_jobs(cmd.root.as_deref(), cmd.app, jobs_cmd).await
}
}
}
fn run_tunnel(project_root: &Path, cmd: CloudflareTunnelCmd) -> Result<(), String> {
let args = match cmd.command {
CloudflareTunnelSubCommand::Create(CloudflareTunnelCreateCmd { name }) => {
vec!["tunnel".to_string(), "create".to_string(), name]
}
CloudflareTunnelSubCommand::Delete(CloudflareTunnelDeleteCmd { tunnel, force }) => {
let mut args = vec!["tunnel".to_string(), "delete".to_string(), tunnel];
if force {
args.push("--force".to_string());
}
args
}
CloudflareTunnelSubCommand::Info(CloudflareTunnelInfoCmd { tunnel }) => {
vec!["tunnel".to_string(), "info".to_string(), tunnel]
}
CloudflareTunnelSubCommand::List(_) => {
vec!["tunnel".to_string(), "list".to_string()]
}
CloudflareTunnelSubCommand::Run(CloudflareTunnelRunCmd {
tunnel,
token,
log_level,
}) => {
let mut args = vec!["tunnel".to_string(), "run".to_string()];
if let Some(tunnel) = tunnel {
args.push(tunnel);
}
if let Some(token) = token {
args.push("--token".to_string());
args.push(token);
}
if let Some(log_level) = log_level {
args.push("--log-level".to_string());
args.push(log_level);
}
args
}
CloudflareTunnelSubCommand::QuickStart(CloudflareTunnelQuickStartCmd { url }) => {
vec!["tunnel".to_string(), "quick-start".to_string(), url]
}
};
wrangler::run_wrangler(project_root, &args)
}
fn run_workflows(project_root: &Path, cmd: CloudflareWorkflowsCmd) -> Result<(), String> {
let CloudflareWorkflowsSubCommand::External(mut workflow_args) = cmd.command;
workflow_args.insert(0, "workflows".to_string());
wrangler::run_wrangler(project_root, &workflow_args)
}
fn run_cloudflare_env(project_root: &Path, cmd: CloudflareEnvCmd) -> Result<(), String> {
let args = match cmd.command {
CloudflareEnvSubCommand::External(args) => args,
};
wrangler::run_wrangler(project_root, &args)
}
async fn run_doctor(app: &ResolvedCloudflareApp, cmd: &CloudflareDoctorCmd) -> Result<(), String> {
println!("Cloudflare app: {}", app.app_name);
println!("Project root: {}", app.project_root.display());
println!("XBP config: {}", app.config_path.display());
println!("Worker root: {}", app.worker_root.display());
println!("Wrangler config: {}", app.wrangler_config.display());
if let Some(script) = app.worker.script_name.as_deref() {
println!("Script name: {script}");
}
match app.container.as_ref() {
Some(_) => println!("Container contract: configured"),
None => println!("Container contract: none (plain Worker)"),
}
let diagnostics = collect_local_diagnostics(app)?;
if !diagnostics.warnings.is_empty() {
for warning in &diagnostics.warnings {
println!("warning: {warning}");
}
}
if !diagnostics.missing_required.is_empty() {
return Err(format!(
"Cloudflare worker contract is incomplete:\n{}",
diagnostics
.missing_required
.iter()
.map(|item| format!("- {item}"))
.collect::<Vec<_>>()
.join("\n")
));
}
if app.container.is_some() {
println!("Local container contract: ok");
} else {
println!("Local worker contract: ok");
}
if cmd.offline {
println!("Wrangler checks skipped (--offline).");
return Ok(());
}
println!(
"{}",
wrangler::describe_wrangler_cloudflare_auth(&app.worker_root)
);
run_wrangler_checked(app, vec!["--version".to_string()], "wrangler --version")?;
run_wrangler_checked(app, vec!["whoami".to_string()], "wrangler whoami")?;
run_wrangler_checked(app, wrangler_check_args(app), "wrangler check")?;
let rollout = app
.container
.as_ref()
.and_then(|container| container.default_rollout.as_deref());
run_wrangler_checked(
app,
wrangler_deploy_args(app, rollout, true),
"wrangler deploy --dry-run",
)?;
println!("Cloudflare doctor: ok");
Ok(())
}
fn run_init_contract(
root_override: Option<&Path>,
app_override: Option<&str>,
cmd: CloudflareInitCmd,
) -> Result<(), String> {
let current_dir =
env::current_dir().map_err(|error| format!("Failed to read current directory: {error}"))?;
let start = root_override.unwrap_or(current_dir.as_path());
let found = find_xbp_config_upwards(start).ok_or_else(|| {
format!(
"Could not find .xbp/xbp.yaml from {}. Run `xbp setup` or pass `--root`.",
start.display()
)
})?;
let content = fs::read_to_string(&found.config_path)
.map_err(|error| format!("Failed to read {}: {error}", found.config_path.display()))?;
let (mut config, healed_content) =
parse_config_with_auto_heal::<XbpConfig>(&content, found.kind)
.map_err(|error| format!("Failed to parse {}: {error}", found.config_path.display()))?;
if let Some(healed_content) = healed_content {
fs::write(&found.config_path, healed_content)
.map_err(|error| format!("Failed to write {}: {error}", found.config_path.display()))?;
}
let app_name = app_override
.map(str::trim)
.filter(|value| !value.is_empty())
.map(ToOwned::to_owned)
.or_else(|| {
cmd.worker_root
.file_name()
.and_then(|value| value.to_str())
.map(ToOwned::to_owned)
})
.ok_or_else(|| "Pass --app or a worker root with a final path component.".to_string())?;
let worker_root = if cmd.worker_root.is_absolute() {
cmd.worker_root.clone()
} else {
found.project_root.join(&cmd.worker_root)
};
let worker_root_string = worker_root.to_string_lossy().to_string();
let workers = config.workers.get_or_insert_with(Vec::new);
let index = workers
.iter()
.position(|worker| {
worker.name == app_name
|| worker.root == worker_root_string
|| found.project_root.join(&worker.root) == worker_root
})
.unwrap_or_else(|| {
workers.push(WorkerConfig {
name: app_name.clone(),
root: worker_root_string.clone(),
script_name: None,
service: None,
deploy: None,
container: None,
containers: None,
durable_objects: None,
});
workers.len() - 1
});
let dockerfile = cmd.dockerfile.to_string_lossy().replace('\\', "/");
let container = WorkerContainerConfig {
class_name: cmd.class_name,
binding: cmd.binding,
dockerfile: Some(dockerfile),
port: Some(cmd.container_port),
application_id: cmd.application_id,
healthcheck_path: Some(normalize_health_path(&cmd.health_path)),
required_secrets: cmd.required_secrets,
version_var: cmd.version_var,
instance_name_template: None,
default_rollout: Some("immediate".to_string()),
health_urls: Vec::new(),
expected_image_name: None,
image_tag_keep_count: None,
allow_unchanged_container_image: false,
image: None,
build_context: None,
instance_type: None,
max_instances: None,
rollout_active_grace_period: None,
rollout_step_percentage: Vec::new(),
regions: Vec::new(),
};
workers[index].root = worker_root_string;
workers[index].container = Some(container);
if !cmd.write {
println!(
"{}",
serde_yaml::to_string(&workers[index]).map_err(|error| error.to_string())?
);
println!("Run again with --write to persist this worker container contract.");
return Ok(());
}
write_worker_configs_to_project(&found.project_root, &found.config_path, &mut config)?;
println!(
"Wrote Cloudflare container contract for `{}` to {}",
app_name,
found.config_path.display()
);
Ok(())
}
async fn run_deploy_workflow(
app: &ResolvedCloudflareApp,
payload: &CloudflareWorkflowPayload,
) -> Result<(), String> {
let _ = app.require_container()?;
validate_payload_for_workflow(payload)?;
run_doctor(app, &CloudflareDoctorCmd { offline: false }).await?;
regenerate_worker_types(app)?;
run_project_build_if_present(app)?;
let before_info = if payload.dry_run {
None
} else {
read_container_app_info(app)?
};
if payload.skip_deploy {
println!("wrangler deploy: skipped (--skip-deploy)");
} else {
let args = wrangler_deploy_args(app, Some(payload.rollout.as_str()), payload.dry_run);
run_wrangler_streaming(app, args, "wrangler deploy")?;
}
if !payload.dry_run {
verify_container_state(
app,
ContainerVerificationContext {
before_info,
expected_version: payload.version.clone(),
skip_deploy: payload.skip_deploy,
allow_unchanged_image: payload.allow_unchanged_container_image,
prune_old_images: payload.prune_old_images,
keep_image_tag_count: payload.keep_image_tag_count,
},
)
.await?;
}
Ok(())
}
async fn run_release_workflow(
app: &mut ResolvedCloudflareApp,
payload: &CloudflareWorkflowPayload,
) -> Result<(), String> {
let _ = app.require_container()?;
validate_payload_for_workflow(payload)?;
let version = payload
.version
.as_deref()
.ok_or_else(|| "Release jobs require --version.".to_string())?;
if let Some(domain) = payload.domain.as_deref() {
check_domain_for_cloudflare_release(
Some(&app.project_root),
domain,
version,
&app.app_name,
)?;
}
sync_version_var(app, version)?;
run_deploy_workflow(app, payload).await?;
Ok(())
}
fn run_containers(app: &ResolvedCloudflareApp, cmd: CloudflareContainersCmd) -> Result<(), String> {
let _ = app.require_container()?;
let args = match cmd.command {
CloudflareContainersSubCommand::List(list_cmd) => {
let mut args = vec!["containers".to_string(), "list".to_string()];
if list_cmd.json {
args.push("--json".to_string());
}
args
}
CloudflareContainersSubCommand::Info(info_cmd) => {
let mut args = vec![
"containers".to_string(),
"info".to_string(),
resolve_application_id(app, info_cmd.application_id.as_deref())?,
];
if info_cmd.json {
args.push("--json".to_string());
}
args
}
CloudflareContainersSubCommand::Instances(instances_cmd) => {
let mut args = vec![
"containers".to_string(),
"instances".to_string(),
resolve_application_id(app, instances_cmd.application_id.as_deref())?,
];
if instances_cmd.json {
args.push("--json".to_string());
}
args
}
CloudflareContainersSubCommand::Push(push_cmd) => {
vec!["containers".to_string(), "push".to_string(), push_cmd.image]
}
CloudflareContainersSubCommand::Ssh(ssh_cmd) => {
let mut args = vec![
"containers".to_string(),
"ssh".to_string(),
ssh_cmd.instance_id,
];
args.extend(ssh_cmd.args);
args
}
};
wrangler::run_wrangler(&app.worker_root, &args)
}
async fn run_jobs(
root_override: Option<&Path>,
app_override: Option<String>,
cmd: CloudflareJobsCmd,
) -> Result<(), String> {
match cmd.command {
CloudflareJobsSubCommand::Enqueue(enqueue_cmd) => {
let project_root = resolve_project_root(root_override)?;
let payload = CloudflareWorkflowPayload {
app: app_override,
token: None,
account_id: None,
workflow: enqueue_cmd.workflow,
rollout: enqueue_cmd.rollout,
version: enqueue_cmd.version,
domain: None,
dry_run: enqueue_cmd.dry_run,
skip_deploy: enqueue_cmd.skip_deploy,
allow_unchanged_container_image: enqueue_cmd.allow_unchanged_container_image,
prune_old_images: enqueue_cmd.prune_old_images,
keep_image_tag_count: enqueue_cmd.keep_image_tag_count,
};
validate_payload_for_workflow(&payload)?;
let job = enqueue_job(&project_root, payload)?;
println!(
"{}",
serde_json::to_string_pretty(&job).map_err(|error| error.to_string())?
);
Ok(())
}
CloudflareJobsSubCommand::Run(run_cmd) => {
let project_root = resolve_project_root(root_override)?;
run_queued_jobs(&project_root, run_cmd.once).await
}
CloudflareJobsSubCommand::Status(status_cmd) => {
let project_root = resolve_project_root(root_override)?;
let job = read_job(&project_root, &status_cmd.job_id)?;
println!(
"{}",
serde_json::to_string_pretty(&job).map_err(|error| error.to_string())?
);
Ok(())
}
CloudflareJobsSubCommand::List(_) => {
let project_root = resolve_project_root(root_override)?;
for job in list_jobs(&project_root)? {
println!(
"{}\t{}\t{}\t{}",
job.id,
job.status.as_str(),
job.payload.workflow.as_str(),
job.updated_at
);
}
Ok(())
}
CloudflareJobsSubCommand::Logs(logs_cmd) => {
let project_root = resolve_project_root(root_override)?;
let job = read_job(&project_root, &logs_cmd.job_id)?;
for line in job.logs {
println!("{line}");
}
Ok(())
}
}
}
async fn run_queued_jobs(project_root: &Path, once: bool) -> Result<(), String> {
let mut jobs = list_jobs(project_root)?
.into_iter()
.filter(|job| job.status == CloudflareJobStatus::Queued)
.collect::<Vec<_>>();
jobs.sort_by_key(|job| job.created_at);
if jobs.is_empty() {
println!("No queued Cloudflare jobs.");
return Ok(());
}
for mut job in jobs {
job.status = CloudflareJobStatus::Running;
job.started_at = Some(Utc::now());
job.updated_at = Utc::now();
job.logs.push("job started".to_string());
write_job(project_root, &job)?;
let result = run_job_payload(project_root, &job.payload).await;
job.updated_at = Utc::now();
job.finished_at = Some(Utc::now());
match result {
Ok(()) => {
job.status = CloudflareJobStatus::Succeeded;
job.logs.push("job succeeded".to_string());
}
Err(error) => {
job.status = CloudflareJobStatus::Failed;
job.error = Some(error.clone());
job.logs.push(format!("job failed: {error}"));
}
}
write_job(project_root, &job)?;
if once {
break;
}
}
Ok(())
}
async fn run_job_payload(
project_root: &Path,
payload: &CloudflareWorkflowPayload,
) -> Result<(), String> {
let app = resolve_cloudflare_app(Some(project_root), payload.app.as_deref())?;
match payload.workflow {
CloudflareJobWorkflow::Doctor => {
run_doctor(&app, &CloudflareDoctorCmd { offline: false }).await
}
CloudflareJobWorkflow::Deploy => run_deploy_workflow(&app, payload).await,
CloudflareJobWorkflow::Release => {
let mut app = app;
run_release_workflow(&mut app, payload).await
}
}
}
fn resolve_cloudflare_app(
root_override: Option<&Path>,
app_override: Option<&str>,
) -> Result<ResolvedCloudflareApp, String> {
let current_dir =
env::current_dir().map_err(|error| format!("Failed to read current directory: {error}"))?;
let mut resolution = resolve_worker_targets(¤t_dir, root_override, app_override, false)?;
let inserted = ensure_selected_workers_configured(&mut resolution)?;
if inserted > 0 {
println!(
"Auto-configured {inserted} Worker entr{} in {} from on-disk Wrangler project(s).",
if inserted == 1 { "y" } else { "ies" },
resolution.config_path.display()
);
}
let selected = resolution
.selected
.first()
.ok_or_else(|| "No Worker app was selected.".to_string())?
.clone();
let mut worker = selected.config.clone().ok_or_else(|| {
format!(
"Worker app `{}` was discovered on disk but could not be auto-configured in {}.",
selected.label,
resolution.config_path.display()
)
})?;
let wrangler_config = resolve_wrangler_config_file(&selected.root)?;
if worker.container.is_none() {
if let Some(inferred) = infer_container_contract_from_wrangler(&wrangler_config) {
worker.container = Some(inferred.clone());
if let Some(workers) = resolution.config.workers.as_mut() {
if let Some(entry) = workers.iter_mut().find(|entry| {
entry.name == worker.name
|| entry.root == worker.root
|| entry.script_name == worker.script_name
}) {
entry.container = Some(inferred);
write_worker_configs_to_project(
&resolution.project_root,
&resolution.config_path,
&mut resolution.config,
)?;
println!(
"Inferred container contract for `{}` from Wrangler and wrote it to {}.",
worker.name,
resolution.config_path.display()
);
}
}
}
}
Ok(ResolvedCloudflareApp {
project_root: resolution.project_root,
config_path: resolution.config_path,
worker_root: selected.root.clone(),
app_name: selected.label.clone(),
container: worker.container.clone(),
worker,
wrangler_config,
})
}
fn infer_container_contract_from_wrangler(wrangler_config: &Path) -> Option<WorkerContainerConfig> {
let config = load_wrangler_config_value(wrangler_config).ok()?;
let containers = config.get("containers")?.as_array()?;
let first = containers.first()?;
let class_name = first
.get("class_name")
.and_then(Value::as_str)
.map(str::to_string)?;
let binding = config
.get("durable_objects")
.and_then(|value| value.get("bindings"))
.and_then(Value::as_array)
.and_then(|bindings| {
bindings.iter().find_map(|binding| {
let binding_class = binding.get("class_name").and_then(Value::as_str)?;
if binding_class == class_name {
binding
.get("name")
.and_then(Value::as_str)
.map(str::to_string)
} else {
None
}
})
});
let port = first
.get("image_vars")
.and_then(|vars| vars.get("PORT"))
.and_then(|port| {
port.as_u64()
.or_else(|| port.as_str().and_then(|text| text.parse::<u64>().ok()))
})
.and_then(|port| u16::try_from(port).ok());
let dockerfile = first
.get("image")
.and_then(Value::as_str)
.filter(|image| image.starts_with("./") || image.starts_with("../") || !image.contains('/'))
.map(str::to_string)
.or_else(|| Some("Dockerfile".to_string()));
Some(WorkerContainerConfig {
class_name: Some(class_name),
binding,
dockerfile,
port,
application_id: None,
healthcheck_path: Some("/health".to_string()),
required_secrets: Vec::new(),
version_var: None,
instance_name_template: None,
default_rollout: Some("immediate".to_string()),
health_urls: Vec::new(),
expected_image_name: None,
image_tag_keep_count: None,
allow_unchanged_container_image: false,
image: None,
build_context: None,
instance_type: None,
max_instances: None,
rollout_active_grace_period: None,
rollout_step_percentage: Vec::new(),
regions: Vec::new(),
})
}
fn resolve_project_root(root_override: Option<&Path>) -> Result<PathBuf, String> {
let current_dir =
env::current_dir().map_err(|error| format!("Failed to read current directory: {error}"))?;
let start = root_override.unwrap_or(current_dir.as_path());
find_xbp_config_upwards(start)
.map(|found| found.project_root)
.ok_or_else(|| format!("Could not find .xbp/xbp.yaml from {}.", start.display()))
}
fn resolve_worker_root(
root_override: Option<&Path>,
app_override: Option<&str>,
) -> Result<PathBuf, String> {
let current_dir =
env::current_dir().map_err(|error| format!("Failed to read current directory: {error}"))?;
resolve_worker_targets(¤t_dir, root_override, app_override, false)?
.selected
.first()
.map(|app| app.root.clone())
.ok_or_else(|| "No Worker app was selected.".to_string())
}
fn resolve_wrangler_config_file(worker_root: &Path) -> Result<PathBuf, String> {
for name in ["wrangler.jsonc", "wrangler.json", "wrangler.toml"] {
let path = worker_root.join(name);
if path.exists() {
return Ok(path);
}
}
Err(format!(
"No Wrangler config found under {}. Expected wrangler.jsonc, wrangler.json, or wrangler.toml.",
worker_root.display()
))
}
#[derive(Debug, Default)]
struct LocalDiagnostics {
missing_required: Vec<String>,
warnings: Vec<String>,
}
fn collect_local_diagnostics(app: &ResolvedCloudflareApp) -> Result<LocalDiagnostics, String> {
let mut diagnostics = LocalDiagnostics::default();
let local_env = load_worker_env(&app.worker_root)?;
let wrangler_config = load_wrangler_config_value(&app.wrangler_config)?;
validate_required_wrangler_secrets(&wrangler_config, &local_env, &mut diagnostics);
audit_wrangler_vars_for_secrets(&wrangler_config, &mut diagnostics);
let Some(container) = app.container.as_ref() else {
return Ok(diagnostics);
};
if container.class_name.as_deref().is_none_or(str::is_empty) {
diagnostics
.missing_required
.push("workers[].container.class_name is required".to_string());
}
if container.binding.as_deref().is_none_or(str::is_empty) {
diagnostics
.missing_required
.push("workers[].container.binding is required".to_string());
}
if container.port.is_none() {
diagnostics
.missing_required
.push("workers[].container.port is required".to_string());
}
let dockerfile = resolve_dockerfile_path(app);
if !dockerfile.exists() {
diagnostics
.missing_required
.push(format!("Dockerfile not found: {}", dockerfile.display()));
}
for secret_name in &container.required_secrets {
let present = local_env
.get(secret_name)
.is_some_and(|value| !value.trim().is_empty())
|| env::var(secret_name).is_ok_and(|value| !value.trim().is_empty());
if !present {
diagnostics.missing_required.push(format!(
"Required secret `{secret_name}` is missing locally"
));
}
}
validate_wrangler_container_shape(app, &wrangler_config, &mut diagnostics);
validate_port_alignment(app, &wrangler_config, &mut diagnostics);
Ok(diagnostics)
}
fn validate_required_wrangler_secrets(
config: &Value,
local_env: &HashMap<String, String>,
diagnostics: &mut LocalDiagnostics,
) {
let Some(required) = config
.get("secrets")
.and_then(|secrets| secrets.get("required"))
.and_then(Value::as_array)
else {
return;
};
for secret_name in required.iter().filter_map(Value::as_str) {
let present = local_env
.get(secret_name)
.is_some_and(|value| !value.trim().is_empty())
|| env::var(secret_name).is_ok_and(|value| !value.trim().is_empty());
if !present {
diagnostics.missing_required.push(format!(
"Wrangler required secret `{secret_name}` is missing locally"
));
}
}
}
fn audit_wrangler_vars_for_secrets(config: &Value, diagnostics: &mut LocalDiagnostics) {
let Some(vars) = config.get("vars").and_then(Value::as_object) else {
return;
};
for key in vars.keys() {
if secret_like_key(key) {
diagnostics.warnings.push(format!(
"wrangler `vars` contains secret-looking key `{key}`. Move real secrets to `wrangler secret put {key}` and keep only documented dummy/local values in .env.local/.env.example."
));
}
}
}
fn secret_like_key(key: &str) -> bool {
let key = key.to_ascii_uppercase();
key.ends_with("_SECRET")
|| key.ends_with("_TOKEN")
|| key.ends_with("_PASSWORD")
|| key.contains("PRIVATE_KEY")
}
fn validate_wrangler_container_shape(
app: &ResolvedCloudflareApp,
config: &Value,
diagnostics: &mut LocalDiagnostics,
) {
let Some(container) = app.container.as_ref() else {
return;
};
let class_name = container.class_name.as_deref().unwrap_or_default();
let binding = container.binding.as_deref().unwrap_or_default();
let has_container = config
.get("containers")
.and_then(Value::as_array)
.is_some_and(|containers| {
containers
.iter()
.any(|item| item.get("class_name").and_then(Value::as_str) == Some(class_name))
});
if !class_name.is_empty() && !has_container {
diagnostics.missing_required.push(format!(
"wrangler config `containers[]` does not include class_name `{class_name}`"
));
}
let has_do_binding = config
.get("durable_objects")
.and_then(|value| value.get("bindings"))
.and_then(Value::as_array)
.is_some_and(|bindings| {
bindings.iter().any(|item| {
item.get("name").and_then(Value::as_str) == Some(binding)
&& item.get("class_name").and_then(Value::as_str) == Some(class_name)
})
});
if !binding.is_empty() && !class_name.is_empty() && !has_do_binding {
diagnostics.missing_required.push(format!(
"wrangler config `durable_objects.bindings[]` does not map `{binding}` to `{class_name}`"
));
}
}
fn validate_port_alignment(
app: &ResolvedCloudflareApp,
config: &Value,
diagnostics: &mut LocalDiagnostics,
) {
let Some(app_container) = app.container.as_ref() else {
return;
};
let Some(expected_port) = app_container.port else {
return;
};
let Some(containers) = config.get("containers").and_then(Value::as_array) else {
return;
};
for container in containers {
if container.get("class_name").and_then(Value::as_str)
!= app_container.class_name.as_deref()
{
continue;
}
let configured_port = container
.get("image_vars")
.and_then(|value| value.get("PORT"))
.and_then(|value| {
value
.as_str()
.and_then(|text| text.parse::<u16>().ok())
.or_else(|| value.as_u64().and_then(|value| u16::try_from(value).ok()))
})
.or_else(|| {
container
.get("port")
.and_then(Value::as_u64)
.and_then(|value| u16::try_from(value).ok())
});
if let Some(configured_port) = configured_port {
if configured_port != expected_port {
diagnostics.warnings.push(format!(
"container port mismatch: xbp.yaml has {expected_port}, wrangler config has {configured_port}"
));
}
}
}
}
fn load_worker_env(worker_root: &Path) -> Result<HashMap<String, String>, String> {
let mut vars = HashMap::new();
for name in [".env", ".env.local", ".dev.vars"] {
let path = worker_root.join(name);
if !path.exists() {
continue;
}
vars.extend(parse_multiline_env_file(&path)?);
}
Ok(vars)
}
fn load_wrangler_config_value(path: &Path) -> Result<Value, String> {
let content = fs::read_to_string(path)
.map_err(|error| format!("Failed to read {}: {error}", path.display()))?;
match path.extension().and_then(|value| value.to_str()) {
Some("toml") => {
let value: toml::Value = toml::from_str(&content)
.map_err(|error| format!("Failed to parse {}: {error}", path.display()))?;
serde_json::to_value(value)
.map_err(|error| format!("Failed to normalize {}: {error}", path.display()))
}
_ => {
let stripped = strip_json_comments(&content);
let normalized = strip_jsonc_trailing_commas(&stripped);
serde_json::from_str(&normalized)
.map_err(|error| format!("Failed to parse {} as JSONC: {error}", path.display()))
}
}
}
fn strip_json_comments(input: &str) -> String {
let mut output = String::with_capacity(input.len());
let mut chars = input.chars().peekable();
let mut in_string = false;
let mut escaped = false;
while let Some(ch) = chars.next() {
if in_string {
output.push(ch);
if escaped {
escaped = false;
} else if ch == '\\' {
escaped = true;
} else if ch == '"' {
in_string = false;
}
continue;
}
if ch == '"' {
in_string = true;
output.push(ch);
continue;
}
if ch == '/' {
match chars.peek().copied() {
Some('/') => {
let _ = chars.next();
for next in chars.by_ref() {
if next == '\n' {
output.push('\n');
break;
}
}
continue;
}
Some('*') => {
let _ = chars.next();
let mut last = '\0';
for next in chars.by_ref() {
if last == '*' && next == '/' {
break;
}
last = next;
}
continue;
}
_ => {}
}
}
output.push(ch);
}
output
}
fn strip_jsonc_trailing_commas(input: &str) -> String {
let mut output = String::with_capacity(input.len());
let chars = input.chars().collect::<Vec<_>>();
let mut in_string = false;
let mut escaped = false;
for (index, ch) in chars.iter().copied().enumerate() {
if in_string {
output.push(ch);
if escaped {
escaped = false;
} else if ch == '\\' {
escaped = true;
} else if ch == '"' {
in_string = false;
}
continue;
}
if ch == '"' {
in_string = true;
output.push(ch);
continue;
}
if ch == ',' {
let next = chars[index + 1..]
.iter()
.copied()
.find(|next| !next.is_whitespace());
if matches!(next, Some('}' | ']')) {
continue;
}
}
output.push(ch);
}
output
}
fn wrangler_config_args(app: &ResolvedCloudflareApp) -> Vec<String> {
vec![
"--config".to_string(),
wrangler::render_path_arg(&app.worker_root, &app.wrangler_config),
]
}
fn wrangler_check_args(app: &ResolvedCloudflareApp) -> Vec<String> {
let mut args = vec!["check".to_string()];
args.extend(wrangler_config_args(app));
args
}
fn wrangler_deploy_args(
app: &ResolvedCloudflareApp,
rollout: Option<&str>,
dry_run: bool,
) -> Vec<String> {
let mut args = vec!["deploy".to_string()];
args.extend(wrangler_config_args(app));
if dry_run {
args.push("--dry-run".to_string());
} else if let Some(rollout) = rollout.filter(|value| !value.trim().is_empty()) {
args.push("--containers-rollout".to_string());
args.push(rollout.to_string());
}
args
}
fn wrangler_types_args(app: &ResolvedCloudflareApp) -> Vec<String> {
let output = if app.worker_root.join("worker").is_dir() {
app.worker_root
.join("worker")
.join("worker-configuration.d.ts")
} else {
app.worker_root.join("worker-configuration.d.ts")
};
let mut args = vec![
"types".to_string(),
wrangler::render_path_arg(&app.worker_root, &output),
];
args.extend(wrangler_config_args(app));
args
}
fn run_wrangler_checked(
app: &ResolvedCloudflareApp,
args: Vec<String>,
label: &str,
) -> Result<WranglerOutput, String> {
let output = wrangler::run_wrangler_capture(&app.worker_root, &args)?;
if output.exit_code != 0 {
let detail = wrangler::format_wrangler_failure(&args, &output);
return Err(format!("{label} failed:\n{}", mask_secrets(&detail)));
}
println!("{label}: ok");
Ok(output)
}
fn run_wrangler_streaming(
app: &ResolvedCloudflareApp,
args: Vec<String>,
label: &str,
) -> Result<(), String> {
println!("{label}: running");
wrangler::run_wrangler_stream(&app.worker_root, &args, |line| {
println!("{}", mask_secrets(line));
Ok(())
})
}
fn regenerate_worker_types(app: &ResolvedCloudflareApp) -> Result<(), String> {
run_wrangler_streaming(app, wrangler_types_args(app), "wrangler types")
}
fn run_project_build_if_present(app: &ResolvedCloudflareApp) -> Result<(), String> {
let package_json = app.worker_root.join("package.json");
if !package_json.exists() {
return Ok(());
}
let content = fs::read_to_string(&package_json)
.map_err(|error| format!("Failed to read {}: {error}", package_json.display()))?;
let value: Value = serde_json::from_str(&content)
.map_err(|error| format!("Failed to parse {}: {error}", package_json.display()))?;
let scripts = value.get("scripts").and_then(Value::as_object);
let script = scripts.and_then(|scripts| {
if scripts.contains_key("build:worker") {
Some("build:worker")
} else if scripts.contains_key("build") {
Some("build")
} else {
None
}
});
let Some(script) = script else {
return Ok(());
};
let pnpm_bin = if cfg!(windows) { "pnpm.cmd" } else { "pnpm" };
let status = Command::new(pnpm_bin)
.args(["run", script])
.current_dir(&app.worker_root)
.stdout(Stdio::inherit())
.stderr(Stdio::inherit())
.status()
.map_err(|error| format!("Failed to run pnpm run {script}: {error}"))?;
if !status.success() {
return Err(format!("pnpm run {script} failed with status {status}."));
}
Ok(())
}
#[derive(Debug, Clone)]
struct ContainerVerificationContext {
before_info: Option<ContainerAppInfo>,
expected_version: Option<String>,
skip_deploy: bool,
allow_unchanged_image: bool,
prune_old_images: bool,
keep_image_tag_count: Option<usize>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct ContainerAppInfo {
image: Option<String>,
image_name: Option<String>,
image_tag: Option<String>,
app_version: Option<i64>,
instance_type: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct ContainerInstance {
name: String,
state: String,
version: Option<i64>,
}
async fn verify_container_state(
app: &ResolvedCloudflareApp,
context: ContainerVerificationContext,
) -> Result<(), String> {
let after_info = read_container_app_info(app)?;
if let Some(after_info) = after_info.as_ref() {
verify_container_app_info(app, context.before_info.as_ref(), after_info, &context)?;
let instances = read_container_instances(app)?;
verify_container_instances(app, after_info, &instances, &context.expected_version)?;
report_or_prune_old_images(app, after_info, &context)?;
} else {
println!("Container app verification skipped: no application_id configured.");
}
verify_health(app, context.expected_version.as_deref()).await?;
Ok(())
}
fn read_container_app_info(
app: &ResolvedCloudflareApp,
) -> Result<Option<ContainerAppInfo>, String> {
let Some(application_id) = app
.container
.as_ref()
.and_then(|container| container.application_id.as_deref())
else {
return Ok(None);
};
let mut args = vec![
"containers".to_string(),
"info".to_string(),
application_id.to_string(),
];
args.extend(wrangler_config_args(app));
let output = run_wrangler_checked(app, args, "wrangler containers info")?;
let value: Value = serde_json::from_str(output.stdout.trim()).map_err(|error| {
format!(
"wrangler containers info did not return JSON that XBP can verify: {error}. Upgrade Wrangler or run `wrangler containers info {application_id}` and confirm it emits JSON."
)
})?;
Ok(Some(parse_container_app_info(&value)))
}
fn read_container_instances(app: &ResolvedCloudflareApp) -> Result<Vec<ContainerInstance>, String> {
let application_id = resolve_application_id(app, None)?;
let mut args = vec![
"containers".to_string(),
"instances".to_string(),
application_id,
];
args.extend(wrangler_config_args(app));
args.extend([
"--per-page".to_string(),
"100".to_string(),
"--json".to_string(),
]);
let output = run_wrangler_checked(app, args, "wrangler containers instances --json")?;
let value: Value = serde_json::from_str(output.stdout.trim()).map_err(|error| {
format!("wrangler containers instances --json returned invalid JSON: {error}")
})?;
parse_container_instances(&value)
}
fn verify_container_app_info(
app: &ResolvedCloudflareApp,
before_info: Option<&ContainerAppInfo>,
after_info: &ContainerAppInfo,
context: &ContainerVerificationContext,
) -> Result<(), String> {
let container = app.require_container()?;
if let Some(expected_image_name) = container.expected_image_name.as_deref() {
let actual_matches = after_info.image_name.as_deref() == Some(expected_image_name)
|| after_info.image.as_deref().is_some_and(|image| {
image == expected_image_name
|| image.starts_with(&format!("{expected_image_name}:"))
|| image.contains(&format!("/{expected_image_name}:"))
});
if !actual_matches {
return Err(format!(
"Container image mismatch: expected image name `{expected_image_name}`, observed `{}`.",
after_info.image.as_deref().unwrap_or("<missing>")
));
}
}
let allow_unchanged =
context.allow_unchanged_image || container.allow_unchanged_container_image;
if !context.skip_deploy && !allow_unchanged {
if let (Some(before), Some(after)) = (
before_info.and_then(|info| info.image.as_deref()),
after_info.image.as_deref(),
) {
if before == after {
return Err(format!(
"Container image did not change after deploy (`{after}`). Pass --allow-unchanged-container-image for deliberate config-only deploys."
));
}
}
if let (Some(before), Some(after)) = (
before_info.and_then(|info| info.app_version),
after_info.app_version,
) {
if after <= before {
return Err(format!(
"Container app version did not advance after deploy (before {before}, after {after})."
));
}
}
}
if let Some(expected_instance_type) = expected_instance_type(app)? {
let Some(actual_instance_type) = after_info.instance_type.as_deref() else {
return Err(format!(
"Could not verify Cloudflare container instance type; expected `{expected_instance_type}` but containers info omitted the type."
));
};
if actual_instance_type != expected_instance_type {
return Err(format!(
"Container instance type mismatch: expected `{expected_instance_type}`, observed `{actual_instance_type}`."
));
}
}
println!("Container app verification: ok");
Ok(())
}
fn verify_container_instances(
app: &ResolvedCloudflareApp,
info: &ContainerAppInfo,
instances: &[ContainerInstance],
expected_version: &Option<String>,
) -> Result<(), String> {
let container = app.require_container()?;
let expected_name = container
.instance_name_template
.as_deref()
.zip(expected_version.as_deref())
.map(|(template, version)| render_instance_name_template(template, version));
if let Some(expected_name) = expected_name.as_deref() {
let has_expected = instances
.iter()
.any(|instance| instance.name == expected_name && instance_is_active(instance));
if !has_expected {
return Err(format!(
"Expected active container instance `{expected_name}` was not found."
));
}
}
let prefix = container
.instance_name_template
.as_deref()
.and_then(instance_name_template_prefix);
for instance in instances {
let matches_release_instance = prefix
.as_deref()
.is_some_and(|prefix| instance.name.starts_with(prefix));
if !matches_release_instance {
continue;
}
let is_expected = expected_name
.as_deref()
.is_some_and(|expected| instance.name == expected);
if is_expected {
continue;
}
if instance_is_active(instance) {
if instance.version == info.app_version {
println!(
"warning: active stale-looking instance `{}` has current app version {:?}",
instance.name, instance.version
);
} else {
return Err(format!(
"Stale active container instance `{}` has version {:?}; expected active version {:?}.",
instance.name, instance.version, info.app_version
));
}
} else {
println!(
"info: inactive stale container instance `{}` in state `{}`",
instance.name, instance.state
);
}
}
println!("Container instance verification: ok");
Ok(())
}
async fn verify_health(
app: &ResolvedCloudflareApp,
expected_version: Option<&str>,
) -> Result<(), String> {
let urls = resolve_health_urls(app)?;
if urls.is_empty() {
println!("Live health verification skipped: no health URL could be inferred.");
return Ok(());
}
let client = reqwest::Client::builder()
.timeout(std::time::Duration::from_secs(20))
.build()
.map_err(|error| format!("Failed to build HTTP client: {error}"))?;
for url in urls {
let response = client
.get(&url)
.send()
.await
.map_err(|error| format!("Health check failed for {url}: {error}"))?;
let status = response.status();
let text = response
.text()
.await
.map_err(|error| format!("Failed to read health response from {url}: {error}"))?;
if !status.is_success() {
return Err(format!("Health check failed for {url}: HTTP {status}"));
}
if let Some(expected_version) = expected_version {
let value: Value = serde_json::from_str(&text)
.map_err(|error| format!("Health check {url} did not return JSON: {error}"))?;
let observed = observed_health_version(&value).ok_or_else(|| {
format!("Health check {url} JSON did not include a string `version` field.")
})?;
if observed != expected_version {
return Err(format!(
"Health check {url} reported version `{observed}`, expected `{expected_version}`."
));
}
}
println!("Health check: {url} ok");
}
Ok(())
}
fn observed_health_version(value: &Value) -> Option<&str> {
value
.get("version")
.and_then(Value::as_str)
.or_else(|| value.get("appVersion").and_then(Value::as_str))
}
fn resolve_health_urls(app: &ResolvedCloudflareApp) -> Result<Vec<String>, String> {
let configured = app
.container
.as_ref()
.map(|container| {
container
.health_urls
.iter()
.map(|url| url.trim())
.filter(|url| !url.is_empty())
.map(ToOwned::to_owned)
.collect::<Vec<_>>()
})
.unwrap_or_default();
if !configured.is_empty() {
return Ok(configured);
}
Ok(resolve_inferred_health_url(app)?.into_iter().collect())
}
fn resolve_inferred_health_url(app: &ResolvedCloudflareApp) -> Result<Option<String>, String> {
let path = app
.container
.as_ref()
.and_then(|container| container.healthcheck_path.as_deref())
.map(normalize_health_path)
.unwrap_or_else(|| "/health".to_string());
let config = load_wrangler_config_value(&app.wrangler_config)?;
if let Some(pattern) = config
.get("routes")
.and_then(Value::as_array)
.and_then(|routes| routes.first())
.and_then(|route| {
route
.get("pattern")
.and_then(Value::as_str)
.or_else(|| route.as_str())
})
{
let host = pattern
.trim()
.trim_start_matches("https://")
.trim_start_matches("http://")
.trim_end_matches("/*")
.trim_end_matches('/');
if !host.is_empty() {
return Ok(Some(format!("https://{host}{path}")));
}
}
if config
.get("workers_dev")
.and_then(Value::as_bool)
.unwrap_or(true)
{
let worker_name = config
.get("name")
.and_then(Value::as_str)
.unwrap_or(&app.worker.name);
return Ok(Some(format!("https://{worker_name}.workers.dev{path}")));
}
Ok(None)
}
fn parse_container_app_info(value: &Value) -> ContainerAppInfo {
let image = find_string_key(value, &["image", "current_image", "image_ref"]);
let image_name = find_string_key(value, &["image_name", "imageName"]);
let image_tag = image
.as_deref()
.and_then(|image| image.rsplit_once(':').map(|(_, tag)| tag.to_string()))
.or_else(|| find_string_key(value, &["image_tag", "tag"]));
ContainerAppInfo {
image,
image_name,
image_tag,
app_version: find_i64_key(value, &["version", "app_version"]),
instance_type: find_string_key(value, &["instance_type", "instanceType"])
.or_else(|| instance_type_from_resource_shape(value)),
}
}
fn parse_container_instances(value: &Value) -> Result<Vec<ContainerInstance>, String> {
let array = value
.as_array()
.or_else(|| value.get("instances").and_then(Value::as_array))
.or_else(|| value.get("items").and_then(Value::as_array))
.or_else(|| value.get("result").and_then(Value::as_array))
.or_else(|| value.get("value").and_then(Value::as_array))
.ok_or_else(|| "wrangler containers instances JSON was not an array.".to_string())?;
Ok(array
.iter()
.filter_map(|item| {
let name = find_string_key(item, &["name", "id"])?;
let state = find_string_key(item, &["state", "status"]).unwrap_or_default();
Some(ContainerInstance {
name,
state,
version: find_i64_key(item, &["version"]),
})
})
.collect())
}
fn expected_instance_type(app: &ResolvedCloudflareApp) -> Result<Option<String>, String> {
let config = load_wrangler_config_value(&app.wrangler_config)?;
let Some(container) = find_wrangler_container(app, &config) else {
return Ok(None);
};
if let Some(instance_type) = container.get("instance_type").and_then(Value::as_str) {
return Ok(Some(instance_type.to_string()));
}
Ok(instance_type_from_resource_shape(container))
}
fn find_wrangler_container<'a>(
app: &ResolvedCloudflareApp,
config: &'a Value,
) -> Option<&'a Value> {
let class_name = app.container.as_ref()?.class_name.as_deref()?;
config
.get("containers")
.and_then(Value::as_array)?
.iter()
.find(|container| container.get("class_name").and_then(Value::as_str) == Some(class_name))
}
fn instance_type_from_resource_shape(value: &Value) -> Option<String> {
let vcpu = find_f64_key(value, &["vcpu", "vcpu_count", "cpu"]);
let memory_mib = find_i64_key(value, &["memory_mib", "memoryMiB", "memory"]);
let disk_mb = find_i64_key(value, &["disk_mb", "diskMB"]).or_else(|| {
value
.get("disk")
.and_then(|disk| find_i64_key(disk, &["size_mb", "sizeMB"]))
});
match (vcpu, memory_mib, disk_mb) {
(Some(vcpu), Some(memory), Some(disk)) => Some(match (vcpu, memory, disk) {
(vcpu, 256, 2000) if approx_eq(vcpu, 0.0625) => "lite".to_string(),
(vcpu, 1024, 4000) if approx_eq(vcpu, 0.25) => "basic".to_string(),
(vcpu, 4096, 8000) if approx_eq(vcpu, 0.5) => "standard-1".to_string(),
(vcpu, 6144, 12000) if approx_eq(vcpu, 1.0) => "standard-2".to_string(),
(vcpu, 8192, 16000) if approx_eq(vcpu, 2.0) => "standard-3".to_string(),
(vcpu, 12288, 20000) if approx_eq(vcpu, 4.0) => "standard-4".to_string(),
_ => format!("custom(vcpu={vcpu},memory_mib={memory},disk_mb={disk})"),
}),
_ => None,
}
}
fn approx_eq(left: f64, right: f64) -> bool {
(left - right).abs() < f64::EPSILON
}
fn find_string_key(value: &Value, keys: &[&str]) -> Option<String> {
match value {
Value::Object(map) => {
for key in keys {
if let Some(text) = map.get(*key).and_then(Value::as_str) {
return Some(text.to_string());
}
}
map.values().find_map(|item| find_string_key(item, keys))
}
Value::Array(items) => items.iter().find_map(|item| find_string_key(item, keys)),
_ => None,
}
}
fn find_i64_key(value: &Value, keys: &[&str]) -> Option<i64> {
match value {
Value::Object(map) => {
for key in keys {
if let Some(number) = map.get(*key).and_then(Value::as_i64) {
return Some(number);
}
if let Some(number) = map
.get(*key)
.and_then(Value::as_str)
.and_then(|text| text.parse::<i64>().ok())
{
return Some(number);
}
}
map.values().find_map(|item| find_i64_key(item, keys))
}
Value::Array(items) => items.iter().find_map(|item| find_i64_key(item, keys)),
_ => None,
}
}
fn find_f64_key(value: &Value, keys: &[&str]) -> Option<f64> {
match value {
Value::Object(map) => {
for key in keys {
if let Some(number) = map.get(*key).and_then(Value::as_f64) {
return Some(number);
}
if let Some(number) = map
.get(*key)
.and_then(Value::as_str)
.and_then(|text| text.parse::<f64>().ok())
{
return Some(number);
}
}
map.values().find_map(|item| find_f64_key(item, keys))
}
Value::Array(items) => items.iter().find_map(|item| find_f64_key(item, keys)),
_ => None,
}
}
fn render_instance_name_template(template: &str, version: &str) -> String {
let raw_segment = version
.chars()
.map(|ch| {
if ch.is_ascii_alphanumeric() {
ch.to_ascii_lowercase()
} else {
'-'
}
})
.collect::<String>();
let segment = raw_segment
.split('-')
.filter(|part| !part.is_empty())
.collect::<Vec<_>>()
.join("-");
template.replace("${VERSION}", &segment)
}
fn instance_name_template_prefix(template: &str) -> Option<String> {
template
.split_once("${VERSION}")
.map(|(prefix, _)| prefix.to_string())
.filter(|prefix| !prefix.is_empty())
}
fn instance_is_active(instance: &ContainerInstance) -> bool {
!matches!(
instance.state.to_ascii_lowercase().as_str(),
"stopped" | "inactive" | "terminated" | "failed" | "exited"
)
}
fn report_or_prune_old_images(
app: &ResolvedCloudflareApp,
after_info: &ContainerAppInfo,
context: &ContainerVerificationContext,
) -> Result<(), String> {
let Some(container) = app.container.as_ref() else {
return Ok(());
};
let Some(expected_image_name) = container.expected_image_name.as_deref() else {
return Ok(());
};
let Some(active_tag) = after_info.image_tag.as_deref() else {
return Ok(());
};
let keep_count = context
.keep_image_tag_count
.or(container.image_tag_keep_count)
.unwrap_or(20);
let mut args = vec![
"containers".to_string(),
"images".to_string(),
"list".to_string(),
"--filter".to_string(),
format!("^{expected_image_name}$"),
"--json".to_string(),
];
args.extend(wrangler_config_args(app));
let output = run_wrangler_checked(app, args, "wrangler containers images list --json");
let Ok(output) = output else {
println!("warning: image cleanup report skipped; Wrangler could not list images.");
return Ok(());
};
let value: Value = serde_json::from_str(output.stdout.trim()).map_err(|error| {
format!("wrangler containers images list --json returned invalid JSON: {error}")
})?;
let tags = parse_image_tags(&value);
let to_delete = select_image_tags_to_delete(&tags, active_tag, keep_count);
if to_delete.is_empty() {
return Ok(());
}
if !context.prune_old_images {
println!(
"Image cleanup available: {} old tags can be pruned with --prune-old-images.",
to_delete.len()
);
return Ok(());
}
for tag in to_delete {
let mut args = vec![
"containers".to_string(),
"images".to_string(),
"delete".to_string(),
format!("{expected_image_name}:{tag}"),
];
args.extend(wrangler_config_args(app));
run_wrangler_checked(app, args, "wrangler containers images delete")?;
}
Ok(())
}
fn parse_image_tags(value: &Value) -> Vec<String> {
match value {
Value::Array(items) => items
.iter()
.filter_map(|item| {
item.as_str()
.map(ToOwned::to_owned)
.or_else(|| {
item.get("tag")
.and_then(Value::as_str)
.map(ToOwned::to_owned)
})
.or_else(|| {
item.get("name")
.and_then(Value::as_str)
.and_then(|name| name.rsplit_once(':').map(|(_, tag)| tag.to_string()))
})
})
.collect(),
Value::Object(map) => map
.get("images")
.or_else(|| map.get("items"))
.or_else(|| map.get("result"))
.map(parse_image_tags)
.unwrap_or_default(),
_ => Vec::new(),
}
}
fn select_image_tags_to_delete(
tags: &[String],
active_tag: &str,
keep_count: usize,
) -> Vec<String> {
let mut kept = HashSet::from([active_tag.to_string()]);
for tag in tags.iter().take(keep_count) {
kept.insert(tag.clone());
}
tags.iter()
.filter(|tag| !kept.contains(*tag))
.cloned()
.collect()
}
fn sync_version_var(app: &mut ResolvedCloudflareApp, version: &str) -> Result<(), String> {
let version_var = app
.require_container()?
.version_var
.as_deref()
.ok_or_else(|| "workers[].container.version_var is required for release.".to_string())?
.to_string();
let version_var = version_var.as_str();
let content = fs::read_to_string(&app.wrangler_config)
.map_err(|error| format!("Failed to read {}: {error}", app.wrangler_config.display()))?;
let pattern = regex::Regex::new(&format!(
r#""{}"\s*:\s*"[^"]*""#,
regex::escape(version_var)
))
.map_err(|error| error.to_string())?;
if !pattern.is_match(&content) {
return Err(format!(
"Could not find version var `{version_var}` in {}. Add it under `vars` before release.",
app.wrangler_config.display()
));
}
let replacement = format!("\"{version_var}\": \"{version}\"");
let updated = pattern.replace(&content, replacement).to_string();
fs::write(&app.wrangler_config, updated)
.map_err(|error| format!("Failed to write {}: {error}", app.wrangler_config.display()))?;
println!("Updated {version_var} in {}", app.wrangler_config.display());
Ok(())
}
fn resolve_dockerfile_path(app: &ResolvedCloudflareApp) -> PathBuf {
let dockerfile = app
.container
.as_ref()
.and_then(|container| container.dockerfile.as_deref())
.unwrap_or("Dockerfile");
let path = PathBuf::from(dockerfile);
if path.is_absolute() {
path
} else {
app.worker_root.join(path)
}
}
fn resolve_application_id(
app: &ResolvedCloudflareApp,
override_id: Option<&str>,
) -> Result<String, String> {
override_id
.map(str::trim)
.filter(|value| !value.is_empty())
.map(ToOwned::to_owned)
.or_else(|| {
app.container
.as_ref()
.and_then(|container| container.application_id.clone())
})
.ok_or_else(|| {
format!(
"No container application id configured for `{}`. Set workers[].container.application_id or pass --application-id.",
app.app_name
)
})
}
fn normalize_health_path(path: &str) -> String {
let trimmed = path.trim();
if trimmed.is_empty() {
return "/health".to_string();
}
if trimmed.starts_with('/') {
trimmed.to_string()
} else {
format!("/{trimmed}")
}
}
fn mask_secrets(input: &str) -> String {
let mut output = input.to_string();
for (key, value) in env::vars() {
let key_upper = key.to_ascii_uppercase();
let looks_secret = key_upper.contains("SECRET")
|| key_upper.contains("TOKEN")
|| key_upper.contains("PRIVATE_KEY")
|| key_upper.contains("PASSWORD");
if looks_secret && value.len() >= 8 {
output = output.replace(&value, "***");
}
}
output
}
fn validate_payload_for_workflow(payload: &CloudflareWorkflowPayload) -> Result<(), String> {
if payload.workflow == CloudflareJobWorkflow::Release && payload.version.is_none() {
return Err("Release jobs require --version.".to_string());
}
if payload.dry_run && payload.skip_deploy {
return Err("--dry-run and --skip-deploy are mutually exclusive.".to_string());
}
Ok(())
}
fn enqueue_job(
project_root: &Path,
payload: CloudflareWorkflowPayload,
) -> Result<CloudflareJobRecord, String> {
let now = Utc::now();
let job = CloudflareJobRecord {
id: Uuid::new_v4().to_string(),
status: CloudflareJobStatus::Queued,
payload,
created_at: now,
updated_at: now,
started_at: None,
finished_at: None,
error: None,
logs: vec!["job queued".to_string()],
};
write_job(project_root, &job)?;
Ok(job)
}
fn job_dir(project_root: &Path) -> PathBuf {
project_root.join(CLOUDFLARE_JOB_DIR)
}
fn job_path(project_root: &Path, job_id: &str) -> PathBuf {
job_dir(project_root).join(format!("{job_id}.json"))
}
fn write_job(project_root: &Path, job: &CloudflareJobRecord) -> Result<(), String> {
let dir = job_dir(project_root);
fs::create_dir_all(&dir)
.map_err(|error| format!("Failed to create {}: {error}", dir.display()))?;
let path = job_path(project_root, &job.id);
let tmp_path = dir.join(format!("{}.{}.tmp", job.id, Uuid::new_v4()));
let content = format!(
"{}\n",
serde_json::to_string_pretty(job)
.map_err(|error| format!("Failed to encode job {}: {error}", job.id))?
);
fs::write(&tmp_path, content)
.map_err(|error| format!("Failed to write {}: {error}", tmp_path.display()))?;
match fs::rename(&tmp_path, &path) {
Ok(()) => Ok(()),
Err(first_error) => {
if path.exists() {
fs::remove_file(&path).map_err(|error| {
format!(
"Failed to replace {}: {first_error}; remove failed: {error}",
path.display()
)
})?;
fs::rename(&tmp_path, &path)
.map_err(|error| format!("Failed to replace {}: {error}", path.display()))
} else {
Err(format!("Failed to write {}: {first_error}", path.display()))
}
}
}
}
fn read_job(project_root: &Path, job_id: &str) -> Result<CloudflareJobRecord, String> {
let path = job_path(project_root, job_id);
let content = fs::read_to_string(&path)
.map_err(|error| format!("Failed to read {}: {error}", path.display()))?;
serde_json::from_str(&content)
.map_err(|error| format!("Failed to parse {}: {error}", path.display()))
}
fn list_jobs(project_root: &Path) -> Result<Vec<CloudflareJobRecord>, String> {
let dir = job_dir(project_root);
if !dir.exists() {
return Ok(Vec::new());
}
let mut jobs = Vec::new();
for entry in
fs::read_dir(&dir).map_err(|error| format!("Failed to read {}: {error}", dir.display()))?
{
let entry = entry.map_err(|error| format!("Failed to read job entry: {error}"))?;
if entry.path().extension().and_then(|value| value.to_str()) != Some("json") {
continue;
}
let content = fs::read_to_string(entry.path())
.map_err(|error| format!("Failed to read {}: {error}", entry.path().display()))?;
let job: CloudflareJobRecord = serde_json::from_str(&content)
.map_err(|error| format!("Failed to parse {}: {error}", entry.path().display()))?;
jobs.push(job);
}
jobs.sort_by_key(|job| job.created_at);
Ok(jobs)
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
use std::time::{SystemTime, UNIX_EPOCH};
fn temp_dir(label: &str) -> PathBuf {
let nanos = SystemTime::now()
.duration_since(UNIX_EPOCH)
.expect("clock")
.as_nanos();
let dir = env::temp_dir().join(format!("xbp-cloudflare-{label}-{nanos}"));
fs::create_dir_all(&dir).expect("temp dir");
dir
}
fn test_app(root: &Path) -> ResolvedCloudflareApp {
ResolvedCloudflareApp {
project_root: root.to_path_buf(),
config_path: root.join(".xbp").join("xbp.yaml"),
worker_root: root.join("apps").join("auth"),
app_name: "auth".to_string(),
worker: WorkerConfig {
name: "auth".to_string(),
root: root.join("apps").join("auth").to_string_lossy().to_string(),
script_name: Some("auth-worker".to_string()),
service: None,
deploy: None,
container: None,
containers: None,
durable_objects: None,
},
container: Some(WorkerContainerConfig {
class_name: Some("AuthContainer".to_string()),
binding: Some("AUTH_CONTAINER".to_string()),
dockerfile: Some("Dockerfile".to_string()),
port: Some(8787),
application_id: Some("app_123".to_string()),
healthcheck_path: Some("/health".to_string()),
required_secrets: vec!["AUTH_SECRET".to_string()],
version_var: Some("AUTH_VERSION".to_string()),
instance_name_template: None,
default_rollout: Some("immediate".to_string()),
health_urls: Vec::new(),
expected_image_name: None,
image_tag_keep_count: None,
allow_unchanged_container_image: false,
image: None,
build_context: None,
instance_type: None,
max_instances: None,
rollout_active_grace_period: None,
rollout_step_percentage: Vec::new(),
regions: Vec::new(),
}),
wrangler_config: root.join("apps").join("auth").join("wrangler.jsonc"),
}
}
#[test]
fn strip_json_comments_preserves_urls_inside_strings() {
let input = r#"{
// comment
"url": "https://example.com/a//b",
"value": 1 /* block */
}"#;
let value: Value = serde_json::from_str(&strip_json_comments(input)).expect("json");
assert_eq!(value["url"], "https://example.com/a//b");
assert_eq!(value["value"], 1);
}
#[test]
fn load_wrangler_jsonc_accepts_trailing_commas_and_nested_environments() {
let root = temp_dir("jsonc-environment");
let path = root.join("wrangler.jsonc");
fs::write(
&path,
r#"{
"$schema": "./node_modules/wrangler/config-schema.json",
// Top-level configuration
"name": "my-worker",
"main": "src/index.js",
"compatibility_date": "2026-07-12",
"workers_dev": false,
"route": {
"pattern": "example.org/*",
"zone_name": "example.org",
},
"kv_namespaces": [
{
"binding": "MY_NAMESPACE",
"id": "kv-id",
},
],
"env": {
"staging": {
"name": "my-worker-staging",
"route": {
"pattern": "staging.example.org/*",
"zone_name": "example.org",
},
"kv_namespaces": [
{
"binding": "MY_NAMESPACE",
"id": "staging-kv-id",
},
],
},
},
}"#,
)
.expect("wrangler jsonc");
let config = load_wrangler_config_value(&path).expect("parse JSONC");
assert_eq!(
config["env"]["staging"]["route"]["pattern"],
"staging.example.org/*"
);
assert_eq!(
config["env"]["staging"]["kv_namespaces"][0]["id"],
"staging-kv-id"
);
let _ = fs::remove_dir_all(root);
}
#[test]
fn wrangler_deploy_args_include_rollout_or_dry_run() {
let root = temp_dir("args");
fs::create_dir_all(root.join("apps/auth")).expect("worker dir");
let app = test_app(&root);
let dry = wrangler_deploy_args(&app, Some("gradual"), true);
assert!(dry.contains(&"--dry-run".to_string()));
assert!(!dry.contains(&"--containers-rollout".to_string()));
let live = wrangler_deploy_args(&app, Some("gradual"), false);
assert!(live
.windows(2)
.any(|pair| pair == ["--containers-rollout", "gradual"]));
let _ = fs::remove_dir_all(root);
}
#[test]
fn diagnostics_detect_missing_secret_without_printing_value() {
let root = temp_dir("diagnostics");
let worker_root = root.join("apps/auth");
fs::create_dir_all(&worker_root).expect("worker dir");
fs::write(worker_root.join("Dockerfile"), "FROM node:22\n").expect("dockerfile");
fs::write(
worker_root.join("wrangler.jsonc"),
r#"{
"name": "auth-worker",
"containers": [{"class_name": "AuthContainer", "image_vars": {"PORT": "8787"}}],
"durable_objects": {"bindings": [{"name": "AUTH_CONTAINER", "class_name": "AuthContainer"}]}
}"#,
)
.expect("wrangler");
let app = test_app(&root);
let diagnostics = collect_local_diagnostics(&app).expect("diagnostics");
assert!(diagnostics
.missing_required
.iter()
.any(|item| item.contains("AUTH_SECRET")));
assert!(!diagnostics
.missing_required
.iter()
.any(|item| item.contains("secret-value")));
let _ = fs::remove_dir_all(root);
}
#[test]
fn local_jobs_round_trip_status_transitions() {
let root = temp_dir("jobs");
let payload = CloudflareWorkflowPayload {
app: Some("auth".to_string()),
token: None,
account_id: None,
workflow: CloudflareJobWorkflow::Deploy,
rollout: CloudflareRollout::Immediate,
version: None,
domain: None,
dry_run: true,
skip_deploy: false,
allow_unchanged_container_image: false,
prune_old_images: false,
keep_image_tag_count: None,
};
let mut job = enqueue_job(&root, payload).expect("enqueue");
assert_eq!(job.status, CloudflareJobStatus::Queued);
job.status = CloudflareJobStatus::Running;
write_job(&root, &job).expect("write running");
let loaded = read_job(&root, &job.id).expect("read");
assert_eq!(loaded.status, CloudflareJobStatus::Running);
let jobs = list_jobs(&root).expect("list");
assert_eq!(jobs.len(), 1);
let _ = fs::remove_dir_all(root);
}
#[test]
fn fake_container_instances_json_parses() {
let output = json!([
{"id": "abc", "status": "running"},
{"id": "def", "status": "stopped"}
]);
let running = output
.as_array()
.expect("array")
.iter()
.filter(|item| item.get("status").and_then(Value::as_str) == Some("running"))
.count();
assert_eq!(running, 1);
}
#[test]
fn renders_release_derived_instance_name() {
assert_eq!(
render_instance_name_template("athena-auth-${VERSION}", "1.14.2-canary.1"),
"athena-auth-1-14-2-canary-1"
);
}
#[test]
fn maps_named_and_resource_shaped_instance_types() {
let named = json!({"configuration": {"instance_type": "basic"}});
let parsed = parse_container_app_info(&named);
assert_eq!(parsed.instance_type.as_deref(), Some("basic"));
let shaped = json!({"vcpu": 0.25, "memory_mib": 1024, "disk": {"size_mb": 4000}});
assert_eq!(
instance_type_from_resource_shape(&shaped).as_deref(),
Some("basic")
);
}
#[test]
fn detects_unchanged_container_image_without_override() {
let root = temp_dir("unchanged-image");
fs::create_dir_all(root.join("apps/auth")).expect("worker dir");
fs::write(
root.join("apps/auth/wrangler.jsonc"),
r#"{"containers":[{"class_name":"AuthContainer","instance_type":"basic"}]}"#,
)
.expect("wrangler");
let mut app = test_app(&root);
app.container
.as_mut()
.expect("container")
.expected_image_name = Some("demo-auth".to_string());
let before = ContainerAppInfo {
image: Some("registry/demo-auth:old".to_string()),
image_name: None,
image_tag: Some("old".to_string()),
app_version: Some(1),
instance_type: Some("basic".to_string()),
};
let err = verify_container_app_info(
&app,
Some(&before),
&before,
&ContainerVerificationContext {
before_info: None,
expected_version: None,
skip_deploy: false,
allow_unchanged_image: false,
prune_old_images: false,
keep_image_tag_count: None,
},
)
.expect_err("unchanged image should fail");
assert!(err.contains("did not change"));
let _ = fs::remove_dir_all(root);
}
#[test]
fn allows_unchanged_container_image_with_override() {
let root = temp_dir("unchanged-image-allowed");
fs::create_dir_all(root.join("apps/auth")).expect("worker dir");
fs::write(
root.join("apps/auth/wrangler.jsonc"),
r#"{"containers":[{"class_name":"AuthContainer","instance_type":"basic"}]}"#,
)
.expect("wrangler");
let mut app = test_app(&root);
app.container
.as_mut()
.expect("container")
.expected_image_name = Some("demo-auth".to_string());
let before = ContainerAppInfo {
image: Some("registry/demo-auth:old".to_string()),
image_name: None,
image_tag: Some("old".to_string()),
app_version: Some(1),
instance_type: Some("basic".to_string()),
};
verify_container_app_info(
&app,
Some(&before),
&before,
&ContainerVerificationContext {
before_info: None,
expected_version: None,
skip_deploy: false,
allow_unchanged_image: true,
prune_old_images: false,
keep_image_tag_count: None,
},
)
.expect("override should allow unchanged image");
let _ = fs::remove_dir_all(root);
}
#[test]
fn stale_active_wrong_version_instance_fails() {
let root = temp_dir("stale-instance");
let mut app = test_app(&root);
app.container
.as_mut()
.expect("container")
.instance_name_template = Some("auth-${VERSION}".to_string());
let info = ContainerAppInfo {
image: Some("registry/demo-auth:new".to_string()),
image_name: None,
image_tag: Some("new".to_string()),
app_version: Some(5),
instance_type: Some("basic".to_string()),
};
let instances = vec![
ContainerInstance {
name: "auth-1-2-3".to_string(),
state: "running".to_string(),
version: Some(5),
},
ContainerInstance {
name: "auth-1-2-2".to_string(),
state: "running".to_string(),
version: Some(4),
},
];
let err = verify_container_instances(&app, &info, &instances, &Some("1.2.3".to_string()))
.expect_err("wrong-version stale active instance should fail");
assert!(err.contains("Stale active"));
let _ = fs::remove_dir_all(root);
}
#[test]
fn same_version_stale_active_instance_warns_only() {
let root = temp_dir("same-version-stale");
let mut app = test_app(&root);
app.container
.as_mut()
.expect("container")
.instance_name_template = Some("auth-${VERSION}".to_string());
let info = ContainerAppInfo {
image: Some("registry/demo-auth:new".to_string()),
image_name: None,
image_tag: Some("new".to_string()),
app_version: Some(5),
instance_type: Some("basic".to_string()),
};
let instances = vec![
ContainerInstance {
name: "auth-1-2-3".to_string(),
state: "running".to_string(),
version: Some(5),
},
ContainerInstance {
name: "auth-1-2-2".to_string(),
state: "running".to_string(),
version: Some(5),
},
];
verify_container_instances(&app, &info, &instances, &Some("1.2.3".to_string()))
.expect("same-version stale active instance should warn only");
let _ = fs::remove_dir_all(root);
}
#[test]
fn parses_health_json_version_for_multiple_hosts() {
let bodies = [json!({"version": "1.2.3"}), json!({"appVersion": "1.2.3"})];
assert!(bodies
.iter()
.all(|body| observed_health_version(body) == Some("1.2.3")));
}
#[test]
fn selects_old_image_tags_to_prune_with_retention() {
let tags = vec![
"new".to_string(),
"1-2-3".to_string(),
"1-2-2".to_string(),
"old".to_string(),
];
assert_eq!(
select_image_tags_to_delete(&tags, "new", 2),
vec!["1-2-2".to_string(), "old".to_string()]
);
}
#[test]
fn detects_secret_like_wrangler_vars_without_value_leakage() {
let mut diagnostics = LocalDiagnostics::default();
let config = json!({"vars": {"API_TOKEN": "super-secret-token", "PUBLIC_URL": "https://example.com"}});
audit_wrangler_vars_for_secrets(&config, &mut diagnostics);
let rendered = diagnostics.warnings.join("\n");
assert!(rendered.contains("API_TOKEN"));
assert!(!rendered.contains("super-secret-token"));
assert!(!rendered.contains("PUBLIC_URL"));
}
}