use std::collections::HashMap;
use std::sync::atomic::AtomicU64;
use std::sync::{Arc, Mutex};
use async_trait::async_trait;
use serde::{Deserialize, Serialize};
use crate::error::Result;
use crate::graph::status::GraphRunStatus;
use crate::graph::stream::{GraphEvent, GraphEventSink};
use crate::harness::ids::{CheckpointId, EventId, GraphId, NodeId, RunId, ThreadId};
use crate::harness::observability::AppendWorker;
use crate::harness::store::AppendStore;
#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
pub struct GraphObservation {
pub event_id: EventId,
pub run_id: RunId,
pub root_run_id: RunId,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub parent_run_id: Option<RunId>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub thread_id: Option<ThreadId>,
pub graph_id: GraphId,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub checkpoint_id: Option<CheckpointId>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub namespace: Vec<String>,
pub step: usize,
pub offset: u64,
pub ts_ms: u64,
pub event: GraphEvent,
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub struct GraphStepLatency {
pub step: usize,
pub elapsed_ms: u64,
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub struct GraphNodeLatency {
pub node: NodeId,
pub step: usize,
pub elapsed_ms: u64,
pub failed: bool,
}
#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct GraphLatencyMetrics {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub run_elapsed_ms: Option<u64>,
#[serde(default)]
pub steps: Vec<GraphStepLatency>,
#[serde(default)]
pub nodes: Vec<GraphNodeLatency>,
pub total_step_ms: u64,
pub max_step_ms: u64,
pub total_node_ms: u64,
pub max_node_ms: u64,
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub struct GraphNodeHealth {
pub node: NodeId,
pub started: u64,
pub completed: u64,
pub failed: u64,
}
impl GraphNodeHealth {
pub fn attempts(&self) -> u64 {
self.completed.saturating_add(self.failed)
}
pub fn failure_rate(&self) -> f64 {
let attempts = self.attempts();
if attempts == 0 {
0.0
} else {
self.failed as f64 / attempts as f64
}
}
pub fn is_healthy(&self) -> bool {
self.failed == 0
}
}
#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct GraphHealthSummary {
#[serde(default)]
pub nodes: Vec<GraphNodeHealth>,
pub total_started: u64,
pub total_completed: u64,
pub total_failed: u64,
pub run_failed: bool,
}
impl GraphHealthSummary {
pub fn total_attempts(&self) -> u64 {
self.total_completed.saturating_add(self.total_failed)
}
pub fn failure_rate(&self) -> f64 {
let attempts = self.total_attempts();
if attempts == 0 {
0.0
} else {
self.total_failed as f64 / attempts as f64
}
}
pub fn is_healthy(&self) -> bool {
self.total_failed == 0 && !self.run_failed
}
pub fn unhealthy_nodes(&self) -> impl Iterator<Item = &GraphNodeHealth> {
self.nodes.iter().filter(|n| !n.is_healthy())
}
}
#[async_trait]
pub trait GraphEventJournal: Send + Sync {
async fn append(&self, obs: GraphObservation) -> Result<u64>;
async fn read_from(&self, run_id: &str, offset: u64) -> Result<Vec<GraphObservation>>;
}
#[derive(Clone, Debug, Default)]
pub struct InMemoryGraphEventJournal {
pub(crate) runs: Arc<Mutex<HashMap<String, Vec<GraphObservation>>>>,
}
#[derive(Clone, Debug)]
pub struct StoreGraphEventJournal<A: AppendStore> {
pub(crate) store: A,
}
#[async_trait]
pub trait GraphStatusStore: Send + Sync {
async fn put_status(&self, status: GraphRunStatus) -> Result<()>;
async fn get_status(&self, run_id: &str) -> Result<Option<GraphRunStatus>>;
async fn list_by_thread(&self, thread_id: &str) -> Result<Vec<GraphRunStatus>>;
}
#[derive(Clone, Debug, Default)]
pub struct InMemoryGraphStatusStore {
pub(crate) state: Arc<Mutex<StatusStoreState>>,
pub(crate) max_runs: Option<usize>,
}
#[derive(Debug, Default)]
pub(crate) struct StatusStoreState {
pub(crate) statuses: HashMap<String, GraphRunStatus>,
pub(crate) by_thread: HashMap<String, Vec<String>>,
pub(crate) order: std::collections::VecDeque<String>,
}
#[derive(Clone)]
pub struct JournalGraphSink {
pub(crate) worker: Arc<AppendWorker<GraphObservation>>,
pub(crate) inner: Option<Arc<dyn GraphEventSink>>,
pub(crate) run_id: RunId,
pub(crate) root_run_id: RunId,
pub(crate) parent_run_id: Option<RunId>,
pub(crate) thread_id: Option<ThreadId>,
pub(crate) graph_id: GraphId,
pub(crate) namespace: Vec<String>,
pub(crate) offset: Arc<AtomicU64>,
pub(crate) step: Arc<AtomicU64>,
}