use std::sync::{Arc, Mutex, Weak};
use std::time::Duration;
use async_trait::async_trait;
use crate::orchestration::operator_command::{is_active_status, is_final_status};
pub const MARK_STABILITY_WINDOW: Duration = Duration::from_secs(10);
pub const MARK_SETTLEMENT_ATTEMPTS: u32 = 3;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum MarkSettlementExclusion {
NotLoadable,
Terminal,
Active,
Waiting,
AlreadyQueued,
AlreadyNotQueued,
Unavailable,
}
impl MarkSettlementExclusion {
pub fn as_str(self) -> &'static str {
match self {
Self::NotLoadable => "not_loadable",
Self::Terminal => "terminal",
Self::Active => "active",
Self::Waiting => "waiting",
Self::AlreadyQueued => "already_queued",
Self::AlreadyNotQueued => "already_not_queued",
Self::Unavailable => "unavailable",
}
}
pub fn is_stable(self) -> bool {
!matches!(self, Self::NotLoadable)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum MarkSettlementFailure {
RuntimeUnbound,
RuntimeGone,
NoTaskRuntime,
UnreconciledBatch,
}
impl MarkSettlementFailure {
pub fn as_str(self) -> &'static str {
match self {
Self::RuntimeUnbound => "runtime_unbound",
Self::RuntimeGone => "runtime_gone",
Self::NoTaskRuntime => "no_task_runtime",
Self::UnreconciledBatch => "unreconciled_batch",
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum MarkSettlementAction {
Add,
Remove,
}
#[derive(Debug, Clone, Copy)]
pub struct MarkSettlementRow<'a> {
pub change_id: &'a str,
pub display_status: &'a str,
pub tracked: bool,
pub parallel_eligible: bool,
pub marked: bool,
}
#[derive(Debug, Default, Clone, PartialEq, Eq)]
pub struct MarkSettlementPlan {
pub additions: Vec<String>,
pub removals: Vec<String>,
pub excluded: Vec<(String, MarkSettlementExclusion)>,
}
impl MarkSettlementPlan {
#[cfg_attr(not(test), allow(dead_code))]
pub fn is_empty(&self) -> bool {
self.additions.is_empty() && self.removals.is_empty()
}
}
pub fn classify_mark_settlement_row(
row: &MarkSettlementRow<'_>,
) -> std::result::Result<MarkSettlementAction, MarkSettlementExclusion> {
if !row.tracked {
return Err(MarkSettlementExclusion::NotLoadable);
}
if is_final_status(row.display_status) || matches!(row.display_status, "error" | "stopped") {
return Err(MarkSettlementExclusion::Terminal);
}
if is_active_status(row.display_status) {
return Err(MarkSettlementExclusion::Active);
}
if matches!(
row.display_status,
"merge wait" | "resolve pending" | "reject pending" | "blocked" | "stalled"
) {
return Err(MarkSettlementExclusion::Waiting);
}
match row.display_status {
"queued" if row.marked => Err(MarkSettlementExclusion::AlreadyQueued),
"queued" => Ok(MarkSettlementAction::Remove),
"not queued" if !row.marked => Err(MarkSettlementExclusion::AlreadyNotQueued),
"not queued" if !row.parallel_eligible => Err(MarkSettlementExclusion::Unavailable),
"not queued" => Ok(MarkSettlementAction::Add),
_ => Err(MarkSettlementExclusion::Waiting),
}
}
pub fn plan_mark_settlement(rows: &[MarkSettlementRow<'_>]) -> MarkSettlementPlan {
let mut plan = MarkSettlementPlan::default();
for row in rows {
match classify_mark_settlement_row(row) {
Ok(MarkSettlementAction::Add) => plan.additions.push(row.change_id.to_string()),
Ok(MarkSettlementAction::Remove) => plan.removals.push(row.change_id.to_string()),
Err(reason) => plan.excluded.push((row.change_id.to_string(), reason)),
}
}
plan
}
#[async_trait]
pub trait MarkSettlementRuntime: Send + Sync {
fn admits_dynamic_queue(&self) -> bool;
async fn settle_marks(&self, targets: Vec<String>) -> MarkSettlementPlan;
async fn report_abandoned_settlement(&self, pending: Vec<String>);
async fn report_settlement_failure(
&self,
_failure: MarkSettlementFailure,
_targets: Vec<String>,
) {
}
}
pub struct MarkSettlementCoordinator {
inner: Mutex<CoordinatorInner>,
window: Duration,
passes: tokio::sync::watch::Sender<u64>,
}
#[derive(Default)]
struct CoordinatorInner {
runtime: Option<Weak<dyn MarkSettlementRuntime>>,
generation: u64,
pending: Option<Vec<String>>,
settled: u64,
abandoned: u64,
last_plan: Option<MarkSettlementPlan>,
attempts: u32,
last_failure: Option<MarkSettlementFailure>,
}
impl std::fmt::Debug for MarkSettlementCoordinator {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
let guard = self.lock();
f.debug_struct("MarkSettlementCoordinator")
.field("armed", &guard.pending.is_some())
.field("generation", &guard.generation)
.field("settled", &guard.settled)
.field("abandoned", &guard.abandoned)
.finish()
}
}
impl Default for MarkSettlementCoordinator {
fn default() -> Self {
Self::new()
}
}
impl MarkSettlementCoordinator {
pub fn new() -> Self {
Self::with_window(MARK_STABILITY_WINDOW)
}
pub fn with_window(window: Duration) -> Self {
Self {
inner: Mutex::new(CoordinatorInner::default()),
window,
passes: tokio::sync::watch::channel(0).0,
}
}
fn lock(&self) -> std::sync::MutexGuard<'_, CoordinatorInner> {
self.inner
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
}
pub fn bind_runtime(&self, runtime: Weak<dyn MarkSettlementRuntime>) {
self.lock().runtime = Some(runtime);
}
#[cfg_attr(not(test), allow(dead_code))]
pub fn window(&self) -> Duration {
self.window
}
pub fn notify(self: &Arc<Self>, changed: Vec<String>) -> bool {
let runtime = match self.runtime_or_failure() {
Ok(runtime) => runtime,
Err(failure) => {
self.record_failure(failure, &changed);
return false;
}
};
if !runtime.admits_dynamic_queue() {
return false;
}
let Ok(handle) = tokio::runtime::Handle::try_current() else {
self.record_failure(MarkSettlementFailure::NoTaskRuntime, &changed);
return false;
};
let generation = {
let mut guard = self.lock();
guard.generation += 1;
guard.attempts = 0;
let batch = guard.pending.get_or_insert_with(Vec::new);
for change_id in changed {
if !batch.iter().any(|existing| existing == &change_id) {
batch.push(change_id);
}
}
guard.generation
};
self.arm(&handle, generation);
true
}
fn arm(self: &Arc<Self>, handle: &tokio::runtime::Handle, generation: u64) {
let coordinator = self.clone();
let window = self.window;
handle.spawn(async move {
tokio::time::sleep(window).await;
coordinator.settle(generation).await;
});
}
#[cfg_attr(not(test), allow(dead_code))]
pub fn pending_snapshot(&self) -> Option<Vec<String>> {
self.lock().pending.clone()
}
#[cfg_attr(not(test), allow(dead_code))]
pub fn is_armed(&self) -> bool {
self.lock().pending.is_some()
}
#[cfg_attr(not(test), allow(dead_code))]
pub fn settled_count(&self) -> u64 {
self.lock().settled
}
#[cfg_attr(not(test), allow(dead_code))]
pub fn abandoned_count(&self) -> u64 {
self.lock().abandoned
}
#[cfg_attr(not(test), allow(dead_code))]
pub fn last_plan(&self) -> Option<MarkSettlementPlan> {
self.lock().last_plan.clone()
}
#[cfg_attr(not(test), allow(dead_code))]
pub fn passes(&self) -> tokio::sync::watch::Receiver<u64> {
self.passes.subscribe()
}
#[cfg_attr(not(test), allow(dead_code))]
pub fn last_failure(&self) -> Option<MarkSettlementFailure> {
self.lock().last_failure
}
fn record_pass(&self) {
self.passes.send_modify(|passes| *passes += 1);
}
fn record_failure(&self, failure: MarkSettlementFailure, targets: &[String]) {
self.lock().last_failure = Some(failure);
let reason = failure.as_str();
let targets = if targets.is_empty() {
"no marked change".to_string()
} else {
targets.join(", ")
};
match failure {
MarkSettlementFailure::RuntimeUnbound => {
tracing::debug!("Mark settlement is mark-only (reason={reason}): {targets}");
}
_ => tracing::warn!("Mark settlement could not complete (reason={reason}): {targets}"),
}
}
fn runtime(&self) -> Option<Arc<dyn MarkSettlementRuntime>> {
self.lock().runtime.as_ref().and_then(Weak::upgrade)
}
fn runtime_or_failure(
&self,
) -> std::result::Result<Arc<dyn MarkSettlementRuntime>, MarkSettlementFailure> {
let bound = self.lock().runtime.clone();
match bound {
None => Err(MarkSettlementFailure::RuntimeUnbound),
Some(weak) => weak.upgrade().ok_or(MarkSettlementFailure::RuntimeGone),
}
}
async fn settle(self: Arc<Self>, generation: u64) {
let batch = {
let mut guard = self.lock();
if guard.generation != generation {
return;
}
guard.pending.take()
};
let batch = batch.unwrap_or_default();
let Some(runtime) = self.runtime() else {
self.record_failure(MarkSettlementFailure::RuntimeGone, &batch);
self.lock().abandoned += 1;
self.record_pass();
return;
};
if !runtime.admits_dynamic_queue() {
self.lock().abandoned += 1;
runtime.report_abandoned_settlement(batch).await;
self.record_pass();
return;
}
let plan = runtime.settle_marks(batch).await;
let unreconciled: Vec<String> = plan
.excluded
.iter()
.filter(|(_, reason)| !reason.is_stable())
.map(|(change_id, _)| change_id.clone())
.collect();
let retry = {
let mut guard = self.lock();
guard.settled += 1;
guard.last_plan = Some(plan);
if unreconciled.is_empty() {
None
} else {
guard.attempts += 1;
if guard.attempts >= MARK_SETTLEMENT_ATTEMPTS {
None
} else {
guard.generation += 1;
let pending = guard.pending.get_or_insert_with(Vec::new);
for change_id in &unreconciled {
if !pending.iter().any(|existing| existing == change_id) {
pending.push(change_id.clone());
}
}
Some(guard.generation)
}
}
};
match retry {
Some(generation) => match tokio::runtime::Handle::try_current() {
Ok(handle) => self.arm(&handle, generation),
Err(_) => {
self.record_failure(MarkSettlementFailure::NoTaskRuntime, &unreconciled);
runtime
.report_settlement_failure(
MarkSettlementFailure::NoTaskRuntime,
unreconciled,
)
.await;
}
},
None if !unreconciled.is_empty() => {
self.record_failure(MarkSettlementFailure::UnreconciledBatch, &unreconciled);
runtime
.report_settlement_failure(
MarkSettlementFailure::UnreconciledBatch,
unreconciled,
)
.await;
}
None => {}
}
self.record_pass();
}
}
#[cfg(test)]
mod tests;