pe-graph 0.1.0

Graph execution engine for Potential Expectations — state graphs, Pregel model, ReAct topology, and builder DSL
Documentation
//! Pending writes tracking for fault tolerance within a superstep.
//!
//! When multiple nodes execute in parallel, some may succeed while others fail.
//! `PendingWrites` records which nodes produced valid updates so that on resume,
//! successful nodes don't re-execute.

use pe_core::error::PeError;
use pe_core::state::StateUpdate;

/// Tracks node write outcomes within a single superstep.
///
/// # Example
///
/// ```ignore
/// let mut writes = PendingWrites::new();
/// writes.record_success("chat", update);
/// writes.record_failure("tools", &PeError::ToolExecution { .. });
///
/// if writes.has_failures() {
///     // Only re-run failed nodes on resume
/// }
/// ```
#[derive(Debug, Clone)]
pub struct PendingWrites<U: StateUpdate> {
    successes: Vec<(String, U)>,
    failures: Vec<(String, String)>,
}

impl<U: StateUpdate> PendingWrites<U> {
    /// Create an empty write tracker.
    pub fn new() -> Self {
        Self {
            successes: Vec::new(),
            failures: Vec::new(),
        }
    }

    /// Record a successful node execution with its update.
    pub fn record_success(&mut self, node_name: &str, update: U) {
        self.successes.push((node_name.to_string(), update));
    }

    /// Record a failed node execution.
    pub fn record_failure(&mut self, node_name: &str, error: &PeError) {
        self.failures
            .push((node_name.to_string(), error.to_string()));
    }

    /// Returns true if any node failed in this superstep.
    pub fn has_failures(&self) -> bool {
        !self.failures.is_empty()
    }

    /// View successful node writes.
    pub fn successes(&self) -> &[(String, U)] {
        &self.successes
    }

    /// View failed nodes and their error messages.
    pub fn failures(&self) -> &[(String, String)] {
        &self.failures
    }

    /// Drain all successful writes, leaving the tracker empty.
    pub fn drain_successes(&mut self) -> Vec<(String, U)> {
        std::mem::take(&mut self.successes)
    }
}

impl<U: StateUpdate> Default for PendingWrites<U> {
    fn default() -> Self {
        Self::new()
    }
}

#[cfg(test)]
mod tests {
    use super::*;
    use serde::{Deserialize, Serialize};

    #[derive(Debug, Clone, Default, Serialize, Deserialize)]
    struct FakeUpdate {
        value: i32,
    }
    impl StateUpdate for FakeUpdate {}

    #[test]
    fn test_record_and_access() {
        let mut writes = PendingWrites::new();
        writes.record_success("node_a", FakeUpdate { value: 1 });
        writes.record_success("node_b", FakeUpdate { value: 2 });

        assert_eq!(writes.successes().len(), 2);
        assert_eq!(writes.successes()[0].0, "node_a");
        assert_eq!(writes.successes()[1].1.value, 2);
        assert!(!writes.has_failures());
    }

    #[test]
    fn test_record_failure() {
        let mut writes: PendingWrites<FakeUpdate> = PendingWrites::new();
        writes.record_failure(
            "bad_node",
            &PeError::Internal {
                details: "boom".into(),
            },
        );

        assert!(writes.has_failures());
        assert_eq!(writes.failures().len(), 1);
        assert_eq!(writes.failures()[0].0, "bad_node");
    }

    #[test]
    fn test_drain_successes() {
        let mut writes = PendingWrites::new();
        writes.record_success("a", FakeUpdate { value: 10 });
        writes.record_success("b", FakeUpdate { value: 20 });

        let drained = writes.drain_successes();
        assert_eq!(drained.len(), 2);
        assert!(writes.successes().is_empty());
    }
}