use crate::error::Result;
use crate::events::LogEntry;
use std::collections::HashSet;
use std::sync::atomic::Ordering;
use tokio::task::JoinSet;
use tracing::{error, info, warn};
use super::cleanup::WorkspaceCleanupGuard;
use super::dynamic_queue::ReanalysisReason;
use super::events::send_event;
use super::queue_state::{
QueueReconciliationOutcome, ReanalysisDispatchContext, RetryEdgeConsumption,
};
use super::types::WorkspaceResult;
use super::work_snapshot::ReducerWorkSnapshot;
use super::ParallelEvent;
use super::ParallelExecutor;
use super::SchedulerLifetime;
use super::SchedulerRunReport;
use crate::upstream::coordinator::SchedulerOutcome;
pub(super) const CANCELLATION_MERGE_DRAIN_DEADLINE: std::time::Duration =
std::time::Duration::from_secs(90);
pub(super) const RUN_COMMAND_CLEANUP_DEADLINE: std::time::Duration =
crate::ai_command_runner::RUN_COMMAND_CLEANUP_DEADLINE;
pub(super) fn remaining_cleanup_budget(
started: std::time::Instant,
outer: std::time::Duration,
cap: std::time::Duration,
) -> std::time::Duration {
cap.min(outer.saturating_sub(started.elapsed()))
}
impl ParallelExecutor {
pub(super) fn is_fully_drained(
&self,
join_set_empty: bool,
queued_empty: bool,
in_flight_empty: bool,
) -> bool {
join_set_empty
&& queued_empty
&& in_flight_empty
&& self.resolve_wait_changes.is_empty()
&& self.reject_wait_changes.is_empty()
&& self.manual_resolve_active() == 0
&& self.pending_merge_count.load(Ordering::Relaxed) == 0
}
pub(super) async fn should_exit_when_idle(
&self,
join_set_empty: bool,
queued: &[crate::openspec::Change],
in_flight: &HashSet<String>,
work_snapshot: Option<&ReducerWorkSnapshot>,
) -> bool {
if self.scheduler_lifetime != SchedulerLifetime::Finite || !join_set_empty {
return false;
}
let captured;
let work_snapshot = match work_snapshot {
Some(snapshot) => snapshot,
None => {
captured = self.capture_reducer_work_snapshot().await;
&captured
}
};
if !work_snapshot.is_complete() {
return false;
}
self.is_fully_drained(join_set_empty, queued.is_empty(), in_flight.is_empty())
|| self
.is_blocked_only_scheduler_state_with_snapshot(queued, in_flight, work_snapshot)
.await
}
pub(super) async fn should_enter_persistent_idle_wait(
&self,
join_set_empty: bool,
queued: &[crate::openspec::Change],
in_flight: &HashSet<String>,
work_snapshot: Option<&ReducerWorkSnapshot>,
) -> bool {
if self.scheduler_lifetime != SchedulerLifetime::Persistent || !join_set_empty {
return false;
}
let captured;
let work_snapshot = match work_snapshot {
Some(snapshot) => snapshot,
None => {
captured = self.capture_reducer_work_snapshot().await;
&captured
}
};
if !work_snapshot.is_complete() {
return false;
}
self.is_fully_drained(join_set_empty, queued.is_empty(), in_flight.is_empty())
|| (queued.is_empty()
&& in_flight.is_empty()
&& (!self.resolve_wait_changes.is_empty() || !self.reject_wait_changes.is_empty())
&& self.manual_resolve_active() == 0
&& self.pending_merge_count.load(Ordering::Relaxed) == 0)
|| self
.is_blocked_only_scheduler_state_with_snapshot(queued, in_flight, work_snapshot)
.await
}
pub(super) async fn admit_persistent_idle_wait(
&self,
join_set_empty: bool,
queued: &[crate::openspec::Change],
in_flight: &HashSet<String>,
work_snapshot: Option<&ReducerWorkSnapshot>,
) -> bool {
if !self
.should_enter_persistent_idle_wait(join_set_empty, queued, in_flight, work_snapshot)
.await
{
return false;
}
let captured;
let work_snapshot = match work_snapshot {
Some(snapshot) => snapshot,
None => {
captured = self.capture_reducer_work_snapshot().await;
&captured
}
};
self.record_persistent_idle_baseline(work_snapshot);
if self.latch_persistent_idle() {
send_event(&self.event_tx, ParallelEvent::PersistentSchedulerIdle).await;
}
true
}
pub(super) fn graceful_stop_requested(&self) -> bool {
self.graceful_stop
.as_ref()
.is_some_and(|flag| flag.load(std::sync::atomic::Ordering::SeqCst))
}
pub(super) async fn settle_graceful_stop_when_no_work(
&self,
join_set_empty: bool,
queued: &[crate::openspec::Change],
in_flight: &HashSet<String>,
work_snapshot: &ReducerWorkSnapshot,
) -> bool {
if !self.graceful_stop_requested() || !work_snapshot.is_complete() {
return false;
}
if !self.is_fully_drained(join_set_empty, queued.is_empty(), in_flight.is_empty()) {
return false;
}
info!(
"Graceful stop requested with no work remaining; settling the scheduler at this boundary"
);
send_event(&self.event_tx, ParallelEvent::Stopped).await;
true
}
#[cfg(test)]
pub(crate) async fn settle_graceful_stop_at_idle_boundary(
&self,
queued: &[crate::openspec::Change],
) -> bool {
let work_snapshot = self.capture_reducer_work_snapshot().await;
self.settle_graceful_stop_when_no_work(true, queued, &HashSet::new(), &work_snapshot)
.await
}
pub(super) fn record_persistent_idle_baseline(&self, work_snapshot: &ReducerWorkSnapshot) {
let baseline: HashSet<String> = work_snapshot.queued_intent_ids().iter().cloned().collect();
*self
.persistent_idle_baseline
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner()) = Some(baseline);
}
pub(super) fn rearm_persistent_idle_from_observed_intent(
&self,
work_snapshot: &ReducerWorkSnapshot,
explicit_retry_edge: bool,
) {
if !work_snapshot.is_complete() {
return;
}
let observed_new_intent = {
let baseline = self
.persistent_idle_baseline
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
match baseline.as_ref() {
None => return,
Some(parked) => work_snapshot
.queued_intent_ids()
.iter()
.any(|id| !parked.contains(id)),
}
};
if observed_new_intent || explicit_retry_edge {
self.rearm_persistent_idle();
}
}
pub(super) fn latch_persistent_idle(&self) -> bool {
!self
.persistent_idle_latched
.swap(true, std::sync::atomic::Ordering::SeqCst)
}
pub(super) fn rearm_persistent_idle(&self) {
self.persistent_idle_latched
.store(false, std::sync::atomic::Ordering::SeqCst);
*self
.persistent_idle_baseline
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner()) = None;
}
#[cfg(test)]
pub(super) fn persistent_idle_is_latched(&self) -> bool {
self.persistent_idle_latched
.load(std::sync::atomic::Ordering::SeqCst)
}
pub(super) fn had_change_failures(&self) -> bool {
!self.change_failures_this_run.is_empty()
}
pub(super) fn derive_pass_reanalysis_reason(
current: ReanalysisReason,
retry_edges: RetryEdgeConsumption,
reconciliation: QueueReconciliationOutcome,
dynamic_queue_added: bool,
) -> ReanalysisReason {
if reconciliation.has_queued_additions() || retry_edges.bypass_armed {
ReanalysisReason::QueueNotification
} else if reconciliation.has_repair_additions() {
ReanalysisReason::RepairCandidate
} else if matches!(current, ReanalysisReason::QueueNotification) && !dynamic_queue_added {
ReanalysisReason::Initial
} else {
current
}
}
pub async fn execute_with_order_based_reanalysis<F>(
&mut self,
changes: Vec<crate::openspec::Change>,
analyzer: F,
) -> Result<SchedulerRunReport>
where
for<'a> F: Fn(
&'a [crate::openspec::Change],
&'a [String],
u32,
) -> std::pin::Pin<
Box<dyn std::future::Future<Output = crate::analyzer::AnalysisOutcome> + Send + 'a>,
> + Send
+ Sync,
{
let must_reach_upstream_boundary =
self.upstream_enabled() || self.explicit_target_plan.is_some();
if changes.is_empty() && !must_reach_upstream_boundary {
let startup_snapshot = self.capture_reducer_work_snapshot().await;
let reducer_has_queued_intent = !startup_snapshot.queued_intent_ids().is_empty();
let reducer_has_lane_wait = !startup_snapshot.resolve_wait_ids().is_empty()
|| !startup_snapshot.reject_wait_ids().is_empty();
if startup_snapshot.is_complete()
&& !reducer_has_queued_intent
&& !reducer_has_lane_wait
{
return Ok(self.finish_completed_run().await);
}
if reducer_has_lane_wait {
info!(
"Starting scheduler loop with reducer-visible base-lane wait retry intent and empty local queue"
);
} else {
info!(
"Starting scheduler loop with reducer-visible queued intent and empty local queue"
);
}
}
info!(
"Starting order-based execution with re-analysis for {} changes",
changes.len()
);
if let Some(token) = &self.cancel_token {
self.run_command_scope.link_cancellation(token.clone());
}
info!("Preparing for parallel execution...");
match self.workspace_manager.prepare_for_parallel().await {
Ok(Some(warning)) => {
warn!("{}", warning.message);
send_event(
&self.event_tx,
ParallelEvent::Warning {
title: warning.title,
message: warning.message,
},
)
.await;
}
Ok(None) => {}
Err(e) => {
let error_msg = format!("Failed to prepare for parallel execution: {}", e);
error!("{}", error_msg);
send_event(&self.event_tx, ParallelEvent::Error { message: error_msg }).await;
return Err(e.into());
}
}
info!("Preparation complete");
if self.upstream_enabled() {
if let Err(err) = self
.run_upstream_checkpoint(
crate::upstream::checkpoint::CheckpointTrigger::BeforeFirstDispatch,
None,
false,
)
.await
{
let error_msg = format!("Upstream pre-dispatch checkpoint failed: {}", err);
error!("{}", error_msg);
send_event(&self.event_tx, ParallelEvent::Error { message: error_msg }).await;
return Err(err);
}
let _ = self.resume_pending_publications().await;
}
let changes = match self.apply_explicit_target_plan(changes).await {
Ok(changes) => changes,
Err(err) => {
let error_msg = format!("Explicit target resolution failed: {}", err);
error!("{}", error_msg);
send_event(&self.event_tx, ParallelEvent::Error { message: error_msg }).await;
return Err(err);
}
};
let max_parallelism = self.workspace_manager.max_concurrent();
self.lifecycle_slots.ensure_capacity(max_parallelism);
let mut join_set: JoinSet<WorkspaceResult> = JoinSet::new();
let (merge_result_tx, mut merge_result_rx) = self.take_merge_result_channel();
let mut in_flight: HashSet<String> = HashSet::new();
let mut queued: Vec<crate::openspec::Change> = changes;
let mut iteration = 1u32;
let mut cleanup_guard = WorkspaceCleanupGuard::new(
self.workspace_manager.backend_type(),
self.repo_root.clone(),
);
let mut reanalysis_reason = ReanalysisReason::Initial;
let mut cancelled = false;
let mut graceful_stop_settled = false;
let mut blocked_exit = false;
let mut scheduler_outcome = SchedulerOutcome::BlockedOrStalled;
loop {
if self.run_fatal_abort.is_some() {
let remaining: Vec<String> = queued.iter().map(|c| c.id.clone()).collect();
error!(
queued = remaining.len(),
in_flight = in_flight.len(),
"Run-fatal base-lane outcome; stopping dispatch and draining owned work"
);
let shutdown_started = std::time::Instant::now();
self.run_command_scope.close();
join_set.abort_all();
while let Some(result) = join_set.join_next().await {
if let Err(err) = result {
if !err.is_cancelled() {
warn!(error = %err, "In-flight workspace task failed while draining after run-fatal abort");
}
}
}
self.await_run_command_quiescence(shutdown_started, "run-fatal abort")
.await;
self.release_execution_handles_after_cancellation().await;
self.clear_preparation_for_aborted_changes(&in_flight).await;
in_flight.clear();
queued.clear();
self.drain_pending_merge_results_after_cancellation(
&merge_result_tx,
&mut merge_result_rx,
remaining_cleanup_budget(
shutdown_started,
crate::tui::orchestrator::PARALLEL_CANCELLATION_CLEANUP_DEADLINE,
CANCELLATION_MERGE_DRAIN_DEADLINE,
),
)
.await;
self.lifecycle_slots.clear();
break;
}
if self.is_cancelled() {
let remaining_changes: Vec<String> = queued.iter().map(|c| c.id.clone()).collect();
let cancel_msg = format!(
"Cancelled parallel execution ({} queued, {} in-flight: queued=[{}], in-flight=[{}])",
remaining_changes.len(),
in_flight.len(),
remaining_changes.join(", "),
in_flight.iter().cloned().collect::<Vec<_>>().join(", ")
);
send_event(
&self.event_tx,
ParallelEvent::Log(LogEntry::warn(&cancel_msg)),
)
.await;
cancelled = true;
scheduler_outcome = SchedulerOutcome::Cancelled;
let shutdown_started = std::time::Instant::now();
self.run_command_scope.close();
join_set.abort_all();
while let Some(result) = join_set.join_next().await {
if let Err(err) = result {
if !err.is_cancelled() {
warn!(error = %err, "In-flight workspace task failed while draining after cancellation");
}
}
}
self.await_run_command_quiescence(shutdown_started, "operator cancellation")
.await;
self.release_execution_handles_after_cancellation().await;
self.clear_preparation_for_aborted_changes(&in_flight).await;
in_flight.clear();
self.drain_pending_merge_results_after_cancellation(
&merge_result_tx,
&mut merge_result_rx,
remaining_cleanup_budget(
shutdown_started,
crate::tui::orchestrator::PARALLEL_CANCELLATION_CLEANUP_DEADLINE,
CANCELLATION_MERGE_DRAIN_DEADLINE,
),
)
.await;
self.lifecycle_slots.clear();
break;
}
let retry_edges = self.consume_explicit_retry_edges().await;
let work_snapshot = self.capture_reducer_work_snapshot().await;
if self.is_cancelled() {
continue;
}
let dynamic_queue_added = self
.check_dynamic_queue_and_add_changes_with_snapshot(
&mut queued,
&in_flight,
&mut reanalysis_reason,
&work_snapshot,
)
.await;
self.sync_resolve_wait_from_snapshot(&work_snapshot);
self.reconcile_retained_lifecycle_slots(&work_snapshot);
self.maybe_dispatch_resolve_wait_retry_with_tx(&merge_result_tx)
.await;
let reconciliation = self
.reconcile_queued_candidates_with_snapshot(&mut queued, &in_flight, &work_snapshot)
.await;
self.rearm_persistent_idle_from_observed_intent(
&work_snapshot,
retry_edges.newly_drained > 0,
);
reanalysis_reason = Self::derive_pass_reanalysis_reason(
reanalysis_reason,
retry_edges,
reconciliation,
dynamic_queue_added,
);
let work_drained = work_snapshot.is_complete()
&& queued.is_empty()
&& in_flight.is_empty()
&& self.resolve_wait_changes.is_empty()
&& self.reject_wait_changes.is_empty()
&& self.manual_resolve_active() == 0
&& self.pending_merge_count.load(Ordering::Relaxed) == 0;
if work_drained && self.scheduler_lifetime == SchedulerLifetime::Finite {
info!(
"All changes completed (queued/in-flight/resolve_wait/manual_resolve empty), stopping"
);
scheduler_outcome = SchedulerOutcome::DrainedSuccessfully;
break;
}
if let Some((should_break, new_iteration)) = self
.evaluate_queued_reanalysis_and_dispatch(
ReanalysisDispatchContext {
queued: &mut queued,
in_flight: &mut in_flight,
max_parallelism,
iteration,
reanalysis_reason,
analyzer: &analyzer,
join_set: &mut join_set,
cleanup_guard: &mut cleanup_guard,
work_snapshot: Some(&work_snapshot),
},
&mut reanalysis_reason,
)
.await?
{
iteration = new_iteration;
if should_break {
break;
}
}
if self
.should_exit_when_idle(
join_set.is_empty(),
&queued,
&in_flight,
Some(&work_snapshot),
)
.await
{
info!(
"All automatic scheduler work completed or blocked-only, exiting scheduler loop"
);
scheduler_outcome = if self.is_fully_drained(
join_set.is_empty(),
queued.is_empty(),
in_flight.is_empty(),
) {
SchedulerOutcome::DrainedSuccessfully
} else {
blocked_exit = true;
SchedulerOutcome::BlockedOrStalled
};
break;
}
if self
.settle_graceful_stop_when_no_work(
join_set.is_empty(),
&queued,
&in_flight,
&work_snapshot,
)
.await
{
graceful_stop_settled = true;
break;
}
if self
.admit_persistent_idle_wait(
join_set.is_empty(),
&queued,
&in_flight,
Some(&work_snapshot),
)
.await
{
self.wait_for_persistent_idle_wake_with_tx(
&mut reanalysis_reason,
&merge_result_tx,
&mut merge_result_rx,
)
.await;
continue;
}
self.wait_for_scheduler_event(
&mut join_set,
&mut in_flight,
max_parallelism,
&merge_result_tx,
&mut merge_result_rx,
&mut reanalysis_reason,
)
.await;
}
drop(cleanup_guard);
if cancelled {
send_event(&self.event_tx, ParallelEvent::Stopped).await;
return Ok(SchedulerRunReport::Stopped);
}
if graceful_stop_settled {
return Ok(SchedulerRunReport::Stopped);
}
if let Some(detail) = self.run_fatal_abort.clone() {
return Err(crate::error::OrchestratorError::GitCommand(detail));
}
if self.upstream_enabled() {
let stranded = self.resume_pending_publications().await;
if !stranded.is_empty() {
send_event(
&self.event_tx,
ParallelEvent::Error {
message: format!(
"Upstream publication is still owed for {}; cumulative base was not published",
stranded.join(", ")
),
},
)
.await;
return Ok(self.terminal_report());
}
}
if !self.finalize_upstream(scheduler_outcome).await {
send_event(
&self.event_tx,
ParallelEvent::Error {
message: format!(
"Upstream integration did not complete ({:?}); cumulative base was not published",
scheduler_outcome
),
},
)
.await;
return Ok(self.terminal_report());
}
if self.upstream_enabled() {
let stranded: Vec<String> = self
.pending_publications()
.await
.into_iter()
.map(|evidence| evidence.trailers.change_id)
.collect();
if !stranded.is_empty() {
send_event(
&self.event_tx,
ParallelEvent::Error {
message: format!(
"Upstream publication is still owed for {}; the run is not complete",
stranded.join(", ")
),
},
)
.await;
return Ok(self.terminal_report());
}
}
if blocked_exit {
return Ok(self.finish_blocked_run(&queued).await);
}
Ok(self.finish_completed_run().await)
}
async fn finish_blocked_run(&self, queued: &[crate::openspec::Change]) -> SchedulerRunReport {
let mut blocked: Vec<String> = queued.iter().map(|change| change.id.clone()).collect();
blocked.sort();
let message = format!(
"Processing stopped with blocked work remaining; no dispatchable candidate is available: {}",
blocked.join(", ")
);
warn!("{}", message);
send_event(&self.event_tx, ParallelEvent::Log(LogEntry::warn(&message))).await;
SchedulerRunReport::BlockedOrStalled
}
async fn finish_completed_run(&self) -> SchedulerRunReport {
let report = self.terminal_report();
if report == SchedulerRunReport::CompletedWithErrors {
let mut failed: Vec<String> = self.change_failures_this_run.iter().cloned().collect();
failed.sort();
let message = format!(
"Processing completed with errors; unresolved change-local failures preserved for explicit retry: {}",
failed.join(", ")
);
warn!("{}", message);
send_event(&self.event_tx, ParallelEvent::Log(LogEntry::warn(&message))).await;
}
send_event(&self.event_tx, ParallelEvent::AllCompleted).await;
report
}
fn terminal_report(&self) -> SchedulerRunReport {
if self.had_change_failures() {
SchedulerRunReport::CompletedWithErrors
} else {
SchedulerRunReport::Completed
}
}
fn take_merge_result_channel(
&mut self,
) -> (
tokio::sync::mpsc::Sender<super::MergeResult>,
tokio::sync::mpsc::Receiver<super::MergeResult>,
) {
#[cfg(test)]
{
if let Some(channel) = self.merge_result_channel_override.take() {
return channel;
}
}
tokio::sync::mpsc::channel(64)
}
async fn release_execution_handles_after_cancellation(&self) {
let Some(queue) = self.dynamic_queue.as_ref() else {
return;
};
let scope = self.run_command_scope.clone();
let release = queue
.release_all_execution_handles(|change_id| scope.change_is_quiescent(change_id))
.await;
if release.confirmed > 0 {
info!(
confirmed = release.confirmed,
"Released registered execution handles whose run-owned commands reached confirmed cleanup"
);
}
for change_id in &release.unconfirmed {
warn!(
change_id = %change_id,
"Execution handle released without confirmed command cleanup; the completion \
handshake stays unfired and its waiter times out truthfully"
);
}
}
async fn await_run_command_quiescence(
&self,
shutdown_started: std::time::Instant,
reason: &str,
) {
#[allow(unused_mut)]
let mut cap = RUN_COMMAND_CLEANUP_DEADLINE;
#[cfg(test)]
if let Some(override_budget) = self.run_command_cleanup_budget_override {
cap = override_budget;
}
let budget = remaining_cleanup_budget(
shutdown_started,
crate::tui::orchestrator::PARALLEL_CANCELLATION_CLEANUP_DEADLINE,
cap,
);
let cleanup = self.run_command_scope.shutdown(budget).await;
if cleanup.is_quiescent() {
info!(
reason,
escalated = cleanup.escalated,
"Run-owned AI commands reached process quiescence"
);
return;
}
let message = format!(
"Run-owned command cleanup could not be fully proven while stopping ({}): {}",
reason,
cleanup.diagnostics()
);
warn!("{}", message);
send_event(&self.event_tx, ParallelEvent::Log(LogEntry::warn(&message))).await;
}
async fn clear_preparation_for_aborted_changes(&self, in_flight: &HashSet<String>) {
for change_id in in_flight {
send_event(
&self.event_tx,
ParallelEvent::WorkspacePreparationEnded {
change_id: change_id.clone(),
},
)
.await;
}
}
async fn drain_pending_merge_results_after_cancellation(
&mut self,
merge_result_tx: &tokio::sync::mpsc::Sender<super::MergeResult>,
merge_result_rx: &mut tokio::sync::mpsc::Receiver<super::MergeResult>,
deadline: std::time::Duration,
) {
let pending = self.pending_merge_count.load(Ordering::Relaxed);
if pending == 0 {
return;
}
let waiting_msg = format!(
"Waiting for {} pending background merge/base-lane task(s) to reach a safe boundary before stopping",
pending
);
info!("{}", waiting_msg);
send_event(
&self.event_tx,
ParallelEvent::Log(LogEntry::info(&waiting_msg)),
)
.await;
let drained = tokio::time::timeout(deadline, async {
while self.pending_merge_count.load(Ordering::Relaxed) > 0 {
let Some(merge_result) = merge_result_rx.recv().await else {
break;
};
self.handle_merge_result_with_tx(merge_result, merge_result_tx)
.await;
}
})
.await;
if drained.is_err() {
let timeout_msg = format!(
"Pending background merge/base-lane task(s) did not report within {}s while stopping; continuing shutdown",
deadline.as_secs()
);
warn!("{}", timeout_msg);
send_event(
&self.event_tx,
ParallelEvent::Log(LogEntry::warn(&timeout_msg)),
)
.await;
} else {
info!("Pending background merge/base-lane tasks reached a safe boundary; stopping");
}
}
pub(super) async fn evaluate_queued_reanalysis_and_dispatch<F>(
&mut self,
ctx: ReanalysisDispatchContext<'_, F>,
reanalysis_reason: &mut ReanalysisReason,
) -> Result<Option<(bool, u32)>>
where
for<'a> F: Fn(
&'a [crate::openspec::Change],
&'a [String],
u32,
) -> std::pin::Pin<
Box<dyn std::future::Future<Output = crate::analyzer::AnalysisOutcome> + Send + 'a>,
> + Send
+ Sync,
{
if ctx.queued.is_empty() {
return Ok(None);
}
let evaluated_reason = ctx.reanalysis_reason;
let result = self.perform_reanalysis_and_dispatch(ctx).await?;
let evaluated = !self.analyzer_capacity_suppressed();
if evaluated
&& *reanalysis_reason == evaluated_reason
&& evaluated_reason.is_one_shot_edge_trigger()
{
*reanalysis_reason = ReanalysisReason::Initial;
}
Ok(Some(result))
}
async fn wait_for_scheduler_event(
&mut self,
join_set: &mut JoinSet<WorkspaceResult>,
in_flight: &mut HashSet<String>,
max_parallelism: usize,
merge_result_tx: &tokio::sync::mpsc::Sender<super::MergeResult>,
merge_result_rx: &mut tokio::sync::mpsc::Receiver<super::MergeResult>,
reanalysis_reason: &mut ReanalysisReason,
) {
tokio::select! {
Some(result) = join_set.join_next() => {
match result {
Ok(workspace_result) => {
self.handle_workspace_completion(workspace_result, max_parallelism, in_flight, merge_result_tx).await;
let manual_resolves_active = self
.manual_resolve_count
.as_ref()
.map(|counter| counter.load(std::sync::atomic::Ordering::Relaxed))
.unwrap_or(0);
*reanalysis_reason = if manual_resolves_active == 0 {
ReanalysisReason::ResolveCompletion
} else {
ReanalysisReason::Completion
};
self.trigger_resolve_wait_retry_dispatch();
}
Err(e) => {
error!("Task panicked: {:?}", e);
}
}
}
Some(merge_result) = merge_result_rx.recv() => {
if self.handle_merge_result_with_tx(merge_result, merge_result_tx).await.is_merged() {
self.trigger_resolve_wait_retry_dispatch();
*reanalysis_reason = ReanalysisReason::ResolveCompletion;
}
}
Some(_) = self.wait_for_dynamic_queue_notification() => {
info!("Queue notification received, will check queue on next iteration");
self.trigger_resolve_wait_retry_dispatch();
*reanalysis_reason = ReanalysisReason::QueueNotification;
}
_ = self.wait_for_cancellation(), if self.cancel_token.is_some() => {
info!("Cancellation received while scheduler is waiting for events");
}
_ = tokio::time::sleep(std::time::Duration::from_millis(500)) => {
}
}
}
#[allow(dead_code)]
pub(super) async fn wait_for_persistent_idle_wake(
&mut self,
reanalysis_reason: &mut ReanalysisReason,
merge_result_rx: &mut tokio::sync::mpsc::Receiver<super::MergeResult>,
) {
let (merge_result_tx, _merge_result_rx) = tokio::sync::mpsc::channel(1);
self.wait_for_persistent_idle_wake_with_tx(
reanalysis_reason,
&merge_result_tx,
merge_result_rx,
)
.await;
}
pub(super) async fn wait_for_persistent_idle_wake_with_tx(
&mut self,
reanalysis_reason: &mut ReanalysisReason,
merge_result_tx: &tokio::sync::mpsc::Sender<super::MergeResult>,
merge_result_rx: &mut tokio::sync::mpsc::Receiver<super::MergeResult>,
) {
info!(
"Scheduler idle with no work; waiting for dynamic queue notifications (persistent lifetime)"
);
tokio::select! {
Some(merge_result) = merge_result_rx.recv() => {
if self.handle_merge_result_with_tx(merge_result, merge_result_tx).await.is_merged() {
self.trigger_resolve_wait_retry_dispatch();
*reanalysis_reason = ReanalysisReason::ResolveCompletion;
}
}
Some(_) = self.wait_for_dynamic_queue_notification() => {
info!("Queue notification received while scheduler idle; resuming scheduler loop");
self.trigger_resolve_wait_retry_dispatch();
*reanalysis_reason = ReanalysisReason::QueueNotification;
}
_ = self.wait_for_cancellation(), if self.cancel_token.is_some() => {
info!("Cancellation received while scheduler idle; resuming scheduler loop");
}
}
}
async fn wait_for_dynamic_queue_notification(&self) -> Option<()> {
if let Some(queue) = &self.dynamic_queue {
queue.notified().await;
Some(())
} else {
std::future::pending().await
}
}
async fn wait_for_cancellation(&self) {
if let Some(token) = &self.cancel_token {
token.cancelled().await;
} else {
std::future::pending::<()>().await;
}
}
}