use std::collections::{HashMap, VecDeque};
use std::sync::Mutex;
use async_trait::async_trait;
use super::*;
use crate::upstream::classify::MergeRepositoryState;
use crate::upstream::ports::{
MergeCommandResult, PushCommandResult, RecoveryCommit, RepairAttemptResult, VerificationOutcome,
};
use crate::upstream::publication::format_publication_marker_message;
use crate::upstream::spine::{CommitTreeEvidence, SpineCommit};
use crate::upstream::trailers::format_upstream_merge_message;
#[derive(Debug, Clone)]
struct FakeCommit {
message: String,
parents: Vec<String>,
tree_evidence: CommitTreeEvidence,
}
#[derive(Debug, Clone, PartialEq, Eq)]
enum MergeBehavior {
Succeed,
Conflict,
CommandFailure,
}
#[derive(Debug, Clone)]
enum PushBehavior {
Succeed,
Reject { porcelain: String },
FailWithoutRefStatus,
}
#[derive(Debug, Default)]
struct FakeGitInner {
commits: HashMap<String, FakeCommit>,
head: String,
remote_refs: HashMap<String, String>,
tracking_refs: HashMap<String, String>,
configured_remotes: Vec<String>,
current_branch: Option<String>,
merge_state: MergeRepositoryState,
working_tree_clean: bool,
porcelain_v2: String,
merge_behaviors: VecDeque<MergeBehavior>,
push_behaviors: VecDeque<PushBehavior>,
fetch_error: Option<String>,
fetch_calls: usize,
merge_calls: Vec<(String, String)>,
push_calls: usize,
successful_pushes: usize,
empty_commits: Vec<String>,
next_sha: usize,
evidence_walk_calls: usize,
metadata_walk_calls: usize,
deny_evidence_walk: bool,
}
#[derive(Default)]
struct FakeGit {
inner: Mutex<FakeGitInner>,
}
fn sha(seed: &str) -> String {
let mut s: String = seed.bytes().map(|b| format!("{:02x}", b)).collect();
while s.len() < 40 {
s.push('0');
}
s.truncate(40);
s
}
impl FakeGit {
fn new_linear() -> Self {
let mut inner = FakeGitInner {
working_tree_clean: true,
current_branch: Some("main".to_string()),
configured_remotes: vec!["origin".to_string()],
next_sha: 100,
..Default::default()
};
inner.commits.insert(
sha("root"),
FakeCommit {
message: "root\n".into(),
parents: vec![],
tree_evidence: CommitTreeEvidence::default(),
},
);
inner.commits.insert(
sha("local"),
FakeCommit {
message: "Merge change: a\n".into(),
parents: vec![sha("root"), sha("wt")],
tree_evidence: CommitTreeEvidence::new(["a".to_string()], []),
},
);
inner.commits.insert(
sha("wt"),
FakeCommit {
message: "work\n".into(),
parents: vec![sha("root")],
tree_evidence: CommitTreeEvidence::default(),
},
);
inner.head = sha("local");
inner.remote_refs.insert("origin/main".into(), sha("root"));
Self {
inner: Mutex::new(inner),
}
}
fn lock(&self) -> std::sync::MutexGuard<'_, FakeGitInner> {
self.inner.lock().unwrap_or_else(|e| e.into_inner())
}
fn advance_remote(&self, seed: &str, parent: &str) -> String {
let mut inner = self.lock();
let new = sha(seed);
inner.commits.insert(
new.clone(),
FakeCommit {
message: format!("{}\n", seed),
parents: vec![parent.to_string()],
tree_evidence: CommitTreeEvidence::default(),
},
);
inner.remote_refs.insert("origin/main".into(), new.clone());
new
}
fn set_remote_ref(&self, value: &str) {
self.lock()
.remote_refs
.insert("origin/main".into(), value.to_string());
}
fn head(&self) -> String {
self.lock().head.clone()
}
fn merge_calls(&self) -> Vec<(String, String)> {
self.lock().merge_calls.clone()
}
fn fetch_calls(&self) -> usize {
self.lock().fetch_calls
}
fn successful_pushes(&self) -> usize {
self.lock().successful_pushes
}
fn push_calls(&self) -> usize {
self.lock().push_calls
}
fn empty_commits(&self) -> Vec<String> {
self.lock().empty_commits.clone()
}
fn deny_evidence_walk(&self) {
self.lock().deny_evidence_walk = true;
}
fn evidence_walk_calls(&self) -> usize {
self.lock().evidence_walk_calls
}
fn metadata_walk_calls(&self) -> usize {
self.lock().metadata_walk_calls
}
fn extend_linear_history(&self, count: usize) {
let mut inner = self.lock();
for _ in 0..count {
inner.next_sha += 1;
let new = sha(&format!("filler{}", inner.next_sha));
let head = inner.head.clone();
inner.commits.insert(
new.clone(),
FakeCommit {
message: "filler\n".into(),
parents: vec![head],
tree_evidence: CommitTreeEvidence::default(),
},
);
inner.head = new;
}
}
fn integrate_change(&self, change_id: &str) -> String {
let mut inner = self.lock();
inner.next_sha += 1;
let branch = sha(&format!("wt-{}", change_id));
let branch_parent = inner.head.clone();
inner.commits.insert(
branch.clone(),
FakeCommit {
message: format!("work for {}\n", change_id),
parents: vec![branch_parent],
tree_evidence: CommitTreeEvidence::default(),
},
);
let new = sha(&format!("base{}", inner.next_sha));
let head = inner.head.clone();
inner.commits.insert(
new.clone(),
FakeCommit {
message: format!("Merge change: {}\n", change_id),
parents: vec![head, branch],
tree_evidence: CommitTreeEvidence::new([change_id.to_string()], []),
},
);
inner.head = new.clone();
new
}
fn walk_first_parent(
inner: &FakeGitInner,
from_exclusive: Option<&str>,
to: &str,
limit: Option<usize>,
) -> Vec<(String, FakeCommit)> {
let mut collected = Vec::new();
let mut current = Some(to.to_string());
while let Some(sha_value) = current {
if Some(sha_value.as_str()) == from_exclusive {
break;
}
if let Some(max) = limit {
if collected.len() >= max {
break;
}
}
let Some(commit) = inner.commits.get(&sha_value) else {
break;
};
collected.push((sha_value.clone(), commit.clone()));
current = commit.parents.first().cloned();
}
collected.reverse();
collected
}
fn is_ancestor_locked(inner: &FakeGitInner, ancestor: &str, descendant: &str) -> bool {
if ancestor == descendant {
return true;
}
let mut stack = vec![descendant.to_string()];
let mut seen = std::collections::HashSet::new();
while let Some(current) = stack.pop() {
if !seen.insert(current.clone()) {
continue;
}
let Some(commit) = inner.commits.get(¤t) else {
continue;
};
for parent in &commit.parents {
if parent == ancestor {
return true;
}
stack.push(parent.clone());
}
}
false
}
}
#[async_trait]
impl UpstreamGit for FakeGit {
async fn remote_configured(&self, remote: &str) -> PortResult<bool> {
Ok(self.lock().configured_remotes.iter().any(|r| r == remote))
}
async fn current_branch(&self) -> PortResult<Option<String>> {
Ok(self.lock().current_branch.clone())
}
async fn fetch(&self, remote: &str, branch: &str) -> PortResult<()> {
let mut inner = self.lock();
inner.fetch_calls += 1;
if let Some(err) = inner.fetch_error.clone() {
return Err(UpstreamPortError::new("git fetch", err));
}
let key = format!("{}/{}", remote, branch);
if let Some(value) = inner.remote_refs.get(&key).cloned() {
inner.tracking_refs.insert(key, value);
} else {
inner.tracking_refs.remove(&key);
}
Ok(())
}
async fn fetched_sha(&self, remote: &str, branch: &str) -> PortResult<Option<String>> {
Ok(self
.lock()
.tracking_refs
.get(&format!("{}/{}", remote, branch))
.cloned())
}
async fn head_sha(&self) -> PortResult<String> {
Ok(self.lock().head.clone())
}
async fn is_ancestor(&self, ancestor: &str, descendant: &str) -> PortResult<bool> {
let inner = self.lock();
Ok(FakeGit::is_ancestor_locked(&inner, ancestor, descendant))
}
async fn merge_base(&self, a: &str, b: &str) -> PortResult<String> {
let inner = self.lock();
if FakeGit::is_ancestor_locked(&inner, a, b) {
return Ok(a.to_string());
}
if FakeGit::is_ancestor_locked(&inner, b, a) {
return Ok(b.to_string());
}
Ok(sha("root"))
}
async fn merge_no_ff(&self, target: &str, message: &str) -> PortResult<MergeCommandResult> {
let mut inner = self.lock();
inner
.merge_calls
.push((target.to_string(), message.to_string()));
let behavior = inner
.merge_behaviors
.pop_front()
.unwrap_or(MergeBehavior::Succeed);
match behavior {
MergeBehavior::Succeed => {
inner.next_sha += 1;
let new = sha(&format!("merge{}", inner.next_sha));
let head = inner.head.clone();
inner.commits.insert(
new.clone(),
FakeCommit {
message: message.to_string(),
parents: vec![head, target.to_string()],
tree_evidence: CommitTreeEvidence::default(),
},
);
inner.head = new;
inner.merge_state = MergeRepositoryState {
merge_head_present: false,
has_unmerged_entries: false,
};
Ok(MergeCommandResult {
exit_success: true,
state: inner.merge_state,
})
}
MergeBehavior::Conflict => {
inner.merge_state = MergeRepositoryState {
merge_head_present: true,
has_unmerged_entries: true,
};
inner.working_tree_clean = false;
Ok(MergeCommandResult {
exit_success: false,
state: inner.merge_state,
})
}
MergeBehavior::CommandFailure => {
inner.merge_state = MergeRepositoryState {
merge_head_present: false,
has_unmerged_entries: false,
};
Ok(MergeCommandResult {
exit_success: false,
state: inner.merge_state,
})
}
}
}
async fn commit_empty(&self, message: &str) -> PortResult<String> {
let mut inner = self.lock();
inner.next_sha += 1;
let new = sha(&format!("marker{}", inner.next_sha));
let head = inner.head.clone();
inner.commits.insert(
new.clone(),
FakeCommit {
message: message.to_string(),
parents: vec![head],
tree_evidence: CommitTreeEvidence::default(),
},
);
inner.head = new.clone();
inner.empty_commits.push(message.to_string());
Ok(new)
}
async fn merge_repository_state(&self) -> PortResult<MergeRepositoryState> {
Ok(self.lock().merge_state)
}
async fn is_working_tree_clean(&self) -> PortResult<bool> {
Ok(self.lock().working_tree_clean)
}
async fn status_porcelain_v2(&self) -> PortResult<String> {
Ok(self.lock().porcelain_v2.clone())
}
async fn commit_message(&self, commit: &str) -> PortResult<String> {
self.lock()
.commits
.get(commit)
.map(|c| c.message.clone())
.ok_or_else(|| UpstreamPortError::new("git log", format!("unknown commit {}", commit)))
}
async fn commit_parents(&self, commit: &str) -> PortResult<Vec<String>> {
self.lock()
.commits
.get(commit)
.map(|c| c.parents.clone())
.ok_or_else(|| {
UpstreamPortError::new("git rev-list", format!("unknown commit {}", commit))
})
}
async fn first_parent_recovery_metadata(
&self,
to: &str,
limit: Option<usize>,
) -> PortResult<Vec<RecoveryCommit>> {
let mut inner = self.lock();
inner.metadata_walk_calls += 1;
Ok(Self::walk_first_parent(&inner, None, to, limit)
.into_iter()
.map(|(sha_value, commit)| RecoveryCommit {
sha: sha_value,
message: commit.message,
parents: commit.parents,
})
.collect())
}
async fn first_parent_commits(
&self,
from_exclusive: Option<&str>,
to: &str,
limit: Option<usize>,
) -> PortResult<Vec<SpineCommit>> {
let mut inner = self.lock();
inner.evidence_walk_calls += 1;
if inner.deny_evidence_walk {
return Err(UpstreamPortError::new(
"git ls-tree",
"evidence-bearing first-parent walk is forbidden on this path",
));
}
Ok(Self::walk_first_parent(&inner, from_exclusive, to, limit)
.into_iter()
.map(|(sha_value, commit)| SpineCommit {
sha: sha_value,
message: commit.message,
parents: commit.parents,
tree_evidence: commit.tree_evidence,
})
.collect())
}
async fn local_ref_sha(&self, reference: &str) -> PortResult<Option<String>> {
let key = reference
.strip_prefix("refs/remotes/")
.unwrap_or(reference)
.to_string();
Ok(self.lock().tracking_refs.get(&key).cloned())
}
async fn push_porcelain(&self, remote: &str, branch: &str) -> PortResult<PushCommandResult> {
let mut inner = self.lock();
inner.push_calls += 1;
let behavior = inner
.push_behaviors
.pop_front()
.unwrap_or(PushBehavior::Succeed);
match behavior {
PushBehavior::Succeed => {
inner.successful_pushes += 1;
let head = inner.head.clone();
inner
.remote_refs
.insert(format!("{}/{}", remote, branch), head);
Ok(PushCommandResult {
exit_success: true,
porcelain_stdout:
"To fake\n \trefs/heads/main:refs/heads/main\tabc..def\nDone\n".to_string(),
})
}
PushBehavior::Reject { porcelain } => Ok(PushCommandResult {
exit_success: false,
porcelain_stdout: porcelain,
}),
PushBehavior::FailWithoutRefStatus => Ok(PushCommandResult {
exit_success: false,
porcelain_stdout: String::new(),
}),
}
}
async fn ls_remote_sha(&self, remote: &str, branch: &str) -> PortResult<Option<String>> {
Ok(self
.lock()
.remote_refs
.get(&format!("{}/{}", remote, branch))
.cloned())
}
}
#[derive(Default)]
struct FakeVerifier {
results: Mutex<VecDeque<bool>>,
calls: Mutex<usize>,
}
impl FakeVerifier {
fn always_pass() -> Self {
Self::default()
}
fn scripted(results: impl IntoIterator<Item = bool>) -> Self {
Self {
results: Mutex::new(results.into_iter().collect()),
calls: Mutex::new(0),
}
}
fn calls(&self) -> usize {
*self.calls.lock().unwrap()
}
}
#[async_trait]
impl UpstreamVerifier for FakeVerifier {
async fn verify(&self) -> PortResult<VerificationOutcome> {
*self.calls.lock().unwrap() += 1;
let pass = self.results.lock().unwrap().pop_front().unwrap_or(true);
Ok(if pass {
VerificationOutcome::passed()
} else {
VerificationOutcome::failed("assertion failed")
})
}
}
type RepairAction = Box<dyn Fn(&FakeGit) + Send>;
#[derive(Default)]
struct FakeRepairAgent {
max_attempts: u32,
calls: Mutex<Vec<RepairRequest>>,
on_repair: Mutex<Option<RepairAction>>,
git: Mutex<Option<std::sync::Weak<FakeGit>>>,
}
impl FakeRepairAgent {
fn with_attempts(max_attempts: u32) -> Self {
Self {
max_attempts,
..Default::default()
}
}
fn calls(&self) -> Vec<RepairRequest> {
self.calls.lock().unwrap().clone()
}
fn call_count(&self) -> usize {
self.calls.lock().unwrap().len()
}
fn bind(&self, git: &std::sync::Arc<FakeGit>, action: impl Fn(&FakeGit) + Send + 'static) {
*self.git.lock().unwrap() = Some(std::sync::Arc::downgrade(git));
*self.on_repair.lock().unwrap() = Some(Box::new(action));
}
}
#[async_trait]
impl UpstreamRepairAgent for FakeRepairAgent {
fn max_attempts(&self) -> u32 {
self.max_attempts
}
async fn repair(&self, request: &RepairRequest) -> PortResult<RepairAttemptResult> {
self.calls.lock().unwrap().push(request.clone());
let git = self.git.lock().unwrap().clone();
let action = self.on_repair.lock().unwrap();
if let (Some(weak), Some(action)) = (git, action.as_ref()) {
if let Some(git) = weak.upgrade() {
action(&git);
}
}
Ok(RepairAttemptResult {
command_success: true,
})
}
}
#[derive(Default)]
struct RecordingObserver {
events: Mutex<Vec<UpstreamEvent>>,
}
impl RecordingObserver {
fn events(&self) -> Vec<UpstreamEvent> {
self.events.lock().unwrap().clone()
}
}
#[async_trait]
impl UpstreamObserver for RecordingObserver {
async fn observe(&self, event: UpstreamEvent) {
self.events.lock().unwrap().push(event);
}
}
struct Harness {
git: Arc<FakeGit>,
verifier: Arc<FakeVerifier>,
repair: Arc<FakeRepairAgent>,
observer: Arc<RecordingObserver>,
coordinator: UpstreamCoordinator,
}
fn harness_with(verifier: FakeVerifier, repair: FakeRepairAgent) -> Harness {
let git = Arc::new(FakeGit::new_linear());
let verifier = Arc::new(verifier);
let repair = Arc::new(repair);
let observer = Arc::new(RecordingObserver::default());
let coordinator = UpstreamCoordinator::new(
UpstreamIntegrationConfig::new("origin", "cargo test"),
"main",
git.clone(),
verifier.clone(),
repair.clone(),
observer.clone(),
);
Harness {
git,
verifier,
repair,
observer,
coordinator,
}
}
fn harness() -> Harness {
harness_with(
FakeVerifier::always_pass(),
FakeRepairAgent::with_attempts(2),
)
}
#[tokio::test]
async fn upstream_integration_initial_fetch_rejects_missing_remote_branch() {
let h = harness();
h.git.lock().remote_refs.clear();
let result = h.coordinator.validate_initial_fetch().await.unwrap();
assert_eq!(
result,
Err(UpstreamOptionError::RemoteBranchMissing {
remote: "origin".into(),
branch: "main".into()
})
);
}
#[tokio::test]
async fn upstream_integration_initial_fetch_rejects_unconfigured_remote_and_detached_head() {
let h = harness();
h.git.lock().configured_remotes.clear();
assert_eq!(
h.coordinator.validate_initial_fetch().await.unwrap(),
Err(UpstreamOptionError::RemoteNotConfigured("origin".into()))
);
let h = harness();
h.git.lock().current_branch = None;
assert_eq!(
h.coordinator.validate_initial_fetch().await.unwrap(),
Err(UpstreamOptionError::DetachedHead)
);
}
#[tokio::test]
async fn upstream_integration_initial_fetch_accepts_cumulative_change_merge_history() {
let h = harness();
let validation = h
.coordinator
.validate_initial_fetch()
.await
.unwrap()
.unwrap();
assert_eq!(validation.branch, "main");
assert!(validation.spine.is_publishable());
assert_eq!(
validation.spine.integrated_change_ids,
vec!["a".to_string()]
);
}
#[tokio::test]
async fn upstream_integration_initial_fetch_rejects_unrelated_local_history() {
let h = harness();
{
let mut inner = h.git.lock();
let head = inner.head.clone();
inner.commits.insert(
sha("hack"),
FakeCommit {
message: "hotfix: local only\n".into(),
parents: vec![head],
tree_evidence: CommitTreeEvidence::default(),
},
);
inner.head = sha("hack");
}
let result = h.coordinator.validate_initial_fetch().await.unwrap();
assert!(matches!(
result,
Err(UpstreamOptionError::UnrelatedLocalHistory { .. })
));
}
#[tokio::test]
async fn upstream_integration_initial_fetch_accepts_upstream_recovery_commit() {
let h = harness();
{
let mut inner = h.git.lock();
let head = inner.head.clone();
let fetched = sha("root");
inner.commits.insert(
sha("upmerge"),
FakeCommit {
message: format_upstream_merge_message("origin", "main", &fetched),
parents: vec![head, fetched],
tree_evidence: CommitTreeEvidence::default(),
},
);
inner.head = sha("upmerge");
}
let validation = h
.coordinator
.validate_initial_fetch()
.await
.unwrap()
.unwrap();
assert!(validation.spine.is_publishable());
assert_eq!(validation.spine.upstream_merges.len(), 1);
}
#[tokio::test]
async fn upstream_integration_fetch_failure_leaves_worktree_untouched() {
let mut h = harness();
h.git.lock().fetch_error = Some("network unreachable".into());
let err = h
.coordinator
.checkpoint(
CheckpointTrigger::BeforeFirstDispatch,
&BaseLaneState::clean(),
None,
)
.await
.unwrap_err();
assert_eq!(err.operation, "git fetch");
assert!(h.git.merge_calls().is_empty());
assert_eq!(h.verifier.calls(), 0);
}
#[tokio::test]
async fn upstream_integration_already_integrated_revision_is_a_noop() {
let mut h = harness();
let outcome = h
.coordinator
.checkpoint(
CheckpointTrigger::BeforeFirstDispatch,
&BaseLaneState::clean(),
None,
)
.await
.unwrap();
assert_eq!(
outcome,
UpstreamStepOutcome::NoOp {
fetched_sha: sha("root")
}
);
assert!(h.git.merge_calls().is_empty());
assert_eq!(h.verifier.calls(), 0, "no-op must not reverify");
assert_eq!(h.repair.call_count(), 0, "no-op must not invoke an agent");
}
#[tokio::test]
async fn upstream_integration_remote_ahead_uses_no_ff_merge_with_trailers() {
let mut h = harness();
let advanced = h.git.advance_remote("remoteahead", &h.git.head());
let outcome = h
.coordinator
.checkpoint(CheckpointTrigger::AfterDrain, &BaseLaneState::clean(), None)
.await
.unwrap();
let calls = h.git.merge_calls();
assert_eq!(calls.len(), 1, "strictly remote-ahead history still merges");
assert_eq!(calls[0].0, advanced);
assert_eq!(
calls[0].1,
format_upstream_merge_message("origin", "main", &advanced)
);
assert!(matches!(outcome, UpstreamStepOutcome::Integrated { .. }));
assert_eq!(
h.repair.call_count(),
0,
"conflict-free merge starts no agent"
);
assert_eq!(h.verifier.calls(), 1, "changed tree runs full verification");
}
#[tokio::test]
async fn upstream_integration_diverged_revision_is_integrated() {
let mut h = harness();
let diverged = h.git.advance_remote("diverged", &sha("root"));
h.coordinator
.checkpoint(CheckpointTrigger::AfterDrain, &BaseLaneState::clean(), None)
.await
.unwrap();
assert_eq!(h.git.merge_calls()[0].0, diverged);
}
#[tokio::test]
async fn upstream_integration_command_failure_does_not_start_an_agent() {
let mut h = harness();
h.git.advance_remote("remoteahead", &h.git.head());
h.git
.lock()
.merge_behaviors
.push_back(MergeBehavior::CommandFailure);
let outcome = h
.coordinator
.checkpoint(CheckpointTrigger::AfterDrain, &BaseLaneState::clean(), None)
.await
.unwrap();
assert!(matches!(outcome, UpstreamStepOutcome::Stalled { .. }));
assert_eq!(h.repair.call_count(), 0);
assert_eq!(h.verifier.calls(), 0);
}
#[tokio::test]
async fn upstream_integration_dirty_base_defers_before_fetch() {
let mut h = harness();
let lane = BaseLaneState {
base_dirty_reason: Some("uncommitted changes".into()),
lane_busy_reason: None,
};
let outcome = h
.coordinator
.checkpoint(CheckpointTrigger::BeforeBaseIntegration, &lane, Some("c1"))
.await
.unwrap();
assert_eq!(
outcome,
UpstreamStepOutcome::Deferred {
reason: "uncommitted changes".into()
}
);
assert_eq!(
h.git.fetch_calls(),
0,
"deferred checkpoint performs no fetch"
);
assert_eq!(h.coordinator.queued_results(), vec!["c1".to_string()]);
}
#[tokio::test]
async fn upstream_integration_textual_conflict_converges_through_repair() {
let git = Arc::new(FakeGit::new_linear());
let advanced = git.advance_remote("remoteahead", &git.head());
git.lock()
.merge_behaviors
.push_back(MergeBehavior::Conflict);
let verifier = Arc::new(FakeVerifier::always_pass());
let repair = Arc::new(FakeRepairAgent::with_attempts(2));
let observer = Arc::new(RecordingObserver::default());
let advanced_for_repair = advanced.clone();
repair.bind(&git, move |git| {
let mut inner = git.lock();
let head = inner.head.clone();
let message = format_upstream_merge_message("origin", "main", &advanced_for_repair);
inner.commits.insert(
sha("repaired"),
FakeCommit {
message,
parents: vec![head, advanced_for_repair.clone()],
tree_evidence: CommitTreeEvidence::default(),
},
);
inner.head = sha("repaired");
inner.merge_state = MergeRepositoryState {
merge_head_present: false,
has_unmerged_entries: false,
};
inner.working_tree_clean = true;
});
let mut coordinator = UpstreamCoordinator::new(
UpstreamIntegrationConfig::new("origin", "cargo test"),
"main",
git.clone(),
verifier.clone(),
repair.clone(),
observer.clone(),
);
let outcome = coordinator
.checkpoint(CheckpointTrigger::AfterDrain, &BaseLaneState::clean(), None)
.await
.unwrap();
assert_eq!(
outcome,
UpstreamStepOutcome::Integrated {
merge_sha: sha("repaired")
}
);
assert_eq!(repair.call_count(), 1);
assert_eq!(repair.calls()[0].cause, RepairCause::TextualConflict);
assert_eq!(repair.calls()[0].fetched_sha, advanced);
assert_eq!(verifier.calls(), 1);
}
#[tokio::test]
async fn upstream_integration_textual_repair_exhaustion_stalls() {
let mut h = harness();
h.git.advance_remote("remoteahead", &h.git.head());
h.git
.lock()
.merge_behaviors
.push_back(MergeBehavior::Conflict);
let outcome = h
.coordinator
.checkpoint(CheckpointTrigger::AfterDrain, &BaseLaneState::clean(), None)
.await
.unwrap();
assert!(matches!(outcome, UpstreamStepOutcome::Stalled { .. }));
assert_eq!(h.repair.call_count(), 2, "retry budget is Conflux-owned");
assert_eq!(h.verifier.calls(), 0, "unconverged merge never reverifies");
}
#[tokio::test]
async fn upstream_integration_conflict_free_merge_with_semantic_failure_repairs_then_reruns() {
let git = Arc::new(FakeGit::new_linear());
git.advance_remote("remoteahead", &git.head());
let verifier = Arc::new(FakeVerifier::scripted([false, true]));
let repair = Arc::new(FakeRepairAgent::with_attempts(2));
let observer = Arc::new(RecordingObserver::default());
repair.bind(&git, |git| {
let mut inner = git.lock();
let head = inner.head.clone();
inner.next_sha += 1;
let new = sha(&format!("fix{}", inner.next_sha));
inner.commits.insert(
new.clone(),
FakeCommit {
message: "fix: repair semantic breakage\n".into(),
parents: vec![head],
tree_evidence: CommitTreeEvidence::default(),
},
);
inner.head = new;
inner.working_tree_clean = true;
});
let mut coordinator = UpstreamCoordinator::new(
UpstreamIntegrationConfig::new("origin", "cargo test"),
"main",
git.clone(),
verifier.clone(),
repair.clone(),
observer.clone(),
);
let outcome = coordinator
.checkpoint(CheckpointTrigger::AfterDrain, &BaseLaneState::clean(), None)
.await
.unwrap();
assert!(matches!(outcome, UpstreamStepOutcome::Integrated { .. }));
assert_eq!(repair.call_count(), 1);
assert_eq!(repair.calls()[0].cause, RepairCause::SemanticVerification);
assert_eq!(repair.calls()[0].verify_command, "cargo test");
assert!(!repair.calls()[0].verify_output_tail.is_empty());
assert_eq!(
verifier.calls(),
2,
"every repair attempt is followed by a mandatory rerun"
);
}
#[tokio::test]
async fn upstream_integration_semantic_repair_rejects_history_rewrite() {
let git = Arc::new(FakeGit::new_linear());
git.advance_remote("remoteahead", &git.head());
let verifier = Arc::new(FakeVerifier::scripted([false, true]));
let repair = Arc::new(FakeRepairAgent::with_attempts(2));
let observer = Arc::new(RecordingObserver::default());
repair.bind(&git, |git| {
let mut inner = git.lock();
inner.head = sha("root");
});
let mut coordinator = UpstreamCoordinator::new(
UpstreamIntegrationConfig::new("origin", "cargo test"),
"main",
git.clone(),
verifier.clone(),
repair.clone(),
observer.clone(),
);
let outcome = coordinator
.checkpoint(CheckpointTrigger::AfterDrain, &BaseLaneState::clean(), None)
.await
.unwrap();
match outcome {
UpstreamStepOutcome::Stalled { reason } => {
assert!(reason.contains("forward-only"), "reason: {}", reason)
}
other => panic!("expected stall, got {:?}", other),
}
assert_eq!(
verifier.calls(),
1,
"a rewritten history is never granted a verification rerun"
);
}
#[tokio::test]
async fn upstream_integration_semantic_repair_exhaustion_blocks_base_lane() {
let git = Arc::new(FakeGit::new_linear());
git.advance_remote("remoteahead", &git.head());
let verifier = Arc::new(FakeVerifier::scripted([false, false, false, false]));
let repair = Arc::new(FakeRepairAgent::with_attempts(2));
let observer = Arc::new(RecordingObserver::default());
repair.bind(&git, |git| {
let mut inner = git.lock();
let head = inner.head.clone();
inner.next_sha += 1;
let new = sha(&format!("fix{}", inner.next_sha));
inner.commits.insert(
new.clone(),
FakeCommit {
message: "fix attempt\n".into(),
parents: vec![head],
tree_evidence: CommitTreeEvidence::default(),
},
);
inner.head = new;
});
let mut coordinator = UpstreamCoordinator::new(
UpstreamIntegrationConfig::new("origin", "cargo test"),
"main",
git.clone(),
verifier.clone(),
repair.clone(),
observer.clone(),
);
let outcome = coordinator
.checkpoint(CheckpointTrigger::AfterDrain, &BaseLaneState::clean(), None)
.await
.unwrap();
assert!(matches!(outcome, UpstreamStepOutcome::Stalled { .. }));
assert_eq!(repair.call_count(), 2);
assert_eq!(verifier.calls(), 3);
}
#[tokio::test]
async fn upstream_integration_completed_result_runs_full_verification() {
let mut h = harness();
let outcome = h.coordinator.verify_base_result("change-a").await.unwrap();
assert!(matches!(outcome, UpstreamStepOutcome::Integrated { .. }));
assert_eq!(h.verifier.calls(), 1);
assert!(h.git.merge_calls().is_empty());
}
#[tokio::test]
async fn upstream_integration_completed_result_verification_failure_blocks_dispatch() {
let mut h = harness_with(
FakeVerifier::scripted([false]),
FakeRepairAgent::with_attempts(0),
);
let outcome = h.coordinator.verify_base_result("change-a").await.unwrap();
match outcome {
UpstreamStepOutcome::Stalled { reason } => assert!(reason.contains("change-a")),
other => panic!("expected stall, got {:?}", other),
}
assert_eq!(h.repair.call_count(), 0);
}
#[tokio::test]
async fn upstream_integration_scan_finds_unpushed_trailer_identified_merge() {
let git = FakeGit::new_linear();
{
let mut inner = git.lock();
let head = inner.head.clone();
let fetched = sha("root");
inner.commits.insert(
sha("upmerge"),
FakeCommit {
message: format_upstream_merge_message("origin", "main", &fetched),
parents: vec![head, fetched],
tree_evidence: CommitTreeEvidence::default(),
},
);
inner.head = sha("upmerge");
inner
.tracking_refs
.insert("origin/main".into(), sha("root"));
}
let evidence = scan_unpushed_upstream_merges(&git).await.unwrap();
assert_eq!(evidence.len(), 1);
assert_eq!(evidence[0].merge_sha, sha("upmerge"));
assert_eq!(evidence[0].trailers.remote, "origin");
let refusal = upstream_recovery_refusal(&evidence).unwrap();
let message = refusal.to_string();
assert!(message.contains("--integrate-upstream=origin"));
assert!(message.contains("--upstream-verify-command"));
}
#[tokio::test]
async fn upstream_integration_scan_ignores_published_and_untrailered_merges() {
let git = FakeGit::new_linear();
{
let mut inner = git.lock();
let head = inner.head.clone();
let fetched = sha("root");
inner.commits.insert(
sha("upmerge"),
FakeCommit {
message: format_upstream_merge_message("origin", "main", &fetched),
parents: vec![head, fetched],
tree_evidence: CommitTreeEvidence::default(),
},
);
inner.head = sha("upmerge");
inner
.tracking_refs
.insert("origin/main".into(), sha("upmerge"));
}
assert!(scan_unpushed_upstream_merges(&git)
.await
.unwrap()
.is_empty());
let plain = FakeGit::new_linear();
plain
.lock()
.tracking_refs
.insert("origin/main".into(), sha("root"));
assert!(scan_unpushed_upstream_merges(&plain)
.await
.unwrap()
.is_empty());
}
fn attach_unpushed_upstream_merge(git: &FakeGit, merged_parent: &str, recorded_parent: &str) {
let mut inner = git.lock();
let head = inner.head.clone();
inner.commits.insert(
sha("upmerge"),
FakeCommit {
message: format_upstream_merge_message("origin", "main", recorded_parent),
parents: vec![head, merged_parent.to_string()],
tree_evidence: CommitTreeEvidence::default(),
},
);
inner.head = sha("upmerge");
inner
.tracking_refs
.insert("origin/main".into(), sha("root"));
}
#[tokio::test]
async fn upstream_integration_recovery_discovery_never_reads_commit_trees() {
let git = FakeGit::new_linear();
attach_unpushed_upstream_merge(&git, &sha("root"), &sha("root"));
{
let mut inner = git.lock();
let head = inner.head.clone();
inner.commits.insert(
sha("marker"),
FakeCommit {
message: format_publication_marker_message("alpha", "origin", "main"),
parents: vec![head],
tree_evidence: CommitTreeEvidence::default(),
},
);
inner.head = sha("marker");
}
git.deny_evidence_walk();
let publications = scan_pending_publications(&git).await.unwrap();
assert_eq!(publications.len(), 1);
assert_eq!(publications[0].trailers.change_id, "alpha");
let merges = scan_unpushed_upstream_merges(&git).await.unwrap();
assert_eq!(merges.len(), 1);
assert_eq!(merges[0].merge_sha, sha("upmerge"));
assert_eq!(merges[0].trailers.remote, "origin");
assert_eq!(
git.evidence_walk_calls(),
0,
"recovery discovery must not request commit-tree evidence"
);
assert_eq!(git.metadata_walk_calls(), 2);
}
#[tokio::test]
async fn upstream_integration_recovery_rejects_contradicted_upstream_trailer() {
let git = FakeGit::new_linear();
attach_unpushed_upstream_merge(&git, &sha("wt"), &sha("root"));
git.deny_evidence_walk();
assert!(scan_unpushed_upstream_merges(&git)
.await
.unwrap()
.is_empty());
assert_eq!(git.evidence_walk_calls(), 0);
}
#[tokio::test]
async fn upstream_integration_recovery_ignores_evidence_incorporated_by_the_remote() {
let git = FakeGit::new_linear();
attach_unpushed_upstream_merge(&git, &sha("root"), &sha("root"));
{
let mut inner = git.lock();
let head = inner.head.clone();
inner.commits.insert(
sha("marker"),
FakeCommit {
message: format_publication_marker_message("alpha", "origin", "main"),
parents: vec![head.clone()],
tree_evidence: CommitTreeEvidence::default(),
},
);
inner.head = sha("marker");
inner
.tracking_refs
.insert("origin/main".into(), sha("marker"));
}
git.deny_evidence_walk();
assert!(scan_pending_publications(&git).await.unwrap().is_empty());
assert!(scan_unpushed_upstream_merges(&git)
.await
.unwrap()
.is_empty());
assert_eq!(git.evidence_walk_calls(), 0);
}
#[tokio::test]
async fn upstream_integration_recovery_discovery_stops_at_the_bounded_limit() {
let inside = FakeGit::new_linear();
attach_unpushed_upstream_merge(&inside, &sha("root"), &sha("root"));
inside.extend_linear_history(RECOVERY_SCAN_LIMIT - 1);
inside.deny_evidence_walk();
assert_eq!(
scan_unpushed_upstream_merges(&inside).await.unwrap().len(),
1,
"the oldest commit within the bound must still be scanned"
);
let outside = FakeGit::new_linear();
attach_unpushed_upstream_merge(&outside, &sha("root"), &sha("root"));
outside.extend_linear_history(RECOVERY_SCAN_LIMIT);
outside.deny_evidence_walk();
assert!(
scan_unpushed_upstream_merges(&outside)
.await
.unwrap()
.is_empty(),
"the bounded walk must not grow past its limit"
);
assert_eq!(outside.evidence_walk_calls(), 0);
}
#[tokio::test]
async fn upstream_integration_spine_validation_keeps_commit_tree_evidence() {
let publishable = FakeGit::new_linear();
let validation = validate_initial_fetch(&publishable, "origin", "main")
.await
.unwrap()
.expect("the linear fixture's change merge carries archive evidence");
assert_eq!(
validation.spine.integrated_change_ids,
vec!["a".to_string()]
);
assert_eq!(publishable.evidence_walk_calls(), 1);
assert_eq!(
publishable.metadata_walk_calls(),
0,
"spine validation must not fall back to the metadata-only walk"
);
let unproven = FakeGit::new_linear();
unproven
.lock()
.commits
.get_mut(&sha("local"))
.unwrap()
.tree_evidence = CommitTreeEvidence::default();
match validate_initial_fetch(&unproven, "origin", "main")
.await
.unwrap()
{
Err(UpstreamOptionError::UnrelatedLocalHistory { commit, .. }) => {
assert_eq!(commit, sha("local"));
}
other => panic!(
"expected rejection without archive evidence, got {:?}",
other
),
}
let still_active = FakeGit::new_linear();
still_active
.lock()
.commits
.get_mut(&sha("local"))
.unwrap()
.tree_evidence = CommitTreeEvidence::new(["a".to_string()], ["a".to_string()]);
assert!(matches!(
validate_initial_fetch(&still_active, "origin", "main")
.await
.unwrap(),
Err(UpstreamOptionError::UnrelatedLocalHistory { .. })
));
}
#[tokio::test]
async fn upstream_integration_restart_reruns_verification_for_unpushed_merge() {
let git = Arc::new(FakeGit::new_linear());
{
let mut inner = git.lock();
let head = inner.head.clone();
let fetched = sha("root");
inner.commits.insert(
sha("upmerge"),
FakeCommit {
message: format_upstream_merge_message("origin", "main", &fetched),
parents: vec![head, fetched],
tree_evidence: CommitTreeEvidence::default(),
},
);
inner.head = sha("upmerge");
}
let verifier = Arc::new(FakeVerifier::always_pass());
let mut coordinator = UpstreamCoordinator::new(
UpstreamIntegrationConfig::new("origin", "cargo test"),
"main",
git.clone(),
verifier.clone(),
Arc::new(FakeRepairAgent::with_attempts(1)),
Arc::new(RecordingObserver::default()),
);
let outcome = coordinator
.finalize(SchedulerOutcome::DrainedSuccessfully)
.await
.unwrap();
assert_eq!(
outcome,
FinalizeOutcome::Completed {
pushed_head: sha("upmerge")
}
);
assert!(verifier.calls() >= 1);
assert_eq!(git.successful_pushes(), 1);
}
#[tokio::test]
async fn upstream_integration_blocked_and_cancelled_outcomes_never_push() {
for outcome in [
SchedulerOutcome::BlockedOrStalled,
SchedulerOutcome::Cancelled,
] {
let mut h = harness();
let result = h.coordinator.finalize(outcome).await.unwrap();
assert!(matches!(result, FinalizeOutcome::Skipped { .. }));
assert_eq!(h.git.push_calls(), 0);
assert_eq!(h.verifier.calls(), 0);
assert_eq!(h.git.fetch_calls(), 0);
}
}
async fn mark_publication_required(h: &mut Harness, change_id: &str) -> String {
h.coordinator
.record_publication_intent(change_id)
.await
.unwrap();
h.git.head()
}
#[tokio::test]
async fn upstream_integration_successful_drain_pushes_once_and_confirms() {
let mut h = harness();
let head = mark_publication_required(&mut h, "a").await;
let outcome = h
.coordinator
.finalize(SchedulerOutcome::DrainedSuccessfully)
.await
.unwrap();
assert_eq!(
outcome,
FinalizeOutcome::Completed {
pushed_head: head.clone()
}
);
assert_eq!(h.git.successful_pushes(), 1);
assert_eq!(h.coordinator.pushed_head(), Some(head.as_str()));
assert!(h
.observer
.events()
.iter()
.any(|e| matches!(e, UpstreamEvent::PushConfirmed { .. })));
let repeat = h
.coordinator
.finalize(SchedulerOutcome::DrainedSuccessfully)
.await
.unwrap();
assert!(matches!(repeat, FinalizeOutcome::Completed { .. }));
assert_eq!(h.git.successful_pushes(), 1);
}
#[tokio::test]
async fn upstream_integration_zero_change_fresh_run_manufactures_no_history() {
let mut h = harness();
h.git.set_remote_ref(&sha("local"));
let outcome = h
.coordinator
.finalize(SchedulerOutcome::DrainedSuccessfully)
.await
.unwrap();
assert_eq!(outcome, FinalizeOutcome::NoWork);
assert!(h.git.merge_calls().is_empty());
assert_eq!(h.verifier.calls(), 0);
assert_eq!(h.git.push_calls(), 0);
}
#[tokio::test]
async fn upstream_integration_remote_only_advance_creates_no_synthetic_merge() {
let mut h = harness();
let advanced = h.git.advance_remote("ahead", &sha("local"));
let outcome = h
.coordinator
.finalize(SchedulerOutcome::DrainedSuccessfully)
.await
.unwrap();
assert_eq!(outcome, FinalizeOutcome::NoWork);
assert!(h.git.merge_calls().is_empty());
assert_eq!(h.git.push_calls(), 0);
assert_ne!(advanced, sha("local"));
}
#[tokio::test]
async fn upstream_integration_noop_run_still_verifies_before_push() {
let mut h = harness();
mark_publication_required(&mut h, "a").await;
let outcome = h
.coordinator
.finalize(SchedulerOutcome::DrainedSuccessfully)
.await
.unwrap();
assert!(matches!(outcome, FinalizeOutcome::Completed { .. }));
assert!(h.git.merge_calls().is_empty(), "no-op performs no merge");
assert_eq!(
h.verifier.calls(),
1,
"final cumulative HEAD is always verified before push"
);
}
#[tokio::test]
async fn upstream_integration_remote_advance_before_push_returns_to_integration() {
let mut h = harness();
mark_publication_required(&mut h, "a").await;
let advanced = h.git.advance_remote("racer", &sha("root"));
let outcome = h
.coordinator
.finalize(SchedulerOutcome::DrainedSuccessfully)
.await
.unwrap();
assert!(matches!(outcome, FinalizeOutcome::Completed { .. }));
let calls = h.git.merge_calls();
assert_eq!(calls.len(), 1);
assert_eq!(calls[0].0, advanced);
assert_eq!(h.git.successful_pushes(), 1);
}
#[tokio::test]
async fn upstream_integration_race_rejection_returns_to_checkpoint_then_succeeds() {
let mut h = harness();
mark_publication_required(&mut h, "a").await;
h.git.lock().push_behaviors.push_back(PushBehavior::Reject {
porcelain:
"To fake\n!\trefs/heads/main:refs/heads/main\t[rejected]\t(non-fast-forward)\nDone\n"
.to_string(),
});
let outcome = h
.coordinator
.finalize(SchedulerOutcome::DrainedSuccessfully)
.await
.unwrap();
assert!(matches!(outcome, FinalizeOutcome::Completed { .. }));
assert_eq!(
h.git.push_calls(),
2,
"one failed race attempt, then success"
);
assert_eq!(h.git.successful_pushes(), 1, "at most one successful push");
assert_eq!(h.repair.call_count(), 0, "a race never invokes an agent");
}
#[tokio::test]
async fn upstream_integration_non_repairable_push_failure_stalls_without_agent() {
let mut h = harness();
mark_publication_required(&mut h, "a").await;
for _ in 0..MAX_FINALIZE_ATTEMPTS {
h.git
.lock()
.push_behaviors
.push_back(PushBehavior::FailWithoutRefStatus);
}
let outcome = h
.coordinator
.finalize(SchedulerOutcome::DrainedSuccessfully)
.await
.unwrap();
assert!(matches!(outcome, FinalizeOutcome::Stalled { .. }));
assert_eq!(h.git.successful_pushes(), 0);
assert_eq!(h.repair.call_count(), 0);
assert!(h
.observer
.events()
.iter()
.any(|e| matches!(e, UpstreamEvent::PushFailed { .. })));
}
#[tokio::test]
async fn upstream_integration_repairable_push_failure_repairs_then_conflux_pushes() {
let git = Arc::new(FakeGit::new_linear());
git.lock()
.push_behaviors
.push_back(PushBehavior::FailWithoutRefStatus);
git.lock().porcelain_v2 = "1 .M N... 100644 100644 100644 aaa bbb src/lib.rs\n".to_string();
let verifier = Arc::new(FakeVerifier::always_pass());
let repair = Arc::new(FakeRepairAgent::with_attempts(2));
let observer = Arc::new(RecordingObserver::default());
repair.bind(&git, |git| {
git.lock().porcelain_v2 = String::new();
});
let mut coordinator = UpstreamCoordinator::new(
UpstreamIntegrationConfig::new("origin", "cargo test"),
"main",
git.clone(),
verifier.clone(),
repair.clone(),
observer.clone(),
);
coordinator.record_publication_intent("a").await.unwrap();
let outcome = coordinator
.finalize(SchedulerOutcome::DrainedSuccessfully)
.await
.unwrap();
assert!(matches!(outcome, FinalizeOutcome::Completed { .. }));
assert_eq!(repair.call_count(), 1);
assert_eq!(repair.calls()[0].cause, RepairCause::PushRepository);
assert_eq!(
git.successful_pushes(),
1,
"Conflux performs the retry itself"
);
assert!(
verifier.calls() >= 2,
"post-repair convergence reruns the complete command"
);
}
#[tokio::test]
async fn upstream_integration_final_verification_failure_never_pushes() {
let mut h = harness_with(
FakeVerifier::scripted([false]),
FakeRepairAgent::with_attempts(0),
);
mark_publication_required(&mut h, "a").await;
let outcome = h
.coordinator
.finalize(SchedulerOutcome::DrainedSuccessfully)
.await
.unwrap();
assert!(matches!(outcome, FinalizeOutcome::Stalled { .. }));
assert_eq!(h.git.push_calls(), 0);
}
#[tokio::test]
async fn upstream_integration_reports_lifecycle_without_becoming_routing_authority() {
let mut h = harness();
h.git.advance_remote("remoteahead", &h.git.head());
h.coordinator
.checkpoint(CheckpointTrigger::AfterDrain, &BaseLaneState::clean(), None)
.await
.unwrap();
let events = h.observer.events();
assert!(events
.iter()
.any(|e| matches!(e, UpstreamEvent::CheckpointStarted { .. })));
assert!(events.iter().any(|e| matches!(
e,
UpstreamEvent::FetchCompleted { remote, branch, .. } if remote == "origin" && branch == "main"
)));
assert!(events
.iter()
.any(|e| matches!(e, UpstreamEvent::IntegrationCompleted { .. })));
assert!(events
.iter()
.any(|e| matches!(e, UpstreamEvent::Reverifying { .. })));
let repeat = h
.coordinator
.checkpoint(CheckpointTrigger::AfterDrain, &BaseLaneState::clean(), None)
.await
.unwrap();
assert!(matches!(repeat, UpstreamStepOutcome::NoOp { .. }));
}
async fn publish(h: &mut Harness, change_id: &str) -> PublicationOutcome {
h.coordinator
.record_publication_intent(change_id)
.await
.unwrap();
h.coordinator.publish_change(change_id).await.unwrap()
}
#[tokio::test]
async fn per_change_upstream_publishes_one_change_with_one_native_push() {
let mut h = harness();
h.git.integrate_change("alpha");
let outcome = publish(&mut h, "alpha").await;
let head = h.git.head();
assert_eq!(
outcome,
PublicationOutcome::Published { head: head.clone() }
);
assert_eq!(h.git.successful_pushes(), 1, "exactly one native push");
assert_eq!(h.coordinator.pushed_head(), Some(head.as_str()));
let markers = h.git.empty_commits();
assert_eq!(markers.len(), 1);
let trailers = crate::upstream::publication::parse_publication_trailers(&markers[0]).unwrap();
assert_eq!(trailers.change_id, "alpha");
assert_eq!(trailers.remote, "origin");
assert_eq!(trailers.branch, "main");
assert!(h
.observer
.events()
.iter()
.any(|e| matches!(e, UpstreamEvent::PushConfirmed { head: h2, .. } if h2 == &head)));
}
#[tokio::test]
async fn per_change_upstream_does_not_push_an_already_confirmed_head() {
let mut h = harness();
h.git.integrate_change("alpha");
publish(&mut h, "alpha").await;
let head = h.git.head();
let repeat = h.coordinator.publish_change("alpha").await.unwrap();
assert_eq!(repeat, PublicationOutcome::AlreadyConfirmed { head });
assert_eq!(
h.git.successful_pushes(),
1,
"a confirmed revision must never be pushed again"
);
}
#[tokio::test]
async fn per_change_upstream_publishes_multiple_successive_changes() {
let mut h = harness();
h.git.integrate_change("alpha");
let first = publish(&mut h, "alpha").await;
let alpha_head = h.git.head();
h.git.integrate_change("beta");
let second = publish(&mut h, "beta").await;
let beta_head = h.git.head();
assert_eq!(
first,
PublicationOutcome::Published {
head: alpha_head.clone()
}
);
assert_eq!(
second,
PublicationOutcome::Published {
head: beta_head.clone()
}
);
assert_ne!(
alpha_head, beta_head,
"each change advances cumulative HEAD"
);
assert_eq!(
h.git.successful_pushes(),
2,
"one confirmed push per advancing cumulative HEAD"
);
assert_eq!(h.coordinator.pushed_head(), Some(beta_head.as_str()));
}
#[tokio::test]
async fn per_change_upstream_confirms_an_interrupted_push_without_pushing_again() {
let mut h = harness();
h.git.integrate_change("alpha");
h.coordinator
.record_publication_intent("alpha")
.await
.unwrap();
h.git.set_remote_ref(&h.git.head());
let outcome = h.coordinator.publish_change("alpha").await.unwrap();
assert_eq!(
outcome,
PublicationOutcome::AlreadyConfirmed { head: h.git.head() }
);
assert_eq!(
h.git.push_calls(),
0,
"retry must observe the remote before pushing again"
);
}
#[tokio::test]
async fn per_change_upstream_verification_failure_suppresses_push() {
let mut h = harness_with(
FakeVerifier::scripted([false, false, false, false, false, false]),
FakeRepairAgent::with_attempts(1),
);
h.git.integrate_change("alpha");
let outcome = publish(&mut h, "alpha").await;
assert!(
matches!(outcome, PublicationOutcome::Stalled { ref reason } if reason.contains("alpha")),
"outcome: {:?}",
outcome
);
assert_eq!(h.git.push_calls(), 0, "a failed verification must not push");
assert_eq!(h.coordinator.pushed_head(), None);
assert!(!h
.observer
.events()
.iter()
.any(|e| matches!(e, UpstreamEvent::PushConfirmed { .. })));
}
#[tokio::test]
async fn per_change_upstream_non_repairable_push_failure_stalls_without_agent() {
let mut h = harness();
{
let mut inner = h.git.lock();
for _ in 0..MAX_FINALIZE_ATTEMPTS {
inner
.push_behaviors
.push_back(PushBehavior::FailWithoutRefStatus);
}
}
h.git.integrate_change("alpha");
let outcome = publish(&mut h, "alpha").await;
assert!(matches!(outcome, PublicationOutcome::Stalled { .. }));
assert_eq!(h.coordinator.pushed_head(), None);
assert_eq!(
h.repair.call_count(),
0,
"a non-repairable push failure must not invoke the repair agent"
);
}
#[tokio::test]
async fn per_change_upstream_remote_race_returns_to_integration_before_publishing() {
let mut h = harness();
h.git.integrate_change("alpha");
h.coordinator
.record_publication_intent("alpha")
.await
.unwrap();
h.git.advance_remote("advanced", &sha("root"));
let outcome = h.coordinator.publish_change("alpha").await.unwrap();
assert!(matches!(outcome, PublicationOutcome::Published { .. }));
assert!(
h.git
.merge_calls()
.iter()
.any(|(target, _)| target == &sha("advanced")),
"the remote advance must be integrated with a real merge"
);
assert_eq!(h.git.successful_pushes(), 1);
}
#[tokio::test]
async fn per_change_upstream_scan_reports_only_unpublished_markers() {
let mut h = harness();
h.git.integrate_change("alpha");
h.coordinator
.record_publication_intent("alpha")
.await
.unwrap();
let pending = scan_pending_publications(h.git.as_ref()).await.unwrap();
assert_eq!(pending.len(), 1);
assert_eq!(pending[0].trailers.change_id, "alpha");
h.git.set_remote_ref(&h.git.head());
h.git.fetch("origin", "main").await.unwrap();
let published = scan_pending_publications(h.git.as_ref()).await.unwrap();
assert!(published.is_empty(), "pending: {:?}", published);
}
#[tokio::test]
async fn per_change_upstream_scan_ignores_disabled_mode_history() {
let h = harness();
h.git.integrate_change("alpha");
h.git.integrate_change("beta");
let pending = scan_pending_publications(h.git.as_ref()).await.unwrap();
assert!(pending.is_empty(), "pending: {:?}", pending);
}
#[tokio::test]
async fn per_change_upstream_option_less_refusal_names_the_unpublished_change() {
let mut h = harness();
h.git.integrate_change("alpha");
h.coordinator
.record_publication_intent("alpha")
.await
.unwrap();
let pending = scan_pending_publications(h.git.as_ref()).await.unwrap();
let refusal = publication_recovery_refusal(&pending).expect("refusal");
match &refusal {
UpstreamOptionError::PublicationRecoveryRequiresOption {
change_id,
remote,
branch,
..
} => {
assert_eq!(change_id, "alpha");
assert_eq!(remote, "origin");
assert_eq!(branch, "main");
}
other => panic!("unexpected refusal: {:?}", other),
}
let message = refusal.to_string();
assert!(
message.contains("--integrate-upstream=origin"),
"{}",
message
);
assert!(message.contains("--upstream-verify-command"), "{}", message);
assert!(
message.contains("it is not merged"),
"the refusal must state plainly that the change is not merged: {}",
message
);
}
#[tokio::test]
async fn per_change_upstream_zero_change_recovery_requires_explicit_evidence() {
let mut h = harness();
h.git.integrate_change("alpha");
let outcome = h
.coordinator
.finalize(SchedulerOutcome::DrainedSuccessfully)
.await
.unwrap();
assert_eq!(outcome, FinalizeOutcome::NoWork);
assert_eq!(h.git.push_calls(), 0);
assert_eq!(h.verifier.calls(), 0);
assert!(h.git.merge_calls().is_empty());
}
#[tokio::test]
async fn per_change_upstream_zero_change_recovery_publishes_marked_history() {
let mut h = harness();
h.git.integrate_change("alpha");
h.coordinator
.record_publication_intent("alpha")
.await
.unwrap();
let outcome = h
.coordinator
.finalize(SchedulerOutcome::DrainedSuccessfully)
.await
.unwrap();
assert_eq!(
outcome,
FinalizeOutcome::Completed {
pushed_head: h.git.head()
}
);
assert_eq!(h.git.successful_pushes(), 1);
}
#[tokio::test]
async fn per_change_upstream_finalization_is_a_noop_after_change_publication() {
let mut h = harness();
h.git.integrate_change("alpha");
publish(&mut h, "alpha").await;
let head = h.git.head();
let outcome = h
.coordinator
.finalize(SchedulerOutcome::DrainedSuccessfully)
.await
.unwrap();
assert_eq!(outcome, FinalizeOutcome::Completed { pushed_head: head });
assert_eq!(
h.git.successful_pushes(),
1,
"finalization must not republish an already confirmed HEAD"
);
}
#[tokio::test]
async fn per_change_upstream_cancelled_or_blocked_run_never_publishes() {
for outcome in [
SchedulerOutcome::Cancelled,
SchedulerOutcome::BlockedOrStalled,
] {
let mut h = harness();
h.git.integrate_change("alpha");
h.coordinator
.record_publication_intent("alpha")
.await
.unwrap();
let finalized = h.coordinator.finalize(outcome).await.unwrap();
assert!(matches!(finalized, FinalizeOutcome::Skipped { .. }));
assert_eq!(h.git.push_calls(), 0);
assert_eq!(h.coordinator.pushed_head(), None);
}
}