use std::path::PathBuf;
use tracing::warn;
use vtcode_memory::progress::{Milestone, MilestoneStatus, ProgressLedger, load_progress, save_progress};
pub trait ProgressLedgerSink: Send + Sync {
fn persist(&self, ledger: &ProgressLedger);
fn checkpoint(&self, _ledger: &ProgressLedger) {}
}
#[derive(Debug, Default, Clone, Copy)]
pub struct NullProgressSink;
impl ProgressLedgerSink for NullProgressSink {
fn persist(&self, _ledger: &ProgressLedger) {}
}
#[derive(Debug, Clone)]
pub struct SessionProgressSink {
workspace: PathBuf,
}
impl SessionProgressSink {
#[must_use]
pub fn new(workspace: PathBuf) -> Self {
Self { workspace }
}
}
impl ProgressLedgerSink for SessionProgressSink {
fn persist(&self, ledger: &ProgressLedger) {
if let Err(e) = save_progress(&self.workspace, &ledger.session_id, ledger) {
warn!(session = %ledger.session_id, error = %e, "failed to persist progress ledger");
}
}
fn checkpoint(&self, ledger: &ProgressLedger) {
let path = self.workspace.join("memories").join("progress.md");
if let Some(parent) = path.parent() {
if let Err(e) = std::fs::create_dir_all(parent) {
warn!(path = %parent.display(), error = %e, "failed to create memories dir");
return;
}
}
if let Err(e) = std::fs::write(&path, ledger.to_markdown()) {
warn!(path = %path.display(), error = %e, "failed to write progress memory");
}
}
}
pub struct ProgressMonitor {
ledger: ProgressLedger,
sink: Box<dyn ProgressLedgerSink>,
consecutive_stalls: u32,
}
impl ProgressMonitor {
#[must_use]
pub fn new(session_id: &str, goal: &str) -> Self {
Self::with_sink(ProgressLedger::new(session_id, goal), Box::new(NullProgressSink))
}
#[must_use]
pub fn with_sink(ledger: ProgressLedger, sink: Box<dyn ProgressLedgerSink>) -> Self {
Self { ledger, sink, consecutive_stalls: 0 }
}
pub fn with_persistence(workspace: PathBuf, session_id: &str, goal: &str) -> Self {
let ledger = load_progress(&workspace, session_id)
.ok()
.flatten()
.unwrap_or_else(|| ProgressLedger::new(session_id, goal));
Self::with_sink(ledger, Box::new(SessionProgressSink::new(workspace)))
}
#[must_use]
pub fn ledger(&self) -> &ProgressLedger {
&self.ledger
}
#[must_use]
pub fn is_stalled(&self) -> bool {
self.ledger.is_stalled()
}
#[must_use]
pub fn completion_ratio(&self) -> f32 {
self.ledger.completion_ratio()
}
#[must_use]
pub fn is_complete(&self) -> bool {
self.ledger.is_complete()
}
pub fn set_goal(&mut self, goal: &str) {
self.ledger.set_goal(goal);
self.persist();
}
pub fn set_milestones(&mut self, milestones: Vec<Milestone>) {
self.ledger.set_milestones(milestones);
self.persist();
}
pub fn record_advance(&mut self) {
self.ledger.note_advance();
self.consecutive_stalls = 0;
self.persist();
}
pub fn record_stall(&mut self) {
self.ledger.note_stall();
self.consecutive_stalls = self.consecutive_stalls.saturating_add(1);
self.persist();
}
#[must_use]
pub fn consecutive_stalls(&self) -> u32 {
self.consecutive_stalls
}
fn persist(&self) {
self.sink.persist(&self.ledger);
}
pub fn checkpoint(&self) {
self.sink.checkpoint(&self.ledger);
}
}
#[must_use]
pub fn milestone_status_from_str(status: &str) -> MilestoneStatus {
match status.trim().to_ascii_lowercase().as_str() {
"done" | "complete" | "completed" | "pass" | "passed" | "success" => MilestoneStatus::Done,
"blocked" | "stuck" | "waiting" => MilestoneStatus::Blocked,
"in_progress" | "in progress" | "active" | "running" => MilestoneStatus::InProgress,
_ => MilestoneStatus::Pending,
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
#[derive(Default)]
struct CountingSink {
persists: Arc<AtomicUsize>,
checkpoints: Arc<AtomicUsize>,
}
impl ProgressLedgerSink for CountingSink {
fn persist(&self, _ledger: &ProgressLedger) {
self.persists.fetch_add(1, Ordering::Relaxed);
}
fn checkpoint(&self, _ledger: &ProgressLedger) {
self.checkpoints.fetch_add(1, Ordering::Relaxed);
}
}
#[test]
fn mutations_persist_through_injected_sink() {
let persists = Arc::new(AtomicUsize::new(0));
let checkpoints = Arc::new(AtomicUsize::new(0));
let sink = CountingSink {
persists: persists.clone(),
checkpoints: checkpoints.clone(),
};
let mut monitor = ProgressMonitor::with_sink(ProgressLedger::new("s1", "goal"), Box::new(sink));
monitor.set_milestones(vec![Milestone {
id: "1".into(),
description: "step".into(),
status: MilestoneStatus::InProgress,
}]);
monitor.record_advance();
monitor.record_stall();
monitor.checkpoint();
assert_eq!(persists.load(Ordering::Relaxed), 3);
assert_eq!(checkpoints.load(Ordering::Relaxed), 1);
}
#[test]
fn null_sink_monitor_is_pure_in_memory() {
let mut monitor = ProgressMonitor::new("s2", "goal");
assert!(monitor.is_complete()); monitor.set_milestones(vec![Milestone {
id: "1".into(),
description: "step".into(),
status: MilestoneStatus::Pending,
}]);
assert!(!monitor.is_complete());
assert!((monitor.completion_ratio() - 0.0).abs() < f32::EPSILON);
}
#[test]
fn consecutive_stalls_increment_and_reset() {
let mut monitor = ProgressMonitor::new("s3", "goal");
assert_eq!(monitor.consecutive_stalls(), 0);
monitor.record_stall();
assert_eq!(monitor.consecutive_stalls(), 1);
monitor.record_stall();
assert_eq!(monitor.consecutive_stalls(), 2);
monitor.record_advance();
assert_eq!(monitor.consecutive_stalls(), 0);
monitor.record_stall();
assert_eq!(monitor.consecutive_stalls(), 1);
}
}