use std::sync::Arc;
use lunaris_consolidate::{ActRConsolidator, Consolidator, NoopConsolidator};
use lunaris_core::{LunarisError, StorageError};
use parking_lot::{Mutex, RwLock};
pub const ENABLED_ENV_VAR: &str = "LUNARIS_CONSOLIDATE_ENABLED";
pub const BACKEND_ENV_VAR: &str = "LUNARIS_CONSOLIDATOR_BACKEND";
static BACKEND_LOG_ONCE: std::sync::OnceLock<()> = std::sync::OnceLock::new();
pub struct ConsolidatorPipelineHandle {
enabled: RwLock<bool>,
consolidator: RwLock<Option<Arc<dyn Consolidator>>>,
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>>>,
scope_prefix: RwLock<Option<String>>,
}
impl ConsolidatorPipelineHandle {
pub fn new(initial_enabled: bool, consolidator: Arc<dyn Consolidator>) -> Self {
Self {
enabled: RwLock::new(initial_enabled),
consolidator: RwLock::new(Some(consolidator)),
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),
scope_prefix: 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(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!("consolidator_pipeline_enable_without_storage; worker not spawned");
return;
}
};
let consolidator = self
.snapshot_consolidator()
.unwrap_or_else(|| Arc::new(NoopConsolidator) as Arc<dyn Consolidator>);
let shutdown = self.shutdown.clone();
let source_prefix = self.scope_prefix.read().clone();
let handle = tokio::spawn(async move {
#[allow(deprecated)]
match lunaris_consolidate::run_consolidate_worker(
storage,
consolidator,
shutdown,
source_prefix,
)
.await
{
Ok(jh) => {
if let Err(e) = jh.await {
tracing::warn!(err = %e, "consolidator_pipeline_inner_worker_join_failed");
}
}
Err(e) => {
tracing::error!(err = %e, "consolidator_pipeline_worker_spawn_failed");
}
}
});
*wh = Some(handle);
}
pub fn enable(&self) {
*self.scope_prefix.write() = None;
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",
scope = tracing::field::Empty,
"consolidator_pipeline_state_changed"
);
drop(w);
self.spawn_worker_if_idle();
}
}
pub fn enable_for_scope(&self, prefix: impl Into<String>) {
let prefix: String = prefix.into();
if prefix.is_empty() {
tracing::warn!(
"consolidator_pipeline_enable_for_scope_empty_prefix; \
degrading to system-wide enable (T-12-02-05 hygiene signal)"
);
self.enable();
return;
}
let current = self.scope_prefix.read().clone();
let already_scoped_same = matches!(¤t, Some(p) if p == &prefix);
let is_enabled = *self.enabled.read();
if already_scoped_same && is_enabled {
return;
}
let needs_rotation = is_enabled && matches!(¤t, Some(p) if p != &prefix);
if needs_rotation {
self.disable();
}
*self.scope_prefix.write() = Some(prefix.clone());
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",
scope = %prefix,
"consolidator_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", "consolidator_pipeline_state_changed");
drop(w);
self.shutdown.notify_one();
}
*self.scope_prefix.write() = None;
}
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, "consolidator_pipeline_worker_join_failed");
}
}
pub fn is_enabled(&self) -> bool {
*self.enabled.read()
}
pub fn scope_prefix(&self) -> Option<String> {
self.scope_prefix.read().clone()
}
pub fn state_change_count(&self) -> u64 {
self.state_change_count.load(std::sync::atomic::Ordering::SeqCst)
}
pub fn set_consolidator(&self, consolidator: Arc<dyn Consolidator>) {
*self.consolidator.write() = Some(consolidator);
tracing::info!("consolidator_pipeline_consolidator_replaced");
}
pub fn snapshot_consolidator(&self) -> Option<Arc<dyn Consolidator>> {
self.consolidator.read().clone()
}
pub fn with_noop() -> Self {
Self::new(false, Arc::new(NoopConsolidator) as Arc<dyn Consolidator>)
}
pub fn with_actr() -> Self {
Self::new(false, Arc::new(ActRConsolidator::default()) as Arc<dyn Consolidator>)
}
pub fn backend_from_env() -> Result<Arc<dyn Consolidator>, LunarisError> {
let raw = std::env::var(BACKEND_ENV_VAR).ok();
let trimmed = raw.as_deref().map(|v| v.trim());
let (backend, name): (Arc<dyn Consolidator>, &'static str) = match trimmed {
None | Some("") => {
(Arc::new(ActRConsolidator::default()) as Arc<dyn Consolidator>, "ActRConsolidator")
}
Some(v) if v.eq_ignore_ascii_case("actr") => {
(Arc::new(ActRConsolidator::default()) as Arc<dyn Consolidator>, "ActRConsolidator")
}
Some(v) if v.eq_ignore_ascii_case("noop") => {
(Arc::new(NoopConsolidator) as Arc<dyn Consolidator>, "NoopConsolidator")
}
Some(other) => {
return Err(LunarisError::Storage(StorageError::Backend(format!(
"{BACKEND_ENV_VAR}={other:?} is not one of [actr, noop]; \
unset for the ActR default, or set to \"noop\" to pin the \
no-op backend. CONSOL-V1-01."
))));
}
};
BACKEND_LOG_ONCE.get_or_init(|| {
tracing::info!(
target: "lunaris::consolidator",
consolidator_backend = name,
"consolidator_backend_resolved"
);
});
Ok(backend)
}
}
impl std::fmt::Debug for ConsolidatorPipelineHandle {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("ConsolidatorPipelineHandle")
.field("enabled", &*self.enabled.read())
.field("has_consolidator", &self.consolidator.read().is_some())
.field("has_storage", &self.storage.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!(!ConsolidatorPipelineHandle::initial_state_from_value(None));
}
#[test]
fn value_one_enables_initial_state() {
assert!(ConsolidatorPipelineHandle::initial_state_from_value(Some("1")));
assert!(ConsolidatorPipelineHandle::initial_state_from_value(Some("true")));
assert!(ConsolidatorPipelineHandle::initial_state_from_value(Some("TRUE")));
assert!(ConsolidatorPipelineHandle::initial_state_from_value(Some("on")));
assert!(ConsolidatorPipelineHandle::initial_state_from_value(Some("ON")));
}
#[test]
fn value_off_disables_initial_state() {
assert!(!ConsolidatorPipelineHandle::initial_state_from_value(Some("0")));
assert!(!ConsolidatorPipelineHandle::initial_state_from_value(Some("false")));
assert!(!ConsolidatorPipelineHandle::initial_state_from_value(Some("")));
assert!(!ConsolidatorPipelineHandle::initial_state_from_value(Some("yes")));
assert!(!ConsolidatorPipelineHandle::initial_state_from_value(Some("True")));
}
#[tokio::test]
async fn enable_disable_is_observable_and_idempotent() {
let h = ConsolidatorPipelineHandle::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_consolidator_returns_arc_clone() {
let h = ConsolidatorPipelineHandle::new(true, Arc::new(NoopConsolidator));
let snap1 = h.snapshot_consolidator();
let snap2 = h.snapshot_consolidator();
assert!(snap1.is_some());
assert!(snap2.is_some());
assert!(Arc::ptr_eq(snap1.as_ref().unwrap(), snap2.as_ref().unwrap()));
}
#[test]
fn set_consolidator_replaces_handle_preserving_toggle() {
let h = ConsolidatorPipelineHandle::with_noop();
h.enable();
assert_eq!(h.state_change_count(), 1);
assert!(h.is_enabled());
let replacement: Arc<dyn Consolidator> = Arc::new(NoopConsolidator);
h.set_consolidator(replacement);
assert!(h.is_enabled(), "set_consolidator must not flip the toggle");
assert_eq!(h.state_change_count(), 1, "set_consolidator must not increment state changes");
}
#[test]
fn debug_impl_is_safe_to_format() {
let h = ConsolidatorPipelineHandle::with_noop();
let dbg = format!("{:?}", h);
assert!(dbg.contains("enabled"));
assert!(dbg.contains("has_consolidator"));
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 = ConsolidatorPipelineHandle::new(false, Arc::new(NoopConsolidator));
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 = ConsolidatorPipelineHandle::new(false, Arc::new(NoopConsolidator));
assert!(!h.is_enabled(), "enabled bit");
assert!(h.snapshot_consolidator().is_some(), "consolidator 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");
assert!(h.scope_prefix().is_none(), "scope_prefix None by default (v0.1.0 parity)");
}
#[tokio::test]
async fn enable_for_scope_sets_prefix() {
let h = ConsolidatorPipelineHandle::with_noop();
assert!(!h.is_enabled());
assert!(h.scope_prefix().is_none());
h.enable_for_scope("helios:fs/");
assert!(h.is_enabled(), "enable_for_scope turns the pipeline ON");
assert_eq!(
h.scope_prefix(),
Some("helios:fs/".to_string()),
"scope_prefix stored verbatim"
);
assert_eq!(h.state_change_count(), 1, "one state transition");
}
#[tokio::test]
async fn enable_clears_scope_prefix() {
let h = ConsolidatorPipelineHandle::with_noop();
h.enable_for_scope("helios:fs/");
assert_eq!(h.scope_prefix(), Some("helios:fs/".to_string()));
h.disable();
h.enable();
assert!(h.scope_prefix().is_none(), "enable() (system-wide) clears any stale scope_prefix");
}
#[tokio::test]
async fn enable_for_scope_idempotent_same_prefix() {
let h = ConsolidatorPipelineHandle::with_noop();
h.enable_for_scope("helios:fs/");
assert_eq!(h.state_change_count(), 1);
h.enable_for_scope("helios:fs/");
assert_eq!(
h.state_change_count(),
1,
"idempotent: same prefix must not bump state_change_count"
);
assert_eq!(h.scope_prefix(), Some("helios:fs/".to_string()));
}
#[tokio::test]
async fn enable_for_scope_rotate_prefix_reconfigures_worker() {
let h = ConsolidatorPipelineHandle::with_noop();
h.enable_for_scope("helios:fs/");
assert_eq!(h.state_change_count(), 1);
assert_eq!(h.scope_prefix(), Some("helios:fs/".to_string()));
h.enable_for_scope("other:tenant/");
assert_eq!(
h.state_change_count(),
3,
"rotation cycles the worker: disable + re-enable bumps counter by 2"
);
assert_eq!(h.scope_prefix(), Some("other:tenant/".to_string()));
assert!(h.is_enabled());
}
#[tokio::test]
async fn enable_for_scope_empty_prefix_degrades_to_system_wide() {
let h = ConsolidatorPipelineHandle::with_noop();
h.enable_for_scope("");
assert!(h.is_enabled());
assert!(
h.scope_prefix().is_none(),
"empty prefix MUST degrade to system-wide (scope_prefix=None), \
not leak a `Some(\"\")` value that would match everything"
);
}
#[tokio::test]
async fn disable_clears_scope_prefix() {
let h = ConsolidatorPipelineHandle::with_noop();
h.enable_for_scope("helios:fs/");
assert_eq!(h.scope_prefix(), Some("helios:fs/".to_string()));
h.disable();
assert!(h.scope_prefix().is_none(), "disable MUST clear scope_prefix");
}
}