use crate::ai_command_runner::{AiCommandRunner, RunCommandScope, SharedStaggerState};
use crate::analyzer::{ParallelGroup, ParallelizationAnalyzer};
use crate::config::OrchestratorConfig;
use crate::dependency_targets::union_metadata_dependencies;
use crate::error::Result;
use crate::hooks::HookRunner;
use crate::openspec::Change;
use crate::parallel::dedup::{DiagnosticDeduplicationKey, DiagnosticDeduplicationStore};
use crate::parallel::{ParallelEvent, ParallelExecutor, PostArchiveAction, SchedulerRunReport};
use std::collections::{HashMap, HashSet};
use std::path::PathBuf;
use std::sync::Arc;
use tokio::sync::{mpsc, Mutex};
use tokio_util::sync::CancellationToken;
use tracing::{debug, error, info, warn};
type AnalysisDiagnosticStore = Arc<Mutex<DiagnosticDeduplicationStore<DiagnosticDeduplicationKey>>>;
pub struct ParallelRunService {
config: OrchestratorConfig,
repo_root: PathBuf,
no_resume: bool,
shared_stagger_state: SharedStaggerState,
post_archive_action: PostArchiveAction,
shared_orchestrator_state:
Arc<tokio::sync::RwLock<crate::orchestration::state::OrchestratorState>>,
ai_runner: AiCommandRunner,
run_command_scope: RunCommandScope,
diagnostic_dedup: AnalysisDiagnosticStore,
upstream_integration: Option<crate::upstream::UpstreamRuntime>,
explicit_target_plan: Option<crate::orchestration::target_resolution::ExplicitTargetPlan>,
graceful_stop: Option<Arc<std::sync::atomic::AtomicBool>>,
}
impl ParallelRunService {
pub fn new(repo_root: PathBuf, config: OrchestratorConfig) -> Self {
let shared_stagger_state: SharedStaggerState = Arc::new(Mutex::new(None));
let run_command_scope = RunCommandScope::new();
let ai_runner = AiCommandRunner::for_run(
&config,
shared_stagger_state.clone(),
run_command_scope.clone(),
);
let shared_orchestrator_state = Arc::new(tokio::sync::RwLock::new(
crate::orchestration::state::OrchestratorState::new(Vec::new(), 1),
));
Self {
config,
repo_root,
no_resume: false,
shared_stagger_state,
post_archive_action: PostArchiveAction::MergeToBase,
shared_orchestrator_state,
ai_runner,
run_command_scope,
diagnostic_dedup: Arc::new(Mutex::new(DiagnosticDeduplicationStore::new())),
upstream_integration: None,
explicit_target_plan: None,
graceful_stop: None,
}
}
pub fn new_with_shared_state(
repo_root: PathBuf,
config: OrchestratorConfig,
shared_stagger_state: SharedStaggerState,
) -> Self {
let run_command_scope = RunCommandScope::new();
let ai_runner = AiCommandRunner::for_run(
&config,
shared_stagger_state.clone(),
run_command_scope.clone(),
);
let shared_orchestrator_state = Arc::new(tokio::sync::RwLock::new(
crate::orchestration::state::OrchestratorState::new(Vec::new(), 1),
));
Self {
config,
repo_root,
no_resume: false,
shared_stagger_state,
post_archive_action: PostArchiveAction::MergeToBase,
shared_orchestrator_state,
ai_runner,
run_command_scope,
diagnostic_dedup: Arc::new(Mutex::new(DiagnosticDeduplicationStore::new())),
upstream_integration: None,
explicit_target_plan: None,
graceful_stop: None,
}
}
pub fn set_no_resume(&mut self, no_resume: bool) {
self.no_resume = no_resume;
}
pub fn set_post_archive_action(&mut self, action: PostArchiveAction) {
self.post_archive_action = action;
}
pub fn set_upstream_integration(&mut self, runtime: crate::upstream::UpstreamRuntime) {
self.upstream_integration = Some(runtime);
}
#[cfg(test)]
pub fn upstream_integration(&self) -> Option<&crate::upstream::UpstreamRuntime> {
self.upstream_integration.as_ref()
}
pub fn set_graceful_stop_flag(&mut self, graceful_stop: Arc<std::sync::atomic::AtomicBool>) {
self.graceful_stop = Some(graceful_stop);
}
pub fn set_explicit_target_plan(
&mut self,
plan: crate::orchestration::target_resolution::ExplicitTargetPlan,
) {
self.explicit_target_plan = Some(plan);
}
fn install_upstream_integration(&self, executor: &mut ParallelExecutor) {
if let Some(runtime) = &self.upstream_integration {
executor.set_upstream_integration(runtime.clone());
}
if let Some(plan) = &self.explicit_target_plan {
executor.set_explicit_target_plan(plan.clone());
}
}
#[cfg(test)]
pub fn post_archive_action(&self) -> &PostArchiveAction {
&self.post_archive_action
}
pub fn set_shared_orchestrator_state(
&mut self,
shared_state: Arc<tokio::sync::RwLock<crate::orchestration::state::OrchestratorState>>,
) {
self.shared_orchestrator_state = shared_state;
}
#[allow(dead_code)] pub fn run_command_scope(&self) -> RunCommandScope {
self.run_command_scope.clone()
}
pub fn set_run_command_scope(&mut self, scope: RunCommandScope) {
self.ai_runner.set_run_command_scope(scope.clone());
self.run_command_scope = scope;
}
pub async fn check_vcs_available(&self) -> Result<()> {
if !crate::cli::check_git_workspace_usable() {
return Err(crate::error::OrchestratorError::GitCommand(
"Git repository not available for worktree execution".to_string(),
));
}
Ok(())
}
pub fn create_executor_with_queue_state(
&self,
event_tx: Option<mpsc::Sender<ParallelEvent>>,
cancel_token: Option<CancellationToken>,
shared_queue_change: Option<std::sync::Arc<tokio::sync::Mutex<Option<std::time::Instant>>>>,
dynamic_queue: Option<std::sync::Arc<crate::tui::queue::DynamicQueue>>,
manual_resolve_counter: Option<std::sync::Arc<std::sync::atomic::AtomicUsize>>,
shared_orchestrator_state: Option<
std::sync::Arc<tokio::sync::RwLock<crate::orchestration::state::OrchestratorState>>,
>,
) -> ParallelExecutor {
let vcs_backend = self.config.get_vcs_backend();
let hooks = if let Some(ref tx) = event_tx {
HookRunner::with_event_tx(self.config.get_hooks(), &self.repo_root, tx.clone())
} else {
HookRunner::new(self.config.get_hooks(), &self.repo_root)
};
let has_dynamic_queue = dynamic_queue.is_some();
let mut executor = ParallelExecutor::with_backend_and_queue_and_stagger(
self.repo_root.clone(),
self.config.clone(),
event_tx,
vcs_backend,
shared_queue_change,
Some(self.shared_stagger_state.clone()),
);
executor.set_run_command_scope(self.run_command_scope.clone());
executor.set_no_resume(self.no_resume);
executor.set_post_archive_action(self.post_archive_action.clone());
self.install_upstream_integration(&mut executor);
if has_dynamic_queue {
executor.set_persistent_lifetime();
}
if let Some(graceful_stop) = &self.graceful_stop {
executor.set_graceful_stop_flag(graceful_stop.clone());
}
executor.set_hooks(hooks);
if let Some(token) = cancel_token {
executor.set_cancel_token(token);
}
if let Some(queue) = dynamic_queue {
executor.set_dynamic_queue(queue);
}
if let Some(counter) = manual_resolve_counter {
executor.set_manual_resolve_counter(counter);
}
if let Some(shared_state) = shared_orchestrator_state {
executor.set_shared_orchestrator_state(shared_state);
}
executor
}
async fn filter_committed_changes(
&self,
changes: Vec<Change>,
) -> Result<(Vec<Change>, Vec<String>)> {
filter_committed_changes_at(&self.repo_root, changes).await
}
async fn prepare_parallel_execution(
&self,
changes: Vec<Change>,
event_tx: &mpsc::Sender<ParallelEvent>,
allow_empty_when_resolve_wait: bool,
) -> Result<Option<Vec<Change>>> {
let (changes, skipped) = self.filter_committed_changes(changes).await?;
if !skipped.is_empty() {
let message = format!("Skipping uncommitted changes: {}", skipped.join(", "));
warn!("{}", message);
let _ = event_tx
.send(ParallelEvent::Warning {
title: "Uncommitted changes skipped".to_string(),
message,
})
.await;
let _ = event_tx
.send(ParallelEvent::ParallelStartRejected {
change_ids: skipped.clone(),
reason: "uncommitted or not in HEAD".to_string(),
})
.await;
}
if changes.is_empty() {
if allow_empty_when_resolve_wait {
info!(
"No committed changes available, but scheduler-owned ResolveWait retry is present; continuing with empty queue"
);
return Ok(Some(changes));
}
info!("No committed changes available for parallel execution");
return Ok(None);
}
Ok(Some(changes))
}
pub async fn run_parallel<F>(
&self,
changes: Vec<Change>,
cancel_token: Option<CancellationToken>,
event_handler: F,
) -> Result<()>
where
F: Fn(ParallelEvent) + Send + Sync + 'static,
{
let (event_tx, mut event_rx) = mpsc::channel::<ParallelEvent>(100);
{
let mut guard = self.shared_orchestrator_state.write().await;
for change in &changes {
guard.add_dynamic_change(change.id.clone());
}
}
let changes = match self
.prepare_parallel_execution(changes, &event_tx, true)
.await?
{
Some(changes) => changes,
None => {
drop(event_tx);
while let Some(event) = event_rx.recv().await {
event_handler(event);
}
return Ok(());
}
};
{
let mut guard = self.shared_orchestrator_state.write().await;
for change in &changes {
guard.apply_command(crate::orchestration::state::ReducerCommand::AddToQueue(
change.id.clone(),
));
}
}
let forward_handle = tokio::spawn(async move {
while let Some(event) = event_rx.recv().await {
let is_completed =
matches!(event, ParallelEvent::AllCompleted | ParallelEvent::Stopped);
event_handler(event);
if is_completed {
break;
}
}
});
let executor = self.create_executor_with_queue_state(
Some(event_tx.clone()),
cancel_token,
None,
None,
None,
Some(self.shared_orchestrator_state.clone()),
);
let result = self
.run_parallel_order_based_with_executor(executor, changes, event_tx)
.await;
let _ = forward_handle.await;
result.map(|_report| ())
}
#[allow(clippy::too_many_arguments)]
pub async fn run_parallel_with_channel_and_queue_state(
&self,
changes: Vec<Change>,
event_tx: mpsc::Sender<ParallelEvent>,
cancel_token: Option<CancellationToken>,
shared_queue_change: Option<std::sync::Arc<tokio::sync::Mutex<Option<std::time::Instant>>>>,
dynamic_queue: Option<std::sync::Arc<crate::tui::queue::DynamicQueue>>,
manual_resolve_counter: Option<std::sync::Arc<std::sync::atomic::AtomicUsize>>,
shared_orchestrator_state: Option<
std::sync::Arc<tokio::sync::RwLock<crate::orchestration::state::OrchestratorState>>,
>,
explicit_retry: bool,
) -> Result<SchedulerRunReport> {
let mut executor = self.create_executor_with_queue_state(
Some(event_tx.clone()),
cancel_token,
shared_queue_change,
dynamic_queue,
manual_resolve_counter,
shared_orchestrator_state,
);
executor.set_explicit_retry(explicit_retry);
self.run_parallel_order_based_with_executor(executor, changes, event_tx)
.await
}
pub async fn run_parallel_order_based_with_executor(
&self,
mut executor: ParallelExecutor,
changes: Vec<Change>,
event_tx: mpsc::Sender<ParallelEvent>,
) -> Result<SchedulerRunReport> {
executor.ensure_shared_orchestrator_state(self.shared_orchestrator_state.clone());
let allow_empty_queue = changes.is_empty()
&& (executor.has_resolve_wait().await || executor.has_upstream_integration());
let changes = match self
.prepare_parallel_execution(changes, &event_tx, allow_empty_queue)
.await?
{
Some(changes) => changes,
None => return Ok(SchedulerRunReport::Completed),
};
info!(
"Starting order-based parallel execution with re-analysis for {} changes",
changes.len()
);
let config = self.config.clone();
let repo_root = self.repo_root.clone();
let shared_stagger_state = self.shared_stagger_state.clone();
executor
.execute_with_order_based_reanalysis(
changes,
move |remaining, in_flight_ids, iteration| {
let config = config.clone();
let repo_root = repo_root.clone();
let event_tx = event_tx.clone();
let shared_stagger_state = shared_stagger_state.clone();
Box::pin(async move {
let service = ParallelRunService::new_with_shared_state(
repo_root,
config,
shared_stagger_state,
);
service
.analyze_order_with_sender(
remaining,
in_flight_ids,
Some(&event_tx),
iteration,
)
.await
})
},
)
.await
}
pub async fn analyze_and_group_public(&self, changes: &[Change]) -> Vec<ParallelGroup> {
self.analyze_and_group(changes).await
}
async fn analyze_and_group(&self, changes: &[Change]) -> Vec<ParallelGroup> {
self.analyze_and_group_with_sender(changes, None, 1).await
}
async fn analyze_order_with_sender(
&self,
changes: &[Change],
in_flight_ids: &[String],
event_tx: Option<&mpsc::Sender<ParallelEvent>>,
iteration: u32,
) -> crate::analyzer::AnalysisOutcome {
if self.config.use_llm_analysis() {
info!("Using LLM analysis for parallelization (analyze_command)");
match self
.analyze_order_with_llm_streaming(changes, in_flight_ids, event_tx, iteration)
.await
{
Ok(result) => {
info!(
"LLM analysis successful: {} changes in order",
result.order.len()
);
return crate::analyzer::AnalysisOutcome::healthy(result);
}
Err(e) => {
self.emit_recoverable_analysis_fallback_diagnostic_once(
changes,
in_flight_ids,
event_tx,
&e.to_string(),
)
.await;
return crate::analyzer::AnalysisOutcome::recoverable_failure_fallback(
Self::metadata_dependency_analysis_result(changes),
);
}
}
}
info!("LLM analysis disabled, using metadata-dependency-only analysis");
crate::analyzer::AnalysisOutcome::intentional_metadata_only(
Self::metadata_dependency_analysis_result(changes),
)
}
fn log_recoverable_analysis_fallback(error: &dyn std::fmt::Display) {
warn!(
error = %error,
"LLM analysis failed; falling back to metadata-dependency-only analysis"
);
}
pub(crate) fn recoverable_analysis_fallback_diagnostic(
changes: &[Change],
in_flight_ids: &[String],
error: &str,
) -> (DiagnosticDeduplicationKey, String) {
let mut queued_ids: Vec<String> = changes.iter().map(|change| change.id.clone()).collect();
queued_ids.sort();
let mut in_flight = in_flight_ids.to_vec();
in_flight.sort();
let normalized_error = error.trim().to_string();
let key = DiagnosticDeduplicationKey::AnalysisFailure {
queued_ids: queued_ids.clone(),
in_flight_ids: in_flight.clone(),
error: normalized_error.clone(),
};
let message = format!(
"{}: error={}, queued={:?}, in_flight={:?}",
crate::events::RECOVERABLE_ANALYSIS_FALLBACK_MARKER,
normalized_error,
queued_ids,
in_flight
);
(key, message)
}
async fn emit_recoverable_analysis_fallback_diagnostic_once(
&self,
changes: &[Change],
in_flight_ids: &[String],
event_tx: Option<&mpsc::Sender<ParallelEvent>>,
error: &str,
) {
let (key, message) =
Self::recoverable_analysis_fallback_diagnostic(changes, in_flight_ids, error);
let error = error.trim().to_string();
let mut dedup = self.diagnostic_dedup.lock().await;
dedup
.emit_or_suppress(
key,
move || async move {
Self::log_recoverable_analysis_fallback(&error);
if let Some(tx) = event_tx {
let _ = tx
.send(ParallelEvent::Log(crate::events::LogEntry::warn(&message)))
.await;
}
},
|| {
debug!("Suppressing repeated analysis fallback diagnostic");
},
)
.await;
}
fn metadata_dependency_analysis_result(changes: &[Change]) -> crate::analyzer::AnalysisResult {
let mut dependencies = HashMap::new();
for change in changes {
union_metadata_dependencies(&mut dependencies, &change.id, &change.dependencies);
}
crate::analyzer::AnalysisResult {
order: changes.iter().map(|c| c.id.clone()).collect(),
dependencies,
groups: None,
}
}
async fn analyze_and_group_with_sender(
&self,
changes: &[Change],
event_tx: Option<&mpsc::Sender<ParallelEvent>>,
iteration: u32,
) -> Vec<ParallelGroup> {
if self.config.use_llm_analysis() {
info!("Using LLM analysis for parallelization (analyze_command)");
match self
.analyze_with_llm_streaming(changes, event_tx, iteration)
.await
{
Ok(groups) => {
info!("LLM analysis successful: {} groups", groups.len());
return groups;
}
Err(e) => {
error!("LLM analysis failed: {}", e);
warn!(
"Falling back to running all changes in parallel (no dependency analysis)"
);
}
}
} else {
info!("LLM analysis disabled, running all changes in parallel");
}
Self::all_parallel(changes)
}
fn all_parallel(changes: &[Change]) -> Vec<ParallelGroup> {
if changes.is_empty() {
return Vec::new();
}
vec![ParallelGroup {
id: 1,
changes: changes.iter().map(|c| c.id.clone()).collect(),
depends_on: Vec::new(),
}]
}
async fn analyze_order_with_llm_streaming(
&self,
changes: &[Change],
in_flight_ids: &[String],
event_tx: Option<&mpsc::Sender<ParallelEvent>>,
iteration: u32,
) -> Result<crate::analyzer::AnalysisResult> {
let analyzer = ParallelizationAnalyzer::new(
self.ai_runner.clone(),
self.config.clone(),
self.repo_root.clone(),
);
if let Some(tx) = event_tx {
let tx = tx.clone();
analyzer
.analyze_with_callback(changes, in_flight_ids, move |output| {
let _ = tx.try_send(ParallelEvent::AnalysisOutput {
output: output.clone(),
iteration,
});
})
.await
} else {
analyzer.analyze_with_inflight(changes, in_flight_ids).await
}
}
async fn analyze_with_llm_streaming(
&self,
changes: &[Change],
event_tx: Option<&mpsc::Sender<ParallelEvent>>,
iteration: u32,
) -> Result<Vec<ParallelGroup>> {
let analyzer = ParallelizationAnalyzer::new(
self.ai_runner.clone(),
self.config.clone(),
self.repo_root.clone(),
);
if let Some(tx) = event_tx {
let tx = tx.clone();
analyzer
.analyze_groups_with_callback(changes, move |output| {
let _ = tx.try_send(ParallelEvent::AnalysisOutput {
output: output.clone(),
iteration,
});
})
.await
} else {
analyzer.analyze_groups(changes).await
}
}
}
pub(crate) async fn filter_committed_changes_at(
repo_root: &std::path::Path,
changes: Vec<Change>,
) -> Result<(Vec<Change>, Vec<String>)> {
let committed_change_ids: HashSet<String> =
match crate::vcs::git::commands::list_changes_in_head(repo_root).await {
Ok(ids) => ids.into_iter().collect(),
Err(err) => {
warn!(
"Failed to load committed change snapshot; assuming all changes are committed: {}",
err
);
return Ok((changes, Vec::new()));
}
};
let uncommitted_file_change_ids: HashSet<String> =
match crate::vcs::git::commands::list_changes_with_uncommitted_files(repo_root).await {
Ok(ids) => ids.into_iter().collect(),
Err(err) => {
warn!(
"Failed to detect uncommitted files in changes; assuming no uncommitted files: {}",
err
);
HashSet::new()
}
};
let mut committed = Vec::new();
let mut skipped = Vec::new();
for change in changes {
if !committed_change_ids.contains(&change.id)
|| uncommitted_file_change_ids.contains(&change.id)
{
skipped.push(change.id);
} else {
committed.push(change);
}
}
skipped.sort();
Ok((committed, skipped))
}
#[cfg(test)]
mod tests {
use super::*;
use crate::openspec::ProposalMetadata;
use tempfile::TempDir;
use tokio::process::Command;
fn create_test_change(id: &str, dependencies: Vec<&str>) -> Change {
Change {
id: id.to_string(),
completed_tasks: 0,
total_tasks: 5,
last_modified: "1m ago".to_string(),
dependencies: dependencies.into_iter().map(String::from).collect(),
metadata: ProposalMetadata::default(),
}
}
fn create_test_config() -> OrchestratorConfig {
OrchestratorConfig {
apply_command: Some("echo apply {change_id}".to_string()),
archive_command: Some("echo archive {change_id}".to_string()),
analyze_command: Some("echo '{\"order\":[\"route\",\"policy\"],\"dependencies\":{\"route\":[\"ghost\"]}}'".to_string()),
acceptance_command: Some("echo acceptance".to_string()),
resolve_command: Some("echo resolve".to_string()),
..Default::default()
}
}
#[derive(Clone, Default)]
struct CaptureLayer(std::sync::Arc<std::sync::Mutex<Vec<(tracing::Level, String)>>>);
impl CaptureLayer {
fn records(&self) -> Vec<(tracing::Level, String)> {
self.0.lock().expect("capture layer mutex").clone()
}
fn warnings_containing(&self, needle: &str) -> usize {
self.records()
.iter()
.filter(|(level, fields)| *level == tracing::Level::WARN && fields.contains(needle))
.count()
}
}
impl<S> tracing_subscriber::Layer<S> for CaptureLayer
where
S: tracing::Subscriber,
{
fn on_event(
&self,
event: &tracing::Event<'_>,
_ctx: tracing_subscriber::layer::Context<'_, S>,
) {
struct Visitor {
fields: String,
}
impl tracing::field::Visit for Visitor {
fn record_debug(
&mut self,
field: &tracing::field::Field,
value: &dyn std::fmt::Debug,
) {
self.fields
.push_str(&format!("{}={:?};", field.name(), value));
}
}
let mut visitor = Visitor {
fields: String::new(),
};
event.record(&mut visitor);
self.0
.lock()
.expect("capture layer mutex")
.push((*event.metadata().level(), visitor.fields));
}
}
struct TracingCapture {
_exclusive: tokio::sync::MutexGuard<'static, ()>,
_subscriber: tracing::subscriber::DefaultGuard,
}
async fn capture_tracing() -> (CaptureLayer, TracingCapture) {
use tracing_subscriber::filter::LevelFilter;
use tracing_subscriber::layer::SubscriberExt;
use tracing_subscriber::Layer;
let exclusive = crate::test_support::tracing_capture_lock().lock().await;
let capture = CaptureLayer::default();
let subscriber =
tracing_subscriber::registry().with(capture.clone().with_filter(LevelFilter::WARN));
let subscriber = tracing::subscriber::set_default(subscriber);
crate::test_support::refresh_tracing_interest();
(
capture,
TracingCapture {
_exclusive: exclusive,
_subscriber: subscriber,
},
)
}
fn upstream_test_config() -> OrchestratorConfig {
OrchestratorConfig {
apply_command: Some("echo apply {change_id}".to_string()),
archive_command: Some("echo archive {change_id}".to_string()),
resolve_command: Some("echo resolve".to_string()),
..Default::default()
}
}
#[test]
fn upstream_integration_is_absent_by_default() {
let service = ParallelRunService::new(PathBuf::from("/tmp"), upstream_test_config());
assert!(service.upstream_integration().is_none());
let mut executor = crate::parallel::ParallelExecutor::new(
PathBuf::from("/tmp"),
upstream_test_config(),
None,
);
service.install_upstream_integration(&mut executor);
assert!(!executor.has_upstream_integration());
}
#[test]
fn upstream_integration_propagates_to_service_and_executor() {
let runtime = crate::upstream::UpstreamRuntime {
config: crate::upstream::UpstreamIntegrationConfig::new("upstream", "cargo test"),
branch: "develop".to_string(),
};
let mut service = ParallelRunService::new(PathBuf::from("/tmp"), upstream_test_config());
service.set_upstream_integration(runtime.clone());
let stored = service.upstream_integration().expect("runtime installed");
assert_eq!(stored.config.remote, "upstream");
assert_eq!(stored.config.verify_command, "cargo test");
assert_eq!(stored.branch, "develop");
let mut executor = crate::parallel::ParallelExecutor::new(
PathBuf::from("/tmp"),
upstream_test_config(),
None,
);
service.install_upstream_integration(&mut executor);
assert!(executor.has_upstream_integration());
}
#[test]
fn post_archive_action_propagates_to_service() {
let temp_dir = TempDir::new().unwrap();
let mut service =
ParallelRunService::new(temp_dir.path().to_path_buf(), create_test_config());
service.set_post_archive_action(PostArchiveAction::PushToRemote {
remote: "upstream".to_string(),
});
assert_eq!(
service.post_archive_action(),
&PostArchiveAction::PushToRemote {
remote: "upstream".to_string()
}
);
}
fn headless_executor(service: &ParallelRunService) -> crate::parallel::ParallelExecutor {
service.create_executor_with_queue_state(None, None, None, None, None, None)
}
#[test]
fn create_executor_with_queue_state_carries_configured_post_archive_action() {
let temp_dir = TempDir::new().unwrap();
let mut service =
ParallelRunService::new(temp_dir.path().to_path_buf(), create_test_config());
service.set_post_archive_action(PostArchiveAction::PushToRemote {
remote: "origin".to_string(),
});
let executor = headless_executor(&service);
assert_eq!(
executor.post_archive_action_for_test(),
&PostArchiveAction::PushToRemote {
remote: "origin".to_string()
},
"a headless run's configured push action must reach the executor"
);
}
#[test]
fn create_executor_with_queue_state_defaults_to_merge_to_base() {
let temp_dir = TempDir::new().unwrap();
let service = ParallelRunService::new(temp_dir.path().to_path_buf(), create_test_config());
assert_eq!(
headless_executor(&service).post_archive_action_for_test(),
&PostArchiveAction::MergeToBase
);
}
#[test]
fn create_executor_with_queue_state_binds_run_owner_graceful_stop_flag() {
let temp_dir = TempDir::new().unwrap();
let mut service =
ParallelRunService::new(temp_dir.path().to_path_buf(), create_test_config());
let graceful_stop = Arc::new(std::sync::atomic::AtomicBool::new(false));
service.set_graceful_stop_flag(graceful_stop.clone());
let executor = headless_executor(&service);
let bound = executor
.graceful_stop_flag_for_test()
.expect("a bound graceful-stop request must reach the executor");
assert!(
Arc::ptr_eq(bound, &graceful_stop),
"the executor must observe the owner's own flag, not a copy"
);
}
#[test]
fn create_executor_with_queue_state_binds_no_graceful_stop_flag_by_default() {
let temp_dir = TempDir::new().unwrap();
let service = ParallelRunService::new(temp_dir.path().to_path_buf(), create_test_config());
assert!(headless_executor(&service)
.graceful_stop_flag_for_test()
.is_none());
}
#[test]
fn create_executor_with_queue_state_keeps_finite_lifetime_without_dynamic_queue() {
use crate::parallel::SchedulerLifetime;
let temp_dir = TempDir::new().unwrap();
let service = ParallelRunService::new(temp_dir.path().to_path_buf(), create_test_config());
assert_eq!(
headless_executor(&service).scheduler_lifetime_for_test(),
SchedulerLifetime::Finite,
"a CLI run supplies no dynamic queue and must stay finite"
);
let loop_based = service.create_executor_with_queue_state(
None,
None,
None,
Some(Arc::new(crate::tui::queue::DynamicQueue::new())),
None,
None,
);
assert_eq!(
loop_based.scheduler_lifetime_for_test(),
SchedulerLifetime::Persistent,
"a loop-based frontend still gets the persistent lifetime"
);
}
async fn init_git_repo(temp_dir: &TempDir) -> bool {
let init_result = Command::new("git")
.args(["init"])
.current_dir(temp_dir.path())
.output()
.await;
let init_ok = init_result
.as_ref()
.map(|output| output.status.success())
.unwrap_or(false);
if !init_ok {
return false;
}
let _ = Command::new("git")
.args(["config", "user.email", "test@example.com"])
.current_dir(temp_dir.path())
.output()
.await;
let _ = Command::new("git")
.args(["config", "user.name", "Test User"])
.current_dir(temp_dir.path())
.output()
.await;
true
}
#[tokio::test]
async fn test_filter_committed_changes_skips_uncommitted() {
let temp_dir = TempDir::new().expect("tempdir");
if !init_git_repo(&temp_dir).await {
return;
}
let base_dir = temp_dir.path().join("openspec/changes");
std::fs::create_dir_all(base_dir.join("change-a")).unwrap();
std::fs::write(base_dir.join("change-a/proposal.md"), "test").unwrap();
let _ = Command::new("git")
.args(["add", "."])
.current_dir(temp_dir.path())
.output()
.await;
let _ = Command::new("git")
.args(["commit", "-m", "add change-a"])
.current_dir(temp_dir.path())
.output()
.await;
std::fs::create_dir_all(base_dir.join("change-b")).unwrap();
std::fs::write(base_dir.join("change-b/proposal.md"), "test").unwrap();
let service = ParallelRunService::new(temp_dir.path().to_path_buf(), create_test_config());
let changes = vec![
create_test_change("change-a", vec![]),
create_test_change("change-b", vec![]),
];
let (committed, skipped) = service
.filter_committed_changes(changes)
.await
.expect("filter changes");
let committed_ids: Vec<String> = committed.into_iter().map(|change| change.id).collect();
assert_eq!(committed_ids, vec!["change-a".to_string()]);
assert_eq!(skipped, vec!["change-b".to_string()]);
}
#[tokio::test]
async fn test_filter_committed_changes_skips_partially_uncommitted() {
let temp_dir = TempDir::new().expect("tempdir");
if !init_git_repo(&temp_dir).await {
return;
}
let base_dir = temp_dir.path().join("openspec/changes");
std::fs::create_dir_all(base_dir.join("change-a")).unwrap();
std::fs::write(base_dir.join("change-a/proposal.md"), "test").unwrap();
std::fs::create_dir_all(base_dir.join("change-b")).unwrap();
std::fs::write(base_dir.join("change-b/proposal.md"), "test").unwrap();
let _ = Command::new("git")
.args(["add", "."])
.current_dir(temp_dir.path())
.output()
.await;
let _ = Command::new("git")
.args(["commit", "-m", "add changes"])
.current_dir(temp_dir.path())
.output()
.await;
std::fs::write(base_dir.join("change-a/tasks.md"), "new task").unwrap();
let service = ParallelRunService::new(temp_dir.path().to_path_buf(), create_test_config());
let changes = vec![
create_test_change("change-a", vec![]),
create_test_change("change-b", vec![]),
];
let (committed, skipped) = service
.filter_committed_changes(changes)
.await
.expect("filter changes");
let committed_ids: Vec<String> = committed.into_iter().map(|change| change.id).collect();
assert_eq!(committed_ids, vec!["change-b".to_string()]);
assert_eq!(skipped, vec!["change-a".to_string()]);
}
#[tokio::test]
async fn test_prepare_parallel_execution_emits_rejection_event() {
let temp_dir = TempDir::new().expect("tempdir");
if !init_git_repo(&temp_dir).await {
return;
}
let base_dir = temp_dir.path().join("openspec/changes");
std::fs::create_dir_all(base_dir.join("change-a")).unwrap();
std::fs::write(base_dir.join("change-a/proposal.md"), "test").unwrap();
let _ = Command::new("git")
.args(["add", "."])
.current_dir(temp_dir.path())
.output()
.await;
let _ = Command::new("git")
.args(["commit", "-m", "add change-a"])
.current_dir(temp_dir.path())
.output()
.await;
std::fs::create_dir_all(base_dir.join("change-b")).unwrap();
std::fs::write(base_dir.join("change-b/proposal.md"), "test").unwrap();
let service = ParallelRunService::new(temp_dir.path().to_path_buf(), create_test_config());
let changes = vec![
create_test_change("change-a", vec![]),
create_test_change("change-b", vec![]),
];
let (event_tx, mut event_rx) = tokio::sync::mpsc::channel::<ParallelEvent>(32);
let result = service
.prepare_parallel_execution(changes, &event_tx, false)
.await
.expect("prepare_parallel_execution");
assert!(result.is_some(), "change-a should still be eligible");
let committed = result.unwrap();
assert_eq!(committed.len(), 1);
assert_eq!(committed[0].id, "change-a");
drop(event_tx);
let mut got_rejection_event = false;
while let Some(event) = event_rx.recv().await {
if let ParallelEvent::ParallelStartRejected { change_ids, .. } = event {
assert!(
change_ids.contains(&"change-b".to_string()),
"rejection event should include change-b"
);
got_rejection_event = true;
}
}
assert!(
got_rejection_event,
"expected a ParallelStartRejected event for the uncommitted change"
);
}
#[tokio::test]
async fn test_prepare_parallel_execution_all_rejected_emits_rejection_event() {
let temp_dir = TempDir::new().expect("tempdir");
if !init_git_repo(&temp_dir).await {
return;
}
let base_dir = temp_dir.path().join("openspec/changes");
let placeholder = base_dir.join("placeholder");
std::fs::create_dir_all(&placeholder).unwrap();
std::fs::write(placeholder.join("proposal.md"), "placeholder").unwrap();
let _ = Command::new("git")
.args(["add", "."])
.current_dir(temp_dir.path())
.output()
.await;
let _ = Command::new("git")
.args(["commit", "-m", "initial commit"])
.current_dir(temp_dir.path())
.output()
.await;
for id in &["change-a", "change-b"] {
std::fs::create_dir_all(base_dir.join(id)).unwrap();
std::fs::write(base_dir.join(id).join("proposal.md"), "test").unwrap();
}
let service = ParallelRunService::new(temp_dir.path().to_path_buf(), create_test_config());
let changes = vec![
create_test_change("change-a", vec![]),
create_test_change("change-b", vec![]),
];
let (event_tx, mut event_rx) = tokio::sync::mpsc::channel::<ParallelEvent>(32);
let result = service
.prepare_parallel_execution(changes, &event_tx, false)
.await
.expect("prepare_parallel_execution");
assert!(
result.is_none(),
"all changes were uncommitted so result should be None"
);
drop(event_tx);
let mut got_rejection_event = false;
let mut rejected_ids: Vec<String> = Vec::new();
while let Some(event) = event_rx.recv().await {
if let ParallelEvent::ParallelStartRejected { change_ids, .. } = event {
rejected_ids = change_ids;
got_rejection_event = true;
}
}
assert!(
got_rejection_event,
"expected a ParallelStartRejected event even when all changes are rejected"
);
rejected_ids.sort();
assert_eq!(rejected_ids, vec!["change-a", "change-b"]);
}
#[tokio::test]
async fn test_prepare_parallel_execution_allows_empty_when_resolve_wait_requested() {
let temp_dir = TempDir::new().expect("tempdir");
if !init_git_repo(&temp_dir).await {
return;
}
let service = ParallelRunService::new(temp_dir.path().to_path_buf(), create_test_config());
let changes = Vec::new();
let (event_tx, _event_rx) = tokio::sync::mpsc::channel::<ParallelEvent>(32);
let result = service
.prepare_parallel_execution(changes, &event_tx, true)
.await
.expect("prepare_parallel_execution");
assert!(
result.is_some(),
"empty startup should continue when reducer-owned ResolveWait exists"
);
assert!(
result.expect("checked is_some").is_empty(),
"no committed changes should still produce an empty queue"
);
}
#[tokio::test]
async fn test_prepare_parallel_execution_empty_parallel_without_resolve_wait_is_noop() {
let temp_dir = TempDir::new().expect("tempdir");
if !init_git_repo(&temp_dir).await {
return;
}
let service = ParallelRunService::new(temp_dir.path().to_path_buf(), create_test_config());
let changes = Vec::new();
let (event_tx, _event_rx) = tokio::sync::mpsc::channel::<ParallelEvent>(32);
let result = service
.prepare_parallel_execution(changes, &event_tx, false)
.await
.expect("prepare_parallel_execution");
assert!(
result.is_none(),
"empty startup without reducer-owned ResolveWait must remain a safe no-op"
);
}
#[tokio::test]
async fn test_analyze_order_fallback_preserves_metadata_dependencies_when_llm_disabled() {
let temp_dir = TempDir::new().expect("tempdir");
let mut config = create_test_config();
config.use_llm_analysis = Some(false);
let service = ParallelRunService::new(temp_dir.path().to_path_buf(), config);
let changes = vec![
create_test_change("route", vec!["policy"]),
create_test_change("policy", vec![]),
];
let outcome = service
.analyze_order_with_sender(&changes, &[], None, 1)
.await;
assert_eq!(
outcome.result.order,
vec!["route".to_string(), "policy".to_string()]
);
assert_eq!(
outcome.result.dependencies.get("route"),
Some(&vec!["policy".to_string()])
);
assert_eq!(
outcome.provenance,
crate::analyzer::AnalysisProvenance::IntentionalMetadataOnly,
"configured metadata-only analysis is the intended result, not a failure fallback"
);
assert!(
!outcome.provenance.is_degraded(),
"intentional metadata-only analysis must not be suppressed on a degraded interval"
);
}
#[cfg(unix)]
#[tokio::test]
async fn judge_outcomes_never_change_analysis_provenance() {
use std::os::unix::fs::PermissionsExt;
const MODEL: &str = "test-judge-1.0.0";
let temp_dir = TempDir::new().expect("tempdir");
let repo = temp_dir.path();
for change_id in ["route", "policy"] {
let dir = repo.join("openspec/changes").join(change_id);
std::fs::create_dir_all(&dir).expect("create change dir");
std::fs::write(dir.join("proposal.md"), "proposal body").expect("write proposal");
}
let judge_path = repo.join("fake-judge");
std::fs::write(
&judge_path,
format!(
"#!/bin/sh\ncat > /dev/null\nprintf '%s' '{{\"model\":\"{MODEL}\",\"answers\":{{\"q_0000_0001\":{{\"type\":\"noul\",\"noul\":0.99}},\"q_0001_0000\":{{\"type\":\"noul\",\"noul\":0.99}}}},\"usage\":{{\"input_tokens\":1,\"output_tokens\":1}}}}'\n: > '{}'\n",
repo.join("judge-done").display()
),
)
.expect("write fake judge");
std::fs::set_permissions(&judge_path, std::fs::Permissions::from_mode(0o755))
.expect("chmod fake judge");
let judge_entry = |command: Vec<String>| crate::config::JudgeCommandsConfig {
parallel_dependency: Some(crate::config::ParallelDependencyJudgeConfig {
command,
model: MODEL.to_string(),
timeout_ms: Some(300_000),
max_input_bytes: None,
max_output_bytes: None,
yes_threshold: None,
mode: None,
evaluation: None,
}),
};
let cases = vec![
(
"valid judge",
judge_entry(vec![judge_path.display().to_string()]),
format!(
"while [ ! -f '{}' ]; do sleep 0.01; done; printf '%s' '{{\"order\":[\"route\",\"policy\"],\"dependencies\":{{}}}}'",
repo.join("judge-done").display()
),
),
(
"unavailable judge",
judge_entry(vec![repo.join("not-installed").display().to_string()]),
r#"printf '%s' '{"order":["route","policy"],"dependencies":{}}'"#.to_string(),
),
];
for (label, judge_commands, analyze_command) in cases {
let config = OrchestratorConfig {
use_llm_analysis: Some(true),
analyze_command: Some(analyze_command),
judge_commands: Some(judge_commands),
command_queue_stagger_delay_ms: Some(0),
command_queue_max_retries: Some(0),
..create_test_config()
};
let service = ParallelRunService::new(repo.to_path_buf(), config);
let changes = vec![
create_test_change("route", vec![]),
create_test_change("policy", vec![]),
];
let outcome = tokio::time::timeout(
std::time::Duration::from_secs(60),
service.analyze_order_with_sender(&changes, &[], None, 1),
)
.await
.expect("analysis must not wait for the judge's own timeout");
assert_eq!(
outcome.provenance,
crate::analyzer::AnalysisProvenance::HealthyLlm,
"{label} must leave the conventional analyzer's provenance untouched"
);
assert!(
!outcome.provenance.is_degraded(),
"{label} must not degrade the analysis"
);
assert_eq!(
outcome.result.order,
vec!["route".to_string(), "policy".to_string()],
"{label} must not reorder the authoritative result"
);
assert!(
outcome.result.dependencies.is_empty(),
"{label}: a judge `yes` never adds an authoritative dependency"
);
}
}
#[tokio::test]
async fn test_analyze_order_recoverable_fallback_preserves_metadata_dependencies() {
let temp_dir = TempDir::new().expect("tempdir");
let mut config = create_test_config();
config.use_llm_analysis = Some(true);
config.analyze_command = Some("printf 'not json'".to_string());
let service = ParallelRunService::new(temp_dir.path().to_path_buf(), config);
let changes = vec![
create_test_change("route", vec!["policy"]),
create_test_change("policy", vec![]),
];
let outcome = service
.analyze_order_with_sender(&changes, &[], None, 1)
.await;
assert_eq!(
outcome.result.order,
vec!["route".to_string(), "policy".to_string()]
);
assert_eq!(
outcome.result.dependencies.get("route"),
Some(&vec!["policy".to_string()])
);
assert!(
!outcome.result.dependencies.is_empty(),
"recoverable fallback must not degrade to dependency-free analysis"
);
assert_eq!(
outcome.provenance,
crate::analyzer::AnalysisProvenance::RecoverableFailureFallback,
"a failed LLM command must be distinguishable from a configured metadata-only result"
);
assert!(
outcome.provenance.is_degraded(),
"recoverable-failure fallback must use bounded suppression so one retry stays possible"
);
}
#[tokio::test]
async fn test_recoverable_analysis_fallback_emits_warning_without_terminal_error() {
let temp_dir = TempDir::new().expect("tempdir");
let mut config = create_test_config();
config.use_llm_analysis = Some(true);
config.analyze_command =
Some("echo '{\"order\":[\"route\",\"policy\"],\"dependencies\":{}}'".to_string());
let service = ParallelRunService::new(temp_dir.path().to_path_buf(), config);
let changes = vec![
create_test_change("route", vec!["policy"]),
create_test_change("policy", vec![]),
create_test_change("gateway", vec![]),
];
let (event_tx, mut event_rx) = tokio::sync::mpsc::channel::<ParallelEvent>(32);
let outcome = service
.analyze_order_with_sender(&changes, &[], Some(&event_tx), 1)
.await;
assert_eq!(
outcome.result.order,
vec![
"route".to_string(),
"policy".to_string(),
"gateway".to_string()
],
"fallback must represent every queued change exactly once"
);
assert_eq!(
outcome.result.dependencies.get("route"),
Some(&vec!["policy".to_string()]),
"declared metadata dependencies must survive fallback"
);
assert_eq!(
outcome.provenance,
crate::analyzer::AnalysisProvenance::RecoverableFailureFallback,
"a rejected LLM response is a recoverable failure, not healthy analyzer output"
);
drop(event_tx);
let mut warnings = Vec::new();
let mut errors = Vec::new();
while let Some(event) = event_rx.recv().await {
match event {
ParallelEvent::Log(entry) if entry.level == crate::events::LogLevel::Warn => {
warnings.push(entry.message)
}
ParallelEvent::Error { message } => errors.push(message),
_ => {}
}
}
assert!(
errors.is_empty(),
"successful metadata fallback must not emit a terminal error event: {errors:?}"
);
assert_eq!(
warnings.len(),
1,
"successful metadata fallback should emit exactly one warning: {warnings:?}"
);
assert!(
warnings[0].contains("metadata-dependency-only"),
"warning must name the fallback mode: {}",
warnings[0]
);
assert!(
warnings[0].contains("Missing change IDs in response"),
"warning must preserve the original analysis rejection reason: {}",
warnings[0]
);
}
#[test]
fn test_recoverable_analysis_fallback_diagnostic_message_names_fallback_mode() {
let changes = vec![
create_test_change("route", vec!["policy"]),
create_test_change("policy", vec![]),
];
let (key, message) = ParallelRunService::recoverable_analysis_fallback_diagnostic(
&changes,
&["beta".to_string(), "alpha".to_string()],
" Missing change IDs in response: [\"gateway\"] ",
);
assert!(
matches!(
key,
DiagnosticDeduplicationKey::AnalysisFailure {
ref queued_ids,
ref in_flight_ids,
ref error,
} if queued_ids == &["policy".to_string(), "route".to_string()]
&& in_flight_ids == &["alpha".to_string(), "beta".to_string()]
&& error == "Missing change IDs in response: [\"gateway\"]"
),
"dedup identity must stay stable and order-independent: {key:?}"
);
assert!(
message.contains("metadata-dependency-only"),
"message must name the fallback mode: {message}"
);
assert!(
message.contains("Missing change IDs in response"),
"message must preserve the original reason: {message}"
);
assert!(
!message.contains("Dependency analysis failed"),
"recoverable fallback must not be phrased as a terminal failure: {message}"
);
}
#[tokio::test]
async fn test_recoverable_analysis_fallback_diagnostic_dedupes_by_signature() {
let temp_dir = TempDir::new().expect("tempdir");
let service = ParallelRunService::new(temp_dir.path().to_path_buf(), create_test_config());
let changes = vec![create_test_change("route", vec!["policy"])];
let other_changes = vec![
create_test_change("route", vec!["policy"]),
create_test_change("gateway", vec![]),
];
let (event_tx, mut event_rx) = tokio::sync::mpsc::channel::<ParallelEvent>(32);
let (capture, _tracing_guard) = capture_tracing().await;
for _ in 0..2 {
service
.emit_recoverable_analysis_fallback_diagnostic_once(
&changes,
&[],
Some(&event_tx),
"Missing change IDs in response: [\"gateway\"]",
)
.await;
}
service
.emit_recoverable_analysis_fallback_diagnostic_once(
&changes,
&[],
Some(&event_tx),
"Duplicate change ID in order: route",
)
.await;
service
.emit_recoverable_analysis_fallback_diagnostic_once(
&other_changes,
&["alpha".to_string()],
Some(&event_tx),
"Missing change IDs in response: [\"gateway\"]",
)
.await;
drop(event_tx);
let mut warnings = Vec::new();
let mut errors = Vec::new();
while let Some(event) = event_rx.recv().await {
match event {
ParallelEvent::Log(entry) if entry.level == crate::events::LogLevel::Warn => {
warnings.push(entry.message)
}
ParallelEvent::Error { message } => errors.push(message),
_ => {}
}
}
assert!(
errors.is_empty(),
"fallback diagnostics must never emit terminal error events: {errors:?}"
);
assert_eq!(
warnings.len(),
3,
"equivalent diagnostics collapse to one while distinct contexts stay visible: {warnings:?}"
);
assert_eq!(
warnings
.iter()
.filter(|message| message.contains("Duplicate change ID in order: route"))
.count(),
1,
"a different error must remain visible: {warnings:?}"
);
assert_eq!(
warnings
.iter()
.filter(|message| message.contains("\"gateway\", \"route\""))
.count(),
1,
"a different queued set must remain visible: {warnings:?}"
);
let records = capture.records();
assert!(
records
.iter()
.all(|(level, _)| *level != tracing::Level::ERROR),
"recoverable fallback must not emit ERROR-level tracing records: {records:?}"
);
assert_eq!(
records
.iter()
.filter(|(level, _)| *level == tracing::Level::WARN)
.count(),
3,
"tracing records must collapse equivalent signatures exactly like runtime events: {records:?}"
);
assert_eq!(
capture.warnings_containing("Missing change IDs in response"),
2,
"two repeats of one signature plus one changed queued context: {records:?}"
);
assert_eq!(
capture.warnings_containing("Duplicate change ID in order: route"),
1,
"a different rejection reason must emit its own tracing record: {records:?}"
);
}
#[tokio::test]
async fn test_recoverable_analysis_fallback_dedupes_tracing_without_event_sender() {
let temp_dir = TempDir::new().expect("tempdir");
let service = ParallelRunService::new(temp_dir.path().to_path_buf(), create_test_config());
let changes = vec![create_test_change("route", vec!["policy"])];
let (capture, _tracing_guard) = capture_tracing().await;
for _ in 0..3 {
service
.emit_recoverable_analysis_fallback_diagnostic_once(
&changes,
&[],
None,
"Missing change IDs in response: [\"gateway\"]",
)
.await;
}
let records = capture.records();
assert_eq!(
records
.iter()
.filter(|(level, _)| *level == tracing::Level::WARN)
.count(),
1,
"tracing-only callers must still get exactly one record per signature: {records:?}"
);
}
#[tokio::test]
async fn test_recoverable_fallback_log_uses_warn_level_only() {
use tracing::Level;
let (capture, guard) = capture_tracing().await;
ParallelRunService::log_recoverable_analysis_fallback(&"invalid dependency graph");
drop(guard);
let events = capture.records();
assert_eq!(events.len(), 1);
assert_eq!(events[0].0, Level::WARN);
assert!(
events[0]
.1
.contains("falling back to metadata-dependency-only analysis"),
"fallback diagnostic should remain operator-visible"
);
assert!(
events[0].1.contains("invalid dependency graph"),
"original LLM analysis failure should remain visible as warning context"
);
assert!(
events.iter().all(|(level, _)| *level != Level::ERROR),
"recoverable fallback must not emit ERROR-level records"
);
}
#[tokio::test]
async fn test_order_based_empty_resolve_wait_shared_state_enters_scheduler_path() {
use crate::orchestration::state::{
OrchestratorState, ReducerCommand, WorkspaceObservation,
};
use crate::parallel::ParallelExecutor;
use std::sync::Arc;
use tokio::sync::RwLock;
let temp_dir = TempDir::new().expect("tempdir");
if !init_git_repo(&temp_dir).await {
return;
}
let _ = Command::new("git")
.args(["commit", "--allow-empty", "-m", "initial commit"])
.current_dir(temp_dir.path())
.output()
.await;
let service = ParallelRunService::new(temp_dir.path().to_path_buf(), create_test_config());
let (event_tx, mut event_rx) = tokio::sync::mpsc::channel::<ParallelEvent>(32);
let shared = Arc::new(RwLock::new(OrchestratorState::new(
vec!["alpha".to_string()],
3,
)));
{
let mut state = shared.write().await;
state.apply_observation("alpha", WorkspaceObservation::WorkspaceArchived);
state.apply_command(ReducerCommand::ResolveMerge("alpha".to_string()));
}
let mut executor = ParallelExecutor::new(
temp_dir.path().to_path_buf(),
create_test_config(),
Some(event_tx.clone()),
);
executor.set_shared_orchestrator_state(shared.clone());
service
.run_parallel_order_based_with_executor(executor, Vec::new(), event_tx.clone())
.await
.expect("empty ResolveWait scheduler path should run");
drop(event_tx);
let mut rejected_empty_start = false;
while let Some(event) = event_rx.recv().await {
if matches!(event, ParallelEvent::ParallelStartRejected { .. }) {
rejected_empty_start = true;
}
}
assert!(
!rejected_empty_start,
"empty ResolveWait startup must not be treated as a zero-change start rejection"
);
}
#[tokio::test]
async fn test_run_parallel_all_rejected_forwards_event_to_callback() {
let temp_dir = TempDir::new().expect("tempdir");
if !init_git_repo(&temp_dir).await {
return;
}
let base_dir = temp_dir.path().join("openspec/changes");
let placeholder = base_dir.join("placeholder");
std::fs::create_dir_all(&placeholder).unwrap();
std::fs::write(placeholder.join("proposal.md"), "placeholder").unwrap();
let _ = Command::new("git")
.args(["add", "."])
.current_dir(temp_dir.path())
.output()
.await;
let _ = Command::new("git")
.args(["commit", "-m", "initial commit"])
.current_dir(temp_dir.path())
.output()
.await;
for id in &["change-a", "change-b"] {
std::fs::create_dir_all(base_dir.join(id)).unwrap();
std::fs::write(base_dir.join(id).join("proposal.md"), "test").unwrap();
}
let service = ParallelRunService::new(temp_dir.path().to_path_buf(), create_test_config());
let changes = vec![
create_test_change("change-a", vec![]),
create_test_change("change-b", vec![]),
];
let collected_events: std::sync::Arc<std::sync::Mutex<Vec<ParallelEvent>>> =
std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
let collected_events_clone = collected_events.clone();
service
.run_parallel(changes, None, move |event| {
collected_events_clone.lock().unwrap().push(event);
})
.await
.expect("run_parallel should succeed even when all changes are rejected");
let events = collected_events.lock().unwrap();
let got_rejection = events
.iter()
.any(|e| matches!(e, ParallelEvent::ParallelStartRejected { .. }));
assert!(
got_rejection,
"ParallelStartRejected must be forwarded to the callback when all changes are rejected at start time"
);
}
fn per_change_upstream_git(cwd: &std::path::Path, args: &[&str]) -> Option<String> {
let output = std::process::Command::new("git")
.args(args)
.current_dir(cwd)
.output()
.ok()?;
output
.status
.success()
.then(|| String::from_utf8_lossy(&output.stdout).trim().to_string())
}
fn per_change_upstream_unpublished_repo() -> Option<(TempDir, PathBuf)> {
let dir = TempDir::new().ok()?;
let root = dir.path().join("repo");
let remote = dir.path().join("remote.git");
std::fs::create_dir_all(&root).ok()?;
per_change_upstream_git(dir.path(), &["init", "--bare", "-b", "main", "remote.git"])?;
per_change_upstream_git(&root, &["init", "-b", "main"])?;
per_change_upstream_git(&root, &["config", "user.email", "test@example.com"])?;
per_change_upstream_git(&root, &["config", "user.name", "Test User"])?;
per_change_upstream_git(&root, &["config", "commit.gpgsign", "false"])?;
std::fs::write(root.join("README.md"), "# base\n").ok()?;
per_change_upstream_git(&root, &["add", "-A"])?;
per_change_upstream_git(&root, &["commit", "-m", "Initial commit"])?;
per_change_upstream_git(&root, &["remote", "add", "origin", remote.to_str()?])?;
per_change_upstream_git(&root, &["push", "-u", "origin", "main"])?;
let archive = root.join("openspec/changes/archive/alpha");
std::fs::create_dir_all(&archive).ok()?;
std::fs::write(archive.join("proposal.md"), "# archived alpha\n").ok()?;
per_change_upstream_git(&root, &["add", "-A"])?;
per_change_upstream_git(&root, &["commit", "-m", "Archive: alpha"])?;
let marker = crate::upstream::publication::format_publication_marker_message(
"alpha", "origin", "main",
);
per_change_upstream_git(&root, &["commit", "--allow-empty", "-m", &marker])?;
Some((dir, root))
}
async fn per_change_upstream_finite_run(
root: &std::path::Path,
verify_command: &str,
) -> Vec<ParallelEvent> {
let mut service = ParallelRunService::new(root.to_path_buf(), upstream_test_config());
service.set_upstream_integration(crate::upstream::UpstreamRuntime {
config: crate::upstream::UpstreamIntegrationConfig::new("origin", verify_command),
branch: "main".to_string(),
});
let (event_tx, mut event_rx) = mpsc::channel::<ParallelEvent>(256);
let collector = tokio::spawn(async move {
let mut events = Vec::new();
while let Some(event) = event_rx.recv().await {
events.push(event);
}
events
});
service
.run_parallel_with_channel_and_queue_state(
Vec::new(),
event_tx,
None,
None,
None,
None,
None,
true,
)
.await
.expect("finite opted-in run");
collector.await.expect("event collector")
}
#[tokio::test]
async fn per_change_upstream_finite_run_completes_only_after_confirmation() {
let _serialize = crate::parallel::merge_lock_test_mutex().lock().await;
let Some((_dir, root)) = per_change_upstream_unpublished_repo() else {
println!("Skipping test: git not available");
return;
};
let head = per_change_upstream_git(&root, &["rev-parse", "HEAD"]).expect("head");
let events = per_change_upstream_finite_run(&root, "exit 0").await;
let pushed = events
.iter()
.position(|event| {
matches!(event, ParallelEvent::PushCompleted { change_id, .. } if change_id == "alpha")
})
.expect("the targeted change must reach confirmed publication");
let completed = events
.iter()
.position(|event| matches!(event, ParallelEvent::AllCompleted))
.expect("a successful finite run reports completion");
assert!(
pushed < completed,
"AllCompleted must follow remote confirmation, never precede it"
);
assert_eq!(
per_change_upstream_git(&root, &["ls-remote", "origin", "refs/heads/main"])
.expect("ls-remote")
.split_whitespace()
.next()
.expect("remote head"),
head,
"completion is reported against a remotely observed cumulative HEAD"
);
}
#[tokio::test]
async fn per_change_upstream_finite_run_withholds_completion_when_publication_fails() {
let _serialize = crate::parallel::merge_lock_test_mutex().lock().await;
let Some((_dir, root)) = per_change_upstream_unpublished_repo() else {
println!("Skipping test: git not available");
return;
};
let remote_before =
per_change_upstream_git(&root, &["ls-remote", "origin", "refs/heads/main"])
.expect("ls-remote");
let events = per_change_upstream_finite_run(&root, "exit 1").await;
assert!(
!events
.iter()
.any(|event| matches!(event, ParallelEvent::AllCompleted)),
"a finite run with an unpublished change must not report completion"
);
assert!(
!events.iter().any(|event| matches!(
event,
ParallelEvent::PushCompleted { change_id, .. } if change_id == "alpha"
)),
"failed verification must suppress confirmed publication"
);
assert_eq!(
per_change_upstream_git(&root, &["ls-remote", "origin", "refs/heads/main"])
.expect("ls-remote"),
remote_before,
"nothing may reach the remote when verification fails"
);
}
}