use std::sync::Arc;
use lunaris_core::LunarisError;
use lunaris_extract::{Extractor, NoopExtractor};
use parking_lot::{Mutex, RwLock};
pub const ENABLED_ENV_VAR: &str = "LUNARIS_GRAPH_ENABLED";
pub struct GraphPipelineHandle {
enabled: RwLock<bool>,
extractor: RwLock<Option<Arc<dyn Extractor>>>,
state_changes: Mutex<u64>,
}
impl GraphPipelineHandle {
pub fn new(initial_enabled: bool, extractor: Arc<dyn Extractor>) -> Self {
Self {
enabled: RwLock::new(initial_enabled),
extractor: RwLock::new(Some(extractor)),
state_changes: Mutex::new(0),
}
}
pub fn initial_state_from_value(raw: Option<&str>) -> bool {
matches!(raw, Some("1" | "true" | "TRUE" | "on" | "ON"))
}
pub fn initial_state_from_env() -> bool {
Self::initial_state_from_value(std::env::var(ENABLED_ENV_VAR).ok().as_deref())
}
pub fn enable(&self) {
let mut w = self.enabled.write();
if !*w {
*w = true;
*self.state_changes.lock() += 1;
tracing::info!(state = "enabled", "graph_pipeline_state_changed");
}
}
pub fn disable(&self) {
let mut w = self.enabled.write();
if *w {
*w = false;
*self.state_changes.lock() += 1;
tracing::info!(state = "disabled", "graph_pipeline_state_changed");
}
}
pub fn is_enabled(&self) -> bool {
*self.enabled.read()
}
pub fn state_change_count(&self) -> u64 {
*self.state_changes.lock()
}
pub async fn force_reload(&self) -> Result<(), LunarisError> {
tracing::info!("graph_pipeline_extractor_reloaded (noop — extractor is remote-only)");
Ok(())
}
pub fn set_extractor(&self, extractor: Arc<dyn Extractor>) {
*self.extractor.write() = Some(extractor);
tracing::info!("graph_pipeline_extractor_replaced");
}
pub fn snapshot_extractor(&self) -> Option<Arc<dyn Extractor>> {
self.extractor.read().clone()
}
pub fn with_noop() -> Self {
Self::new(false, Arc::new(NoopExtractor) as Arc<dyn Extractor>)
}
}
impl std::fmt::Debug for GraphPipelineHandle {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("GraphPipelineHandle")
.field("enabled", &*self.enabled.read())
.field("has_extractor", &self.extractor.read().is_some())
.field("state_change_count", &*self.state_changes.lock())
.finish()
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn default_state_is_off_when_value_none() {
assert!(!GraphPipelineHandle::initial_state_from_value(None));
}
#[test]
fn value_one_enables_initial_state() {
assert!(GraphPipelineHandle::initial_state_from_value(Some("1")));
assert!(GraphPipelineHandle::initial_state_from_value(Some("true")));
assert!(GraphPipelineHandle::initial_state_from_value(Some("TRUE")));
assert!(GraphPipelineHandle::initial_state_from_value(Some("on")));
assert!(GraphPipelineHandle::initial_state_from_value(Some("ON")));
}
#[test]
fn value_off_disables_initial_state() {
assert!(!GraphPipelineHandle::initial_state_from_value(Some("0")));
assert!(!GraphPipelineHandle::initial_state_from_value(Some("false")));
assert!(!GraphPipelineHandle::initial_state_from_value(Some("")));
assert!(!GraphPipelineHandle::initial_state_from_value(Some("yes"))); assert!(!GraphPipelineHandle::initial_state_from_value(Some("True"))); }
#[test]
fn enable_disable_is_observable_and_idempotent() {
let h = GraphPipelineHandle::with_noop();
assert!(!h.is_enabled());
assert_eq!(h.state_change_count(), 0);
h.enable();
assert!(h.is_enabled());
assert_eq!(h.state_change_count(), 1);
h.enable();
assert_eq!(h.state_change_count(), 1);
h.disable();
assert!(!h.is_enabled());
assert_eq!(h.state_change_count(), 2);
h.disable();
assert_eq!(h.state_change_count(), 2);
h.enable();
assert!(h.is_enabled());
assert_eq!(h.state_change_count(), 3);
}
#[test]
fn snapshot_extractor_returns_arc_clone() {
let h = GraphPipelineHandle::new(true, Arc::new(NoopExtractor));
let snap1 = h.snapshot_extractor();
let snap2 = h.snapshot_extractor();
assert!(snap1.is_some());
assert!(snap2.is_some());
assert!(Arc::ptr_eq(snap1.as_ref().unwrap(), snap2.as_ref().unwrap()));
}
#[tokio::test]
async fn force_reload_without_candle_is_noop() {
let h = GraphPipelineHandle::with_noop();
let before = h.is_enabled();
let _ = h.force_reload().await; assert_eq!(h.is_enabled(), before, "force_reload MUST NOT change toggle state");
}
#[test]
fn set_extractor_replaces_handle() {
let h = GraphPipelineHandle::with_noop();
let original = h.snapshot_extractor().unwrap();
let replacement: Arc<dyn Extractor> = Arc::new(NoopExtractor);
h.set_extractor(replacement.clone());
let after = h.snapshot_extractor().unwrap();
assert!(!Arc::ptr_eq(&original, &after), "set_extractor must swap the Arc");
assert!(Arc::ptr_eq(&replacement, &after));
}
#[test]
fn set_extractor_preserves_toggle_state_and_counter() {
let h = GraphPipelineHandle::with_noop();
h.enable();
assert_eq!(h.state_change_count(), 1);
assert!(h.is_enabled());
let replacement: Arc<dyn Extractor> = Arc::new(NoopExtractor);
h.set_extractor(replacement);
assert!(h.is_enabled(), "set_extractor must not flip the toggle");
assert_eq!(h.state_change_count(), 1, "set_extractor must not increment state changes");
}
#[test]
fn debug_impl_is_safe_to_format() {
let h = GraphPipelineHandle::with_noop();
h.enable();
let dbg = format!("{:?}", h);
assert!(dbg.contains("enabled"));
assert!(dbg.contains("has_extractor"));
assert!(dbg.contains("state_change_count"));
}
}