use std::sync::Arc;
use lunaris_core::HlcClock;
use lunaris_verify::{NoopVerifier, Verifier};
use parking_lot::{Mutex, RwLock};
pub const ENABLED_ENV_VAR: &str = "LUNARIS_VERIFY_ENABLED";
pub struct VerifierPipelineHandle {
enabled: RwLock<bool>,
verifier: RwLock<Option<Arc<dyn Verifier>>>,
state_change_count: std::sync::atomic::AtomicU64,
shutdown: Arc<tokio::sync::Notify>,
worker_handle: Mutex<Option<tokio::task::JoinHandle<()>>>,
storage: RwLock<Option<Arc<dyn lunaris_core::StoragePort>>>,
clock: RwLock<Option<Arc<HlcClock>>>,
}
impl VerifierPipelineHandle {
pub fn new(initial_enabled: bool, verifier: Arc<dyn Verifier>) -> Self {
Self {
enabled: RwLock::new(initial_enabled),
verifier: RwLock::new(Some(verifier)),
state_change_count: std::sync::atomic::AtomicU64::new(0),
shutdown: Arc::new(tokio::sync::Notify::new()),
worker_handle: Mutex::new(None),
storage: RwLock::new(None),
clock: RwLock::new(None),
}
}
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 bind_storage(&self, storage: Arc<dyn lunaris_core::StoragePort>) {
*self.storage.write() = Some(storage);
}
pub fn bind_clock(&self, clock: Arc<HlcClock>) {
*self.clock.write() = Some(clock);
}
pub(crate) fn spawn_worker_if_idle(&self) {
let mut wh = self.worker_handle.lock();
if wh.is_some() {
return;
}
let storage = match self.storage.read().clone() {
Some(s) => s,
None => {
tracing::warn!("verify_pipeline_enable_without_storage; worker not spawned");
return;
}
};
let clock = self.clock.read().clone().unwrap_or_else(|| HlcClock::new(0));
let verifier =
self.snapshot_verifier().unwrap_or_else(|| Arc::new(NoopVerifier) as Arc<dyn Verifier>);
let shutdown = self.shutdown.clone();
let handle = tokio::spawn(async move {
#[allow(deprecated)]
match lunaris_verify::run_verify_worker(storage, verifier, shutdown, clock).await {
Ok(jh) => {
if let Err(e) = jh.await {
tracing::warn!(err = %e, "verify_pipeline_inner_worker_join_failed");
}
}
Err(e) => {
tracing::error!(err = %e, "verify_pipeline_worker_spawn_failed");
}
}
});
*wh = Some(handle);
}
pub fn enable(&self) {
let mut w = self.enabled.write();
if !*w {
*w = true;
self.state_change_count.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
tracing::info!(state = "enabled", "verify_pipeline_state_changed");
drop(w);
self.spawn_worker_if_idle();
}
}
pub fn disable(&self) {
let mut w = self.enabled.write();
if *w {
*w = false;
self.state_change_count.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
tracing::info!(state = "disabled", "verify_pipeline_state_changed");
drop(w);
self.shutdown.notify_one();
}
}
pub async fn join_worker(&self) {
let handle = self.worker_handle.lock().take();
if let Some(h) = handle
&& let Err(e) = h.await
{
tracing::warn!(err = %e, "verify_pipeline_worker_join_failed");
}
}
pub fn is_enabled(&self) -> bool {
*self.enabled.read()
}
pub fn state_change_count(&self) -> u64 {
self.state_change_count.load(std::sync::atomic::Ordering::SeqCst)
}
pub fn set_verifier(&self, verifier: Arc<dyn Verifier>) {
*self.verifier.write() = Some(verifier);
tracing::info!("verify_pipeline_verifier_replaced");
}
pub fn snapshot_verifier(&self) -> Option<Arc<dyn Verifier>> {
self.verifier.read().clone()
}
pub fn with_noop() -> Self {
Self::new(false, Arc::new(NoopVerifier) as Arc<dyn Verifier>)
}
}
impl std::fmt::Debug for VerifierPipelineHandle {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("VerifierPipelineHandle")
.field("enabled", &*self.enabled.read())
.field("has_verifier", &self.verifier.read().is_some())
.field("has_storage", &self.storage.read().is_some())
.field("has_clock", &self.clock.read().is_some())
.field("has_worker", &self.worker_handle.lock().is_some())
.field("state_change_count", &self.state_change_count())
.finish()
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn default_state_is_off_when_value_none() {
assert!(!VerifierPipelineHandle::initial_state_from_value(None));
}
#[test]
fn value_one_enables_initial_state() {
assert!(VerifierPipelineHandle::initial_state_from_value(Some("1")));
assert!(VerifierPipelineHandle::initial_state_from_value(Some("true")));
assert!(VerifierPipelineHandle::initial_state_from_value(Some("TRUE")));
assert!(VerifierPipelineHandle::initial_state_from_value(Some("on")));
assert!(VerifierPipelineHandle::initial_state_from_value(Some("ON")));
}
#[test]
fn value_off_disables_initial_state() {
assert!(!VerifierPipelineHandle::initial_state_from_value(Some("0")));
assert!(!VerifierPipelineHandle::initial_state_from_value(Some("false")));
assert!(!VerifierPipelineHandle::initial_state_from_value(Some("")));
assert!(!VerifierPipelineHandle::initial_state_from_value(Some("yes")));
assert!(!VerifierPipelineHandle::initial_state_from_value(Some("True")));
}
#[tokio::test]
async fn enable_disable_is_observable_and_idempotent() {
let h = VerifierPipelineHandle::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_verifier_returns_arc_clone() {
let h = VerifierPipelineHandle::new(true, Arc::new(NoopVerifier));
let snap1 = h.snapshot_verifier();
let snap2 = h.snapshot_verifier();
assert!(snap1.is_some());
assert!(snap2.is_some());
assert!(Arc::ptr_eq(snap1.as_ref().unwrap(), snap2.as_ref().unwrap()));
}
#[test]
fn set_verifier_replaces_handle_preserving_toggle() {
let h = VerifierPipelineHandle::with_noop();
h.enable();
assert_eq!(h.state_change_count(), 1);
assert!(h.is_enabled());
let replacement: Arc<dyn Verifier> = Arc::new(NoopVerifier);
h.set_verifier(replacement);
assert!(h.is_enabled(), "set_verifier must not flip the toggle");
assert_eq!(h.state_change_count(), 1, "set_verifier must not increment state changes");
}
#[test]
fn debug_impl_is_safe_to_format() {
let h = VerifierPipelineHandle::with_noop();
let dbg = format!("{:?}", h);
assert!(dbg.contains("enabled"));
assert!(dbg.contains("has_verifier"));
assert!(dbg.contains("has_storage"));
assert!(dbg.contains("state_change_count"));
}
#[tokio::test]
async fn enable_without_bound_storage_does_not_spawn_worker() {
let h = VerifierPipelineHandle::new(false, Arc::new(NoopVerifier));
h.enable();
assert!(h.is_enabled());
assert!(h.worker_handle.lock().is_none(), "no storage bound → no worker spawned (B-10)");
}
#[test]
fn new_initializes_all_six_fields() {
let h = VerifierPipelineHandle::new(false, Arc::new(NoopVerifier));
assert!(!h.is_enabled(), "enabled bit");
assert!(h.snapshot_verifier().is_some(), "verifier slot");
assert_eq!(h.state_change_count(), 0, "state_change_count");
assert!(h.storage.read().is_none(), "storage unbound by default");
assert!(h.worker_handle.lock().is_none(), "worker_handle None by default");
}
}