use crate::archive_layout;
use crate::vcs::git::commands::{self as git_commands, CommitDiffEntry, WorktreeIdentity};
use async_trait::async_trait;
use std::path::{Path, PathBuf};
#[cfg(test)]
pub(crate) mod fixtures;
#[cfg(test)]
mod tests;
pub fn final_merge_subject(change_id: &str) -> String {
format!("Merge change: {}", change_id)
}
pub fn presync_subject(change_id: &str) -> String {
format!("Pre-sync base into {}", change_id)
}
pub fn cleanup_subject(change_id: &str) -> String {
format!("Cleanup resurrected change: {}", change_id)
}
pub type EvidenceResult<T> = std::result::Result<T, String>;
#[async_trait]
pub trait ResolveEvidence: Send + Sync {
async fn validate_worktree(
&self,
supplied_path: &Path,
expected_branch: &str,
) -> WorktreeIdentity;
async fn worktree_merge_in_progress(&self, worktree: &Path) -> EvidenceResult<bool>;
async fn worktree_conflicts(&self, worktree: &Path) -> EvidenceResult<Vec<String>>;
async fn target_head(&self) -> EvidenceResult<String>;
async fn target_merge_head(&self) -> EvidenceResult<Option<String>>;
async fn target_conflict_paths(&self) -> EvidenceResult<Vec<String>>;
async fn target_is_clean(&self) -> EvidenceResult<bool>;
async fn target_index_paths(&self, prefix: &str) -> EvidenceResult<Vec<String>>;
async fn parents_of(&self, commit: &str) -> EvidenceResult<Vec<String>>;
async fn is_ancestor(&self, ancestor: &str, descendant: &str) -> EvidenceResult<bool>;
async fn first_parent_lineage(&self, tip: &str) -> EvidenceResult<Vec<String>>;
async fn merge_base(&self, a: &str, b: &str) -> EvidenceResult<Option<String>>;
async fn commits_with_exact_subject(
&self,
from: Option<&str>,
to: &str,
subject: &str,
) -> EvidenceResult<Vec<String>>;
async fn commit_diff_entries(&self, commit: &str) -> EvidenceResult<Vec<CommitDiffEntry>>;
async fn committed_tree_paths(
&self,
revision: &str,
prefix: &str,
) -> EvidenceResult<Vec<String>>;
async fn committed_file_text(
&self,
revision: &str,
path: &str,
) -> EvidenceResult<Option<String>>;
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum BatchState {
UnsafeEvidence {
reason: String,
},
PreSyncUnfinished {
change_id: String,
worktree: PathBuf,
reason: String,
},
PreSyncInvalid {
change_id: String,
worktree: PathBuf,
required_target_state: String,
reason: String,
},
TargetMergeUnfinished {
change_id: String,
worktree: PathBuf,
requires_live_removal: bool,
reason: String,
},
MergeNotAuthorized {
change_id: String,
reason: String,
},
FinalMergeMissing {
change_id: String,
worktree: PathBuf,
required_target_state: String,
},
ResurrectionCleanupRequired {
change_id: String,
reason: String,
},
Complete,
}
impl BatchState {
pub fn phase(&self) -> &'static str {
match self {
Self::UnsafeEvidence { .. } => "unsafe_evidence",
Self::PreSyncUnfinished { .. } => "presync_unfinished",
Self::PreSyncInvalid { .. } => "presync_invalid",
Self::TargetMergeUnfinished { .. } => "target_merge_unfinished",
Self::MergeNotAuthorized { .. } => "merge_not_authorized",
Self::FinalMergeMissing { .. } => "final_merge_missing",
Self::ResurrectionCleanupRequired { .. } => "resurrection_cleanup_required",
Self::Complete => "complete",
}
}
pub fn is_complete(&self) -> bool {
matches!(self, Self::Complete)
}
pub fn allows_agent_action(&self) -> bool {
!matches!(
self,
Self::UnsafeEvidence { .. } | Self::MergeNotAuthorized { .. } | Self::Complete
)
}
pub fn diagnosis(&self) -> String {
const FIELD_MAX_BYTES: usize = 384;
let mut lines = vec![format!("phase: {}", self.phase())];
match self {
Self::UnsafeEvidence { reason } => {
lines.push(format!("detail: {}", reason));
lines.push(
"required_action: none; preserve repository state and do not commit"
.to_string(),
);
}
Self::PreSyncUnfinished {
change_id,
worktree,
reason,
} => {
lines.push(format!("change_id: {}", change_id));
lines.push(format!("worktree: {}", worktree.display()));
lines.push(format!("detail: {}", reason));
lines.push(format!(
"required_action: finish the in-progress worktree merge and commit it with subject '{}'",
presync_subject(change_id)
));
}
Self::PreSyncInvalid {
change_id,
worktree,
required_target_state,
reason,
} => {
lines.push(format!("change_id: {}", change_id));
lines.push(format!("worktree: {}", worktree.display()));
lines.push(format!("required_target_state: {}", required_target_state));
lines.push(format!("detail: {}", reason));
lines.push(format!(
"required_action: in the worktree run 'git merge --no-ff -m \"{}\" {}' so the pre-sync commit has exactly two parents with non-first parent {}",
presync_subject(change_id),
required_target_state,
required_target_state
));
}
Self::TargetMergeUnfinished {
change_id,
worktree,
requires_live_removal,
reason,
} => {
lines.push(format!("change_id: {}", change_id));
lines.push(format!("worktree: {}", worktree.display()));
lines.push(format!("requires_live_removal: {}", requires_live_removal));
lines.push(format!("detail: {}", reason));
if *requires_live_removal {
lines.push(format!(
"required_action: run 'git rm -r --cached -f openspec/changes/{}' and remove the directory, then commit the in-progress merge with subject '{}'",
change_id,
final_merge_subject(change_id)
));
} else {
lines.push(format!(
"required_action: commit the in-progress target merge with subject '{}'",
final_merge_subject(change_id)
));
}
}
Self::MergeNotAuthorized { change_id, reason } => {
lines.push(format!("change_id: {}", change_id));
lines.push(format!("detail: {}", reason));
lines.push(
"required_action: none; the change is not proven complete, so do not integrate it and do not mutate any repository state; operator action is required"
.to_string(),
);
}
Self::FinalMergeMissing {
change_id,
worktree,
required_target_state,
} => {
lines.push(format!("change_id: {}", change_id));
lines.push(format!("worktree: {}", worktree.display()));
lines.push(format!("required_target_state: {}", required_target_state));
lines.push(format!(
"required_action: in the repo root run 'git merge --no-ff --no-commit <branch>' and commit with subject '{}'",
final_merge_subject(change_id)
));
}
Self::ResurrectionCleanupRequired { change_id, reason } => {
lines.push(format!("change_id: {}", change_id));
lines.push(format!("detail: {}", reason));
lines.push(format!(
"required_action: remove only openspec/changes/{} and record a forward commit with subject '{}'; never amend or rewrite history",
change_id,
cleanup_subject(change_id)
));
}
Self::Complete => {
lines.push(
"detail: every batch item is integrated and the target is clean".to_string(),
);
lines.push("required_action: none".to_string());
}
}
lines
.into_iter()
.map(|line| crate::history::bounded_head(&line, FIELD_MAX_BYTES))
.collect::<Vec<_>>()
.join("\n")
}
}
pub use super::merge::SequentialMergeItem;
#[derive(Debug, Clone)]
pub struct ValidatedItem {
pub change_id: String,
pub revision: String,
pub worktree: PathBuf,
pub tip: String,
pub branch_base: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ItemIntegration {
Exact {
commit: String,
target_state: String,
},
Historical,
NotIntegrated,
}
fn unsafe_evidence(context: &str, error: String) -> BatchState {
BatchState::UnsafeEvidence {
reason: format!("{}: {}", context, error),
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum TaskCompletion {
Complete {
sources: Vec<String>,
total: u32,
},
Incomplete {
source: String,
completed: u32,
total: u32,
},
Unestablished {
reason: String,
},
}
impl TaskCompletion {
fn detail(&self, change_id: &str) -> Option<String> {
match self {
Self::Complete { .. } => None,
Self::Incomplete {
source,
completed,
total,
} => Some(format!(
"Change '{}' records {}/{} tasks complete in {}; an unfinished change is never integrated",
change_id, completed, total, source
)),
Self::Unestablished { reason } => Some(format!(
"Change '{}' has no task evidence that proves completion: {}",
change_id, reason
)),
}
}
}
enum TaskEvidence {
Coherent {
format: crate::task_file::TaskFileFormat,
paths: Vec<String>,
},
Absent,
Ambiguous { paths: Vec<String> },
}
fn task_evidence_paths(tree_paths: &[String], change_id: &str) -> TaskEvidence {
let mut found: Vec<(crate::task_file::TaskFileFormat, String)> = tree_paths
.iter()
.filter_map(|path| {
archive_layout::active_change_tasks_format(path, change_id)
.or_else(|| archive_layout::archive_tasks_format(path, change_id))
.map(|format| (format, path.clone()))
})
.collect();
found.sort_by(|left, right| left.1.cmp(&right.1));
found.dedup_by(|left, right| left.1 == right.1);
let Some(format) = found.first().map(|(format, _)| *format) else {
return TaskEvidence::Absent;
};
if found.iter().any(|(candidate, _)| *candidate != format) {
return TaskEvidence::Ambiguous {
paths: found.into_iter().map(|(_, path)| path).collect(),
};
}
TaskEvidence::Coherent {
format,
paths: found.into_iter().map(|(_, path)| path).collect(),
}
}
pub async fn read_task_completion(
evidence: &dyn ResolveEvidence,
change_id: &str,
revision: &str,
tree_paths: &[String],
) -> TaskCompletion {
let (format, paths) = match task_evidence_paths(tree_paths, change_id) {
TaskEvidence::Coherent { format, paths } => (format, paths),
TaskEvidence::Absent => {
return TaskCompletion::Unestablished {
reason: format!(
"no active or archived task artifact ({}) exists for '{}' at {}",
archive_layout::active_change_tasks_paths(change_id)
.into_iter()
.map(|(_, path)| path)
.collect::<Vec<_>>()
.join(" or "),
change_id,
revision
),
}
}
TaskEvidence::Ambiguous { paths } => {
return TaskCompletion::Unestablished {
reason: format!(
"task evidence for '{}' at {} mixes task-file formats ({}); a change speaks in exactly one",
change_id,
revision,
paths.join(", ")
),
}
}
};
let mut total = 0u32;
for path in &paths {
let content = match evidence.committed_file_text(revision, path).await {
Ok(Some(content)) => content,
Ok(None) => {
return TaskCompletion::Unestablished {
reason: format!("{} is not readable at {}", path, revision),
}
}
Err(error) => {
return TaskCompletion::Unestablished {
reason: format!("failed to read {} at {}: {}", path, revision, error),
}
}
};
let progress = match crate::task_file::parse_progress(format, &content, None) {
Ok(progress) => progress,
Err(error) => {
return TaskCompletion::Unestablished {
reason: format!(
"{} is not valid task evidence at {}: {}",
path, revision, error
),
}
}
};
if progress.total == 0 {
return TaskCompletion::Unestablished {
reason: format!("{} records no tasks at all", path),
};
}
if progress.completed < progress.total {
return TaskCompletion::Incomplete {
source: path.clone(),
completed: progress.completed,
total: progress.total,
};
}
total = total.saturating_add(progress.total);
}
TaskCompletion::Complete {
sources: paths,
total,
}
}
#[derive(Debug, Default)]
pub struct MergeAuthorizationLatch {
refusals: std::sync::Mutex<std::collections::BTreeMap<String, String>>,
}
impl MergeAuthorizationLatch {
pub fn new() -> Self {
Self::default()
}
fn guard(&self) -> std::sync::MutexGuard<'_, std::collections::BTreeMap<String, String>> {
self.refusals
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
}
pub fn refuse(&self, change_id: &str, reason: &str) {
self.guard()
.entry(change_id.to_string())
.or_insert_with(|| reason.to_string());
}
pub fn refusal(&self, change_id: &str) -> Option<String> {
self.guard().get(change_id).cloned()
}
#[allow(dead_code)]
pub fn refused_count(&self) -> usize {
self.guard().len()
}
}
pub async fn validate_items(
evidence: &dyn ResolveEvidence,
items: &[SequentialMergeItem],
) -> std::result::Result<Vec<ValidatedItem>, BatchState> {
let mut validated = Vec::with_capacity(items.len());
for item in items {
match evidence
.validate_worktree(&item.archive_path, &item.revision)
.await
{
WorktreeIdentity::Supplied { path, tip }
| WorktreeIdentity::Rediscovered { path, tip } => {
validated.push(ValidatedItem {
change_id: item.change_id.clone(),
revision: item.revision.clone(),
worktree: path,
tip,
branch_base: item.branch_base.clone(),
});
}
WorktreeIdentity::Unsafe { reason } => {
return Err(BatchState::UnsafeEvidence { reason });
}
}
}
Ok(validated)
}
pub async fn classify_item_integration(
evidence: &dyn ResolveEvidence,
item: &ValidatedItem,
base_revision: &str,
target_head: &str,
) -> std::result::Result<ItemIntegration, BatchState> {
let subject = final_merge_subject(&item.change_id);
let candidates = evidence
.commits_with_exact_subject(Some(base_revision), target_head, &subject)
.await
.map_err(|error| unsafe_evidence("Failed to enumerate final merge candidates", error))?;
match candidates.len() {
0 => {
let integrated =
evidence
.is_ancestor(&item.tip, target_head)
.await
.map_err(|error| {
unsafe_evidence("Failed to check historical integration", error)
})?;
if integrated {
Ok(ItemIntegration::Historical)
} else {
Ok(ItemIntegration::NotIntegrated)
}
}
1 => {
let commit = candidates[0].clone();
let parents = evidence
.parents_of(&commit)
.await
.map_err(|error| unsafe_evidence("Failed to read merge parents", error))?;
if parents.len() != 2 {
return Err(BatchState::UnsafeEvidence {
reason: format!(
"Final merge commit {} for '{}' has {} parent(s); exactly two are required",
commit,
item.change_id,
parents.len()
),
});
}
if parents[1] != item.tip {
return Err(BatchState::UnsafeEvidence {
reason: format!(
"Final merge commit {} for '{}' has non-first parent {} but branch '{}' tip is {}",
commit, item.change_id, parents[1], item.revision, item.tip
),
});
}
let contained = evidence
.is_ancestor(&commit, target_head)
.await
.map_err(|error| unsafe_evidence("Failed to check merge containment", error))?;
if !contained {
return Err(BatchState::UnsafeEvidence {
reason: format!(
"Final merge commit {} for '{}' is not contained by target HEAD",
commit, item.change_id
),
});
}
Ok(ItemIntegration::Exact {
commit,
target_state: parents[0].clone(),
})
}
count => Err(BatchState::UnsafeEvidence {
reason: format!(
"Found {} commits with subject '{}'; exactly one is required",
count, subject
),
}),
}
}
async fn validate_presync(
evidence: &dyn ResolveEvidence,
item: &ValidatedItem,
required_target_state: &str,
) -> std::result::Result<Option<BatchState>, BatchState> {
let lineage = evidence
.first_parent_lineage(&item.tip)
.await
.map_err(|error| unsafe_evidence("Failed to read worktree first-parent lineage", error))?;
if lineage.iter().any(|commit| commit == required_target_state) {
return Ok(None);
}
let subject = presync_subject(&item.change_id);
let from = match &item.branch_base {
Some(base) => Some(base.clone()),
None => evidence
.merge_base(required_target_state, &item.tip)
.await
.map_err(|error| unsafe_evidence("Failed to compute pre-sync search base", error))?,
};
let candidates = evidence
.commits_with_exact_subject(from.as_deref(), &item.tip, &subject)
.await
.map_err(|error| unsafe_evidence("Failed to enumerate pre-sync candidates", error))?;
let invalid = |reason: String| {
Ok(Some(BatchState::PreSyncInvalid {
change_id: item.change_id.clone(),
worktree: item.worktree.clone(),
required_target_state: required_target_state.to_string(),
reason,
}))
};
match candidates.len() {
0 => invalid(format!(
"Required target state {} is not on branch '{}' first-parent lineage and no '{}' commit exists",
required_target_state, item.revision, subject
)),
1 => {
let commit = &candidates[0];
let parents = evidence
.parents_of(commit)
.await
.map_err(|error| unsafe_evidence("Failed to read pre-sync parents", error))?;
if parents.len() != 2 {
return invalid(format!(
"Pre-sync commit {} has {} parent(s); exactly two are required",
commit,
parents.len()
));
}
if parents[1] != required_target_state {
return invalid(format!(
"Pre-sync commit {} has non-first parent {} but required target state is {}",
commit, parents[1], required_target_state
));
}
let contained = evidence
.is_ancestor(commit, &item.tip)
.await
.map_err(|error| unsafe_evidence("Failed to check pre-sync containment", error))?;
if !contained {
return invalid(format!(
"Pre-sync commit {} is not contained by branch '{}' tip {}",
commit, item.revision, item.tip
));
}
Ok(None)
}
count => invalid(format!(
"Found {} commits with subject '{}'; exactly one is required",
count, subject
)),
}
}
async fn validate_cleanup_commits(
evidence: &dyn ResolveEvidence,
item: &ValidatedItem,
integration: &ItemIntegration,
base_revision: &str,
target_head: &str,
target_lineage: &[String],
) -> std::result::Result<(), BatchState> {
let change_id = item.change_id.as_str();
let subject = cleanup_subject(change_id);
let candidates = evidence
.commits_with_exact_subject(Some(base_revision), target_head, &subject)
.await
.map_err(|error| unsafe_evidence("Failed to enumerate cleanup candidates", error))?;
let commit = match candidates.as_slice() {
[] => return Ok(()),
[commit] => commit.clone(),
_ => {
return Err(BatchState::UnsafeEvidence {
reason: format!(
"Found {} commits with subject '{}'; exactly one is allowed",
candidates.len(),
subject
),
})
}
};
if !target_lineage.iter().any(|candidate| candidate == &commit) {
return Err(BatchState::UnsafeEvidence {
reason: format!(
"Cleanup commit {} for '{}' is not on target HEAD {} first-parent lineage; only forward cleanup on the target is accepted",
commit, change_id, target_head
),
});
}
let parents = evidence
.parents_of(&commit)
.await
.map_err(|error| unsafe_evidence("Failed to read cleanup commit parents", error))?;
if parents.len() != 1 {
return Err(BatchState::UnsafeEvidence {
reason: format!(
"Cleanup commit {} has {} parent(s); a forward cleanup commit has exactly one",
commit,
parents.len()
),
});
}
let predecessor = parents[0].clone();
let entries = evidence
.commit_diff_entries(&commit)
.await
.map_err(|error| unsafe_evidence("Failed to read cleanup commit diff", error))?;
if entries.is_empty() {
return Err(BatchState::UnsafeEvidence {
reason: format!("Cleanup commit {} changes nothing", commit),
});
}
for entry in &entries {
if entry.status != 'D' || !archive_layout::is_active_change_path(&entry.path, change_id) {
return Err(BatchState::UnsafeEvidence {
reason: format!(
"Cleanup commit {} has entry '{}{}'; it may only delete openspec/changes/{}",
commit, entry.status, entry.path, change_id
),
});
}
}
let integrated_at = match integration {
ItemIntegration::Exact { commit, .. } => commit.clone(),
ItemIntegration::Historical => item.tip.clone(),
ItemIntegration::NotIntegrated => {
return Err(BatchState::UnsafeEvidence {
reason: format!(
"Cleanup commit {} exists for '{}' but the change has no committed integration",
commit, change_id
),
})
}
};
let after_integration = evidence
.is_ancestor(&integrated_at, &predecessor)
.await
.map_err(|error| unsafe_evidence("Failed to order cleanup against integration", error))?;
if !after_integration {
return Err(BatchState::UnsafeEvidence {
reason: format!(
"Cleanup commit {} for '{}' has predecessor {} which does not contain its committed integration {}",
commit, change_id, predecessor, integrated_at
),
});
}
let predecessor_paths = evidence
.committed_tree_paths(&predecessor, archive_layout::ACTIVE_CHANGES_PREFIX)
.await
.map_err(|error| unsafe_evidence("Failed to read cleanup predecessor tree", error))?;
if !archive_layout::paths_contain_active_change(&predecessor_paths, change_id)
|| !archive_layout::paths_contain_valid_archive(&predecessor_paths, change_id)
{
return Err(BatchState::UnsafeEvidence {
reason: format!(
"Cleanup commit {} for '{}' has predecessor {} which does not hold the live/archive coexistence it claims to remove",
commit, change_id, predecessor
),
});
}
Ok(())
}
#[allow(dead_code)]
pub async fn classify_batch(
evidence: &dyn ResolveEvidence,
items: &[SequentialMergeItem],
base_revision: &str,
) -> BatchState {
classify_batch_with_latch(
evidence,
items,
base_revision,
&MergeAuthorizationLatch::new(),
)
.await
}
pub async fn classify_batch_with_latch(
evidence: &dyn ResolveEvidence,
items: &[SequentialMergeItem],
base_revision: &str,
latch: &MergeAuthorizationLatch,
) -> BatchState {
match classify_batch_inner(evidence, items, base_revision, latch).await {
Ok(state) | Err(state) => state,
}
}
async fn classify_batch_inner(
evidence: &dyn ResolveEvidence,
items: &[SequentialMergeItem],
base_revision: &str,
latch: &MergeAuthorizationLatch,
) -> std::result::Result<BatchState, BatchState> {
if items.is_empty() {
return Err(BatchState::UnsafeEvidence {
reason: "Sequential resolve received an empty batch".to_string(),
});
}
let validated = validate_items(evidence, items).await?;
let target_head = evidence
.target_head()
.await
.map_err(|error| unsafe_evidence("Failed to read target HEAD", error))?;
let merge_head = evidence
.target_merge_head()
.await
.map_err(|error| unsafe_evidence("Failed to read target MERGE_HEAD", error))?;
let mut integrations = Vec::with_capacity(validated.len());
for item in &validated {
integrations
.push(classify_item_integration(evidence, item, base_revision, &target_head).await?);
}
let owner = match merge_head.as_deref() {
Some(merge_head) => {
let matched: Vec<usize> = validated
.iter()
.enumerate()
.filter(|(_, item)| item.tip == merge_head)
.map(|(index, _)| index)
.collect();
match matched.as_slice() {
[index] => Some(*index),
[] => {
return Err(BatchState::UnsafeEvidence {
reason: format!(
"Target MERGE_HEAD {} does not match any batch branch tip",
merge_head
),
})
}
_ => {
return Err(BatchState::UnsafeEvidence {
reason: format!(
"Target MERGE_HEAD {} matches {} batch branch tips",
merge_head,
matched.len()
),
})
}
}
}
None => None,
};
let first_incomplete = integrations
.iter()
.position(|integration| matches!(integration, ItemIntegration::NotIntegrated));
if let Some(first) = first_incomplete {
if let Some((offset, _)) = integrations
.iter()
.enumerate()
.skip(first + 1)
.find(|(_, integration)| !matches!(integration, ItemIntegration::NotIntegrated))
{
return Err(BatchState::UnsafeEvidence {
reason: format!(
"Item '{}' is integrated while earlier item '{}' is not; declared order was violated",
validated[offset].change_id, validated[first].change_id
),
});
}
}
if let Some(owner) = owner {
if first_incomplete != Some(owner) {
return Err(BatchState::UnsafeEvidence {
reason: format!(
"Target MERGE_HEAD owner '{}' is not the first incomplete batch item",
validated[owner].change_id
),
});
}
}
let target_lineage = evidence
.first_parent_lineage(&target_head)
.await
.map_err(|error| unsafe_evidence("Failed to read target first-parent lineage", error))?;
let mut previous: Option<(&str, String, usize)> = None;
for (item, integration) in validated.iter().zip(integrations.iter()) {
let ItemIntegration::Exact {
commit,
target_state,
} = integration
else {
continue;
};
let Some(position) = target_lineage
.iter()
.position(|candidate| candidate == commit)
else {
return Err(BatchState::UnsafeEvidence {
reason: format!(
"Final merge commit {} for '{}' is not on target HEAD {} first-parent lineage",
commit, item.change_id, target_head
),
});
};
if let Some((previous_id, previous_commit, previous_position)) = &previous {
if position >= *previous_position {
return Err(BatchState::UnsafeEvidence {
reason: format!(
"Item '{}' was integrated before earlier item '{}' on the target first-parent history; declared order was violated",
item.change_id, previous_id
),
});
}
let chained = evidence
.is_ancestor(previous_commit, target_state)
.await
.map_err(|error| {
unsafe_evidence("Failed to chain declared-order integrations", error)
})?;
if !chained {
return Err(BatchState::UnsafeEvidence {
reason: format!(
"Final merge commit {} for '{}' does not build on the committed integration {} of earlier item '{}'; declared order was violated",
commit, item.change_id, previous_commit, previous_id
),
});
}
}
previous = Some((item.change_id.as_str(), commit.clone(), position));
}
for (item, integration) in validated.iter().zip(integrations.iter()) {
match integration {
ItemIntegration::Exact { target_state, .. } => {
if let Some(state) = validate_presync(evidence, item, target_state).await? {
return Ok(state);
}
}
ItemIntegration::Historical => {}
ItemIntegration::NotIntegrated => break,
}
}
if let Some(index) = first_incomplete {
let item = &validated[index];
if evidence
.worktree_merge_in_progress(&item.worktree)
.await
.map_err(|error| unsafe_evidence("Failed to read worktree merge state", error))?
{
return Ok(BatchState::PreSyncUnfinished {
change_id: item.change_id.clone(),
worktree: item.worktree.clone(),
reason: "Worktree merge is in progress (MERGE_HEAD exists)".to_string(),
});
}
let worktree_conflicts = evidence
.worktree_conflicts(&item.worktree)
.await
.map_err(|error| unsafe_evidence("Failed to read worktree conflicts", error))?;
if !worktree_conflicts.is_empty() {
return Ok(BatchState::PreSyncUnfinished {
change_id: item.change_id.clone(),
worktree: item.worktree.clone(),
reason: format!(
"Worktree has unresolved conflicts: {}",
worktree_conflicts.join(", ")
),
});
}
if let Some(state) = validate_presync(evidence, item, &target_head).await? {
return Ok(state);
}
let target_conflicts = evidence
.target_conflict_paths()
.await
.map_err(|error| unsafe_evidence("Failed to read target index stages", error))?;
if owner == Some(index) {
if !target_conflicts.is_empty() {
return Ok(BatchState::TargetMergeUnfinished {
change_id: item.change_id.clone(),
worktree: item.worktree.clone(),
requires_live_removal: false,
reason: format!(
"Target index still has conflict stages: {}",
target_conflicts.join(", ")
),
});
}
let index_paths = evidence
.target_index_paths(archive_layout::ACTIVE_CHANGES_PREFIX)
.await
.map_err(|error| unsafe_evidence("Failed to read target index entries", error))?;
let live = archive_layout::paths_contain_active_change(&index_paths, &item.change_id);
let archived =
archive_layout::paths_contain_valid_archive(&index_paths, &item.change_id);
return Ok(BatchState::TargetMergeUnfinished {
change_id: item.change_id.clone(),
worktree: item.worktree.clone(),
requires_live_removal: live && archived,
reason: if live && archived {
format!(
"Merged stage-0 index holds both openspec/changes/{} and its valid archive entry",
item.change_id
)
} else {
"Target merge is conflict-free and awaiting its commit".to_string()
},
});
}
if !target_conflicts.is_empty() {
return Err(BatchState::UnsafeEvidence {
reason: format!(
"Target index has conflict stages without MERGE_HEAD: {}",
target_conflicts.join(", ")
),
});
}
let worktree_tree_paths = evidence
.committed_tree_paths(&item.tip, archive_layout::ACTIVE_CHANGES_PREFIX)
.await
.map_err(|error| {
unsafe_evidence("Failed to read validated worktree committed tree", error)
})?;
if !archive_layout::paths_contain_valid_archive(&worktree_tree_paths, &item.change_id) {
let nested = archive_layout::paths_contain_invalid_nested_archive(
&worktree_tree_paths,
&item.change_id,
);
return Err(BatchState::UnsafeEvidence {
reason: if nested {
format!(
"Branch '{}' tip {} archives '{}' under an invalid nested layout; expected {}/YYYY-MM-DD-{}/proposal.md",
item.revision,
item.tip,
item.change_id,
archive_layout::ARCHIVE_PREFIX,
item.change_id
)
} else {
format!(
"Branch '{}' tip {} has no valid archive proposal for '{}'; the final merge would integrate an unarchived change",
item.revision, item.tip, item.change_id
)
},
});
}
if let Some(reason) = latch.refusal(&item.change_id) {
return Ok(BatchState::MergeNotAuthorized {
change_id: item.change_id.clone(),
reason,
});
}
let completion =
read_task_completion(evidence, &item.change_id, &item.tip, &worktree_tree_paths).await;
if let Some(reason) = completion.detail(&item.change_id) {
latch.refuse(&item.change_id, &reason);
return Ok(BatchState::MergeNotAuthorized {
change_id: item.change_id.clone(),
reason,
});
}
return Ok(BatchState::FinalMergeMissing {
change_id: item.change_id.clone(),
worktree: item.worktree.clone(),
required_target_state: target_head.clone(),
});
}
let tree_paths = evidence
.committed_tree_paths(&target_head, archive_layout::ACTIVE_CHANGES_PREFIX)
.await
.map_err(|error| unsafe_evidence("Failed to read committed target tree", error))?;
for (item, integration) in validated.iter().zip(integrations.iter()) {
validate_cleanup_commits(
evidence,
item,
integration,
base_revision,
&target_head,
&target_lineage,
)
.await?;
let live = archive_layout::paths_contain_active_change(&tree_paths, &item.change_id);
let archived = archive_layout::paths_contain_valid_archive(&tree_paths, &item.change_id);
if live && archived {
return Ok(BatchState::ResurrectionCleanupRequired {
change_id: item.change_id.clone(),
reason: format!(
"Committed target HEAD {} holds both openspec/changes/{} and its valid archive entry",
target_head, item.change_id
),
});
}
}
let target_conflicts = evidence
.target_conflict_paths()
.await
.map_err(|error| unsafe_evidence("Failed to read target index stages", error))?;
if !target_conflicts.is_empty() {
return Err(BatchState::UnsafeEvidence {
reason: format!(
"Every item is integrated but the target index has conflict stages: {}",
target_conflicts.join(", ")
),
});
}
let clean = evidence
.target_is_clean()
.await
.map_err(|error| unsafe_evidence("Failed to read target cleanliness", error))?;
if !clean {
return Err(BatchState::UnsafeEvidence {
reason: "Every item is integrated but the target index or worktree is not clean"
.to_string(),
});
}
Ok(BatchState::Complete)
}
pub async fn verify_final_integration(
evidence: &dyn ResolveEvidence,
items: &[SequentialMergeItem],
base_revision: &str,
) -> std::result::Result<(), String> {
let validated = validate_items(evidence, items)
.await
.map_err(|state| state.diagnosis())?;
let target_head = evidence
.target_head()
.await
.map_err(|error| format!("Failed to read target HEAD: {}", error))?;
let mut missing = Vec::new();
for item in &validated {
match classify_item_integration(evidence, item, base_revision, &target_head)
.await
.map_err(|state| state.diagnosis())?
{
ItemIntegration::Exact { .. } | ItemIntegration::Historical => {}
ItemIntegration::NotIntegrated => missing.push(item.change_id.clone()),
}
}
if missing.is_empty() {
Ok(())
} else {
Err(format!(
"Missing merge commit message containing change_id(s): {}",
missing.join(", ")
))
}
}
pub struct GitResolveEvidence {
repo_root: PathBuf,
}
impl GitResolveEvidence {
pub fn new(repo_root: impl Into<PathBuf>) -> Self {
Self {
repo_root: repo_root.into(),
}
}
}
#[async_trait]
impl ResolveEvidence for GitResolveEvidence {
async fn validate_worktree(
&self,
supplied_path: &Path,
expected_branch: &str,
) -> WorktreeIdentity {
git_commands::validate_worktree_identity(&self.repo_root, supplied_path, expected_branch)
.await
}
async fn worktree_merge_in_progress(&self, worktree: &Path) -> EvidenceResult<bool> {
git_commands::is_merge_in_progress(worktree)
.await
.map_err(|error| error.to_string())
}
async fn worktree_conflicts(&self, worktree: &Path) -> EvidenceResult<Vec<String>> {
git_commands::get_conflict_files(worktree)
.await
.map_err(|error| error.to_string())
}
async fn target_head(&self) -> EvidenceResult<String> {
git_commands::rev_parse_commit(&self.repo_root, "HEAD")
.await
.map_err(|error| error.to_string())?
.ok_or_else(|| "target repository has no HEAD commit".to_string())
}
async fn target_merge_head(&self) -> EvidenceResult<Option<String>> {
git_commands::merge_head(&self.repo_root)
.await
.map_err(|error| error.to_string())
}
async fn target_conflict_paths(&self) -> EvidenceResult<Vec<String>> {
let entries = git_commands::index_conflict_entries(&self.repo_root)
.await
.map_err(|error| error.to_string())?;
let mut paths: Vec<String> = entries.into_iter().map(|entry| entry.path).collect();
paths.sort();
paths.dedup();
Ok(paths)
}
async fn target_is_clean(&self) -> EvidenceResult<bool> {
git_commands::is_clean_including_untracked(&self.repo_root)
.await
.map_err(|error| error.to_string())
}
async fn target_index_paths(&self, prefix: &str) -> EvidenceResult<Vec<String>> {
git_commands::index_stage0_paths(&self.repo_root, prefix)
.await
.map_err(|error| error.to_string())
}
async fn parents_of(&self, commit: &str) -> EvidenceResult<Vec<String>> {
git_commands::parents_of(&self.repo_root, commit)
.await
.map_err(|error| error.to_string())
}
async fn is_ancestor(&self, ancestor: &str, descendant: &str) -> EvidenceResult<bool> {
git_commands::is_ancestor(&self.repo_root, ancestor, descendant)
.await
.map_err(|error| error.to_string())
}
async fn first_parent_lineage(&self, tip: &str) -> EvidenceResult<Vec<String>> {
git_commands::first_parent_lineage(&self.repo_root, tip)
.await
.map_err(|error| error.to_string())
}
async fn merge_base(&self, a: &str, b: &str) -> EvidenceResult<Option<String>> {
git_commands::merge_base(&self.repo_root, a, b)
.await
.map_err(|error| error.to_string())
}
async fn commits_with_exact_subject(
&self,
from: Option<&str>,
to: &str,
subject: &str,
) -> EvidenceResult<Vec<String>> {
git_commands::commits_with_exact_subject(&self.repo_root, from, to, subject)
.await
.map_err(|error| error.to_string())
}
async fn commit_diff_entries(&self, commit: &str) -> EvidenceResult<Vec<CommitDiffEntry>> {
git_commands::commit_diff_entries(&self.repo_root, commit)
.await
.map_err(|error| error.to_string())
}
async fn committed_tree_paths(
&self,
revision: &str,
prefix: &str,
) -> EvidenceResult<Vec<String>> {
git_commands::committed_tree_paths(&self.repo_root, revision, prefix)
.await
.map_err(|error| error.to_string())
}
async fn committed_file_text(
&self,
revision: &str,
path: &str,
) -> EvidenceResult<Option<String>> {
git_commands::committed_file_text(&self.repo_root, revision, path)
.await
.map_err(|error| error.to_string())
}
}