use std::collections::HashSet;
use std::path::PathBuf;
use tracing::{debug, info, warn};
use crate::dependency_targets::{
classify_dependency_target, collect_active_change_ids, collect_archived_change_ids,
collect_rejected_change_ids, DependencyTargetClass,
};
use crate::error::{OrchestratorError, Result};
use crate::vcs::WorkspaceManager;
use super::work_snapshot::ReducerWorkSnapshot;
use super::{DependencyBlockerFingerprint, ParallelExecutor};
#[derive(Debug, Clone, PartialEq, Eq)]
pub(super) struct EffectiveDependencyBaseEvidence {
pub(super) base_ref: String,
pub(super) revision: String,
}
impl EffectiveDependencyBaseEvidence {
pub(super) fn signature_material(&self) -> String {
format!("{}@{}", self.base_ref, self.revision)
}
}
#[derive(Debug, Clone)]
pub(super) struct DependencyContext {
repo_root: PathBuf,
queued_ids: HashSet<String>,
in_flight_ids: HashSet<String>,
active_ids: HashSet<String>,
archived_ids: HashSet<String>,
rejected_ids: HashSet<String>,
terminal_error_ids: HashSet<String>,
resolving_ids: HashSet<String>,
resolve_wait_ids: HashSet<String>,
ordinary_ineligible_ids: Option<HashSet<String>>,
effective_dependency_base: Option<String>,
}
impl DependencyContext {
pub(super) async fn from_executor(
executor: &ParallelExecutor,
queued_ids: impl IntoIterator<Item = impl AsRef<str>>,
in_flight: &HashSet<String>,
) -> Self {
let snapshot = executor.capture_reducer_work_snapshot().await;
Self::from_snapshot(executor.repo_root.clone(), queued_ids, in_flight, &snapshot)
}
pub(super) fn from_snapshot(
repo_root: PathBuf,
queued_ids: impl IntoIterator<Item = impl AsRef<str>>,
in_flight: &HashSet<String>,
snapshot: &ReducerWorkSnapshot,
) -> Self {
let queued_ids = queued_ids
.into_iter()
.map(|id| id.as_ref().to_string())
.collect::<HashSet<_>>();
let ordinary_ineligible_ids = snapshot.reducer_present().then(|| {
queued_ids
.iter()
.filter(|id| !snapshot.is_ordinary_queue_eligible(id))
.cloned()
.collect::<HashSet<_>>()
});
let terminal_error_ids = snapshot.terminal_error_ids().clone();
let resolving_ids = snapshot.resolving_ids().clone();
let resolve_wait_ids = snapshot.resolve_wait_ids().clone();
let in_flight_ids = in_flight.iter().cloned().collect::<HashSet<_>>();
let active_ids = collect_active_change_ids(&repo_root);
let archived_ids = collect_archived_change_ids(&repo_root);
let rejected_ids = collect_rejected_change_ids(&repo_root);
debug!(
queued = queued_ids.len(),
in_flight = in_flight_ids.len(),
active = active_ids.len(),
archived = archived_ids.len(),
rejected = rejected_ids.len(),
terminal_error = terminal_error_ids.len(),
resolving = resolving_ids.len(),
resolve_wait = resolve_wait_ids.len(),
reducer_present = snapshot.reducer_present(),
reducer_evidence_complete = snapshot.is_complete(),
"Built dependency classification context"
);
Self {
repo_root,
queued_ids,
in_flight_ids,
active_ids,
archived_ids,
rejected_ids,
terminal_error_ids,
resolving_ids,
resolve_wait_ids,
ordinary_ineligible_ids,
effective_dependency_base: None,
}
}
pub(super) fn withholds_ordinary_queue_intent(&self, change_id: &str) -> bool {
self.ordinary_ineligible_ids
.as_ref()
.map(|ineligible| ineligible.contains(change_id))
.unwrap_or(false)
}
pub(super) fn classify(&self, dep_id: &str) -> DependencyTargetClass {
let class = classify_dependency_target(
dep_id,
self.queued_ids.iter().map(String::as_str),
self.in_flight_ids.iter().map(String::as_str),
self.active_ids.iter().map(String::as_str),
&self.archived_ids,
&self.rejected_ids,
);
if matches!(class, DependencyTargetClass::Rejected) {
return class;
}
if self.terminal_error_ids.contains(dep_id) {
DependencyTargetClass::Error
} else if self.resolving_ids.contains(dep_id) || self.resolve_wait_ids.contains(dep_id) {
DependencyTargetClass::Resolving
} else {
class
}
}
pub(super) fn is_terminal_error_change(&self, change_id: &str) -> bool {
self.terminal_error_ids.contains(change_id)
}
pub(super) async fn is_blocked(
&mut self,
dependencies: &[String],
workspace_manager: &dyn WorkspaceManager,
) -> Option<DependencyBlockerFingerprint> {
let mut blockers = Vec::new();
for dep_id in dependencies {
let class = self.classify(dep_id);
if !matches!(class, DependencyTargetClass::Archived) {
blockers.push((dep_id.clone(), class.as_str().to_string()));
continue;
}
let resolved = self
.is_dependency_resolved_with_base(dep_id, workspace_manager)
.await
.map(|(resolved, _)| resolved)
.unwrap_or(false);
if !resolved {
blockers.push((dep_id.clone(), class.as_str().to_string()));
}
}
(!blockers.is_empty()).then_some(blockers)
}
pub(super) async fn effective_dependency_base(
&mut self,
workspace_manager: &dyn WorkspaceManager,
) -> Result<&str> {
if self.effective_dependency_base.is_none() {
let original_branch = workspace_manager
.ensure_original_branch_initialized()
.await
.map_err(OrchestratorError::from_vcs_error)?;
let effective_base = match workspace_manager.current_branch().await {
Ok(Some(current_branch)) if current_branch != original_branch => {
debug!(
original_branch = %original_branch,
effective_dependency_base = %current_branch,
"Using current integration branch as effective dependency base"
);
current_branch
}
Ok(Some(_)) | Ok(None) => original_branch,
Err(err) => {
warn!(
error = %err,
"Failed to determine current branch for effective dependency base"
);
return Err(OrchestratorError::from_vcs_error(err));
}
};
self.effective_dependency_base = Some(effective_base);
}
Ok(self
.effective_dependency_base
.as_deref()
.expect("effective dependency base initialized above"))
}
pub(super) async fn effective_dependency_base_evidence(
&mut self,
workspace_manager: &dyn WorkspaceManager,
) -> Result<EffectiveDependencyBaseEvidence> {
let base_ref = self
.effective_dependency_base(workspace_manager)
.await?
.to_string();
let revision = workspace_manager
.revision_for_ref(&base_ref)
.await
.map_err(OrchestratorError::from_vcs_error)?;
Ok(EffectiveDependencyBaseEvidence { base_ref, revision })
}
pub(super) async fn is_dependency_resolved_with_base(
&mut self,
dep_id: &str,
workspace_manager: &dyn WorkspaceManager,
) -> Result<(bool, String)> {
let effective_base = self
.effective_dependency_base(workspace_manager)
.await?
.to_string();
match crate::execution::state::is_merged_to_base(dep_id, &self.repo_root, &effective_base)
.await
{
Ok(is_merged) => Ok((is_merged, effective_base)),
Err(e) => {
warn!(
dependency = %dep_id,
effective_dependency_base = %effective_base,
error = %e,
"Failed to check if dependency is merged to effective base; assuming not resolved"
);
Ok((false, effective_base))
}
}
}
pub(super) fn blocker_fingerprint(
change_id: &str,
blockers: &[(String, DependencyTargetClass)],
) -> DependencyBlockerFingerprint {
let mut fingerprint = blockers
.iter()
.map(|(dep_id, class)| (dep_id.clone(), class.as_str().to_string()))
.collect::<Vec<_>>();
fingerprint.sort();
fingerprint.insert(0, ("change_id".to_string(), change_id.to_string()));
fingerprint
}
pub(super) fn log_archived_dependency_check(change_id: &str, dep_id: &str) {
debug!(
change_id = %change_id,
dependency = %dep_id,
"Archived dependency evidence found; verifying base-branch merge before dispatch"
);
}
pub(super) fn log_dependency_resolved(
change_id: &str,
dep_id: &str,
class: DependencyTargetClass,
effective_base: &str,
) {
if matches!(class, DependencyTargetClass::Archived) {
debug!(
change_id = %change_id,
dependency = %dep_id,
effective_dependency_base = %effective_base,
"Archived dependency is merged into effective dependency base"
);
}
}
pub(super) fn log_dependency_unresolved(
change_id: &str,
dep_id: &str,
class: DependencyTargetClass,
effective_base: &str,
) {
if matches!(class, DependencyTargetClass::Archived) {
info!(
change_id = %change_id,
dependency = %dep_id,
effective_dependency_base = %effective_base,
"Archived dependency is not merged into effective dependency base; dispatch remains blocked"
);
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::orchestration::state::OrchestratorState;
use std::sync::Arc;
use std::time::Duration;
use tempfile::TempDir;
use tokio::sync::RwLock;
async fn capture(state: &Arc<RwLock<OrchestratorState>>) -> ReducerWorkSnapshot {
let guard = state.read().await;
ReducerWorkSnapshot::from_state(&guard)
}
fn write_change(root: &std::path::Path, id: &str) {
let change_dir = root.join("openspec/changes").join(id);
std::fs::create_dir_all(&change_dir).unwrap();
std::fs::write(change_dir.join("proposal.md"), "# Change\n").unwrap();
}
fn write_change_with_verifications(root: &std::path::Path, id: &str, verifications: &str) {
let change_dir = root.join("openspec/changes").join(id);
std::fs::create_dir_all(&change_dir).unwrap();
std::fs::write(
change_dir.join("proposal.md"),
format!("---\nverifications:\n{verifications}---\n# Change\n"),
)
.unwrap();
}
#[tokio::test]
async fn verification_role_metadata_does_not_change_scheduler_classification() {
let temp_dir = TempDir::new().unwrap();
write_change_with_verifications(
temp_dir.path(),
"queued-gate",
" - id: deployed-smoke\n phase: post-integration\n execution_class: deployed-service\n completion_role: change-blocking\n",
);
write_change_with_verifications(
temp_dir.path(),
"active-observation",
" - id: release-smoke\n phase: post-integration\n execution_class: deployed-service\n completion_role: operational-observation\n",
);
let archive_dir = temp_dir
.path()
.join("openspec/changes/archive/2026-07-21-archived-gate");
std::fs::create_dir_all(&archive_dir).unwrap();
std::fs::write(
archive_dir.join("proposal.md"),
"---\nverifications:\n - id: deployed-smoke\n phase: post-integration\n execution_class: physical-device\n completion_role: change-blocking\n---\n# Archived\n",
)
.unwrap();
let context = DependencyContext::from_snapshot(
temp_dir.path().to_path_buf(),
["queued-gate"],
&HashSet::new(),
&ReducerWorkSnapshot::absent(),
);
assert_eq!(
context.classify("queued-gate"),
DependencyTargetClass::Queued
);
assert_eq!(
context.classify("active-observation"),
DependencyTargetClass::ActiveButNotQueued
);
assert_eq!(
context.classify("archived-gate"),
DependencyTargetClass::Archived
);
}
#[tokio::test]
async fn reducer_snapshot_contention_defers_dependency_context_until_writer_releases() {
let temp_dir = TempDir::new().unwrap();
write_change(temp_dir.path(), "resolving-a");
let state = Arc::new(RwLock::new(OrchestratorState::new(
vec!["resolving-a".to_string()],
1,
)));
{
let mut guard = state.write().await;
guard.apply_execution_event(&crate::events::ExecutionEvent::ResolveStarted {
change_id: "resolving-a".to_string(),
command: "resolve".to_string(),
});
}
let write_guard = state.write().await;
let repo_root = temp_dir.path().to_path_buf();
let contended = Arc::clone(&state);
let mut classification = tokio::spawn(async move {
let snapshot = capture(&contended).await;
DependencyContext::from_snapshot(repo_root, ["dependent"], &HashSet::new(), &snapshot)
.classify("resolving-a")
});
assert!(
tokio::time::timeout(Duration::from_millis(50), &mut classification)
.await
.is_err(),
"classification must suspend behind the writer instead of resolving from partial evidence"
);
drop(write_guard);
let class = tokio::time::timeout(Duration::from_secs(5), classification)
.await
.expect("released writer must let the pending snapshot acquisition finish")
.expect("classification task must not panic");
assert_eq!(
class,
DependencyTargetClass::Resolving,
"the resumed evaluation classifies from real reducer evidence"
);
}
#[tokio::test]
async fn context_classifies_from_single_collected_evidence_snapshot() {
let temp_dir = TempDir::new().unwrap();
write_change(temp_dir.path(), "active-a");
write_change(temp_dir.path(), "resolving-a");
write_change(temp_dir.path(), "resolve-wait-a");
let archive_dir = temp_dir
.path()
.join("openspec/changes/archive/2026-06-17-archived-a");
std::fs::create_dir_all(&archive_dir).unwrap();
std::fs::write(archive_dir.join("proposal.md"), "# Archived\n").unwrap();
write_change(temp_dir.path(), "rejected-a");
std::fs::write(
temp_dir
.path()
.join("openspec/changes/rejected-a/REJECTED.md"),
"# REJECTED\n",
)
.unwrap();
let in_flight = HashSet::from(["flight-a".to_string()]);
let state = Arc::new(RwLock::new(OrchestratorState::new(
vec!["resolving-a".to_string(), "resolve-wait-a".to_string()],
1,
)));
let mut state_guard = state.write().await;
state_guard.apply_execution_event(&crate::events::ExecutionEvent::ResolveStarted {
change_id: "resolving-a".to_string(),
command: "resolve".to_string(),
});
state_guard.apply_command(crate::orchestration::state::ReducerCommand::ResolveMerge(
"resolve-wait-a".to_string(),
));
drop(state_guard);
let context = DependencyContext::from_snapshot(
temp_dir.path().to_path_buf(),
["queued-a"],
&in_flight,
&capture(&state).await,
);
let cases = [
("queued-a", DependencyTargetClass::Queued),
("flight-a", DependencyTargetClass::InFlight),
("active-a", DependencyTargetClass::ActiveButNotQueued),
("archived-a", DependencyTargetClass::Archived),
("resolving-a", DependencyTargetClass::Resolving),
("resolve-wait-a", DependencyTargetClass::Resolving),
("rejected-a", DependencyTargetClass::Rejected),
("missing-a", DependencyTargetClass::Missing),
];
for (target, expected) in cases {
assert_eq!(context.classify(target), expected, "target={target}");
}
}
}