use super::{load_local_env, print_json, resolve_target_script_name};
use crate::cli::commands::{
WorkersDurableObjectsSubCommand, WorkersDurableObjectsSyncCmd, WorkersTargetArgs,
};
use crate::commands::workers::project::WorkerTargetResolution;
use crate::provider_support::models::CloudflareDurableObjectNamespace;
use crate::strategies::{
get_all_workers, WorkerConfig, WorkerDurableObjectNamespace, WorkerDurableObjectsConfig,
XbpConfig,
};
use crate::utils::{find_xbp_config_upwards, serialize_xbp_yaml};
use serde_json::{json, Value};
use std::collections::BTreeSet;
use std::fs;
use std::path::Path;
pub(crate) async fn run(
resolution: WorkerTargetResolution,
worker_root: &Path,
token_override: Option<&str>,
account_id_override: Option<&str>,
target: WorkersTargetArgs,
command: WorkersDurableObjectsSubCommand,
) -> Result<(), String> {
let local_env = load_local_env(worker_root)?;
let script_name = resolve_target_script_name(worker_root, &target, &local_env);
let credentials =
super::resolve_cloudflare_credentials(token_override, account_id_override, &local_env)?;
let account_id = credentials.account_id.clone();
let client =
crate::provider_support::CloudflareClient::new(credentials.token, account_id.clone())?;
let live = client.list_durable_object_namespaces().await?;
let worker = selected_worker(&resolution, worker_root);
match command {
WorkersDurableObjectsSubCommand::List(cmd) => {
let entries = inventory_entries(
&resolution.config,
worker,
&script_name,
&live,
&account_id,
cmd.all,
);
if cmd.json {
print_json(&json!({ "script_name": script_name, "namespaces": entries }))
} else {
for entry in entries {
println!(
"{}\towner={}\tclass={}\tstorage={}\tcontainers={}\tbinding={}",
entry["name"].as_str().unwrap_or(""),
entry["script"].as_str().unwrap_or("unknown"),
entry["class"].as_str().unwrap_or("unknown"),
entry["storage"].as_str().unwrap_or("unknown"),
entry["use_containers"].as_bool().unwrap_or(false),
entry["bindings"].as_array().map_or(0, Vec::len),
);
}
Ok(())
}
}
WorkersDurableObjectsSubCommand::Inspect(cmd) => {
let entries = inventory_entries(
&resolution.config,
worker,
&script_name,
&live,
&account_id,
cmd.all,
);
let configured = worker
.and_then(|item| item.durable_objects.as_ref())
.cloned()
.unwrap_or_default();
let diff = build_diff(&configured, &live, &script_name);
let output = json!({
"script_name": script_name,
"configured": configured,
"namespaces": entries,
"diff": diff,
"secrets": "redacted",
});
if cmd.json {
print_json(&output)
} else {
println!(
"{}",
serde_json::to_string_pretty(&output).map_err(|e| e.to_string())?
);
Ok(())
}
}
WorkersDurableObjectsSubCommand::Sync(cmd) => {
sync_local_config(&resolution, worker, &script_name, &live, &cmd)
}
}
}
fn selected_worker<'a>(
resolution: &'a WorkerTargetResolution,
worker_root: &Path,
) -> Option<&'a WorkerConfig> {
resolution
.selected
.iter()
.find_map(|app| {
if app.root == worker_root {
app.config.as_ref()
} else {
None
}
})
.or_else(|| {
resolution
.selected
.first()
.and_then(|app| app.config.as_ref())
})
}
fn inventory_entries(
config: &XbpConfig,
worker: Option<&WorkerConfig>,
script_name: &str,
live: &[CloudflareDurableObjectNamespace],
account_id: &str,
all: bool,
) -> Vec<Value> {
let configured = worker
.and_then(|item| item.durable_objects.as_ref())
.map(|item| item.namespaces.as_slice())
.unwrap_or(&[]);
let bindings = worker
.and_then(|item| item.durable_objects.as_ref())
.map(|item| item.bindings.as_slice())
.unwrap_or(&[]);
live.iter()
.filter(|namespace| {
all || namespace.script.as_deref().is_none_or(|script| script == script_name)
})
.map(|namespace| {
let owner = get_all_workers(config).into_iter().find(|candidate| {
candidate.script_name.as_deref() == namespace.script.as_deref()
|| candidate.name == namespace.script.as_deref().unwrap_or_default()
});
let owner_configured = owner
.as_ref()
.and_then(|item| item.durable_objects.as_ref())
.map(|item| item.namespaces.as_slice())
.unwrap_or(configured);
let owner_bindings = owner
.as_ref()
.and_then(|item| item.durable_objects.as_ref())
.map(|item| item.bindings.as_slice())
.unwrap_or(bindings);
let local = owner_configured.iter().find(|item| {
item.id.as_deref() == Some(namespace.id.as_str()) || item.name == namespace.name
});
let namespace_class = namespace.class.as_deref().unwrap_or("");
let namespace_bindings = owner_bindings
.iter()
.filter(|binding| binding.class_name == namespace_class)
.map(|binding| binding.name.clone())
.collect::<Vec<_>>();
json!({
"id": namespace.id,
"name": namespace.name,
"script": namespace.script,
"class": namespace.class,
"service": owner.as_ref().and_then(|item| item.service.clone()).or_else(|| worker.and_then(|item| item.service.clone())),
"storage": if namespace.use_sqlite { "sqlite" } else { "durable_object" },
"use_sqlite": namespace.use_sqlite,
"use_containers": namespace.use_containers,
"bindings": namespace_bindings,
"configured": local.is_some(),
"dashboard_url": format!("https://dash.cloudflare.com/{}/workers/durable-objects/view/{}/overview", account_id, namespace.id),
"object_count": Value::Null,
"account_configured_workers": get_all_workers(config).len(),
})
})
.collect()
}
fn build_diff(
configured: &WorkerDurableObjectsConfig,
live: &[CloudflareDurableObjectNamespace],
script_name: &str,
) -> Value {
let local_names = configured
.namespaces
.iter()
.map(|item| item.name.clone())
.collect::<BTreeSet<_>>();
let live_names = live
.iter()
.filter(|item| {
item.script
.as_deref()
.is_none_or(|script| script == script_name)
})
.map(|item| item.name.clone())
.collect::<BTreeSet<_>>();
json!({
"missing_locally": live_names.difference(&local_names).collect::<Vec<_>>(),
"missing_live": local_names.difference(&live_names).collect::<Vec<_>>(),
"mismatches": live.iter().filter_map(|remote| {
let local = configured.namespaces.iter().find(|item| item.id.as_deref() == Some(remote.id.as_str()) || item.name == remote.name)?;
let mut fields = Vec::new();
if local.script.as_deref().is_some_and(|script| remote.script.as_deref() != Some(script)) { fields.push("script"); }
if local.class_name != remote.class.as_deref().unwrap_or("") { fields.push("class"); }
if local.use_sqlite != remote.use_sqlite { fields.push("use_sqlite"); }
if local.use_containers != remote.use_containers { fields.push("use_containers"); }
(!fields.is_empty()).then(|| json!({ "name": remote.name, "fields": fields }))
}).collect::<Vec<_>>(),
})
}
fn sync_local_config(
resolution: &WorkerTargetResolution,
worker: Option<&WorkerConfig>,
script_name: &str,
live: &[CloudflareDurableObjectNamespace],
cmd: &WorkersDurableObjectsSyncCmd,
) -> Result<(), String> {
let selected = live
.iter()
.filter(|item| {
item.script
.as_deref()
.is_none_or(|script| script == script_name)
})
.collect::<Vec<_>>();
let imported = selected
.iter()
.map(|item| WorkerDurableObjectNamespace {
id: Some(item.id.clone()),
name: item.name.clone(),
script: item.script.clone(),
class_name: item.class.clone().unwrap_or_default(),
use_sqlite: item.use_sqlite,
use_containers: item.use_containers,
})
.collect::<Vec<_>>();
let plan = json!({
"script_name": script_name,
"apply": cmd.apply,
"namespaces": imported,
"migrations": "unchanged; namespace sync never creates or mutates migrations",
});
if !cmd.apply {
println!(
"{}",
serde_json::to_string_pretty(&plan).map_err(|e| e.to_string())?
);
return Ok(());
}
let Some(worker_name) = worker.map(|item| item.name.as_str()) else {
return Err(
"The selected Worker is not represented in the local XBP workers[] config.".to_string(),
);
};
let found = find_xbp_config_upwards(&resolution.project_root)
.ok_or_else(|| "Could not locate the project XBP config.".to_string())?;
if found.kind != "yaml" && found.kind != "json" {
return Err(format!(
"Sync currently writes YAML or JSON XBP configs; found {}.",
found.kind
));
}
let content = fs::read_to_string(&found.config_path)
.map_err(|e| format!("Failed to read {}: {}", found.config_path.display(), e))?;
let (mut config, _) =
crate::utils::parse_config_with_auto_heal::<XbpConfig>(&content, found.kind)?;
let configured_worker = config
.workers
.as_mut()
.and_then(|workers| workers.iter_mut().find(|item| item.name == worker_name))
.ok_or_else(|| format!("Worker `{worker_name}` was not found while applying sync."))?;
let durable_objects = configured_worker
.durable_objects
.get_or_insert_with(Default::default);
durable_objects.namespaces = imported;
let rendered = if found.kind == "yaml" {
serialize_xbp_yaml(&serde_json::to_value(&config).map_err(|e| e.to_string())?)?
} else {
serde_json::to_string_pretty(&config).map_err(|e| e.to_string())?
};
fs::write(&found.config_path, rendered)
.map_err(|e| format!("Failed to write {}: {}", found.config_path.display(), e))?;
println!("Updated {}", found.config_path.display());
Ok(())
}
pub(crate) fn validate_worker(
worker: &WorkerConfig,
all_workers: &[WorkerConfig],
) -> Result<(), String> {
let Some(durable_objects) = worker.durable_objects.as_ref() else {
return Ok(());
};
validate_migration_history(&durable_objects.migrations, &durable_objects.migrations)?;
let mut classes = BTreeSet::new();
let mut container_classes = BTreeSet::new();
if let Some(container) = worker
.container
.as_ref()
.and_then(|item| item.class_name.as_ref())
{
classes.insert(container.clone());
container_classes.insert(container.clone());
}
for container in worker.containers.as_deref().unwrap_or(&[]) {
if let Some(class_name) = &container.class_name {
classes.insert(class_name.clone());
container_classes.insert(class_name.clone());
}
}
for migration in &durable_objects.migrations {
for class_name in migration
.new_classes
.iter()
.chain(&migration.new_sqlite_classes)
{
classes.insert(class_name.clone());
}
}
for binding in &durable_objects.bindings {
let external = binding.script_name.as_deref().is_some_and(|script| {
script
!= worker
.script_name
.as_deref()
.unwrap_or(worker.name.as_str())
});
if !external && !classes.contains(&binding.class_name) {
return Err(format!(
"Worker `{}` Durable Object binding `{}` references missing class `{}`.",
worker.name, binding.name, binding.class_name
));
}
if external {
let known = all_workers.iter().any(|candidate| {
candidate.script_name.as_deref() == binding.script_name.as_deref()
|| candidate.name == binding.script_name.as_deref().unwrap_or_default()
});
if !known {
return Err(format!(
"Worker `{}` Durable Object binding `{}` references unknown Worker `{}`.",
worker.name,
binding.name,
binding.script_name.as_deref().unwrap_or("")
));
}
if binding.environment.as_deref().is_some_and(str::is_empty) {
return Err(format!(
"Worker `{}` Durable Object binding `{}` has an empty environment.",
worker.name, binding.name
));
}
}
}
for namespace in &durable_objects.namespaces {
if namespace.use_containers && !container_classes.contains(&namespace.class_name) {
return Err(format!("Durable Object namespace `{}` is Containers-backed but class `{}` has no container definition.", namespace.name, namespace.class_name));
}
}
let mut tags = BTreeSet::new();
for migration in &durable_objects.migrations {
if !tags.insert(&migration.tag) {
return Err(format!(
"Durable Object migration tag `{}` is repeated.",
migration.tag
));
}
}
Ok(())
}
pub(crate) fn validate_selected_workers(config: &XbpConfig) -> Result<(), String> {
let workers = get_all_workers(config);
for worker in &workers {
validate_worker(worker, &workers)?;
}
Ok(())
}
pub(crate) fn validate_live_worker_namespaces(
worker: &WorkerConfig,
script_name: &str,
live: &[CloudflareDurableObjectNamespace],
) -> Result<(), String> {
let Some(configured) = worker.durable_objects.as_ref() else {
return Ok(());
};
for local in &configured.namespaces {
let Some(remote) = live
.iter()
.find(|item| local.id.as_deref() == Some(item.id.as_str()) || local.name == item.name)
else {
continue;
};
if remote
.script
.as_deref()
.is_some_and(|script| script != script_name)
{
return Err(format!("Durable Object namespace `{}` is owned by Worker `{}` but local Worker `{}` declares it.", local.name, remote.script.as_deref().unwrap_or(""), script_name));
}
if local.class_name != remote.class.as_deref().unwrap_or("") {
return Err(format!(
"Durable Object namespace `{}` class differs from live Cloudflare namespace.",
local.name
));
}
if local.use_sqlite != remote.use_sqlite {
return Err(format!(
"Durable Object namespace `{}` SQLite mode differs from live Cloudflare namespace.",
local.name
));
}
}
Ok(())
}
pub(crate) fn validate_migration_history(
previous: &[crate::strategies::WorkerDurableObjectMigration],
current: &[crate::strategies::WorkerDurableObjectMigration],
) -> Result<(), String> {
for old in previous {
let Some(new) = current.iter().find(|item| item.tag == old.tag) else {
return Err(format!(
"Durable Object migration tag `{}` was removed.",
old.tag
));
};
if old.new_classes != new.new_classes || old.new_sqlite_classes != new.new_sqlite_classes {
return Err(format!(
"Durable Object migration tag `{}` was rewritten.",
old.tag
));
}
}
let mut seen = BTreeSet::new();
for migration in current {
if !seen.insert(&migration.tag) {
return Err(format!(
"Durable Object migration tag `{}` is repeated.",
migration.tag
));
}
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::validate_worker;
use super::{build_diff, CloudflareDurableObjectNamespace, WorkerDurableObjectsConfig};
use crate::strategies::{
WorkerConfig, WorkerDurableObjectBinding, WorkerDurableObjectMigration,
};
#[test]
fn rejects_missing_same_worker_class() {
let worker = WorkerConfig {
name: "web".into(),
root: "apps/web".into(),
script_name: Some("web".into()),
service: None,
deploy: None,
container: None,
containers: None,
durable_objects: Some(WorkerDurableObjectsConfig {
bindings: vec![WorkerDurableObjectBinding {
name: "MISSING".into(),
class_name: "Nope".into(),
script_name: None,
environment: None,
}],
..Default::default()
}),
};
assert!(validate_worker(&worker, std::slice::from_ref(&worker)).is_err());
}
#[test]
fn migration_history_is_append_only() {
let previous = vec![WorkerDurableObjectMigration {
tag: "v1".into(),
new_classes: vec!["A".into()],
..Default::default()
}];
let rewritten = vec![WorkerDurableObjectMigration {
tag: "v1".into(),
new_classes: vec!["B".into()],
..Default::default()
}];
assert!(super::validate_migration_history(&previous, &rewritten).is_err());
assert!(super::validate_migration_history(&previous, &[]).is_err());
let appended = vec![
previous[0].clone(),
WorkerDurableObjectMigration {
tag: "v2".into(),
new_sqlite_classes: vec!["B".into()],
..Default::default()
},
];
assert!(super::validate_migration_history(&previous, &appended).is_ok());
}
#[test]
fn inventory_diff_covers_the_current_seven_namespace_shape() {
let names = [
"athena-auth_AthenaAuthContainer",
"xbp_SignalRelay",
"athena_AthenaContainerDev",
"athena_AthenaContainer",
"workflows-starter-template_WorkflowStatusDO",
"athena_AthenaContainerStaging",
"athena_AthenaContainerProduction",
];
let live = names
.iter()
.enumerate()
.map(|(index, name)| CloudflareDurableObjectNamespace {
id: format!("namespace-{index}"),
name: (*name).to_string(),
script: Some(
if name.starts_with("athena-auth") {
"athena-auth"
} else if name.starts_with("xbp_") {
"xbp"
} else if name.starts_with("workflows") {
"workflows-starter-template"
} else {
"athena"
}
.to_string(),
),
class: Some("TestClass".to_string()),
use_sqlite: !name.starts_with("xbp_"),
use_containers: name.contains("Container"),
extra: Default::default(),
})
.collect::<Vec<_>>();
let diff = build_diff(&WorkerDurableObjectsConfig::default(), &live, "athena");
assert_eq!(diff["missing_locally"].as_array().map(Vec::len), Some(4));
assert_eq!(live.len(), 7);
}
}