use std::collections::BTreeMap;
use std::collections::HashSet;
use std::time::Duration;
use serde::Deserialize;
use crate::contract::schema::{Adapter, Ecosystem, Registry};
use crate::protocol::release::{BuildArtifacts, DryRunReport, PlannedCommand, PublishReceipt};
use super::{
hash_file, 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 }
}
fn crates_io_registry(&self, t: &AdapterTarget) -> Result<&'static str, AdapterError> {
match t.target.registry {
Registry::CratesIo => Ok(CRATES_IO_ALIAS),
registry => Err(AdapterError::UnsupportedRegistry {
adapter: self.adapter,
registry,
}),
}
}
}
impl ReleaseAdapter for CargoAdapter {
fn adapter(&self) -> Adapter {
self.adapter
}
fn is_ci_delegated(&self) -> bool {
matches!(self.adapter, Adapter::CargoDist)
}
fn ci_owns_github_release(&self) -> bool {
matches!(self.adapter, Adapter::CargoDist)
}
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 registry = self.crates_io_registry(t)?;
let deps = target_workspace_deps(ctx, t)?;
let deferred = unpublished_workspace_deps(ctx, t.ecosystem(), &deps);
let defer_packaging = !deferred.is_empty();
let (planned_commands, _artifacts) =
cargo_build_gate(registry, &t.package, &t.version, defer_packaging);
run_all(ctx, &planned_commands)?;
let mut notes = vec![format!(
"publishes with `cargo publish --registry {registry} -p {}` in publish-all",
t.package
)];
if defer_packaging {
let chain = deferred
.iter()
.map(|m| format!("{}@{}", m.name, m.version))
.collect::<Vec<_>>()
.join(", ");
notes.push(format!(
"packaging of `{}` is deferred to that publish: it depends on workspace \
crate(s) not yet on the crates.io index, and a dependent cannot be \
`cargo package`d until they are published",
t.package
));
notes.push(format!(
"waits for these workspace dependencies to be crates.io-index-visible \
before publishing `{}`: {chain}",
t.package
));
}
Ok(DryRunReport {
adapter: self.adapter,
planned_commands,
notes,
})
}
fn build(
&self,
ctx: &EffectCtx<'_>,
t: &AdapterTarget,
) -> Result<BuildArtifacts, AdapterError> {
let (cmds, artifacts, notes) = if matches!(self.adapter, Adapter::CargoDist) {
(
vec![PlannedCommand::new("dist", &["build"])],
vec!["dist/".to_string()],
vec![],
)
} else {
let registry = self.crates_io_registry(t)?;
let deps = target_workspace_deps(ctx, t)?;
let deferred = unpublished_workspace_deps(ctx, t.ecosystem(), &deps);
let defer_packaging = !deferred.is_empty();
let (cmds, artifacts) =
cargo_build_gate(registry, &t.package, &t.version, defer_packaging);
let notes = if defer_packaging {
let chain = deferred
.iter()
.map(|m| format!("{}@{}", m.name, m.version))
.collect::<Vec<_>>()
.join(", ");
vec![format!(
"packaging of `{}` deferred to `cargo publish` in publish-all (it \
depends on workspace crate(s) not yet on the crates.io index: {chain})",
t.package
)]
} else {
vec![]
};
(cmds, artifacts, notes)
};
run_all(ctx, &cmds)?;
Ok(BuildArtifacts {
adapter: self.adapter,
artifacts,
notes,
})
}
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 registry = self.crates_io_registry(t)?;
let ecosystem = t.ecosystem();
if is_published(ctx, ecosystem, &t.package, &t.version)? {
return authenticate_skip(ctx, registry, ecosystem, t);
}
for dep in &target_workspace_deps(ctx, t)? {
wait_for_index(ctx, ecosystem, &dep.name, &dep.version)?;
}
run_all(
ctx,
&[PlannedCommand::new(
"cargo",
&["publish", "--registry", registry, "-p", &t.package],
)],
)?;
confirm_self_published(ctx, ecosystem, &t.package, &t.version)?;
Ok(make_receipt(ctx, t, None, Some(remote_url(t))))
}
fn timeout(&self) -> Duration {
Duration::from_secs(600)
}
}
fn cargo_build_gate(
registry: &str,
package: &str,
version: &str,
defer_packaging: bool,
) -> (Vec<PlannedCommand>, Vec<String>) {
let mut cmds = vec![PlannedCommand::new("cargo", &["check", "-p", package])];
if defer_packaging {
return (cmds, Vec::new());
}
cmds.push(PlannedCommand::new(
"cargo",
&[
"package",
"--registry",
registry,
"-p",
package,
"--no-verify",
],
));
(cmds, vec![format!("{package}-{version}.crate")])
}
fn remote_url(t: &AdapterTarget) -> String {
format!("https://crates.io/crates/{}/{}", t.package, t.version)
}
fn target_workspace_deps(
ctx: &EffectCtx<'_>,
t: &AdapterTarget,
) -> Result<Vec<Member>, AdapterError> {
let meta = load_metadata(ctx)?;
let members = publishable_members(&meta);
let by_name: BTreeMap<&str, &Member> = members.iter().map(|m| (m.name.as_str(), m)).collect();
let Some(target) = by_name.get(t.package.as_str()) else {
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:?}); declare each publishable crate as its own \
target and check the contract `package` and each crate's `publish` setting",
t.package
),
});
};
Ok(target
.deps
.iter()
.filter_map(|d| by_name.get(d.as_str()).map(|m| (*m).clone()))
.collect())
}
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 unpublished_workspace_deps(
ctx: &EffectCtx<'_>,
ecosystem: Ecosystem,
deps: &[Member],
) -> Vec<Member> {
deps.iter()
.filter(|d| !matches!(is_published(ctx, ecosystem, &d.name, &d.version), Ok(true)))
.cloned()
.collect()
}
fn is_published(
ctx: &EffectCtx<'_>,
ecosystem: Ecosystem,
package: &str,
version: &str,
) -> Result<bool, AdapterError> {
match ctx.registry.published_versions(ecosystem.as_str(), package) {
Ok(versions) => Ok(versions.iter().any(|v| v == version)),
Err(e) => Err(AdapterError::RegistryUnavailable {
package: package.to_string(),
version: version.to_string(),
source: e.to_string(),
}),
}
}
fn authenticate_skip(
ctx: &EffectCtx<'_>,
registry: &str,
ecosystem: Ecosystem,
t: &AdapterTarget,
) -> Result<PublishReceipt, AdapterError> {
let local = intended_crate_digest(ctx, registry, &t.package, &t.version)?;
let remote = ctx
.registry
.published_checksum(ecosystem.as_str(), &t.package, &t.version)
.map_err(|e| AdapterError::RegistryUnavailable {
package: t.package.clone(),
version: t.version.clone(),
source: e.to_string(),
})?;
if !is_sha256_hex(&remote) {
return Err(AdapterError::RegistryUnavailable {
package: t.package.clone(),
version: t.version.clone(),
source: format!("registry returned a malformed checksum: {remote:?}"),
});
}
if !local.eq_ignore_ascii_case(&remote) {
return Err(AdapterError::DigestMismatch {
package: t.package.clone(),
version: t.version.clone(),
local,
remote,
});
}
Ok(make_receipt(ctx, t, Some(local), Some(remote_url(t))))
}
fn intended_crate_digest(
ctx: &EffectCtx<'_>,
registry: &str,
package: &str,
version: &str,
) -> Result<String, AdapterError> {
let target_dir = load_metadata(ctx)?.target_directory;
run_all(
ctx,
&[PlannedCommand::new(
"cargo",
&[
"package",
"--registry",
registry,
"-p",
package,
"--no-verify",
],
)],
)?;
let crate_path = format!("{target_dir}/package/{package}-{version}.crate");
hash_file(ctx, &crate_path).map_err(|source| AdapterError::Command {
command: format!("sha256 of {crate_path}"),
code: None,
stderr: source,
})
}
fn is_sha256_hex(s: &str) -> bool {
s.len() == 64 && s.bytes().all(|b| b.is_ascii_hexdigit())
}
enum WaitFailure {
Absent {
waited_secs: u64,
},
Unreachable {
source: String,
},
}
fn poll_for_index(
ctx: &EffectCtx<'_>,
ecosystem: Ecosystem,
package: &str,
version: &str,
) -> Result<(), WaitFailure> {
let start = ctx.clock.now_unix();
let mut observed_absent = false;
let mut last_err: Option<String> = None;
loop {
match ctx.registry.published_versions(ecosystem.as_str(), package) {
Ok(versions) => {
if versions.iter().any(|v| v == version) {
return Ok(());
}
observed_absent = true;
}
Err(e) => last_err = Some(e.to_string()),
}
let waited = ctx.clock.now_unix().saturating_sub(start);
if waited >= INDEX_WAIT_TIMEOUT_SECS {
return Err(if observed_absent {
WaitFailure::Absent {
waited_secs: waited,
}
} else {
WaitFailure::Unreachable {
source: last_err.unwrap_or_else(|| {
"the registry never returned a definitive answer".to_string()
}),
}
});
}
ctx.clock.sleep(INDEX_POLL_INTERVAL);
}
}
fn wait_for_index(
ctx: &EffectCtx<'_>,
ecosystem: Ecosystem,
package: &str,
version: &str,
) -> Result<(), AdapterError> {
poll_for_index(ctx, ecosystem, package, version).map_err(|f| match f {
WaitFailure::Absent { waited_secs } => AdapterError::IndexTimeout {
package: package.to_string(),
version: version.to_string(),
waited_secs,
},
WaitFailure::Unreachable { source } => AdapterError::RegistryUnavailable {
package: package.to_string(),
version: version.to_string(),
source,
},
})
}
fn confirm_self_published(
ctx: &EffectCtx<'_>,
ecosystem: Ecosystem,
package: &str,
version: &str,
) -> Result<(), AdapterError> {
poll_for_index(ctx, ecosystem, package, version).map_err(|f| match f {
WaitFailure::Absent { waited_secs } => AdapterError::PublishNotVisible {
package: package.to_string(),
version: version.to_string(),
waited_secs,
},
WaitFailure::Unreachable { source } => AdapterError::RegistryUnavailable {
package: package.to_string(),
version: version.to_string(),
source,
},
})
}
#[derive(Clone)]
struct Member {
name: String,
version: String,
deps: Vec<String>,
}
#[derive(Deserialize)]
struct CargoMetadata {
packages: Vec<MetaPackage>,
workspace_members: Vec<String>,
#[serde(default)]
target_directory: 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>,
}