use std::collections::{BTreeMap, HashSet};
use std::time::Duration;
use serde::Deserialize;
use crate::contract::schema::{Adapter, Ecosystem};
use crate::protocol::release::{BuildArtifacts, DryRunReport, PlannedCommand, PublishReceipt};
use super::{make_receipt, run_all, AdapterError, AdapterTarget, EffectCtx, ReleaseAdapter};
const INDEX_WAIT_TIMEOUT_SECS: u64 = 300;
const INDEX_POLL_INTERVAL: Duration = Duration::from_secs(3);
const CRATES_IO_ALIAS: &str = "crates-io";
pub struct CargoAdapter {
adapter: Adapter,
}
impl CargoAdapter {
#[must_use]
pub fn new(adapter: Adapter) -> Self {
debug_assert!(matches!(
adapter,
Adapter::CargoPublish | Adapter::CargoDist
));
Self { adapter }
}
}
impl ReleaseAdapter for CargoAdapter {
fn adapter(&self) -> Adapter {
self.adapter
}
fn dry_run(
&self,
ctx: &EffectCtx<'_>,
t: &AdapterTarget,
) -> Result<DryRunReport, AdapterError> {
if matches!(self.adapter, Adapter::CargoDist) {
return Ok(DryRunReport {
adapter: self.adapter,
planned_commands: vec![PlannedCommand::new(
"dist",
&["plan", "--output-format=json"],
)],
notes: vec![],
});
}
let order = publish_order(ctx, t)?;
let mut planned_commands = Vec::with_capacity(order.len());
let mut notes = Vec::new();
if order.len() > 1 {
let chain = order
.iter()
.map(|m| m.name.as_str())
.collect::<Vec<_>>()
.join(" → ");
notes.push(format!("workspace publish order: {chain}"));
}
for (i, m) in order.iter().enumerate() {
planned_commands.push(PlannedCommand::new(
"cargo",
&["publish", "-p", &m.name, "--dry-run"],
));
if has_later_dependent(&order, i) {
notes.push(format!(
"then wait for crates.io to index `{}@{}` before publishing dependents",
m.name, m.version
));
}
}
Ok(DryRunReport {
adapter: self.adapter,
planned_commands,
notes,
})
}
fn build(
&self,
ctx: &EffectCtx<'_>,
t: &AdapterTarget,
) -> Result<BuildArtifacts, AdapterError> {
let (cmds, artifacts) = match self.adapter {
Adapter::CargoDist => (
vec![PlannedCommand::new("dist", &["build"])],
vec!["dist/".to_string()],
),
_ => (
vec![PlannedCommand::new("cargo", &["package", "-p", &t.package])],
vec![format!("{}-{}.crate", t.package, t.version)],
),
};
run_all(ctx, &cmds)?;
Ok(BuildArtifacts {
adapter: self.adapter,
artifacts,
notes: vec![],
})
}
fn publish(
&self,
ctx: &EffectCtx<'_>,
t: &AdapterTarget,
) -> Result<PublishReceipt, AdapterError> {
if matches!(self.adapter, Adapter::CargoDist) {
return Err(AdapterError::Unsupported {
adapter: self.adapter,
operation: "publish",
});
}
let order = publish_order(ctx, t)?;
for (i, m) in order.iter().enumerate() {
if !is_published(ctx, &m.name, &m.version) {
run_all(
ctx,
&[PlannedCommand::new("cargo", &["publish", "-p", &m.name])],
)?;
}
if has_later_dependent(&order, i) {
wait_for_index(ctx, &m.name, &m.version)?;
}
}
let remote_url = Some(format!(
"https://crates.io/crates/{}/{}",
t.package, t.version
));
Ok(make_receipt(ctx, t, None, remote_url))
}
fn timeout(&self) -> Duration {
Duration::from_secs(600)
}
}
fn publish_order(ctx: &EffectCtx<'_>, t: &AdapterTarget) -> Result<Vec<Member>, AdapterError> {
let meta = load_metadata(ctx)?;
let members = publishable_members(&meta);
if !members.iter().any(|m| m.name == t.package) {
let available: Vec<&str> = members.iter().map(|m| m.name.as_str()).collect();
return Err(AdapterError::Command {
command: "cargo metadata".to_string(),
code: None,
stderr: format!(
"target package `{}` is not a crates.io-publishable member of this workspace \
(publishable members: {available:?}); check the contract `package` and each \
crate's `publish` setting",
t.package
),
});
}
let closure = dep_closure(members, &t.package);
topo_sort(closure)
}
fn load_metadata(ctx: &EffectCtx<'_>) -> Result<CargoMetadata, AdapterError> {
let cmd = PlannedCommand::new("cargo", &["metadata", "--no-deps", "--format-version", "1"]);
let outputs = run_all(ctx, std::slice::from_ref(&cmd))?;
let stdout = outputs[0].stdout.trim();
if stdout.is_empty() {
return Err(AdapterError::Command {
command: cmd.rendered(),
code: None,
stderr: "`cargo metadata` succeeded but emitted no output — cannot resolve the \
workspace publish set"
.to_string(),
});
}
serde_json::from_str(stdout).map_err(|e| AdapterError::Command {
command: cmd.rendered(),
code: None,
stderr: format!("could not parse `cargo metadata` output: {e}"),
})
}
fn publishable_members(meta: &CargoMetadata) -> Vec<Member> {
let member_ids: HashSet<&str> = meta.workspace_members.iter().map(String::as_str).collect();
let pkgs: Vec<&MetaPackage> = meta
.packages
.iter()
.filter(|p| member_ids.contains(p.id.as_str()))
.filter(|p| publishable_to_crates_io(p.publish.as_deref()))
.collect();
let names: HashSet<&str> = pkgs.iter().map(|p| p.name.as_str()).collect();
pkgs.iter()
.map(|p| {
let mut deps: Vec<String> = p
.dependencies
.iter()
.filter(|d| matches!(d.kind.as_deref(), None | Some("build")))
.filter(|d| d.name != p.name && names.contains(d.name.as_str()))
.map(|d| d.name.clone())
.collect();
deps.sort();
deps.dedup();
Member {
name: p.name.clone(),
version: p.version.clone(),
deps,
}
})
.collect()
}
fn publishable_to_crates_io(publish: Option<&[String]>) -> bool {
match publish {
None => true,
Some(regs) => regs.iter().any(|r| r == CRATES_IO_ALIAS),
}
}
fn dep_closure(members: Vec<Member>, root: &str) -> Vec<Member> {
let by_name: BTreeMap<&str, &Member> = members.iter().map(|m| (m.name.as_str(), m)).collect();
let mut reached: HashSet<String> = HashSet::new();
let mut stack: Vec<String> = vec![root.to_string()];
while let Some(name) = stack.pop() {
if !reached.insert(name.clone()) {
continue;
}
if let Some(m) = by_name.get(name.as_str()) {
for d in &m.deps {
if !reached.contains(d) {
stack.push(d.clone());
}
}
}
}
members
.into_iter()
.filter(|m| reached.contains(&m.name))
.collect()
}
fn topo_sort(members: Vec<Member>) -> Result<Vec<Member>, AdapterError> {
let mut graph: BTreeMap<String, Member> = BTreeMap::new();
for m in members {
graph.insert(m.name.clone(), m);
}
let mut ordered: Vec<Member> = Vec::with_capacity(graph.len());
let mut published: HashSet<String> = HashSet::new();
let mut remaining: Vec<String> = graph.keys().cloned().collect();
while !remaining.is_empty() {
let ready = remaining
.iter()
.find(|n| graph[*n].deps.iter().all(|d| published.contains(d)))
.cloned();
match ready {
Some(n) => {
published.insert(n.clone());
remaining.retain(|x| x != &n);
ordered.push(graph.remove(&n).expect("ready name is a graph key"));
}
None => {
return Err(AdapterError::Command {
command: "cargo metadata".to_string(),
code: None,
stderr: format!(
"workspace publish order has a dependency cycle among: {remaining:?}"
),
});
}
}
}
Ok(ordered)
}
fn has_later_dependent(order: &[Member], i: usize) -> bool {
let name = &order[i].name;
order[i + 1..]
.iter()
.any(|later| later.deps.iter().any(|d| d == name))
}
fn is_published(ctx: &EffectCtx<'_>, package: &str, version: &str) -> bool {
ctx.registry
.published_versions(Ecosystem::Rust.as_str(), package)
.is_ok_and(|versions| versions.iter().any(|v| v == version))
}
fn wait_for_index(ctx: &EffectCtx<'_>, package: &str, version: &str) -> Result<(), AdapterError> {
let start = ctx.clock.now_unix();
loop {
if let Ok(versions) = ctx
.registry
.published_versions(Ecosystem::Rust.as_str(), package)
{
if versions.iter().any(|v| v == version) {
return Ok(());
}
}
if ctx.clock.now_unix().saturating_sub(start) >= INDEX_WAIT_TIMEOUT_SECS {
return Err(AdapterError::IndexTimeout {
package: package.to_string(),
version: version.to_string(),
waited_secs: INDEX_WAIT_TIMEOUT_SECS,
});
}
ctx.clock.sleep(INDEX_POLL_INTERVAL);
}
}
struct Member {
name: String,
version: String,
deps: Vec<String>,
}
#[derive(Deserialize)]
struct CargoMetadata {
packages: Vec<MetaPackage>,
workspace_members: Vec<String>,
}
#[derive(Deserialize)]
struct MetaPackage {
name: String,
version: String,
id: String,
#[serde(default)]
dependencies: Vec<MetaDep>,
#[serde(default)]
publish: Option<Vec<String>>,
}
#[derive(Deserialize)]
struct MetaDep {
name: String,
#[serde(default)]
kind: Option<String>,
}