use std::collections::VecDeque;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum CheckpointTrigger {
BeforeFirstDispatch,
BeforeBaseIntegration,
AfterDrain,
PrePushRemoteAdvance,
PushRaceRejection,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct BaseLaneState {
pub base_dirty_reason: Option<String>,
pub lane_busy_reason: Option<String>,
}
impl BaseLaneState {
pub fn clean() -> Self {
Self {
base_dirty_reason: None,
lane_busy_reason: None,
}
}
fn unsafe_reason(&self) -> Option<String> {
self.base_dirty_reason
.clone()
.or_else(|| self.lane_busy_reason.clone())
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum CheckpointDecision {
Start { generation: u64 },
Batched { generation: u64 },
Deferred { reason: String },
}
#[derive(Debug, Default)]
pub struct CheckpointScheduler {
generation: u64,
active: Option<u64>,
queued_results: VecDeque<String>,
first_dispatch_done: bool,
}
impl CheckpointScheduler {
pub fn new() -> Self {
Self::default()
}
pub fn active_generation(&self) -> Option<u64> {
self.active
}
pub fn is_active(&self) -> bool {
self.active.is_some()
}
pub fn queued_results(&self) -> Vec<String> {
self.queued_results.iter().cloned().collect()
}
pub fn request(
&mut self,
trigger: CheckpointTrigger,
lane: &BaseLaneState,
pending_result: Option<&str>,
) -> CheckpointDecision {
if let Some(result) = pending_result {
if !self.queued_results.iter().any(|queued| queued == result) {
self.queued_results.push_back(result.to_string());
}
}
if let Some(generation) = self.active {
return CheckpointDecision::Batched { generation };
}
if let Some(reason) = lane.unsafe_reason() {
return CheckpointDecision::Deferred { reason };
}
if trigger == CheckpointTrigger::BeforeFirstDispatch && self.first_dispatch_done {
return CheckpointDecision::Deferred {
reason: "pre-dispatch checkpoint already completed for this run".to_string(),
};
}
self.generation += 1;
self.active = Some(self.generation);
if trigger == CheckpointTrigger::BeforeFirstDispatch {
self.first_dispatch_done = true;
}
CheckpointDecision::Start {
generation: self.generation,
}
}
pub fn release(&mut self, generation: u64) -> Option<Vec<String>> {
if self.active != Some(generation) {
return None;
}
self.active = None;
Some(self.queued_results.drain(..).collect())
}
}
#[cfg(test)]
mod tests {
use super::*;
fn busy_lane(reason: &str) -> BaseLaneState {
BaseLaneState {
base_dirty_reason: None,
lane_busy_reason: Some(reason.to_string()),
}
}
#[test]
fn upstream_integration_first_dispatch_checkpoint_starts_once() {
let mut scheduler = CheckpointScheduler::new();
assert_eq!(
scheduler.request(
CheckpointTrigger::BeforeFirstDispatch,
&BaseLaneState::clean(),
None
),
CheckpointDecision::Start { generation: 1 }
);
assert_eq!(scheduler.release(1), Some(Vec::new()));
assert!(matches!(
scheduler.request(
CheckpointTrigger::BeforeFirstDispatch,
&BaseLaneState::clean(),
None
),
CheckpointDecision::Deferred { .. }
));
}
#[test]
fn upstream_integration_checkpoints_do_not_overlap() {
let mut scheduler = CheckpointScheduler::new();
let first = scheduler.request(
CheckpointTrigger::BeforeFirstDispatch,
&BaseLaneState::clean(),
None,
);
assert_eq!(first, CheckpointDecision::Start { generation: 1 });
assert_eq!(
scheduler.request(CheckpointTrigger::AfterDrain, &BaseLaneState::clean(), None),
CheckpointDecision::Batched { generation: 1 }
);
}
#[test]
fn upstream_integration_batches_queued_results_behind_one_fetch() {
let mut scheduler = CheckpointScheduler::new();
assert_eq!(
scheduler.request(
CheckpointTrigger::BeforeBaseIntegration,
&BaseLaneState::clean(),
Some("change-a")
),
CheckpointDecision::Start { generation: 1 }
);
assert_eq!(
scheduler.request(
CheckpointTrigger::BeforeBaseIntegration,
&BaseLaneState::clean(),
Some("change-b")
),
CheckpointDecision::Batched { generation: 1 }
);
assert_eq!(
scheduler.request(
CheckpointTrigger::BeforeBaseIntegration,
&BaseLaneState::clean(),
Some("change-b")
),
CheckpointDecision::Batched { generation: 1 }
);
assert_eq!(
scheduler.queued_results(),
vec!["change-a".to_string(), "change-b".to_string()]
);
assert_eq!(
scheduler.release(1),
Some(vec!["change-a".to_string(), "change-b".to_string()])
);
}
#[test]
fn upstream_integration_dirty_base_defers_all_checkpoint_side_effects() {
let mut scheduler = CheckpointScheduler::new();
let dirty = BaseLaneState {
base_dirty_reason: Some("uncommitted changes".to_string()),
lane_busy_reason: None,
};
assert_eq!(
scheduler.request(CheckpointTrigger::AfterDrain, &dirty, None),
CheckpointDecision::Deferred {
reason: "uncommitted changes".to_string()
}
);
assert!(!scheduler.is_active());
}
#[test]
fn upstream_integration_busy_lane_defers_checkpoint() {
let mut scheduler = CheckpointScheduler::new();
assert_eq!(
scheduler.request(
CheckpointTrigger::BeforeBaseIntegration,
&busy_lane("merge in progress"),
Some("change-a")
),
CheckpointDecision::Deferred {
reason: "merge in progress".to_string()
}
);
assert_eq!(scheduler.queued_results(), vec!["change-a".to_string()]);
}
#[test]
fn upstream_integration_stale_completion_cannot_release_ownership() {
let mut scheduler = CheckpointScheduler::new();
scheduler.request(CheckpointTrigger::AfterDrain, &BaseLaneState::clean(), None);
assert_eq!(scheduler.release(0), None);
assert_eq!(scheduler.release(99), None);
assert!(scheduler.is_active());
assert_eq!(scheduler.release(1), Some(Vec::new()));
assert_eq!(scheduler.release(1), None);
}
#[test]
fn upstream_integration_release_allows_next_checkpoint_generation() {
let mut scheduler = CheckpointScheduler::new();
scheduler.request(
CheckpointTrigger::PrePushRemoteAdvance,
&BaseLaneState::clean(),
None,
);
scheduler.release(1);
assert_eq!(
scheduler.request(
CheckpointTrigger::PushRaceRejection,
&BaseLaneState::clean(),
None
),
CheckpointDecision::Start { generation: 2 }
);
}
}