use std::collections::HashSet;
use crate::ports::{CommandRunner, Tagger};
use crate::protocol::journal::{
EventKind, JournalEvent, Phase, PhaseOutcome, PublishReceipt as JournalReceipt, RunState,
JOURNAL_SCHEMA_VERSION,
};
use crate::protocol::plan::ReleasePlan;
use crate::protocol::reconcile::DelegatedRunStatus;
use crate::protocol::release::{PublishReceipt as AdapterReceipt, VerifyOutcome};
use super::adapters::{
hash_file, observe_cargo_dist_github_release, resolve, verification_artifacts, AdapterTarget,
EcosystemAdapter, EffectCtx, HomebrewAsset, HomebrewFormula, ReleaseAdapter, ReleaseArtifacts,
SourceTarball,
};
use super::journal::Journal;
use super::journal_target_ids;
use crate::contract::schema::{Adapter, Registry, Target};
pub trait ProgressSink {
fn event(&mut self, event: &JournalEvent);
fn verify_wait(&mut self, _progress: &VerifyWaitProgress) {}
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
pub struct VerifyWaitProgress {
pub target: String,
pub destination: String,
pub state: String,
pub elapsed_secs: u64,
pub remaining_secs: u64,
}
pub struct NullSink;
impl ProgressSink for NullSink {
fn event(&mut self, _event: &JournalEvent) {}
}
#[derive(Debug)]
pub enum CutError {
PhaseFailed {
phase: Phase,
target: Option<String>,
message: String,
},
Journal(std::io::Error),
Plan(String),
Checkout(String),
DelegatedRunPending {
target: String,
message: String,
},
DelegatedRunFailed {
target: String,
message: String,
},
}
impl std::fmt::Display for CutError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::PhaseFailed {
phase,
target,
message,
} => match target {
Some(t) => write!(
f,
"{}-phase failed on target `{t}`: {message}",
phase.as_str()
),
None => write!(f, "{}-phase failed: {message}", phase.as_str()),
},
Self::Journal(e) => write!(f, "could not write the release journal: {e}"),
Self::Plan(m) => write!(f, "the sealed plan is not executable: {m}"),
Self::Checkout(m) => {
write!(
f,
"could not check out the sealed commit to publish from: {m}"
)
}
Self::DelegatedRunPending { target, message } => {
write!(f, "verify-phase pending on target `{target}`: {message}")
}
Self::DelegatedRunFailed { target, message } => {
write!(f, "verify-phase failed on target `{target}`: {message}")
}
}
}
}
impl std::error::Error for CutError {}
struct TargetPlan {
id: String,
adapter: EcosystemAdapter,
input: AdapterTarget,
}
pub fn execute(
journal: &mut Journal<'_>,
plan: &ReleasePlan,
ctx: &EffectCtx<'_>,
tagger: &dyn Tagger,
sink: &mut dyn ProgressSink,
) -> Result<(), CutError> {
validate_plan(plan)?;
let targets = resolve_target_plans(plan)?;
let repo_slug = resolve_repo_slug(ctx, &targets);
let checkout_commit = journal
.state()
.bump
.as_ref()
.map_or_else(|| plan.head_sha.clone(), |b| b.commit.clone());
let checkout = SealedCheckout::materialize(ctx, &checkout_commit)?;
let checkout_ctx = ctx.with_repo_root(checkout.path());
let ctx = &checkout_ctx;
bump_phase(journal, sink, ctx, plan)?;
let tag_commit = journal
.state()
.bump
.as_ref()
.map_or_else(|| plan.head_sha.clone(), |b| b.commit.clone());
let source_tarball = repo_slug
.as_deref()
.and_then(|slug| source_tarball(slug, plan, &targets));
let homebrew = homebrew_inputs(plan, &targets);
let pre_artifacts = ReleaseArtifacts {
assets: Vec::new(),
source_tarball: source_tarball.clone(),
repo_slug: repo_slug.clone(),
homebrew: homebrew.clone(),
homebrew_assets: Vec::new(),
};
let pre_ctx = ctx.with_artifacts(&pre_artifacts);
reversible_phase(journal, sink, &pre_ctx, Phase::DryRun, &targets, None)?;
let mut assets = Vec::new();
reversible_phase(
journal,
sink,
&pre_ctx,
Phase::Build,
&targets,
Some(&mut assets),
)?;
let artifacts = ReleaseArtifacts {
assets,
source_tarball,
repo_slug,
homebrew,
homebrew_assets: Vec::new(),
};
publish_phase(journal, sink, &ctx.with_artifacts(&artifacts), &targets)?;
tag_phase(
journal,
sink,
tagger,
plan,
&tag_commit,
release_disposition(&targets),
)?;
dist_then_verify(journal, sink, ctx, &targets, plan, &artifacts)?;
if plan
.phases
.contains(&crate::protocol::plan::PlanPhase::AdvanceBranch)
{
advance_branch_phase(journal, sink, tagger, &tag_commit)
} else {
Ok(())
}
}
fn dist_then_verify(
journal: &mut Journal<'_>,
sink: &mut dyn ProgressSink,
ctx: &EffectCtx<'_>,
targets: &[TargetPlan],
plan: &ReleasePlan,
artifacts: &ReleaseArtifacts,
) -> Result<(), CutError> {
let dist_result = dist_phase(
journal,
sink,
ctx,
targets,
plan,
artifacts.repo_slug.as_deref(),
artifacts.homebrew.as_ref(),
);
match dist_result {
Ok(()) => verify_phase(
journal,
sink,
ctx,
targets,
plan,
artifacts.homebrew.as_ref(),
VerifyMode::CompletionBarrier,
),
Err(dist_error @ CutError::PhaseFailed { .. }) => {
let verify_result = verify_phase(
journal,
sink,
ctx,
targets,
plan,
artifacts.homebrew.as_ref(),
VerifyMode::ObserveAfterDistFailure,
);
match verify_result {
Ok(()) => Err(with_post_failure_verification(
dist_error,
journal.state(),
None,
)),
Err(verify_error @ CutError::PhaseFailed { .. }) => {
Err(with_post_failure_verification(
dist_error,
journal.state(),
Some(&verify_error),
))
}
Err(verify_error) => Err(verify_error),
}
}
Err(other) => Err(other),
}
}
fn advance_branch_phase(
journal: &mut Journal<'_>,
sink: &mut dyn ProgressSink,
tagger: &dyn Tagger,
tag_commit: &str,
) -> Result<(), CutError> {
let phase = Phase::AdvanceBranch;
if phase_completed_ok(journal.state(), phase) {
return Ok(());
}
record(journal, sink, EventKind::PhaseEntered { phase })?;
let branch = if let Some(branch) = journal.state().selected_default_branch.clone() {
branch
} else {
let branch = match tagger.default_branch() {
Ok(branch) => branch,
Err(error) => {
return fail_phase(
journal,
sink,
phase,
None,
format!("resolve remote default branch: {error}"),
);
}
};
record(
journal,
sink,
EventKind::DefaultBranchSelected {
branch: branch.clone(),
},
)?;
branch
};
match journal.state().default_branch.as_ref() {
Some(evidence) if evidence.branch == branch && evidence.commit == tag_commit => {}
Some(evidence) => {
return fail_phase(
journal,
sink,
phase,
None,
format!(
"journal branch evidence conflicts with this release: recorded {} at {}, expected {} at {}",
evidence.branch, evidence.commit, branch, tag_commit
),
);
}
None => {
if let Err(error) = tagger.advance_branch(&branch, tag_commit) {
return fail_phase(
journal,
sink,
phase,
None,
format!("advance origin/{branch} to {tag_commit}: {error}"),
);
}
record(
journal,
sink,
EventKind::DefaultBranchAdvanced {
branch,
commit: tag_commit.to_string(),
},
)?;
}
}
record(
journal,
sink,
EventKind::PhaseCompleted {
phase,
outcome: PhaseOutcome::Ok,
},
)
}
pub fn validate_plan(plan: &ReleasePlan) -> Result<(), CutError> {
resolve_target_plans(plan)?;
if plan
.targets
.iter()
.any(|target| matches!(target.adapter, Adapter::HomebrewTap | Adapter::HomebrewCore))
&& !plan.homebrew_platforms.iter().any(|triple| {
crate::release::adapters::homebrew::homebrew_platform_condition(triple).is_some()
})
{
return Err(CutError::Plan(
"Homebrew formula has no Homebrew-servable cargo-dist platforms; supported platforms are macOS aarch64/x86_64 and Linux musl aarch64/x86_64; refusing to write a formula with no installable archive".into(),
));
}
Ok(())
}
fn resolve_target_plans(plan: &ReleasePlan) -> Result<Vec<TargetPlan>, CutError> {
let ids = journal_target_ids(&plan.targets);
let mut out = Vec::with_capacity(plan.targets.len());
let mut seen: Vec<String> = Vec::new();
for (t, id) in plan.targets.iter().zip(ids) {
let package = t.package.clone().ok_or_else(|| {
CutError::Plan(format!(
"target `{}` has no resolved package name — pin an explicit `package` \
in OSS-RELEASE.md and re-plan",
t.ecosystem.as_str()
))
})?;
if seen.contains(&id) {
return Err(CutError::Plan(format!(
"two targets resolve to the same journal id `{id}` — the plan has two \
identical targets (same ecosystem, package, registry, and adapter); \
remove the duplicate target in OSS-RELEASE.md"
)));
}
seen.push(id.clone());
let input = AdapterTarget {
target: Target {
ecosystem: t.ecosystem,
package: Some(package.clone()),
registry: t.registry,
adapter: t.adapter,
},
package,
version: plan.version.clone(),
};
out.push(TargetPlan {
id,
adapter: resolve(t.adapter),
input,
});
}
Ok(out)
}
struct SealedCheckout<'a> {
runner: &'a dyn CommandRunner,
repo_root: &'a std::path::Path,
path: std::path::PathBuf,
}
impl<'a> SealedCheckout<'a> {
fn materialize(ctx: &EffectCtx<'a>, head_sha: &str) -> Result<SealedCheckout<'a>, CutError> {
let commitish = format!("{head_sha}^{{commit}}");
let probe = ctx
.runner
.run("git", &["cat-file", "-e", &commitish], ctx.repo_root)
.map_err(|e| {
CutError::Checkout(format!(
"cannot probe the sealed commit `{head_sha}` (`git cat-file` failed to run: {e})"
))
})?;
if probe.status != Some(0) {
return Err(CutError::Checkout(format!(
"the sealed commit `{head_sha}` is not present in this repository (never \
committed, not fetched, or garbage-collected). Commit and push the release \
commit, then re-plan/cut — a cut publishes from a clean checkout of the sealed \
HEAD, never the live working tree"
)));
}
let path = checkout_path(head_sha);
let path_str = path.to_string_lossy().to_string();
let out = ctx
.runner
.run(
"git",
&["worktree", "add", "--detach", &path_str, head_sha],
ctx.repo_root,
)
.map_err(|e| {
CutError::Checkout(format!(
"cannot create a clean checkout worktree for `{head_sha}` \
(`git worktree add` failed to run: {e})"
))
})?;
if out.status != Some(0) {
return Err(CutError::Checkout(format!(
"`git worktree add` could not check out the sealed commit `{head_sha}` into a \
clean worktree at `{path_str}`: {}",
out.stderr.trim()
)));
}
Ok(SealedCheckout {
runner: ctx.runner,
repo_root: ctx.repo_root,
path,
})
}
fn path(&self) -> &std::path::Path {
&self.path
}
}
impl Drop for SealedCheckout<'_> {
fn drop(&mut self) {
let path_str = self.path.to_string_lossy().to_string();
let _ = self.runner.run(
"git",
&["worktree", "remove", "--force", &path_str],
self.repo_root,
);
let _ = self
.runner
.run("git", &["worktree", "prune"], self.repo_root);
}
}
fn checkout_path(head_sha: &str) -> std::path::PathBuf {
let nanos = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map_or(0, |d| d.as_nanos());
let short: String = head_sha.chars().take(12).collect();
std::env::temp_dir().join(format!(
"shipshape-cut-{}-{short}-{nanos}",
std::process::id()
))
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum ReleaseDisposition<'a> {
Engine,
DelegatedToCi(&'a str),
TagOnly,
}
fn release_disposition(targets: &[TargetPlan]) -> ReleaseDisposition<'_> {
if targets.is_empty() {
return ReleaseDisposition::TagOnly;
}
match targets
.iter()
.find(|tp| tp.adapter.ci_owns_github_release())
{
Some(tp) => ReleaseDisposition::DelegatedToCi(tp.input.target.adapter.as_str()),
None => ReleaseDisposition::Engine,
}
}
fn needs_github_slug(targets: &[TargetPlan]) -> bool {
targets.iter().any(|tp| {
matches!(
tp.input.target.adapter,
Adapter::Manual | Adapter::HomebrewTap | Adapter::HomebrewCore
)
})
}
fn resolve_repo_slug(ctx: &EffectCtx<'_>, targets: &[TargetPlan]) -> Option<String> {
if !needs_github_slug(targets) {
return None;
}
let out = ctx
.runner
.run("git", &["remote", "get-url", "origin"], ctx.repo_root)
.ok()?;
if out.status != Some(0) {
return None;
}
crate::vcs::parse_github_slug(out.stdout.trim())
}
fn source_tarball(slug: &str, plan: &ReleasePlan, targets: &[TargetPlan]) -> Option<SourceTarball> {
let needed = targets.iter().any(|tp| {
matches!(
tp.input.target.adapter,
Adapter::HomebrewTap | Adapter::HomebrewCore
)
});
if !needed {
return None;
}
let tag = format!("v{}", plan.version);
Some(SourceTarball {
url: format!("https://github.com/{slug}/archive/refs/tags/{tag}.tar.gz"),
sha256: None,
})
}
fn homebrew_inputs(plan: &ReleasePlan, targets: &[TargetPlan]) -> Option<HomebrewFormula> {
let needed = targets.iter().any(|tp| {
matches!(
tp.input.target.adapter,
Adapter::HomebrewTap | Adapter::HomebrewCore
)
});
if !needed {
return None;
}
Some(HomebrewFormula {
tap: plan.homebrew_tap.clone(),
license: plan.license.clone(),
description: plan.description.clone(),
version: plan.version.clone(),
platforms: plan.homebrew_platforms.clone(),
})
}
fn reversible_phase(
journal: &mut Journal<'_>,
sink: &mut dyn ProgressSink,
ctx: &EffectCtx<'_>,
phase: Phase,
targets: &[TargetPlan],
mut assets: Option<&mut Vec<String>>,
) -> Result<(), CutError> {
if phase_completed_ok(journal.state(), phase) {
return Ok(());
}
record(journal, sink, EventKind::PhaseEntered { phase })?;
for tp in targets {
if target_cleared(journal.state(), phase, &tp.id) {
continue;
}
let outcome = match phase {
Phase::DryRun => tp.adapter.dry_run(ctx, &tp.input).map(|_| ()),
Phase::Build => tp.adapter.build(ctx, &tp.input).map(|built| {
if let Some(sink) = assets.as_deref_mut() {
sink.extend(built.artifacts);
}
}),
Phase::Bump
| Phase::Publish
| Phase::Tag
| Phase::Dist
| Phase::Verify
| Phase::AdvanceBranch => unreachable!("reversible_phase only runs dry_run/build"),
};
match outcome {
Ok(()) => {
let ev = match phase {
Phase::DryRun => EventKind::TargetDryRun {
target: tp.id.clone(),
},
Phase::Build => EventKind::TargetBuilt {
target: tp.id.clone(),
},
_ => unreachable!(),
};
record(journal, sink, ev)?;
}
Err(e) => return fail_phase(journal, sink, phase, Some(tp.id.clone()), e.to_string()),
}
}
record(
journal,
sink,
EventKind::PhaseCompleted {
phase,
outcome: PhaseOutcome::Ok,
},
)?;
Ok(())
}
fn publish_phase(
journal: &mut Journal<'_>,
sink: &mut dyn ProgressSink,
ctx: &EffectCtx<'_>,
targets: &[TargetPlan],
) -> Result<(), CutError> {
let phase = Phase::Publish;
if phase_completed_ok(journal.state(), phase) {
return Ok(());
}
record(journal, sink, EventKind::PhaseEntered { phase })?;
for tp in targets {
if journal.state().published.contains_key(&tp.id) {
continue;
}
if needs_post_tag(tp) {
continue;
}
if journal.state().delegated.contains(&tp.id) {
continue;
}
if tp.adapter.is_ci_delegated() {
record(
journal,
sink,
EventKind::TargetDelegated {
target: tp.id.clone(),
adapter: tp.input.target.adapter.as_str().to_string(),
},
)?;
continue;
}
match tp.adapter.publish(ctx, &tp.input) {
Ok(receipt) => {
record(
journal,
sink,
EventKind::TargetPublished {
target: tp.id.clone(),
receipt: to_journal_receipt(&receipt),
},
)?;
}
Err(e) => return fail_phase(journal, sink, phase, Some(tp.id.clone()), e.to_string()),
}
}
record(
journal,
sink,
EventKind::PhaseCompleted {
phase,
outcome: PhaseOutcome::Ok,
},
)?;
Ok(())
}
fn bump_phase(
journal: &mut Journal<'_>,
sink: &mut dyn ProgressSink,
ctx: &EffectCtx<'_>,
plan: &ReleasePlan,
) -> Result<(), CutError> {
let Some(bump) = plan.bump.as_ref() else {
return Ok(());
};
let phase = Phase::Bump;
if phase_completed_ok(journal.state(), phase) {
return Ok(());
}
record(journal, sink, EventKind::PhaseEntered { phase })?;
if journal.state().bump.is_none() {
let effective_date = crate::release::bump_exec::civil_date(ctx.clock.now_unix());
match crate::release::bump_exec::apply_bump(ctx, bump, &effective_date) {
Ok(outcome) => {
record(
journal,
sink,
EventKind::BumpApplied {
commit: outcome.commit,
effective_date: outcome.effective_date,
},
)?;
}
Err(e) => return fail_phase(journal, sink, phase, None, e.to_string()),
}
}
record(
journal,
sink,
EventKind::PhaseCompleted {
phase,
outcome: PhaseOutcome::Ok,
},
)?;
Ok(())
}
fn tag_phase(
journal: &mut Journal<'_>,
sink: &mut dyn ProgressSink,
tagger: &dyn Tagger,
plan: &ReleasePlan,
tag_commit: &str,
disposition: ReleaseDisposition<'_>,
) -> Result<(), CutError> {
let phase = Phase::Tag;
if phase_completed_ok(journal.state(), phase) {
return Ok(());
}
record(journal, sink, EventKind::PhaseEntered { phase })?;
let tag = format!("v{}", plan.version);
let title = format!("Release {}", plan.version);
if !tag_step_done(journal.state(), &tag, |s| s.created_local) {
if let Err(e) = tagger.create_tag(&tag, tag_commit, &title) {
return fail_phase(journal, sink, phase, None, format!("create local tag: {e}"));
}
record(
journal,
sink,
EventKind::TagCreatedLocal { tag: tag.clone() },
)?;
}
if !tag_step_done(journal.state(), &tag, |s| s.pushed_remote) {
if let Err(e) = tagger.push_tag(&tag) {
return fail_phase(journal, sink, phase, None, format!("push tag: {e}"));
}
record(
journal,
sink,
EventKind::TagPushedRemote { tag: tag.clone() },
)?;
}
let contradiction = match disposition {
ReleaseDisposition::DelegatedToCi(adapter) => {
tag_step_done(journal.state(), &tag, |s| s.github_release).then(|| {
format!(
"tag {tag} already has an engine-created GitHub Release, but the plan \
delegates the Release to CI ({adapter}); the adapter's ownership \
classification changed between attempts — reconcile the tag by hand"
)
})
}
ReleaseDisposition::Engine => {
tag_step_done(journal.state(), &tag, |s| s.github_release_delegated).then(|| {
format!(
"tag {tag}'s GitHub Release was already delegated to CI, but the plan now \
has the coordinator create it; the adapter's ownership classification \
changed between attempts — reconcile the tag by hand"
)
})
}
ReleaseDisposition::TagOnly => tag_step_done(journal.state(), &tag, |s| {
s.github_release || s.github_release_delegated
})
.then(|| {
format!(
"tag {tag} already carries a GitHub Release disposition, but this plan has \
no publish targets and must be tag-only — a publish-none run cannot \
complete over a Release that was created or delegated; reconcile the tag \
by hand"
)
}),
};
if let Some(message) = contradiction {
return fail_phase(journal, sink, phase, None, message);
}
github_release_step(journal, sink, tagger, disposition, &tag, &title)?;
record(
journal,
sink,
EventKind::PhaseCompleted {
phase,
outcome: PhaseOutcome::Ok,
},
)?;
Ok(())
}
fn github_release_step(
journal: &mut Journal<'_>,
sink: &mut dyn ProgressSink,
tagger: &dyn Tagger,
disposition: ReleaseDisposition<'_>,
tag: &str,
title: &str,
) -> Result<(), CutError> {
match disposition {
ReleaseDisposition::DelegatedToCi(adapter) => {
if !tag_step_done(journal.state(), tag, |s| s.github_release_delegated) {
record(
journal,
sink,
EventKind::GithubReleaseDelegated {
tag: tag.to_string(),
delegated_to: adapter.to_string(),
},
)?;
}
}
ReleaseDisposition::Engine => {
if !tag_step_done(journal.state(), tag, |s| s.github_release) {
match tagger.create_github_release(tag, title) {
Ok(url) => record(
journal,
sink,
EventKind::GithubReleaseCreated {
tag: tag.to_string(),
url,
},
)?,
Err(e) => {
return fail_phase(
journal,
sink,
Phase::Tag,
None,
format!("create GitHub Release: {e}"),
)
}
}
}
}
ReleaseDisposition::TagOnly => {}
}
Ok(())
}
fn needs_post_tag(tp: &TargetPlan) -> bool {
matches!(
tp.input.target.adapter,
Adapter::HomebrewTap | Adapter::HomebrewCore
)
}
const DIST_PUBLISH_ATTEMPTS: u32 = 3;
const DIST_PUBLISH_BACKOFF: std::time::Duration = std::time::Duration::from_secs(2);
fn dist_phase(
journal: &mut Journal<'_>,
sink: &mut dyn ProgressSink,
ctx: &EffectCtx<'_>,
targets: &[TargetPlan],
plan: &ReleasePlan,
repo_slug: Option<&str>,
homebrew: Option<&HomebrewFormula>,
) -> Result<(), CutError> {
let phase = Phase::Dist;
if phase_completed_ok(journal.state(), phase) {
return Ok(());
}
record(journal, sink, EventKind::PhaseEntered { phase })?;
let post_tag: Vec<&TargetPlan> = targets.iter().filter(|tp| needs_post_tag(tp)).collect();
if !post_tag.is_empty() {
let source_tarball = match repo_slug {
Some(slug) => {
let url = tag_archive_url(slug, &plan.version);
match compute_source_tarball_sha256(ctx, &url) {
Ok(sha256) => Some(SourceTarball {
url,
sha256: Some(sha256),
}),
Err(message) => return fail_phase(journal, sink, phase, None, message),
}
}
None => None,
};
let homebrew_assets = match (repo_slug, homebrew) {
(Some(slug), Some(formula)) => {
match fetch_homebrew_assets(ctx, slug, plan, formula, &post_tag) {
Ok(assets) => assets,
Err(message) => return fail_phase(journal, sink, phase, None, message),
}
}
_ => Vec::new(),
};
let artifacts = ReleaseArtifacts {
assets: Vec::new(),
source_tarball,
repo_slug: repo_slug.map(str::to_string),
homebrew: homebrew.cloned(),
homebrew_assets,
};
let dist_ctx = ctx.with_artifacts(&artifacts);
for tp in post_tag {
if journal.state().published.contains_key(&tp.id) {
continue;
}
match publish_dist_with_retry(&dist_ctx, tp) {
Ok(receipt) => {
record(
journal,
sink,
EventKind::TargetPublished {
target: tp.id.clone(),
receipt: to_journal_receipt(&receipt),
},
)?;
}
Err(e) => {
return fail_phase(journal, sink, phase, Some(tp.id.clone()), e.to_string())
}
}
}
}
record(
journal,
sink,
EventKind::PhaseCompleted {
phase,
outcome: PhaseOutcome::Ok,
},
)?;
Ok(())
}
fn publish_dist_with_retry(
ctx: &EffectCtx<'_>,
target: &TargetPlan,
) -> Result<AdapterReceipt, super::adapters::AdapterError> {
let mut attempt = 1;
loop {
match target.adapter.publish(ctx, &target.input) {
Ok(receipt) => return Ok(receipt),
Err(error)
if error.is_retryable_dist_setup_failure() && attempt < DIST_PUBLISH_ATTEMPTS =>
{
ctx.clock
.sleep(DIST_PUBLISH_BACKOFF.saturating_mul(attempt));
attempt += 1;
}
Err(error) => return Err(error),
}
}
}
fn with_post_failure_verification(
error: CutError,
state: &RunState,
verify_error: Option<&CutError>,
) -> CutError {
let CutError::PhaseFailed {
phase,
target,
mut message,
} = error
else {
return error;
};
let observations = state
.targets
.iter()
.map(|target| {
let outcome = state
.verified
.get(target)
.map_or("not_observed", |outcome| outcome.as_str());
format!("{target}={outcome}")
})
.collect::<Vec<_>>()
.join(", ");
message = format!(
"{message}; post-failure verify ran after the irreversible tag/publishes and observed journal targets: {observations}"
);
if let Some(verify_error) = verify_error {
message = format!("{message}; post-failure verify reported: {verify_error}");
}
CutError::PhaseFailed {
phase,
target,
message,
}
}
const DELEGATED_RELEASE_VERIFY_TIMEOUT_SECS: u64 = 20 * 60;
const DELEGATED_RELEASE_VERIFY_POLL_INTERVAL: std::time::Duration =
std::time::Duration::from_secs(15);
const DELEGATED_RELEASE_VERIFY_PROGRESS_INTERVAL_SECS: u64 = 60;
const DELEGATED_RELEASE_VERIFY_MAX_SLEEPS: u64 =
DELEGATED_RELEASE_VERIFY_TIMEOUT_SECS / DELEGATED_RELEASE_VERIFY_POLL_INTERVAL.as_secs() + 2;
struct DelegatedVerifyWindow {
start: u64,
sleeps: u64,
next_progress_at: std::collections::HashMap<String, u64>,
}
impl DelegatedVerifyWindow {
fn new(ctx: &EffectCtx<'_>) -> Self {
Self {
start: ctx.clock.now_unix(),
sleeps: 0,
next_progress_at: std::collections::HashMap::new(),
}
}
fn snapshot(&self, ctx: &EffectCtx<'_>) -> (u64, u64) {
let elapsed = ctx.clock.now_unix().saturating_sub(self.start);
(
elapsed,
DELEGATED_RELEASE_VERIFY_TIMEOUT_SECS.saturating_sub(elapsed),
)
}
fn sleep(&mut self, ctx: &EffectCtx<'_>) -> bool {
let (_, remaining) = self.snapshot(ctx);
if remaining == 0 || self.sleeps >= DELEGATED_RELEASE_VERIFY_MAX_SLEEPS {
return false;
}
ctx.clock.sleep(std::time::Duration::from_secs(
remaining.min(DELEGATED_RELEASE_VERIFY_POLL_INTERVAL.as_secs()),
));
self.sleeps += 1;
true
}
fn stream_wait(
&mut self,
sink: &mut dyn ProgressSink,
ctx: &EffectCtx<'_>,
target: &str,
destination: String,
state: &str,
) {
let (elapsed_secs, remaining_secs) = self.snapshot(ctx);
if remaining_secs == 0 {
return;
}
let next = self.next_progress_at.entry(target.to_string()).or_default();
if elapsed_secs < *next {
return;
}
sink.verify_wait(&VerifyWaitProgress {
target: target.to_string(),
destination,
state: state.to_string(),
elapsed_secs,
remaining_secs,
});
*next = elapsed_secs.saturating_add(DELEGATED_RELEASE_VERIFY_PROGRESS_INTERVAL_SECS);
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum VerifyMode {
CompletionBarrier,
ObserveAfterDistFailure,
}
#[derive(Debug)]
enum DelegatedVerifyFailure {
Pending(String),
Failed(String),
Unknown(String),
}
fn delegated_failure(
status: DelegatedRunStatus,
detail: Option<String>,
) -> Option<DelegatedVerifyFailure> {
match status {
DelegatedRunStatus::Success => None,
DelegatedRunStatus::Pending => Some(DelegatedVerifyFailure::Pending(
detail.unwrap_or_else(|| "the delegated workflow is still pending".to_string()),
)),
DelegatedRunStatus::Failed => Some(DelegatedVerifyFailure::Failed(detail.unwrap_or_else(
|| "the delegated workflow ended with a terminal failure".to_string(),
))),
DelegatedRunStatus::Unknown => Some(DelegatedVerifyFailure::Unknown(
detail
.unwrap_or_else(|| "the delegated workflow run could not be observed".to_string()),
)),
}
}
fn wait_for_delegated_runs(
ctx: &EffectCtx<'_>,
targets: &[TargetPlan],
delegated: &std::collections::BTreeSet<String>,
version: &str,
window: &mut DelegatedVerifyWindow,
sink: &mut dyn ProgressSink,
) -> Option<Result<(), (String, DelegatedVerifyFailure)>> {
let mut seen = HashSet::new();
let owners: Vec<(String, Adapter)> = targets
.iter()
.filter(|target| delegated.contains(&target.id))
.filter(|target| {
matches!(
target.input.target.adapter,
Adapter::CargoDist | Adapter::CargoPublishCi
)
})
.filter(|target| seen.insert(target.input.target.adapter))
.map(|target| (target.id.clone(), target.input.target.adapter))
.collect();
if owners.is_empty() {
return None;
}
let mut ready = HashSet::new();
loop {
let mut pending = Vec::new();
let mut first_unknown = None;
for (target, adapter) in &owners {
if ready.contains(adapter) {
continue;
}
let Some(run) = super::delegated::observe_github_run(ctx, *adapter, version) else {
continue;
};
match delegated_failure(run.status, run.detail) {
None => {
ready.insert(*adapter);
}
Some(failure @ DelegatedVerifyFailure::Failed(_)) => {
return Some(Err((target.clone(), failure)));
}
Some(failure @ DelegatedVerifyFailure::Unknown(_)) => {
first_unknown.get_or_insert_with(|| (target.clone(), failure));
}
Some(failure @ DelegatedVerifyFailure::Pending(_)) => {
pending.push((target.clone(), *adapter, failure));
}
}
}
if let Some(unknown) = first_unknown {
return Some(Err(unknown));
}
if ready.len() == owners.len() {
return Some(Ok(()));
}
for (target, adapter, _) in &pending {
window.stream_wait(
sink,
ctx,
target,
format!("GitHub Actions workflow ({})", adapter.as_str()),
"pending",
);
}
if !window.sleep(ctx) {
let (target, _, failure) = pending
.into_iter()
.next()
.expect("an unresolved workflow owner is pending");
return Some(Err((target, failure)));
}
}
}
struct DelegatedDestinationState<'a> {
target: &'a TargetPlan,
outcome: VerifyOutcome,
ever_reachable: bool,
settled: bool,
}
fn delegated_destination_label(plan: &ReleasePlan, target: &AdapterTarget) -> String {
match (target.target.adapter, target.target.registry) {
(_, Registry::Homebrew) => format!(
"Homebrew tap {} formula {}@{}",
plan.homebrew_tap.as_deref().unwrap_or("<unconfigured>"),
target.package,
plan.version
),
(Adapter::CargoDist, _) => {
format!(
"GitHub Release v{} assets for {}",
plan.version, target.package
)
}
(Adapter::CargoPublishCi, _) | (_, Registry::Npm | Registry::Pypi | Registry::TestPypi) => {
format!(
"{} registry {}@{}",
target.target.registry.as_str(),
target.package,
target.version
)
}
_ => format!(
"GitHub Release v{} assets for {}",
plan.version, target.package
),
}
}
fn observe_delegated_destination_once(
ctx: &EffectCtx<'_>,
plan: &ReleasePlan,
state: &mut DelegatedDestinationState<'_>,
) {
let target = &state.target.input;
state.outcome = match (target.target.adapter, target.target.registry) {
(_, Registry::Homebrew) => {
let Some(tap) = plan.homebrew_tap.as_deref() else {
debug_assert!(false, "delegated Homebrew target planned without a tap");
state.outcome = VerifyOutcome::Unknown;
state.settled = true;
return;
};
super::adapters::homebrew::verify_tap_formula(
ctx,
tap,
&target.package,
&plan.version,
false,
Some(&plan.homebrew_platforms),
)
}
(Adapter::CargoDist, _) => {
observe_cargo_dist_github_release(ctx, &plan.version, &target.package)
}
(Adapter::CargoPublishCi, _) | (_, Registry::Npm | Registry::Pypi | Registry::TestPypi) => {
match ctx
.registry
.published_versions(target.ecosystem().as_str(), &target.package)
{
Ok(versions) => {
state.ever_reachable = true;
if versions.iter().any(|version| version == &target.version) {
VerifyOutcome::Matches
} else {
VerifyOutcome::Missing
}
}
Err(_) if state.ever_reachable => VerifyOutcome::Missing,
Err(_) => VerifyOutcome::Unknown,
}
}
_ => observe_cargo_dist_github_release(ctx, &plan.version, &target.package),
};
state.settled = matches!(
(target.target.registry, state.outcome),
(_, VerifyOutcome::Matches) | (Registry::GhReleases, VerifyOutcome::Conflicts)
);
}
fn verify_delegated_destinations(
ctx: &EffectCtx<'_>,
plan: &ReleasePlan,
targets: &[TargetPlan],
delegated: &std::collections::BTreeSet<String>,
verified: &std::collections::BTreeMap<String, VerifyOutcome>,
window: &mut DelegatedVerifyWindow,
sink: &mut dyn ProgressSink,
) -> std::collections::HashMap<String, VerifyOutcome> {
let mut states: Vec<DelegatedDestinationState<'_>> = targets
.iter()
.filter(|target| delegated.contains(&target.id))
.filter(|target| verified.get(&target.id) != Some(&VerifyOutcome::Matches))
.map(|target| DelegatedDestinationState {
target,
outcome: VerifyOutcome::Unknown,
ever_reachable: false,
settled: false,
})
.collect();
if states.is_empty() {
return std::collections::HashMap::new();
}
loop {
for state in states.iter_mut().filter(|state| !state.settled) {
observe_delegated_destination_once(ctx, plan, state);
}
if states.iter().all(|state| state.settled) {
break;
}
for state in states.iter().filter(|state| !state.settled) {
window.stream_wait(
sink,
ctx,
&state.target.id,
delegated_destination_label(plan, &state.target.input),
state.outcome.as_str(),
);
}
if !window.sleep(ctx) {
break;
}
}
states
.into_iter()
.map(|state| (state.target.id.clone(), state.outcome))
.collect()
}
#[allow(clippy::too_many_lines)] fn verify_phase(
journal: &mut Journal<'_>,
sink: &mut dyn ProgressSink,
ctx: &EffectCtx<'_>,
targets: &[TargetPlan],
plan: &ReleasePlan,
homebrew: Option<&HomebrewFormula>,
mode: VerifyMode,
) -> Result<(), CutError> {
let phase = Phase::Verify;
if phase_completed_ok(journal.state(), phase) {
return Ok(());
}
let mut verification_artifacts = verification_artifacts(plan);
verification_artifacts.homebrew = homebrew.cloned();
let verify_ctx = ctx.with_artifacts(&verification_artifacts);
record(journal, sink, EventKind::PhaseEntered { phase })?;
let mut window = DelegatedVerifyWindow::new(&verify_ctx);
if mode == VerifyMode::CompletionBarrier {
if let Some(Err((target, failure))) = wait_for_delegated_runs(
&verify_ctx,
targets,
&journal.state().delegated,
&plan.version,
&mut window,
sink,
) {
record(
journal,
sink,
EventKind::TargetVerified {
target: target.clone(),
outcome: VerifyOutcome::Unknown,
},
)?;
record(
journal,
sink,
EventKind::PhaseCompleted {
phase,
outcome: PhaseOutcome::Failed,
},
)?;
return match failure {
DelegatedVerifyFailure::Pending(message) => {
Err(CutError::DelegatedRunPending { target, message })
}
DelegatedVerifyFailure::Failed(message) => {
Err(CutError::DelegatedRunFailed { target, message })
}
DelegatedVerifyFailure::Unknown(message) => Err(CutError::PhaseFailed {
phase,
target: Some(target),
message,
}),
};
}
}
let delegated_outcomes = if mode == VerifyMode::CompletionBarrier {
verify_delegated_destinations(
&verify_ctx,
plan,
targets,
&journal.state().delegated,
&journal.state().verified,
&mut window,
sink,
)
} else {
std::collections::HashMap::new()
};
let mut first_failure: Option<(String, DelegatedVerifyFailure)> = None;
for tp in targets {
if journal.state().verified.get(&tp.id) == Some(&VerifyOutcome::Matches) {
continue;
}
let is_delegated = journal.state().delegated.contains(&tp.id);
let delegated_failure = (mode == VerifyMode::ObserveAfterDistFailure && is_delegated)
.then(|| {
super::delegated::observe_github_run(
&verify_ctx,
tp.input.target.adapter,
&plan.version,
)
})
.flatten()
.and_then(|run| delegated_failure(run.status, run.detail));
let outcome = if delegated_failure.is_some() {
VerifyOutcome::Unknown
} else if is_delegated {
if mode == VerifyMode::CompletionBarrier {
delegated_outcomes
.get(&tp.id)
.copied()
.unwrap_or(VerifyOutcome::Unknown)
} else {
let mut state = DelegatedDestinationState {
target: tp,
outcome: VerifyOutcome::Unknown,
ever_reachable: false,
settled: false,
};
observe_delegated_destination_once(&verify_ctx, plan, &mut state);
state.outcome
}
} else if let Some(receipt) = journal.state().published.get(&tp.id) {
let receipt = AdapterReceipt {
adapter: tp.input.target.adapter,
ecosystem: tp.input.target.ecosystem,
package: tp.input.package.clone(),
version: receipt.version.clone(),
canonical_ref: tp.input.canonical_ref(),
digest: receipt.digest.clone(),
remote_url: receipt.registry_url.clone(),
timestamp: 0,
};
tp.adapter
.verify(&verify_ctx, &receipt)
.unwrap_or(VerifyOutcome::Unknown)
} else {
VerifyOutcome::Missing
};
record(
journal,
sink,
EventKind::TargetVerified {
target: tp.id.clone(),
outcome,
},
)?;
let failure = delegated_failure.or_else(|| match outcome {
VerifyOutcome::Matches => None,
VerifyOutcome::Unknown => Some(DelegatedVerifyFailure::Unknown(format!(
"could not observe {} at its destination",
tp.id
))),
VerifyOutcome::Missing => Some(DelegatedVerifyFailure::Unknown(format!(
"{} is missing at its destination",
tp.id
))),
VerifyOutcome::Conflicts => Some(DelegatedVerifyFailure::Unknown(format!(
"{} conflicts with its recorded receipt",
tp.id
))),
});
if first_failure.is_none() {
if let Some(failure) = failure {
first_failure = Some((tp.id.clone(), failure));
}
}
}
match (mode, first_failure) {
(VerifyMode::CompletionBarrier, None) => {
record(
journal,
sink,
EventKind::PhaseCompleted {
phase,
outcome: PhaseOutcome::Ok,
},
)?;
Ok(())
}
(VerifyMode::ObserveAfterDistFailure, Some((target, failure))) => fail_phase(
journal,
sink,
phase,
Some(target),
match failure {
DelegatedVerifyFailure::Pending(message)
| DelegatedVerifyFailure::Failed(message)
| DelegatedVerifyFailure::Unknown(message) => message,
},
),
(VerifyMode::CompletionBarrier, Some((target, failure))) => {
record(
journal,
sink,
EventKind::PhaseCompleted {
phase,
outcome: PhaseOutcome::Failed,
},
)?;
match failure {
DelegatedVerifyFailure::Pending(message) => {
Err(CutError::DelegatedRunPending { target, message })
}
DelegatedVerifyFailure::Failed(message) => {
Err(CutError::DelegatedRunFailed { target, message })
}
DelegatedVerifyFailure::Unknown(message) => Err(CutError::PhaseFailed {
phase,
target: Some(target),
message,
}),
}
}
(VerifyMode::ObserveAfterDistFailure, None) => fail_phase(
journal,
sink,
phase,
None,
"all declared destinations match, but the preceding dist barrier failed and the run remains resumable".to_string(),
),
}
}
fn tag_archive_url(slug: &str, version: &str) -> String {
format!("https://github.com/{slug}/archive/refs/tags/v{version}.tar.gz")
}
const TAG_ARCHIVE_FETCH_ATTEMPTS: u32 = 5;
const TAG_ARCHIVE_FETCH_BACKOFF: std::time::Duration = std::time::Duration::from_secs(3);
fn compute_source_tarball_sha256(ctx: &EffectCtx<'_>, url: &str) -> Result<String, String> {
let tmp = source_tarball_tmp_path();
let tmp_str = tmp.to_string_lossy().to_string();
let result = fetch_and_hash(ctx, url, &tmp_str);
let _ = ctx.runner.run("rm", &["-f", &tmp_str], ctx.repo_root);
result
}
fn fetch_and_hash(ctx: &EffectCtx<'_>, url: &str, tmp: &str) -> Result<String, String> {
fetch_tag_archive(ctx, url, tmp)?;
hash_file(ctx, tmp)
}
fn fetch_tag_archive(ctx: &EffectCtx<'_>, url: &str, tmp: &str) -> Result<(), String> {
let mut last = String::new();
for attempt in 0..TAG_ARCHIVE_FETCH_ATTEMPTS {
let out = ctx
.runner
.run("curl", &["-sSfL", "-o", tmp, "--", url], ctx.repo_root)
.map_err(|e| format!("cannot run `curl` to fetch the source tarball `{url}`: {e}"))?;
if out.status == Some(0) {
return Ok(());
}
last = format!(
"exit {}: {}",
out.status
.map_or_else(|| "signal".to_string(), |c| c.to_string()),
out.stderr.trim()
);
if attempt + 1 < TAG_ARCHIVE_FETCH_ATTEMPTS {
ctx.clock.sleep(TAG_ARCHIVE_FETCH_BACKOFF);
}
}
Err(format!(
"`curl` could not fetch the source tarball `{url}` after {TAG_ARCHIVE_FETCH_ATTEMPTS} \
attempts ({last}); the tag archive may not be published yet"
))
}
const RELEASE_ASSET_WAIT_TIMEOUT_SECS: u64 = 300;
const RELEASE_ASSET_POLL_INTERVAL: std::time::Duration = std::time::Duration::from_secs(3);
fn fetch_homebrew_assets(
ctx: &EffectCtx<'_>,
slug: &str,
plan: &ReleasePlan,
formula: &HomebrewFormula,
targets: &[&TargetPlan],
) -> Result<Vec<HomebrewAsset>, String> {
let package = targets
.first()
.ok_or_else(|| "homebrew dist has no target".to_string())?
.input
.package
.as_str();
let mut assets = Vec::new();
for triple in formula.platforms.iter().filter(|triple| {
crate::release::adapters::homebrew::homebrew_platform_condition(triple).is_some()
}) {
let filename = format!("{package}-{triple}.tar.xz");
let url = format!(
"https://github.com/{slug}/releases/download/v{}/{filename}",
plan.version
);
let tmp = std::env::temp_dir().join(format!(
"shipshape-homebrew-{filename}-{}",
std::process::id()
));
let tmp_str = tmp.to_string_lossy().to_string();
let start = ctx.clock.now_unix();
#[allow(unused_assignments)]
let mut last = String::new();
let sha256_result = loop {
let out = ctx.runner.run("curl", &["-sSfL", "-o", &tmp_str, "--", &url], ctx.repo_root)
.map_err(|e| format!("cannot run `curl` while waiting for Homebrew release asset `{filename}`: {e}"))?;
if out.status == Some(0) {
break hash_file(ctx, &tmp_str)
.map_err(|e| format!("cannot hash Homebrew release asset `{filename}`: {e}"));
}
last = format!(
"exit {}: {}",
out.status
.map_or_else(|| "signal".to_string(), |c| c.to_string()),
out.stderr.trim()
);
let waited = ctx.clock.now_unix().saturating_sub(start);
if waited >= RELEASE_ASSET_WAIT_TIMEOUT_SECS {
break Err(format!("Homebrew release asset `{filename}` was not visible after {waited}s (bounded release-asset wait; cargo-dist CI may have failed or not uploaded it): {last}. Refusing to write a source-build or unchecked formula"));
}
ctx.clock.sleep(RELEASE_ASSET_POLL_INTERVAL);
};
let _ = ctx.runner.run("rm", &["-f", &tmp_str], ctx.repo_root);
let sha256 = sha256_result?;
assets.push(HomebrewAsset {
triple: triple.clone(),
url,
sha256,
});
}
if assets.is_empty() {
return Err("Homebrew formula has no Homebrew-servable cargo-dist platforms; supported platforms are macOS aarch64/x86_64 and Linux musl aarch64/x86_64; refusing to write a formula with no installable archive".to_string());
}
Ok(assets)
}
fn source_tarball_tmp_path() -> std::path::PathBuf {
let nanos = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map_or(0, |d| d.as_nanos());
std::env::temp_dir().join(format!(
"shipshape-src-tarball-{}-{nanos}.tar.gz",
std::process::id()
))
}
fn fail_phase(
journal: &mut Journal<'_>,
sink: &mut dyn ProgressSink,
phase: Phase,
target: Option<String>,
message: String,
) -> Result<(), CutError> {
record(
journal,
sink,
EventKind::PhaseCompleted {
phase,
outcome: PhaseOutcome::Failed,
},
)?;
Err(CutError::PhaseFailed {
phase,
target,
message,
})
}
fn record(
journal: &mut Journal<'_>,
sink: &mut dyn ProgressSink,
kind: EventKind,
) -> Result<(), CutError> {
let idempotency_key = kind.idempotency_key();
let kind_for_sink = kind.clone();
let state = journal.append(kind).map_err(CutError::Journal)?;
let event = JournalEvent {
schema_version: JOURNAL_SCHEMA_VERSION,
seq: state.applied_seq,
ts: state.updated_ts,
idempotency_key,
kind: kind_for_sink,
};
sink.event(&event);
Ok(())
}
fn phase_completed_ok(state: &RunState, phase: Phase) -> bool {
state
.phases
.iter()
.any(|r| r.phase == phase && r.outcome == PhaseOutcome::Ok)
}
fn target_cleared(state: &RunState, phase: Phase, target: &str) -> bool {
match phase {
Phase::DryRun => state.dry_run.contains(target),
Phase::Build => state.built.contains(target),
Phase::Publish | Phase::Dist => state.published.contains_key(target),
Phase::Verify => state.verified.get(target) == Some(&VerifyOutcome::Matches),
Phase::Bump | Phase::Tag | Phase::AdvanceBranch => false,
}
}
fn tag_step_done(
state: &RunState,
tag: &str,
pick: impl Fn(&crate::protocol::journal::TagState) -> bool,
) -> bool {
state.tags.get(tag).is_some_and(pick)
}
fn to_journal_receipt(r: &AdapterReceipt) -> JournalReceipt {
JournalReceipt {
ecosystem: r.ecosystem.as_str().to_string(),
package: Some(r.package.clone()),
version: r.version.clone(),
registry_url: r.remote_url.clone(),
digest: r.digest.clone(),
}
}
#[cfg(test)]
mod tests;