use crate::calendar::CalendarClient;
use crate::config::RedFolderConfig;
use crate::engine::BlackoutEngine;
use crate::error::Result;
use crate::events::{EventListener, RedFolderEvent};
use crate::types::{BlackoutNotification, BlackoutWindow};
use chrono::{DateTime, Duration, NaiveDate, Utc};
use serde::{Deserialize, Serialize};
use std::collections::HashMap;
use std::path::PathBuf;
use std::sync::Arc;
use tokio::sync::{broadcast, mpsc, Mutex, Notify, RwLock};
use tokio_util::sync::CancellationToken;
use tracing::{debug, info, warn};
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
pub enum ServiceState {
Stopped,
Starting,
Running,
Stopping,
}
struct WorkerState {
config: RedFolderConfig,
legacy_sender: Option<mpsc::UnboundedSender<BlackoutNotification>>,
event_sender: Option<mpsc::UnboundedSender<RedFolderEvent>>,
in_blackout: bool,
active_window: Option<BlackoutWindow>,
last_warned_window_start: Option<DateTime<Utc>>,
}
struct ServiceSyncMeta {
client: CalendarClient,
check_interval: std::time::Duration,
state: ServiceState,
last_fetch_date: Option<NaiveDate>,
last_sync_time: Option<DateTime<Utc>>,
last_sync_error: Option<String>,
}
#[derive(Clone)]
pub struct RedFolderService {
engine: Arc<RwLock<BlackoutEngine>>,
workers: Arc<RwLock<HashMap<String, WorkerState>>>,
listeners: Arc<RwLock<Vec<Arc<dyn EventListener>>>>,
sync_meta: Arc<Mutex<ServiceSyncMeta>>,
refresh_lock: Arc<Mutex<()>>,
broadcast_tx: broadcast::Sender<RedFolderEvent>,
cancel_token: Arc<Mutex<Option<CancellationToken>>>,
task_handles: Arc<Mutex<Vec<tokio::task::JoinHandle<()>>>>,
notify: Arc<Notify>,
}
pub const MIN_CHECK_INTERVAL: std::time::Duration = std::time::Duration::from_millis(100);
impl RedFolderService {
pub fn new(cache_dir: Option<PathBuf>) -> Self {
Self::with_client(CalendarClient::new(cache_dir))
}
pub fn with_client(client: CalendarClient) -> Self {
let (broadcast_tx, _) = broadcast::channel(256);
let notify = Arc::new(Notify::new());
Self {
engine: Arc::new(RwLock::new(BlackoutEngine::new())),
workers: Arc::new(RwLock::new(HashMap::new())),
listeners: Arc::new(RwLock::new(Vec::new())),
sync_meta: Arc::new(Mutex::new(ServiceSyncMeta {
client,
check_interval: std::time::Duration::from_secs(15),
state: ServiceState::Stopped,
last_fetch_date: None,
last_sync_time: None,
last_sync_error: None,
})),
refresh_lock: Arc::new(Mutex::new(())),
broadcast_tx,
cancel_token: Arc::new(Mutex::new(None)),
task_handles: Arc::new(Mutex::new(Vec::new())),
notify,
}
}
pub async fn state(&self) -> ServiceState {
self.sync_meta.lock().await.state
}
pub fn subscribe(&self) -> broadcast::Receiver<RedFolderEvent> {
self.broadcast_tx.subscribe()
}
pub async fn add_listener(&self, listener: Arc<dyn EventListener>) {
self.listeners.write().await.push(listener);
}
pub async fn set_check_interval(&self, interval: std::time::Duration) -> Result<()> {
if interval < MIN_CHECK_INTERVAL {
return Err(crate::error::RedFolderError::Config(format!(
"check interval must be at least {:?} (got {:?})",
MIN_CHECK_INTERVAL, interval
)));
}
self.sync_meta.lock().await.check_interval = interval;
self.notify.notify_waiters();
Ok(())
}
pub async fn evaluate_and_notify(&self) {
self.check_and_notify_workers().await;
}
pub async fn set_engine(&self, engine: BlackoutEngine) {
*self.engine.write().await = engine;
self.notify.notify_waiters();
}
pub async fn is_calendar_stale(&self) -> bool {
self.check_calendar_staleness().await
}
async fn check_calendar_staleness(&self) -> bool {
let (last_sync, max_age) = {
let meta = self.sync_meta.lock().await;
(meta.last_sync_time, meta.client.max_stale_cache_age())
};
let is_empty = self.engine.read().await.is_empty();
let is_stale = if let Some(last_sync) = last_sync {
if let Some(max_age) = max_age {
if let Ok(chrono_dur) = Duration::from_std(max_age) {
Utc::now() - last_sync > chrono_dur
} else {
false
}
} else {
false
}
} else {
is_empty
};
self.engine.write().await.set_stale(is_stale);
is_stale
}
pub async fn check_and_notify_workers(&self) -> Option<DateTime<Utc>> {
self.check_calendar_staleness().await;
let mut workers = self.workers.write().await;
if workers.is_empty() {
return None;
}
let engine = self.engine.read().await;
let now = Utc::now();
let mut events_to_dispatch: Vec<RedFolderEvent> = Vec::new();
let mut next_transition: Option<DateTime<Utc>> = None;
let update_next = |current: &mut Option<DateTime<Utc>>, candidate: DateTime<Utc>| {
if candidate > now {
match current {
Some(ref mut t) if candidate < *t => *t = candidate,
None => *current = Some(candidate),
_ => {}
}
}
};
for (worker_id, ws) in workers.iter_mut() {
if !ws.config.enabled {
continue;
}
let is_active = engine.is_blackout(&ws.config);
if is_active != ws.in_blackout {
ws.in_blackout = is_active;
let window = if is_active {
engine.current_window(&ws.config)
} else {
None
};
if let Some(ref sender) = ws.legacy_sender {
sender
.send(BlackoutNotification {
active: is_active,
window: window.clone(),
})
.ok();
}
if is_active {
ws.active_window = window.clone();
if let Some(w) = window {
let end_str = w.end.format("%H:%M UTC").to_string();
warn!(worker=%worker_id, until=%end_str, "ENTERING news blackout window");
let ev = RedFolderEvent::BlackoutStarted {
window: w,
worker_id: Some(worker_id.clone()),
};
if let Some(ref sender) = ws.event_sender {
sender.send(ev.clone()).ok();
}
events_to_dispatch.push(ev);
}
} else {
info!(worker=%worker_id, "EXITING news blackout window");
let ended_window = ws.active_window.take().unwrap_or_else(|| BlackoutWindow {
start: now,
end: now,
events: vec![],
});
let ev = RedFolderEvent::BlackoutEnded {
window: ended_window,
worker_id: Some(worker_id.clone()),
};
if let Some(ref sender) = ws.event_sender {
sender.send(ev.clone()).ok();
}
events_to_dispatch.push(ev);
}
}
if ws.in_blackout {
if let Some(ref w) = ws.active_window {
update_next(&mut next_transition, w.end);
}
} else {
let upcoming = engine.upcoming_blackouts(&ws.config, 24);
if let Some(next_window) = upcoming.first() {
update_next(&mut next_transition, next_window.start);
if let Some(warn_min) = ws.config.warning_before_min {
let warn_time = next_window.start - Duration::minutes(warn_min);
update_next(&mut next_transition, warn_time);
let mins_until_start = (next_window.start - now).num_minutes();
if mins_until_start > 0 && mins_until_start <= warn_min {
let already_warned = ws
.last_warned_window_start
.map(|t| t == next_window.start)
.unwrap_or(false);
if !already_warned {
ws.last_warned_window_start = Some(next_window.start);
let ev = RedFolderEvent::BlackoutWarning {
window: next_window.clone(),
minutes_until_start: mins_until_start,
worker_id: Some(worker_id.clone()),
};
if let Some(ref sender) = ws.event_sender {
sender.send(ev.clone()).ok();
}
events_to_dispatch.push(ev);
}
}
}
}
}
}
drop(workers);
drop(engine);
for ev in events_to_dispatch {
self.dispatch_event(ev).await;
}
next_transition
}
async fn dispatch_event(&self, event: RedFolderEvent) {
self.broadcast_tx.send(event.clone()).ok();
let listeners_guard = self.listeners.read().await;
for listener in listeners_guard.iter() {
let listener_clone = listener.clone();
let ev_clone = event.clone();
tokio::spawn(async move {
let _ = tokio::time::timeout(
std::time::Duration::from_secs(10),
listener_clone.on_event(&ev_clone),
)
.await;
});
}
}
async fn register_worker_internal<R>(
&self,
worker_id: impl Into<String>,
config: RedFolderConfig,
allow_overwrite: bool,
reregister_method_hint: &str,
setup: impl FnOnce(&str, &RedFolderConfig, bool, Option<BlackoutWindow>) -> (WorkerState, R),
) -> Result<R> {
config.validate()?;
let worker_id = worker_id.into();
let mut workers = self.workers.write().await;
if workers.contains_key(&worker_id) {
if allow_overwrite {
warn!(worker=%worker_id, "re-registering worker: overwriting previous worker state and channels");
} else {
return Err(crate::error::RedFolderError::Service(format!(
"worker '{worker_id}' is already registered; use {reregister_method_hint} to update or unregister first"
)));
}
}
let engine = self.engine.read().await;
let is_active = engine.is_blackout(&config);
let active_window = if is_active {
engine.current_window(&config)
} else {
None
};
drop(engine);
let (state, rx) = setup(&worker_id, &config, is_active, active_window);
workers.insert(worker_id.clone(), state);
drop(workers);
self.notify.notify_waiters();
if allow_overwrite {
debug!(worker=%worker_id, "re-registered worker in RedFolderService");
} else {
debug!(worker=%worker_id, "registered worker in RedFolderService");
}
Ok(rx)
}
pub async fn register_worker(
&self,
worker_id: impl Into<String>,
config: RedFolderConfig,
) -> Result<mpsc::UnboundedReceiver<BlackoutNotification>> {
self.register_worker_internal(
worker_id,
config,
false,
"reregister_worker",
|_, cfg, is_active, active_window| {
let (tx, rx) = mpsc::unbounded_channel();
if is_active {
tx.send(BlackoutNotification {
active: true,
window: active_window.clone(),
})
.ok();
}
let state = WorkerState {
config: cfg.clone(),
legacy_sender: Some(tx),
event_sender: None,
in_blackout: is_active,
active_window,
last_warned_window_start: None,
};
(state, rx)
},
)
.await
}
pub async fn register_worker_events(
&self,
worker_id: impl Into<String>,
config: RedFolderConfig,
) -> Result<mpsc::UnboundedReceiver<RedFolderEvent>> {
self.register_worker_internal(
worker_id,
config,
false,
"reregister_worker_events",
|wid, cfg, is_active, active_window| {
let (tx, rx) = mpsc::unbounded_channel();
if is_active {
if let Some(ref w) = active_window {
tx.send(RedFolderEvent::BlackoutStarted {
window: w.clone(),
worker_id: Some(wid.to_string()),
})
.ok();
}
}
let state = WorkerState {
config: cfg.clone(),
legacy_sender: None,
event_sender: Some(tx),
in_blackout: is_active,
active_window,
last_warned_window_start: None,
};
(state, rx)
},
)
.await
}
pub async fn reregister_worker(
&self,
worker_id: impl Into<String>,
config: RedFolderConfig,
) -> Result<mpsc::UnboundedReceiver<BlackoutNotification>> {
self.register_worker_internal(
worker_id,
config,
true,
"reregister_worker",
|_, cfg, is_active, active_window| {
let (tx, rx) = mpsc::unbounded_channel();
if is_active {
tx.send(BlackoutNotification {
active: true,
window: active_window.clone(),
})
.ok();
}
let state = WorkerState {
config: cfg.clone(),
legacy_sender: Some(tx),
event_sender: None,
in_blackout: is_active,
active_window,
last_warned_window_start: None,
};
(state, rx)
},
)
.await
}
pub async fn reregister_worker_events(
&self,
worker_id: impl Into<String>,
config: RedFolderConfig,
) -> Result<mpsc::UnboundedReceiver<RedFolderEvent>> {
self.register_worker_internal(
worker_id,
config,
true,
"reregister_worker_events",
|wid, cfg, is_active, active_window| {
let (tx, rx) = mpsc::unbounded_channel();
if is_active {
if let Some(ref w) = active_window {
tx.send(RedFolderEvent::BlackoutStarted {
window: w.clone(),
worker_id: Some(wid.to_string()),
})
.ok();
}
}
let state = WorkerState {
config: cfg.clone(),
legacy_sender: None,
event_sender: Some(tx),
in_blackout: is_active,
active_window,
last_warned_window_start: None,
};
(state, rx)
},
)
.await
}
pub async fn unregister_worker(&self, worker_id: &str) {
self.workers.write().await.remove(worker_id);
self.notify.notify_waiters();
}
pub async fn windows_for_worker(&self, worker_id: &str) -> Vec<BlackoutWindow> {
let cfg = {
let workers = self.workers.read().await;
workers.get(worker_id).map(|w| w.config.clone())
};
if let Some(config) = cfg {
let engine = self.engine.read().await;
engine.windows_for_config(&config, Utc::now())
} else {
Vec::new()
}
}
pub async fn is_running(&self) -> bool {
self.sync_meta.lock().await.state == ServiceState::Running
}
pub async fn start(&self) -> Result<()> {
{
let mut meta = self.sync_meta.lock().await;
if meta.state == ServiceState::Running || meta.state == ServiceState::Starting {
return Err(crate::error::RedFolderError::Service(
"service is already running".to_string(),
));
}
let workers = self.workers.read().await;
if workers.is_empty() {
warn!("no workers registered — RedFolderService not starting");
return Ok(());
}
if !workers.values().any(|w| w.config.enabled) {
info!("all registered worker blackout configs are disabled");
return Ok(());
}
meta.state = ServiceState::Starting;
}
if let Err(e) = self.refresh_internal(false, false).await {
let mut meta = self.sync_meta.lock().await;
meta.state = ServiceState::Stopped;
return Err(e);
}
let cancel_token = CancellationToken::new();
*self.cancel_token.lock().await = Some(cancel_token.clone());
let mut handles = Vec::new();
{
let this = self.clone();
let token = cancel_token.child_token();
let handle = tokio::spawn(async move {
loop {
let now = Utc::now();
let next_midnight = (now + Duration::days(1))
.date_naive()
.and_hms_opt(0, 0, 0)
.unwrap()
.and_utc();
let sleep_secs = (next_midnight - now).num_seconds().max(1) as u64;
tokio::select! {
_ = tokio::time::sleep(tokio::time::Duration::from_secs(sleep_secs)) => {
let res = tokio::select! {
res = this.refresh_internal(false, true) => res,
_ = token.cancelled() => {
break;
}
};
if let Err(ref e) = res {
warn!(err = %e, "daily midnight calendar sync failed; will retry every 15 minutes");
loop {
tokio::select! {
_ = tokio::time::sleep(tokio::time::Duration::from_secs(900)) => {
let retry_res = this.refresh_internal(false, true).await;
if retry_res.is_ok() {
info!("calendar sync successfully recovered after retry");
break;
} else if let Err(err) = retry_res {
warn!(err = %err, "calendar sync retry failed; will retry in 15 minutes");
}
}
_ = token.cancelled() => {
break;
}
}
}
}
}
_ = token.cancelled() => {
break;
}
}
}
});
handles.push(handle);
}
{
let this = self.clone();
let token = cancel_token.child_token();
let notify = self.notify.clone();
let handle = tokio::spawn(async move {
loop {
let next_transition = this.check_and_notify_workers().await;
let check_interval = this.sync_meta.lock().await.check_interval;
let now = Utc::now();
let sleep_duration = if let Some(target) = next_transition {
let dur = (target - now).to_std().unwrap_or(std::time::Duration::ZERO);
dur.min(check_interval)
} else {
check_interval
};
tokio::select! {
_ = tokio::time::sleep(sleep_duration) => {}
_ = notify.notified() => {}
_ = token.cancelled() => {
break;
}
}
}
});
handles.push(handle);
}
*self.task_handles.lock().await = handles;
{
let mut meta = self.sync_meta.lock().await;
meta.state = ServiceState::Running;
}
info!("RedFolderService background tasks started successfully");
Ok(())
}
pub async fn stop(&self) {
{
let mut meta = self.sync_meta.lock().await;
if meta.state == ServiceState::Stopped || meta.state == ServiceState::Stopping {
return;
}
meta.state = ServiceState::Stopping;
}
info!("stopping RedFolderService");
if let Some(token) = self.cancel_token.lock().await.take() {
token.cancel();
}
self.notify.notify_waiters();
let mut handles: Vec<_> = self.task_handles.lock().await.drain(..).collect();
for handle in &mut handles {
let res = tokio::time::timeout(std::time::Duration::from_secs(3), &mut *handle).await;
if res.is_err() {
handle.abort();
let _ = handle.await;
}
}
{
let mut meta = self.sync_meta.lock().await;
meta.state = ServiceState::Stopped;
}
}
pub async fn refresh(&self) -> Result<()> {
self.refresh_internal(false, false).await
}
pub async fn force_refresh(&self) -> Result<()> {
self.refresh_internal(true, false).await
}
async fn refresh_internal(
&self,
force_remote: bool,
skip_if_already_fetched_today: bool,
) -> Result<()> {
let _refresh_guard = self.refresh_lock.lock().await;
let today = Utc::now().date_naive();
if skip_if_already_fetched_today {
let meta = self.sync_meta.lock().await;
let engine_empty = self.engine.read().await.is_empty();
if meta.last_fetch_date == Some(today) && !engine_empty {
debug!("calendar already fetched and compiled for today");
return Ok(());
}
}
let (client, client_tz) = {
let meta = self.sync_meta.lock().await;
(meta.client.clone(), meta.client.calendar_timezone())
};
let fetch_res = if force_remote {
client.force_fetch().await
} else {
client.fetch_or_cached().await
};
let raw = match fetch_res {
Ok(raw) => raw,
Err(e) => {
let err_str = e.to_string();
{
let mut meta = self.sync_meta.lock().await;
meta.last_sync_error = Some(err_str.clone());
}
self.check_calendar_staleness().await;
self.check_and_notify_workers().await;
self.notify.notify_waiters();
self.dispatch_event(RedFolderEvent::CalendarSyncFailed { error: err_str })
.await;
return Err(e);
}
};
let configs: Vec<RedFolderConfig> = {
let workers = self.workers.read().await;
workers
.values()
.filter(|w| w.config.enabled)
.map(|w| w.config.clone())
.collect()
};
let config_refs: Vec<&RedFolderConfig> = configs.iter().collect();
let new_engine = BlackoutEngine::compile_with_tz(&raw, &config_refs, Utc::now(), client_tz);
let total_windows = new_engine.windows().len();
{
let mut engine = self.engine.write().await;
*engine = new_engine;
}
{
let mut meta = self.sync_meta.lock().await;
meta.last_fetch_date = Some(today);
meta.last_sync_time = Some(Utc::now());
meta.last_sync_error = None;
}
self.check_and_notify_workers().await;
self.notify.notify_waiters();
info!(windows=%total_windows, "refreshed economic calendar windows");
self.dispatch_event(RedFolderEvent::CalendarUpdated {
total_events: raw.len(),
total_windows,
})
.await;
Ok(())
}
pub async fn last_sync_time(&self) -> Option<DateTime<Utc>> {
self.sync_meta.lock().await.last_sync_time
}
pub async fn last_sync_error(&self) -> Option<String> {
self.sync_meta.lock().await.last_sync_error.clone()
}
pub async fn is_blackout(&self, worker_id: &str) -> bool {
self.check_calendar_staleness().await;
let cfg = {
let workers = self.workers.read().await;
workers.get(worker_id).map(|w| w.config.clone())
};
if let Some(config) = cfg {
let engine = self.engine.read().await;
engine.is_blackout(&config)
} else {
false
}
}
pub async fn current_window(&self, worker_id: &str) -> Option<BlackoutWindow> {
self.check_calendar_staleness().await;
let cfg = {
let workers = self.workers.read().await;
workers.get(worker_id).map(|w| w.config.clone())
};
if let Some(config) = cfg {
let engine = self.engine.read().await;
engine.current_window(&config)
} else {
None
}
}
pub async fn get_upcoming_blackouts(&self, worker_id: &str, hours: u32) -> Vec<BlackoutWindow> {
let cfg = {
let workers = self.workers.read().await;
workers.get(worker_id).map(|w| w.config.clone())
};
if let Some(config) = cfg {
let engine = self.engine.read().await;
engine.upcoming_blackouts(&config, hours)
} else {
Vec::new()
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
async fn test_event_bus_and_warning_dispatch() {
let service = RedFolderService::new(None);
let mut broadcast_rx = service.subscribe();
let config = RedFolderConfig::builder()
.currencies(vec!["USD"])
.impacts(vec!["High"])
.buffer_minutes(0, 15)
.warning_minutes(10)
.build();
let mut worker_events = service
.register_worker_events("worker_usd", config)
.await
.unwrap();
let now = Utc::now();
let raw = vec![crate::calendar::RawCalendarEvent {
title: "US CPI Release".into(),
country: "USD".into(),
date: (now + Duration::minutes(5)).to_rfc3339(),
time: "".into(),
impact: "High".into(),
}];
let cfg = RedFolderConfig {
weekend_enabled: false,
before_min: 0,
after_min: 15,
..Default::default()
};
service
.set_engine(BlackoutEngine::compile(&raw, &[&cfg], now))
.await;
service.evaluate_and_notify().await;
let ev = worker_events
.try_recv()
.expect("should receive warning event");
assert!(ev.is_warning());
assert_eq!(ev.worker_id(), Some("worker_usd"));
let bus_ev = broadcast_rx
.try_recv()
.expect("broadcast should receive warning");
assert!(bus_ev.is_warning());
}
#[tokio::test]
async fn test_set_check_interval_validation() {
let service = RedFolderService::new(None);
let res_zero = service.set_check_interval(std::time::Duration::ZERO).await;
assert!(res_zero.is_err(), "zero interval must be rejected");
assert!(res_zero
.unwrap_err()
.to_string()
.contains("must be at least"));
let res_small = service
.set_check_interval(std::time::Duration::from_millis(50))
.await;
assert!(res_small.is_err(), "sub-100ms interval must be rejected");
let res_valid = service
.set_check_interval(std::time::Duration::from_secs(5))
.await;
assert!(res_valid.is_ok(), "valid interval must be accepted");
}
#[tokio::test]
async fn test_register_worker_rejects_unvalidated_config() {
let service = RedFolderService::new(None);
let invalid_cfg = RedFolderConfig {
before_min: -10,
..Default::default()
};
let res_legacy = service
.register_worker("bad_worker_legacy", invalid_cfg.clone())
.await;
assert!(
res_legacy.is_err(),
"register_worker must reject unvalidated config with negative buffer"
);
let res_events = service
.register_worker_events("bad_worker_events", invalid_cfg)
.await;
assert!(
res_events.is_err(),
"register_worker_events must reject unvalidated config with negative buffer"
);
}
#[tokio::test]
async fn test_duplicate_worker_registration_rejected() {
let service = RedFolderService::new(None);
let cfg = RedFolderConfig::default();
let rx1 = service.register_worker("dup_worker", cfg.clone()).await;
assert!(rx1.is_ok());
let rx2 = service.register_worker("dup_worker", cfg.clone()).await;
assert!(rx2.is_err());
assert!(rx2.unwrap_err().to_string().contains("already registered"));
let rereg_rx = service.reregister_worker("dup_worker", cfg).await;
assert!(rereg_rx.is_ok());
}
#[tokio::test]
async fn test_calendar_sync_failed_event_dispatch() {
let unreachable_client = CalendarClient::with_options(
reqwest::Client::new(),
"http://127.0.0.1:9/unreachable",
None,
std::time::Duration::from_millis(50),
);
let service = RedFolderService::with_client(unreachable_client);
let mut sub = service.subscribe();
let cfg = RedFolderConfig::default();
let _rx = service.register_worker("test_bot", cfg).await.unwrap();
let refresh_res = service.refresh().await;
assert!(refresh_res.is_err());
let ev = sub
.try_recv()
.expect("must receive CalendarSyncFailed event");
assert!(ev.is_sync_failed());
let err = service.last_sync_error().await;
assert!(err.is_some());
}
#[tokio::test]
async fn test_fail_closed_worker_in_service() {
let service = RedFolderService::new(None);
let open_cfg = RedFolderConfig::builder()
.weekend_curfew(false, "20:00", "21:00", "short")
.fail_safe_mode(crate::types::FailSafeMode::FailOpen)
.build();
let closed_cfg = RedFolderConfig::builder()
.weekend_curfew(false, "20:00", "21:00", "short")
.fail_safe_mode(crate::types::FailSafeMode::FailClosed)
.build();
let _open_rx = service
.register_worker("open_worker", open_cfg)
.await
.unwrap();
let mut closed_rx = service
.register_worker("closed_worker", closed_cfg)
.await
.unwrap();
assert!(!service.is_blackout("open_worker").await);
assert!(service.is_blackout("closed_worker").await);
let win = service
.current_window("closed_worker")
.await
.expect("fail-closed worker must have active window");
assert!(win.summary_title().contains("Fail-Closed Safety Blackout"));
let notification = closed_rx
.try_recv()
.expect("must receive initial notification");
assert!(notification.active);
assert!(notification.window.is_some());
}
}