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::release::PublishReceipt as AdapterReceipt;
use super::adapters::{
hash_file, resolve, AdapterTarget, EcosystemAdapter, EffectCtx, HomebrewFormula,
ReleaseAdapter, ReleaseArtifacts, SourceTarball,
};
use super::journal::Journal;
use super::journal_target_ids;
use crate::contract::schema::{Adapter, Target};
pub trait ProgressSink {
fn event(&mut self, event: &JournalEvent);
}
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),
}
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}"
)
}
}
}
}
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> {
let targets = resolve_target_plans(plan)?;
let repo_slug = resolve_repo_slug(ctx, &targets);
let checkout = SealedCheckout::materialize(ctx, &plan.head_sha)?;
let checkout_ctx = ctx.with_repo_root(checkout.path());
let ctx = &checkout_ctx;
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(),
};
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,
};
publish_phase(journal, sink, &ctx.with_artifacts(&artifacts), &targets)?;
let release_owner = github_release_owner(&targets);
tag_phase(journal, sink, tagger, plan, release_owner.as_deref())?;
dist_phase(
journal,
sink,
ctx,
&targets,
plan,
artifacts.repo_slug.as_deref(),
artifacts.homebrew.as_ref(),
)?;
Ok(())
}
pub fn validate_plan(plan: &ReleasePlan) -> Result<(), CutError> {
resolve_target_plans(plan).map(|_| ())
}
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!("ossctl-cut-{}-{short}-{nanos}", std::process::id()))
}
fn github_release_owner(targets: &[TargetPlan]) -> Option<String> {
targets
.iter()
.find(|tp| tp.adapter.ci_owns_github_release())
.map(|tp| tp.input.target.adapter.as_str().to_string())
}
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(),
})
}
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::Publish | Phase::Tag | Phase::Dist => {
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 tag_phase(
journal: &mut Journal<'_>,
sink: &mut dyn ProgressSink,
tagger: &dyn Tagger,
plan: &ReleasePlan,
release_owner: Option<&str>,
) -> 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, &plan.head_sha, &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() },
)?;
}
if let Some(adapter) = release_owner {
if tag_step_done(journal.state(), &tag, |s| s.github_release) {
return fail_phase(
journal,
sink,
phase,
None,
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"
),
);
}
} else if tag_step_done(journal.state(), &tag, |s| s.github_release_delegated) {
return fail_phase(
journal,
sink,
phase,
None,
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"
),
);
}
if let Some(adapter) = release_owner {
if !tag_step_done(journal.state(), &tag, |s| s.github_release_delegated) {
record(
journal,
sink,
EventKind::GithubReleaseDelegated {
tag: tag.clone(),
delegated_to: adapter.to_string(),
},
)?;
}
} else 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.clone(),
url,
},
)?,
Err(e) => {
return fail_phase(
journal,
sink,
phase,
None,
format!("create GitHub Release: {e}"),
)
}
}
}
record(
journal,
sink,
EventKind::PhaseCompleted {
phase,
outcome: PhaseOutcome::Ok,
},
)?;
Ok(())
}
fn needs_post_tag(tp: &TargetPlan) -> bool {
matches!(
tp.input.target.adapter,
Adapter::HomebrewTap | Adapter::HomebrewCore
)
}
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 artifacts = ReleaseArtifacts {
assets: Vec::new(),
source_tarball,
repo_slug: repo_slug.map(str::to_string),
homebrew: homebrew.cloned(),
};
let dist_ctx = ctx.with_artifacts(&artifacts);
for tp in post_tag {
if journal.state().published.contains_key(&tp.id) {
continue;
}
match tp.adapter.publish(&dist_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 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"
))
}
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!(
"ossctl-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::Tag => 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;