use parking_lot::{Mutex, RwLock};
use std::collections::{BTreeMap, HashMap};
use std::future::Future;
use std::pin::Pin;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;
use crate::effect::Disposable;
use crate::service::{CordisError, Service};
use crate::EventId;
#[derive(Debug)]
pub struct AggregateError {
pub errors: Vec<(String, String)>,
}
impl std::fmt::Display for AggregateError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{}", format_listener_errors(&self.errors))
}
}
impl std::error::Error for AggregateError {}
pub fn summarize_listener_errors(errors: Vec<(String, String)>) -> String {
format_listener_errors(&errors)
}
fn format_listener_errors(errors: &[(String, String)]) -> String {
if errors.is_empty() {
return "0 listener failures".to_string();
}
let joined = errors
.iter()
.map(|(name, message)| format!("{name}: {message}"))
.collect::<Vec<_>>()
.join("; ");
if errors.len() == 1 {
format!("1 listener failure: {joined}")
} else {
format!("{} listener failures: {joined}", errors.len())
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum Dispatch {
Emit,
Parallel,
Serial,
Bail,
Waterfall,
}
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
pub struct EventOptions {
pub prepend: bool,
pub global: bool,
}
pub type ListenerFilter = Box<dyn Fn(&EventOptions) -> bool + Send + Sync>;
pub const INTERNAL_GET_EVENT: &str = "internal/get";
pub const INTERNAL_SET_EVENT: &str = "internal/set";
pub const INTERNAL_CONFIG_EVENT: &str = "internal/config";
pub const INTERNAL_UPDATE_EVENT: &str = "internal/update";
pub const INTERNAL_LISTENER_EVENT: &str = "internal/listener";
pub const INTERNAL_DISPATCH_EVENT: &str = "internal/dispatch";
pub fn is_internal_meta_event(event: &str) -> bool {
matches!(
event,
INTERNAL_GET_EVENT
| INTERNAL_SET_EVENT
| INTERNAL_CONFIG_EVENT
| INTERNAL_UPDATE_EVENT
| INTERNAL_LISTENER_EVENT
| INTERNAL_DISPATCH_EVENT
)
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct InternalDispatchPayload {
pub mode: String,
pub name: String,
pub args: serde_json::Value,
}
const MECHANICS_TEST_EVENTS: &[&str] = &[
"test",
"test.event",
"gone",
"gone.wf",
"parallel.result",
"serial.bail",
"serial.identity",
"serial.test",
"bail.test",
"emit.test",
"emit.counter",
"wf.next",
"wf.short",
"wf.empty",
"par.test",
"par2.test",
"par.agg",
"par.solo",
"around.empty",
"around.wrap",
"around.short",
"once.test",
"once.bail",
"prepend.test",
"filtered.test",
"global.test",
"blocked.event",
"allowed.event",
"another.event",
"observed.a",
"observed.b",
"observed.c",
"wf.filtered",
];
fn bypasses_catalog(event: &str) -> bool {
(cfg!(test) && MECHANICS_TEST_EVENTS.contains(&event)) || is_internal_meta_event(event)
}
fn debug_enforce_dispatch(event: &EventId, mode: Dispatch) {
if bypasses_catalog(event) {
return;
}
if let Err(msg) = crate::events_catalog::validate_dispatch(event, mode) {
debug_assert!(false, "{msg}");
}
}
fn debug_enforce_listener(event: &EventId, waterfall_registration: bool) {
if bypasses_catalog(event) {
return;
}
if let Err(msg) = crate::events_catalog::validate_listener(event, waterfall_registration) {
debug_assert!(false, "{msg}");
}
}
impl std::fmt::Display for Dispatch {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
let name = match self {
Dispatch::Emit => "emit",
Dispatch::Parallel => "parallel",
Dispatch::Serial => "serial",
Dispatch::Bail => "bail",
Dispatch::Waterfall => "waterfall",
};
f.write_str(name)
}
}
type Handler = Arc<
dyn Fn(
serde_json::Value,
) -> Pin<Box<dyn Future<Output = Result<serde_json::Value, CordisError>> + Send>>
+ Send
+ Sync,
>;
pub type WaterfallNext = Box<
dyn FnOnce(
serde_json::Value,
)
-> Pin<Box<dyn Future<Output = Result<serde_json::Value, CordisError>> + Send>>
+ Send,
>;
type WaterfallHandler = Arc<
dyn Fn(
serde_json::Value,
WaterfallNext,
) -> Pin<Box<dyn Future<Output = Result<serde_json::Value, CordisError>> + Send>>
+ Send
+ Sync,
>;
#[derive(Clone)]
struct HandlerSlot {
cancelled: Arc<AtomicBool>,
options: EventOptions,
handler: Handler,
}
#[derive(Clone)]
struct WaterfallSlot {
cancelled: Arc<AtomicBool>,
options: EventOptions,
handler: WaterfallHandler,
}
pub struct EventsService {
handlers: RwLock<HashMap<EventId, Vec<HandlerSlot>>>,
waterfall_handlers: RwLock<HashMap<EventId, Vec<WaterfallSlot>>>,
bus: tokio::sync::broadcast::Sender<(EventId, serde_json::Value)>,
dispatch_counts: Mutex<BTreeMap<String, u64>>,
}
impl EventsService {
pub fn new() -> Self {
let (tx, _rx) = tokio::sync::broadcast::channel(32);
let svc = Self {
handlers: RwLock::new(HashMap::new()),
waterfall_handlers: RwLock::new(HashMap::new()),
bus: tx,
dispatch_counts: Mutex::new(BTreeMap::new()),
};
svc.register_default_admit_handler();
svc
}
fn register_default_admit_handler(&self) {
let cancelled = Arc::new(AtomicBool::new(false));
let slot = HandlerSlot {
cancelled,
options: EventOptions::default(),
handler: Arc::new(|payload| Box::pin(async move { Ok(default_agent_admit(payload)) })),
};
self.handlers
.write()
.entry("agent.admit".into())
.or_default()
.push(slot);
}
pub fn listener_count(&self, event: &str) -> usize {
let flat = self
.handlers
.read()
.get(event)
.map(|slots| {
slots
.iter()
.filter(|slot| !slot.cancelled.load(Ordering::SeqCst))
.count()
})
.unwrap_or(0);
let waterfall = self
.waterfall_handlers
.read()
.get(event)
.map(|slots| {
slots
.iter()
.filter(|slot| !slot.cancelled.load(Ordering::SeqCst))
.count()
})
.unwrap_or(0);
flat + waterfall
}
pub fn dispatch_snapshot(&self) -> (u64, Vec<(String, u64)>) {
let map = self.dispatch_counts.lock();
let total = map.values().sum();
(total, map.iter().map(|(k, v)| (k.clone(), *v)).collect())
}
pub fn subscribe(&self) -> tokio::sync::broadcast::Receiver<(EventId, serde_json::Value)> {
self.bus.subscribe()
}
pub fn once<F, Fut>(&self, event: EventId, handler: F) -> Box<dyn Disposable>
where
F: Fn(serde_json::Value) -> Fut + Send + Sync + 'static,
Fut: Future<Output = Result<serde_json::Value, CordisError>> + Send + 'static,
{
self.once_with(event, EventOptions::default(), handler)
}
pub fn once_with<F, Fut>(
&self,
event: EventId,
options: EventOptions,
handler: F,
) -> Box<dyn Disposable>
where
F: Fn(serde_json::Value) -> Fut + Send + Sync + 'static,
Fut: Future<Output = Result<serde_json::Value, CordisError>> + Send + 'static,
{
debug_enforce_listener(&event, false);
if !blocking_listener_veto(self, &event) {
return Box::new(|| {});
}
let claim = Arc::new(AtomicBool::new(false));
let slot_flag = claim.clone();
let handle_flag = claim.clone();
let user = Arc::new(handler);
let once_handler: Handler = {
let user = user.clone();
Arc::new(move |payload: serde_json::Value| {
let claimed = claim.swap(true, Ordering::SeqCst);
let user = user.clone();
Box::pin(async move {
if claimed {
return Ok(payload);
}
user(payload).await
})
})
};
let slot = HandlerSlot {
cancelled: slot_flag,
options,
handler: once_handler,
};
self.insert_handler(event, options.prepend, slot);
Box::new(move || {
handle_flag.store(true, Ordering::SeqCst);
})
}
pub fn on<F, Fut>(&self, event: EventId, handler: F) -> Box<dyn Disposable>
where
F: Fn(serde_json::Value) -> Fut + Send + Sync + 'static,
Fut: Future<Output = Result<serde_json::Value, CordisError>> + Send + 'static,
{
self.on_with(event, EventOptions::default(), handler)
}
pub fn on_with<F, Fut>(
&self,
event: EventId,
options: EventOptions,
handler: F,
) -> Box<dyn Disposable>
where
F: Fn(serde_json::Value) -> Fut + Send + Sync + 'static,
Fut: Future<Output = Result<serde_json::Value, CordisError>> + Send + 'static,
{
debug_enforce_listener(&event, false);
if !blocking_listener_veto(self, &event) {
return Box::new(|| {});
}
let cancelled = Arc::new(AtomicBool::new(false));
let slot = HandlerSlot {
cancelled: cancelled.clone(),
options,
handler: Arc::new(move |v| Box::pin(handler(v))),
};
self.insert_handler(event, options.prepend, slot);
Box::new(move || {
cancelled.store(true, Ordering::SeqCst);
})
}
fn insert_handler(&self, event: EventId, prepend: bool, slot: HandlerSlot) {
let mut handlers = self.handlers.write();
let entry = handlers.entry(event).or_default();
if prepend {
entry.insert(0, slot);
} else {
entry.push(slot);
}
}
pub fn on_waterfall<F, Fut>(&self, event: EventId, handler: F) -> Box<dyn Disposable>
where
F: Fn(serde_json::Value, WaterfallNext) -> Fut + Send + Sync + 'static,
Fut: Future<Output = Result<serde_json::Value, CordisError>> + Send + 'static,
{
debug_enforce_listener(&event, true);
if !blocking_listener_veto(self, &event) {
return Box::new(|| {});
}
let cancelled = Arc::new(AtomicBool::new(false));
let slot = WaterfallSlot {
cancelled: cancelled.clone(),
options: EventOptions::default(),
handler: Arc::new(move |v, next| Box::pin(handler(v, next))),
};
let mut handlers = self.waterfall_handlers.write();
let entry = handlers.entry(event).or_default();
entry.push(slot);
Box::new(move || {
cancelled.store(true, Ordering::SeqCst);
})
}
fn active_handlers_filtered(
&self,
event: &EventId,
filter: Option<&ListenerFilter>,
) -> Vec<Handler> {
let Some(filter) = filter else {
return self.active_handlers(event);
};
let mut handlers = self.handlers.write();
let active = {
let Some(slots) = handlers.get_mut(event) else {
return Vec::new();
};
slots.retain(|slot| !slot.cancelled.load(Ordering::SeqCst));
slots
.iter()
.filter(|slot| slot.options.global || filter(&slot.options))
.map(|slot| slot.handler.clone())
.collect::<Vec<_>>()
};
if active.is_empty() && handlers.get(event).is_some_and(Vec::is_empty) {
handlers.remove(event);
}
active
}
fn active_handlers(&self, event: &EventId) -> Vec<Handler> {
let mut handlers = self.handlers.write();
let active = {
let Some(slots) = handlers.get_mut(event) else {
return Vec::new();
};
slots.retain(|slot| !slot.cancelled.load(Ordering::SeqCst));
slots
.iter()
.map(|slot| slot.handler.clone())
.collect::<Vec<_>>()
};
if active.is_empty() {
handlers.remove(event);
}
active
}
fn active_waterfall(&self, event: &EventId) -> Vec<WaterfallHandler> {
self.active_waterfall_filtered(event, None)
}
fn active_waterfall_filtered(
&self,
event: &EventId,
filter: Option<&ListenerFilter>,
) -> Vec<WaterfallHandler> {
let mut handlers = self.waterfall_handlers.write();
let active = {
let Some(slots) = handlers.get_mut(event) else {
return Vec::new();
};
slots.retain(|slot| !slot.cancelled.load(Ordering::SeqCst));
slots
.iter()
.filter(|slot| match filter {
None => true,
Some(filter) => slot.options.global || filter(&slot.options),
})
.map(|slot| slot.handler.clone())
.collect::<Vec<_>>()
};
if active.is_empty() && handlers.get(event).is_some_and(Vec::is_empty) {
handlers.remove(event);
}
active
}
fn observe_dispatch(&self, mode: Dispatch, name: &EventId, args: &serde_json::Value) {
if is_internal_meta_event(name) || self.listener_count(INTERNAL_DISPATCH_EVENT) == 0 {
return;
}
let Ok(payload) = serde_json::to_value(InternalDispatchPayload {
mode: mode.to_string(),
name: name.clone(),
args: args.clone(),
}) else {
return;
};
for handler in self.active_handlers(&INTERNAL_DISPATCH_EVENT.to_string()) {
let p = payload.clone();
tokio::spawn(async move {
let _ = handler(p).await;
});
}
}
pub fn emit_filtered(
&self,
event: EventId,
args: serde_json::Value,
filter: ListenerFilter,
) -> Result<serde_json::Value, CordisError> {
debug_enforce_dispatch(&event, Dispatch::Emit);
*self
.dispatch_counts
.lock()
.entry(event.to_string())
.or_insert(0) += 1;
self.observe_dispatch(Dispatch::Emit, &event, &args);
let _ = self.bus.send((event.clone(), args.clone()));
for h in self.active_handlers_filtered(&event, Some(&filter)) {
let p = args.clone();
tokio::spawn(async move {
let _ = h(p).await;
});
}
Ok(serde_json::Value::Null)
}
pub async fn bail_from(
&self,
event: EventId,
payload: serde_json::Value,
filter: Option<ListenerFilter>,
) -> Result<serde_json::Value, CordisError> {
debug_enforce_dispatch(&event, Dispatch::Bail);
*self
.dispatch_counts
.lock()
.entry(event.to_string())
.or_insert(0) += 1;
let handlers = match filter {
Some(filter) => self.active_handlers_filtered(&event, Some(&filter)),
None => self.active_handlers(&event),
};
run_bail_handlers(handlers, payload).await
}
pub async fn waterfall_from(
&self,
event: EventId,
payload: serde_json::Value,
filter: Option<ListenerFilter>,
) -> Result<serde_json::Value, CordisError> {
debug_enforce_dispatch(&event, Dispatch::Waterfall);
*self
.dispatch_counts
.lock()
.entry(event.to_string())
.or_insert(0) += 1;
let handlers = self.active_waterfall_filtered(&event, filter.as_ref());
if handlers.is_empty() {
return Ok(payload);
}
run_waterfall_chain(handlers, 0, payload, None).await
}
pub async fn dispatch(
&self,
event: EventId,
payload: serde_json::Value,
mode: Dispatch,
) -> Result<serde_json::Value, CordisError> {
debug_enforce_dispatch(&event, mode);
*self
.dispatch_counts
.lock()
.entry(event.to_string())
.or_insert(0) += 1;
self.observe_dispatch(mode, &event, &payload);
let handlers = self.active_handlers(&event);
match mode {
Dispatch::Waterfall => {
let wf_handlers = self.active_waterfall(&event);
if wf_handlers.is_empty() {
return Ok(payload);
}
run_waterfall_chain(wf_handlers, 0, payload, None).await
}
Dispatch::Emit => {
let _ = self.bus.send((event, payload.clone()));
for h in handlers {
let p = payload.clone();
tokio::spawn(async move {
let _ = h(p).await;
});
}
Ok(serde_json::Value::Null)
}
Dispatch::Parallel => {
let mut set = tokio::task::JoinSet::new();
for (name, h) in handlers.into_iter().enumerate() {
let p = payload.clone();
set.spawn(async move { (name, h(p).await) });
}
let mut failures: Vec<(String, String)> = Vec::new();
while let Some(res) = set.join_next().await {
match res {
Err(join_err) => {
let message = match join_err.try_into_panic() {
Ok(payload) => panic_payload_message(&payload),
Err(cancelled) => cancelled.to_string(),
};
failures.push(("listener-task".to_string(), message));
}
Ok((name, Err(err))) => {
failures.push((format!("listener[{name}]"), err.message()))
}
Ok((_, Ok(_))) => {}
}
}
if failures.is_empty() {
return Ok(serde_json::Value::Null);
}
tracing::warn!(
event = %event,
failures = %format_listener_errors(&failures),
"parallel dispatch collected listener failures"
);
Err(CordisError::Internal(format_listener_errors(&failures)))
}
Dispatch::Serial | Dispatch::Bail => run_bail_handlers(handlers, payload).await,
}
}
pub async fn waterfall_around<F, Fut>(
&self,
event: EventId,
payload: serde_json::Value,
core: F,
) -> Result<serde_json::Value, CordisError>
where
F: FnOnce(serde_json::Value) -> Fut + Send + 'static,
Fut: Future<Output = Result<serde_json::Value, CordisError>> + Send + 'static,
{
debug_enforce_dispatch(&event, Dispatch::Waterfall);
self.observe_dispatch(Dispatch::Waterfall, &event, &payload);
let handlers = self.active_waterfall(&event);
if handlers.is_empty() {
return core(payload).await;
}
let core: WaterfallCore = Box::new(move |p| {
Box::pin(core(p))
as Pin<Box<dyn Future<Output = Result<serde_json::Value, CordisError>> + Send>>
});
run_waterfall_chain(handlers, 0, payload, Some(core)).await
}
pub async fn waterfall_async_from<F, Fut>(
&self,
event: EventId,
payload: serde_json::Value,
filter: Option<ListenerFilter>,
core: F,
) -> Result<serde_json::Value, CordisError>
where
F: FnOnce(serde_json::Value) -> Fut + Send + 'static,
Fut: Future<Output = Result<serde_json::Value, CordisError>> + Send + 'static,
{
debug_enforce_dispatch(&event, Dispatch::Waterfall);
*self
.dispatch_counts
.lock()
.entry(event.to_string())
.or_insert(0) += 1;
let handlers = self.active_waterfall_filtered(&event, filter.as_ref());
if handlers.is_empty() {
return core(payload).await;
}
let core: WaterfallCore = Box::new(move |p| {
Box::pin(core(p))
as Pin<Box<dyn Future<Output = Result<serde_json::Value, CordisError>> + Send>>
});
run_waterfall_chain(handlers, 0, payload, Some(core)).await
}
pub async fn intercept_get(
&self,
service: &str,
ctx_hint: Option<String>,
) -> Result<Option<serde_json::Value>, CordisError> {
if self.listener_count(INTERNAL_GET_EVENT) == 0 {
return Ok(None);
}
let payload = serde_json::json!({ "service": service, "ctx": ctx_hint });
let out = self
.bail_from(INTERNAL_GET_EVENT.into(), payload, None)
.await?;
Ok((!out.is_null()).then_some(out))
}
pub async fn intercept_set(
&self,
service: &str,
ctx_hint: Option<String>,
) -> Result<(), CordisError> {
if self.listener_count(INTERNAL_SET_EVENT) == 0 {
return Ok(());
}
let payload = serde_json::json!({ "service": service, "ctx": ctx_hint });
self.bail_from(INTERNAL_SET_EVENT.into(), payload, None)
.await?;
Ok(())
}
pub async fn intercept_config(
&self,
raw: serde_json::Value,
) -> Result<serde_json::Value, CordisError> {
if self.listener_count(INTERNAL_CONFIG_EVENT) == 0 {
return Ok(raw);
}
self.bail_from(INTERNAL_CONFIG_EVENT.into(), raw, None).await
}
pub async fn intercept_update(&self, service: &str) -> Result<bool, CordisError> {
if self.listener_count(INTERNAL_UPDATE_EVENT) == 0 {
return Ok(true);
}
let payload = serde_json::json!({ "service": service });
let out = self
.bail_from(INTERNAL_UPDATE_EVENT.into(), payload, None)
.await?;
Ok(!(out.is_null() || out.as_bool() == Some(false)))
}
pub async fn intercept_listener(&self, event: &str) -> Result<bool, CordisError> {
if self.listener_count(INTERNAL_LISTENER_EVENT) == 0 {
return Ok(true);
}
let payload = serde_json::json!({ "event": event });
let out = self
.bail_from(INTERNAL_LISTENER_EVENT.into(), payload, None)
.await?;
Ok(out.is_null() || out.as_bool() == Some(true))
}
pub async fn dispatch_typed<E: crate::events_payload::TypedEvent>(
&self,
payload: &E::Payload,
) -> Result<serde_json::Value, CordisError> {
let value =
serde_json::to_value(payload).map_err(|e| CordisError::Configuration(e.to_string()))?;
self.dispatch(E::NAME.to_string(), value, E::MODE).await
}
pub fn on_typed<E, F, Fut>(&self, handler: F) -> Box<dyn Disposable>
where
E: crate::events_payload::TypedEvent,
F: Fn(E::Payload) -> Fut + Send + Sync + 'static,
Fut: Future<Output = Result<serde_json::Value, CordisError>> + Send + 'static,
{
debug_enforce_listener(&E::NAME.to_string(), E::AROUND);
let wrapped = move |v: serde_json::Value| {
let fut = match serde_json::from_value::<E::Payload>(v.clone()) {
Ok(payload) => handler(payload),
Err(err) => {
tracing::warn!(event = E::NAME, error = %err, "typed listener skipped malformed payload");
return Box::pin(async { Ok(v) })
as Pin<
Box<dyn Future<Output = Result<serde_json::Value, CordisError>> + Send>,
>;
}
};
Box::pin(fut)
as Pin<Box<dyn Future<Output = Result<serde_json::Value, CordisError>> + Send>>
};
self.on(E::NAME.to_string(), wrapped)
}
pub fn on_typed_waterfall<E, F, Fut>(&self, handler: F) -> Box<dyn Disposable>
where
E: crate::events_payload::TypedEvent,
F: Fn(E::Payload, WaterfallNext) -> Fut + Send + Sync + 'static,
Fut: Future<Output = Result<serde_json::Value, CordisError>> + Send + 'static,
{
debug_enforce_listener(&E::NAME.to_string(), E::AROUND);
let wrapped = move |v: serde_json::Value, next: WaterfallNext| {
let fut = match serde_json::from_value::<E::Payload>(v.clone()) {
Ok(payload) => handler(payload, next),
Err(err) => {
tracing::warn!(event = E::NAME, error = %err, "typed listener skipped malformed payload");
return Box::pin(async move {
next(v).await
})
as Pin<
Box<dyn Future<Output = Result<serde_json::Value, CordisError>> + Send>,
>;
}
};
Box::pin(fut)
as Pin<Box<dyn Future<Output = Result<serde_json::Value, CordisError>> + Send>>
};
self.on_waterfall(E::NAME.to_string(), wrapped)
}
}
impl Default for EventsService {
fn default() -> Self {
Self::new()
}
}
struct InterceptFence;
impl InterceptFence {
fn enter() -> Option<Self> {
INTERCEPT_FENCE.with(|fence| {
if fence.get() {
None
} else {
fence.set(true);
Some(Self)
}
})
}
}
impl Drop for InterceptFence {
fn drop(&mut self) {
INTERCEPT_FENCE.with(|fence| fence.set(false));
}
}
thread_local! {
static INTERCEPT_FENCE: std::cell::Cell<bool> = const { std::cell::Cell::new(false) };
}
fn bridge_handle() -> Option<tokio::runtime::Handle> {
let handle = tokio::runtime::Handle::try_current().ok()?;
if handle.runtime_flavor() == tokio::runtime::RuntimeFlavor::MultiThread {
Some(handle)
} else {
None
}
}
fn blocking_listener_veto(svc: &EventsService, event: &str) -> bool {
if svc.listener_count(INTERNAL_LISTENER_EVENT) == 0 {
return true;
}
let Some(_fence) = InterceptFence::enter() else {
return true;
};
let Some(handle) = bridge_handle() else {
tracing::warn!(
event = %event,
"internal/listener veto listener present but runtime cannot block in place; allowing registration"
);
return true;
};
let svc: &'static EventsService = unsafe { &*(svc as *const EventsService) };
tokio::task::block_in_place(|| {
handle.block_on(async move {
svc.intercept_listener(event).await.unwrap_or(false)
})
})
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum ReadVerdict {
Pass,
RedirectFrame,
Refuse,
}
pub(crate) fn blocking_intercept_get(events: &EventsService, service: &str) -> ReadVerdict {
if events.listener_count(INTERNAL_GET_EVENT) == 0 {
return ReadVerdict::Pass;
}
let Some(_fence) = InterceptFence::enter() else {
return ReadVerdict::Pass;
};
let Some(handle) = bridge_handle() else {
tracing::warn!(
service,
"internal/get listener present but runtime cannot block in place; passing read through"
);
return ReadVerdict::Pass;
};
let events: &'static EventsService = unsafe { &*(events as *const EventsService) };
let service = service.to_string();
tokio::task::block_in_place(|| {
handle.block_on(async move {
match events.intercept_get(&service, None).await {
Err(_) => ReadVerdict::Refuse,
Ok(None) => ReadVerdict::Pass,
Ok(Some(out)) => {
if out.is_null() {
ReadVerdict::Pass
} else if out.get("refuse").and_then(|v| v.as_bool()) == Some(true) {
ReadVerdict::Refuse
} else {
ReadVerdict::RedirectFrame
}
}
}
})
})
}
pub(crate) fn blocking_intercept_set(events: &EventsService, service: &str) -> Result<(), CordisError> {
if events.listener_count(INTERNAL_SET_EVENT) == 0 {
return Ok(());
}
let Some(_fence) = InterceptFence::enter() else {
return Ok(());
};
let Some(handle) = bridge_handle() else {
tracing::warn!(
service,
"internal/set listener present but runtime cannot block in place; allowing write"
);
return Ok(());
};
let events: &'static EventsService = unsafe { &*(events as *const EventsService) };
let service = service.to_string();
tokio::task::block_in_place(|| {
handle.block_on(async move { events.intercept_set(&service, None).await })
})
}
pub(crate) fn blocking_intercept_config(
events: &EventsService,
raw: serde_json::Value,
) -> Result<serde_json::Value, CordisError> {
if events.listener_count(INTERNAL_CONFIG_EVENT) == 0 {
return Ok(raw);
}
let Some(_fence) = InterceptFence::enter() else {
return Ok(raw);
};
let Some(handle) = bridge_handle() else {
tracing::warn!(
"internal/config listener present but runtime cannot block in place; using raw config"
);
return Ok(raw);
};
let events: &'static EventsService = unsafe { &*(events as *const EventsService) };
tokio::task::block_in_place(|| handle.block_on(events.intercept_config(raw)))
}
fn panic_payload_message(payload: &(dyn std::any::Any + Send)) -> String {
if let Some(s) = payload.downcast_ref::<&str>() {
(*s).to_string()
} else if let Some(s) = payload.downcast_ref::<String>() {
s.clone()
} else {
"listener panicked".to_string()
}
}
fn json_u64(v: &serde_json::Value, key: &str) -> Option<u64> {
v.get(key).and_then(|x| {
x.as_u64()
.or_else(|| x.as_i64().and_then(|n| u64::try_from(n).ok()))
})
}
fn default_agent_admit(payload: serde_json::Value) -> serde_json::Value {
if payload.get("tier").and_then(|v| v.as_str()) == Some("enterprise") {
return serde_json::Value::Null;
}
let monthly = json_u64(&payload, "monthly").unwrap_or(0);
let daily = json_u64(&payload, "daily").unwrap_or(0);
let Some(rpm) = json_u64(&payload, "requests_per_month") else {
return serde_json::Value::Null;
};
let Some(rpd) = json_u64(&payload, "requests_per_day") else {
return serde_json::Value::Null;
};
if monthly >= rpm {
return serde_json::json!({ "deny": "monthly" });
}
if daily >= rpd {
return serde_json::json!({ "deny": "daily" });
}
serde_json::Value::Null
}
impl Service for EventsService {}
async fn run_bail_handlers(
handlers: Vec<Handler>,
payload: serde_json::Value,
) -> Result<serde_json::Value, CordisError> {
for handler in handlers {
let result = handler(payload.clone()).await?;
if !result.is_null() {
return Ok(result);
}
}
Ok(payload)
}
type WaterfallCore = Box<
dyn FnOnce(
serde_json::Value,
)
-> Pin<Box<dyn Future<Output = Result<serde_json::Value, CordisError>> + Send>>
+ Send,
>;
fn run_waterfall_chain(
handlers: Vec<WaterfallHandler>,
index: usize,
payload: serde_json::Value,
core: Option<WaterfallCore>,
) -> Pin<Box<dyn Future<Output = Result<serde_json::Value, CordisError>> + Send>> {
Box::pin(async move {
if index >= handlers.len() {
return match core {
Some(core) => core(payload).await,
None => Ok(payload),
};
}
let handler = handlers[index].clone();
let next = move |p: serde_json::Value| {
let remaining = handlers.clone();
Box::pin(async move { run_waterfall_chain(remaining, index + 1, p, core).await })
as Pin<Box<dyn Future<Output = Result<serde_json::Value, CordisError>> + Send>>
};
handler(payload, Box::new(next)).await
})
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::atomic::AtomicUsize;
#[tokio::test]
async fn on_dispose_unregisters_handler() {
let svc = EventsService::new();
let flag = Arc::new(AtomicBool::new(false));
let f = flag.clone();
let d = svc.on("gone".into(), move |_v| {
let f = f.clone();
async move {
f.store(true, Ordering::SeqCst);
Ok(serde_json::Value::Null)
}
});
d.dispose();
svc.dispatch("gone".into(), serde_json::json!({}), Dispatch::Emit)
.await
.unwrap();
svc.dispatch("gone".into(), serde_json::json!({}), Dispatch::Serial)
.await
.unwrap();
assert!(svc.handlers.read().get("gone").is_none());
tokio::time::sleep(std::time::Duration::from_millis(30)).await;
assert!(
!flag.load(Ordering::SeqCst),
"disposed on() handler must not run for Emit or Serial"
);
}
#[tokio::test]
async fn on_waterfall_dispose_unregisters_handler() {
let svc = EventsService::new();
let flag = Arc::new(AtomicBool::new(false));
let f = flag.clone();
let d = svc.on_waterfall("gone.wf".into(), move |payload, _next| {
let f = f.clone();
async move {
f.store(true, Ordering::SeqCst);
Ok(payload)
}
});
d.dispose();
svc.dispatch(
"gone.wf".into(),
serde_json::json!({ "n": 1 }),
Dispatch::Waterfall,
)
.await
.unwrap();
assert!(
!flag.load(Ordering::SeqCst),
"disposed on_waterfall handler must not run"
);
assert!(svc.waterfall_handlers.read().get("gone.wf").is_none());
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn once_fires_exactly_once_concurrently() {
let svc = std::sync::Arc::new(EventsService::new());
let runs = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let r = runs.clone();
svc.once("once.test".into(), move |payload| {
let r = r.clone();
async move {
r.fetch_add(1, Ordering::SeqCst);
Ok(payload)
}
});
let mut tasks = tokio::task::JoinSet::new();
for _ in 0..16 {
let svc = std::sync::Arc::clone(&svc);
tasks.spawn(async move {
svc.dispatch(
"once.test".into(),
serde_json::json!({}),
Dispatch::Parallel,
)
.await
});
}
while let Some(res) = tasks.join_next().await {
res.expect("dispatch task").expect("parallel dispatch ok");
}
assert_eq!(
runs.load(Ordering::SeqCst),
1,
"exactly one concurrent dispatch may run the once handler"
);
}
#[tokio::test]
async fn once_stays_registered_when_skipped_by_bail() {
let svc = EventsService::new();
let ran = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let first = ran.clone();
let bailer = svc.on("once.bail".into(), move |_payload| {
let first = first.clone();
async move {
first.fetch_add(1, Ordering::SeqCst);
Ok(serde_json::json!({ "handled": true }))
}
});
let second = ran.clone();
svc.once("once.bail".into(), move |payload| {
let second = second.clone();
async move {
second.fetch_add(1, Ordering::SeqCst);
Ok(payload)
}
});
let out = svc
.dispatch("once.bail".into(), serde_json::json!({}), Dispatch::Bail)
.await
.unwrap();
assert_eq!(out, serde_json::json!({ "handled": true }));
assert_eq!(ran.load(Ordering::SeqCst), 1, "only the bail handler ran");
assert!(svc.handlers.read().get("once.bail").is_some());
let _ = svc
.dispatch("once.bail".into(), serde_json::json!({}), Dispatch::Bail)
.await;
assert_eq!(
ran.load(Ordering::SeqCst),
2,
"the bail handler claims every chain"
);
assert!(
svc.handlers.read().get("once.bail").is_some(),
"skipped-by-bail once slot stays registered"
);
bailer.dispose();
let out = svc
.dispatch(
"once.bail".into(),
serde_json::json!({"n": 1}),
Dispatch::Bail,
)
.await
.unwrap();
assert_eq!(out, serde_json::json!({"n": 1}));
assert_eq!(
ran.load(Ordering::SeqCst),
3,
"the surviving once slot fires exactly once"
);
let _ = svc
.dispatch(
"once.bail".into(),
serde_json::json!({"n": 2}),
Dispatch::Bail,
)
.await;
assert_eq!(
ran.load(Ordering::SeqCst),
3,
"spent once slot must not run again"
);
assert!(
svc.handlers.read().get("once.bail").is_none(),
"pruned from the registry after spending"
);
}
#[tokio::test]
async fn parallel_returns_null_after_all_handlers_complete() {
let svc = EventsService::new();
let completed = Arc::new(std::sync::atomic::AtomicUsize::new(0));
for value in [1, 2] {
let completed = completed.clone();
svc.on("parallel.result".into(), move |_payload| {
let completed = completed.clone();
async move {
completed.fetch_add(value, Ordering::SeqCst);
Ok(serde_json::json!({ "value": value }))
}
});
}
let out = svc
.dispatch(
"parallel.result".into(),
serde_json::json!({ "input": true }),
Dispatch::Parallel,
)
.await
.unwrap();
assert_eq!(out, serde_json::Value::Null);
assert_eq!(completed.load(Ordering::SeqCst), 3);
}
#[tokio::test]
async fn serial_stops_at_first_non_null_result() {
let svc = EventsService::new();
let ran = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let first = ran.clone();
svc.on("serial.bail".into(), move |payload| {
let first = first.clone();
async move {
first.fetch_add(1, Ordering::SeqCst);
assert_eq!(payload["input"], true);
Ok(serde_json::Value::Null)
}
});
let second = ran.clone();
svc.on("serial.bail".into(), move |payload| {
let second = second.clone();
async move {
second.fetch_add(1, Ordering::SeqCst);
assert_eq!(payload["input"], true);
Ok(serde_json::json!({ "handled": true }))
}
});
let third = ran.clone();
svc.on("serial.bail".into(), move |_payload| {
let third = third.clone();
async move {
third.fetch_add(1, Ordering::SeqCst);
Ok(serde_json::json!({ "late": true }))
}
});
let out = svc
.dispatch(
"serial.bail".into(),
serde_json::json!({ "input": true }),
Dispatch::Serial,
)
.await
.unwrap();
assert_eq!(out, serde_json::json!({ "handled": true }));
assert_eq!(ran.load(Ordering::SeqCst), 2);
}
#[tokio::test]
async fn serial_preserves_original_payload_when_no_handler_bails() {
let svc = EventsService::new();
svc.on("serial.identity".into(), |_payload| async move {
Ok(serde_json::Value::Null)
});
svc.on("serial.identity".into(), |_payload| async move {
Ok(serde_json::Value::Null)
});
let payload = serde_json::json!({ "input": [1, 2], "nested": { "ok": true } });
let out = svc
.dispatch("serial.identity".into(), payload.clone(), Dispatch::Serial)
.await
.unwrap();
assert_eq!(out, payload);
}
#[tokio::test]
async fn default_agent_admit_handler_denies_monthly() {
let svc = EventsService::new();
let out = svc
.dispatch(
"agent.admit".into(),
serde_json::json!({
"monthly": 10,
"daily": 0,
"requests_per_month": 10,
"requests_per_day": 50,
"tier": "free"
}),
Dispatch::Bail,
)
.await
.unwrap();
assert_eq!(out["deny"], "monthly");
}
#[tokio::test]
async fn default_agent_admit_handler_allows_under_quota() {
let svc = EventsService::new();
let out = svc
.dispatch(
"agent.admit".into(),
serde_json::json!({
"monthly": 0,
"daily": 0,
"requests_per_month": 10,
"requests_per_day": 50,
"tier": "free"
}),
Dispatch::Bail,
)
.await
.unwrap();
assert!(
out.get("deny").is_none(),
"under-quota must not deny, got {out}"
);
}
#[tokio::test]
async fn waterfall_around_no_handlers_runs_core() {
let svc = EventsService::new();
let out = svc
.waterfall_around(
"around.empty".into(),
serde_json::json!({}),
|mut payload| async move {
if let Some(obj) = payload.as_object_mut() {
obj.insert("core".into(), serde_json::json!(true));
}
Ok(payload)
},
)
.await
.unwrap();
assert_eq!(out["core"], true);
}
#[tokio::test]
async fn waterfall_around_handler_calls_next_then_core() {
let svc = EventsService::new();
svc.on_waterfall("around.wrap".into(), |mut payload, next| async move {
if let Some(obj) = payload.as_object_mut() {
obj.insert("wrap".into(), serde_json::json!(true));
}
next(payload).await
});
let out = svc
.waterfall_around(
"around.wrap".into(),
serde_json::json!({}),
|mut payload| async move {
if let Some(obj) = payload.as_object_mut() {
obj.insert("core".into(), serde_json::json!(true));
}
Ok(payload)
},
)
.await
.unwrap();
assert_eq!(out["wrap"], true);
assert_eq!(out["core"], true);
}
#[tokio::test]
async fn waterfall_around_short_circuit_skips_core() {
let svc = EventsService::new();
let flag = Arc::new(AtomicBool::new(false));
svc.on_waterfall("around.short".into(), |payload, _next| async move {
Ok(payload)
});
let f = flag.clone();
let out = svc
.waterfall_around(
"around.short".into(),
serde_json::json!({ "ok": true }),
move |payload| {
let f = f.clone();
async move {
f.store(true, Ordering::SeqCst);
Ok(payload)
}
},
)
.await
.unwrap();
assert_eq!(out["ok"], true);
assert!(
!flag.load(Ordering::SeqCst),
"core must not run when handler short-circuits"
);
}
#[tokio::test]
async fn typed_dispatch_and_listener_round_trip() {
let svc = EventsService::new();
let seen = Arc::new(parking_lot::Mutex::new(
None::<crate::events_payload::AgentUsagePayload>,
));
let slot = seen.clone();
let _d = svc.on_typed::<crate::events_payload::AgentUsageEvent, _, _>(move |p| {
let slot = slot.clone();
async move {
*slot.lock() = Some(p);
Ok(serde_json::Value::Null)
}
});
svc.dispatch_typed::<crate::events_payload::AgentUsageEvent>(
&crate::events_payload::AgentUsagePayload {
tenant: Some("acme".into()),
prompt: 3,
completion: 4,
total: 7,
},
)
.await
.unwrap();
for _ in 0..100 {
if seen.lock().is_some() {
break;
}
tokio::time::sleep(std::time::Duration::from_millis(5)).await;
}
let got = seen.lock().clone().expect("handler must observe payload");
assert_eq!(got.tenant.as_deref(), Some("acme"));
assert_eq!(got.total, 7);
}
#[tokio::test]
async fn typed_listener_skips_malformed_payload_passthrough() {
let svc = EventsService::new();
let ran = Arc::new(AtomicBool::new(false));
let flag = ran.clone();
let _d = svc.on_typed::<crate::events_payload::AgentAdmitEvent, _, _>(move |_p| {
let flag = flag.clone();
async move {
flag.store(true, Ordering::SeqCst);
Ok(serde_json::json!({ "deny": "daily" }))
}
});
let out = svc
.dispatch(
crate::events_catalog::ev::AGENT_ADMIT.into(),
serde_json::json!({ "tenant_id": 42 }), Dispatch::Bail,
)
.await
.unwrap();
assert!(
!ran.load(Ordering::SeqCst),
"malformed payload must skip the typed handler"
);
assert_eq!(out["tenant_id"], 42, "value must pass through unchanged");
}
#[tokio::test]
async fn typed_waterfall_short_circuit() {
use crate::events::WaterfallNext;
let svc = EventsService::new();
let _d = svc.on_typed_waterfall::<crate::events_payload::LlmGetClientEvent, _, _>(
|p, next: WaterfallNext| async move {
if p.capability == "blocked" {
return Ok(serde_json::json!({ "deny": true }));
}
next(serde_json::json!({ "capability": p.capability })).await
},
);
let denied = svc
.dispatch(
crate::events_catalog::ev::LLM_GET_CLIENT.into(),
serde_json::json!({ "capability": "blocked" }),
Dispatch::Waterfall,
)
.await
.unwrap();
assert_eq!(denied["deny"], true);
let passed = svc
.dispatch(
crate::events_catalog::ev::LLM_GET_CLIENT.into(),
serde_json::json!({ "capability": "chat" }),
Dispatch::Waterfall,
)
.await
.unwrap();
assert_eq!(passed["capability"], "chat");
}
#[test]
fn summarize_listener_errors_formats_multiple_failures() {
let summary = crate::events::summarize_listener_errors(vec![
("quota".to_string(), "monthly cap reached".to_string()),
("audit".to_string(), "db write failed".to_string()),
]);
assert_eq!(
summary,
"2 listener failures: quota: monthly cap reached; audit: db write failed"
);
}
#[test]
fn summarize_listener_errors_formats_single_failure() {
let summary = crate::events::summarize_listener_errors(vec![(
"solo".to_string(),
"boom".to_string(),
)]);
assert_eq!(summary, "1 listener failure: solo: boom");
assert_eq!(
crate::events::summarize_listener_errors(Vec::new()),
"0 listener failures"
);
}
#[test]
fn aggregate_error_display_and_error_impl() {
let agg = AggregateError {
errors: vec![
("a".to_string(), "x".to_string()),
("b".to_string(), "y".to_string()),
],
};
assert_eq!(agg.to_string(), "2 listener failures: a: x; b: y");
let dyn_err: &dyn std::error::Error = &agg;
assert!(dyn_err.to_string().contains("b: y"));
}
#[tokio::test]
async fn parallel_dispatch_aggregates_all_listener_failures() {
use crate::CordisError as Err;
let svc = EventsService::new();
let completed = Arc::new(AtomicBool::new(false));
let c = completed.clone();
svc.on("par.agg".into(), move |_p| {
let c = c.clone();
async move {
c.store(true, Ordering::SeqCst);
Ok(serde_json::json!({ "ok": true }))
}
});
svc.on("par.agg".into(), |_p| async move {
Err::<serde_json::Value, _>(Err::Configuration("first failure".into()))
});
svc.on("par.agg".into(), |_p| async move {
Err::<serde_json::Value, _>(Err::Fiber("second failure".into()))
});
let err = svc
.dispatch("par.agg".into(), serde_json::json!({}), Dispatch::Parallel)
.await
.unwrap_err();
let text = err.message();
assert!(
text.starts_with("internal kernel error: 2 listener failures:")
&& text.contains("listener[1]: configuration error: first failure")
&& text.contains("listener[2]: fiber error: second failure"),
"aggregate must list every failing listener, got: {text}"
);
assert!(
completed.load(Ordering::SeqCst),
"healthy listener still ran"
);
}
#[tokio::test]
async fn parallel_dispatch_single_failure_uses_singular_summary() {
let svc = EventsService::new();
svc.on("par.solo".into(), |_p| async move {
Err::<serde_json::Value, _>(crate::CordisError::Internal("only one".into()))
});
let err = svc
.dispatch("par.solo".into(), serde_json::json!({}), Dispatch::Parallel)
.await
.unwrap_err();
assert_eq!(
err.message(),
"internal kernel error: 1 listener failure: listener[0]: internal kernel error: only one"
);
}
#[tokio::test]
async fn prepend_ordering_observed() {
let svc = EventsService::new();
let order = Arc::new(parking_lot::Mutex::<Vec<String>>::new(Vec::new()));
for name in ["first", "second"] {
let slot = order.clone();
svc.on("prepend.test".into(), move |_p| {
let slot = slot.clone();
async move {
slot.lock().push(name.to_string());
Ok(serde_json::Value::Null)
}
});
}
let prepended_slot = order.clone();
svc.on_with(
"prepend.test".into(),
EventOptions {
prepend: true,
global: false,
},
move |_p| {
let prepended_slot = prepended_slot.clone();
async move {
prepended_slot.lock().push("prepended".to_string());
Ok(serde_json::Value::Null)
}
},
);
let once_slot = order.clone();
svc.once_with(
"prepend.test".into(),
EventOptions {
prepend: true,
global: false,
},
move |_p| {
let once_slot = once_slot.clone();
async move {
once_slot.lock().push("once-prepended".to_string());
Ok(serde_json::Value::Null)
}
},
);
svc.dispatch(
"prepend.test".into(),
serde_json::json!({}),
Dispatch::Serial,
)
.await
.unwrap();
assert_eq!(
*order.lock(),
["once-prepended", "prepended", "first", "second"],
"prepend inserts at the dispatch-order front; defaults append"
);
svc.dispatch(
"prepend.test".into(),
serde_json::json!({}),
Dispatch::Serial,
)
.await
.unwrap();
assert_eq!(
order.lock()[4..],
["prepended", "first", "second"],
"spent once slot drops out; relative order is stable"
);
}
#[tokio::test]
async fn filter_excludes_nonmatching_contexts() {
let svc = EventsService::new();
let ran_a = Arc::new(AtomicUsize::new(0));
let ran_b = Arc::new(AtomicUsize::new(0));
let a = ran_a.clone();
svc.on_with("filtered.test".into(), EventOptions::default(), move |p| {
let a = a.clone();
async move {
a.fetch_add(1, Ordering::SeqCst);
Ok(p)
}
});
let b = ran_b.clone();
svc.on_with("filtered.test".into(), EventOptions::default(), move |p| {
let b = b.clone();
async move {
b.fetch_add(1, Ordering::SeqCst);
Ok(p)
}
});
svc.emit_filtered(
"filtered.test".into(),
serde_json::json!({ "tenant": "a" }),
Box::new(|_opts| false),
)
.unwrap();
tokio::time::sleep(std::time::Duration::from_millis(30)).await;
assert_eq!(ran_a.load(Ordering::SeqCst), 0, "rejecting filter excludes a");
assert_eq!(ran_b.load(Ordering::SeqCst), 0, "rejecting filter excludes b");
svc.dispatch("filtered.test".into(), serde_json::json!({}), Dispatch::Emit)
.await
.unwrap();
tokio::time::sleep(std::time::Duration::from_millis(30)).await;
assert_eq!(
ran_a.load(Ordering::SeqCst),
1,
"listener a must be back for unfiltered dispatches"
);
assert_eq!(
ran_b.load(Ordering::SeqCst),
1,
"listener b must be back for unfiltered dispatches"
);
let handlers = svc.handlers.read();
let slots = handlers.get("filtered.test").expect("entry kept");
assert_eq!(
slots.len(),
2,
"filter exclusion must not unregister anyone"
);
}
#[tokio::test]
async fn global_bypasses_filter() {
let svc = EventsService::new();
let ran = Arc::new(AtomicUsize::new(0));
let b = ran.clone();
svc.on_with(
"global.test".into(),
EventOptions::default(),
move |payload| {
let b = b.clone();
async move {
if payload["tenant"] == "b" {
b.fetch_add(1, Ordering::SeqCst);
}
Ok(payload)
}
},
);
let g = ran.clone();
svc.on_with(
"global.test".into(),
EventOptions {
prepend: false,
global: true,
},
move |payload| {
let g = g.clone();
async move {
if payload["tenant"] == "b" {
g.fetch_add(10, Ordering::SeqCst);
}
Ok(payload)
}
},
);
svc.emit_filtered(
"global.test".into(),
serde_json::json!({ "tenant": "b" }),
Box::new(|_opts| false),
)
.unwrap();
tokio::time::sleep(std::time::Duration::from_millis(30)).await;
assert_eq!(
ran.load(Ordering::SeqCst),
10,
"global listener runs despite a rejecting filter; non-global does not"
);
}
#[tokio::test]
async fn get_interceptor_rewrites_read() {
let svc = EventsService::new();
assert_eq!(svc.intercept_get("Svc", None).await.unwrap(), None);
let d = svc.on(INTERNAL_GET_EVENT.into(), |_payload| async move {
Ok(serde_json::json!({ "service": "Svc", "rewritten": true }))
});
let out = svc.intercept_get("Svc", Some("tenant-a".into())).await;
match out {
Ok(Some(value)) => {
assert_eq!(value["rewritten"], serde_json::json!(true));
assert_eq!(value["service"], "Svc");
}
other => panic!("expected rewritten read, got {other:?}"),
}
d.dispose();
assert_eq!(svc.intercept_get("Svc", None).await.unwrap(), None);
}
#[tokio::test]
async fn set_interceptor_vetoes_write_leaves_old_value() {
let svc = EventsService::new();
assert!(svc.intercept_set("Svc", None).await.is_ok());
let d = svc.on(INTERNAL_SET_EVENT.into(), |_payload| async move {
Err::<serde_json::Value, CordisError>(CordisError::Configuration(
"writes are frozen".into(),
))
});
let err = svc.intercept_set("Svc", None).await.unwrap_err();
assert!(
err.to_string().contains("frozen"),
"veto error must surface, got {err}"
);
d.dispose();
assert!(svc.intercept_set("Svc", None).await.is_ok());
}
#[tokio::test]
async fn config_interceptor_rewrites_effective_config() {
let svc = EventsService::new();
let raw = serde_json::json!({ "model": "base" });
assert_eq!(svc.intercept_config(raw.clone()).await.unwrap(), raw);
let d = svc.on(INTERNAL_CONFIG_EVENT.into(), |raw| async move {
let mut effective = raw;
if let Some(obj) = effective.as_object_mut() {
obj.insert("model".into(), serde_json::json!("rewritten"));
obj.insert("seen_by_interceptor".into(), serde_json::json!(true));
}
Ok(effective)
});
let effective = svc.intercept_config(raw).await.unwrap();
assert_eq!(effective["model"], "rewritten");
assert_eq!(effective["seen_by_interceptor"], true);
d.dispose();
svc.on(INTERNAL_CONFIG_EVENT.into(), |_raw| async move {
Ok(serde_json::Value::Null)
});
let raw2 = serde_json::json!({ "keep": 1 });
assert_eq!(svc.intercept_config(raw2.clone()).await.unwrap(), raw2);
}
#[tokio::test]
async fn update_interceptor_veto_skips_restart_keeps_config() {
use crate::{Context, Fiber};
let ctx = Context::new_root();
let events = Arc::new(EventsService::new());
ctx.provide_arc(events.clone());
let fiber = Arc::new(Fiber::new());
fiber.set_reload_context(&ctx);
fiber.set_id(70_100);
let calls = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let c = calls.clone();
fiber.set_reload_runner(Box::new(move |_| {
c.fetch_add(1, Ordering::SeqCst);
Ok(true)
}));
fiber.declare_inject::<crate::ReflectService>();
let _prov = ctx.provide(crate::ReflectService::new());
fiber.refresh(&ctx).await;
assert!(matches!(fiber.state(), crate::FiberState::Active { .. }));
assert_eq!(calls.load(Ordering::SeqCst), 1, "initial apply ran");
let d = events.on(INTERNAL_UPDATE_EVENT.into(), |_payload| async move {
Ok(serde_json::json!({ "veto": "maintenance window" }))
});
fiber.update(&ctx).await.unwrap();
assert!(
matches!(fiber.state(), crate::FiberState::Active { .. }),
"vetoed update must keep the fiber Active, got {:?}",
fiber.state()
);
assert_eq!(
calls.load(Ordering::SeqCst),
1,
"runner must not run again under veto"
);
d.dispose();
fiber.update(&ctx).await.unwrap();
assert!(matches!(fiber.state(), crate::FiberState::Active { .. }));
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn listener_interceptor_bail_cancels_registration_inert_handle() {
let svc = EventsService::new();
let allow_gate = svc.on(INTERNAL_LISTENER_EVENT.into(), |payload| async move {
if payload["event"] == "blocked.event" {
Ok(serde_json::json!("denied"))
} else {
Ok(serde_json::Value::Null)
}
});
let ran = Arc::new(AtomicBool::new(false));
let flag = ran.clone();
let handle = svc.on("blocked.event".into(), move |_p| {
let flag = flag.clone();
async move {
flag.store(true, Ordering::SeqCst);
Ok(serde_json::Value::Null)
}
});
handle.dispose(); svc.dispatch("blocked.event".into(), serde_json::json!({}), Dispatch::Bail)
.await
.unwrap();
tokio::time::sleep(std::time::Duration::from_millis(30)).await;
assert!(
!ran.load(Ordering::SeqCst),
"cancelled registration must never run"
);
let wf_handle = svc.on_waterfall(
"blocked.event".into(),
|_p: serde_json::Value, next| async move { next(_p).await },
);
wf_handle.dispose();
assert_eq!(
svc.listener_count("blocked.event"),
0,
"neither registry may hold a cancelled registration"
);
let allowed = svc.on("allowed.event".into(), |_p| async move {
Ok(serde_json::Value::Null)
});
allowed.dispose();
assert_eq!(svc.listener_count("allowed.event"), 0, "dispose works");
let fail_gate = svc.on(INTERNAL_LISTENER_EVENT.into(), |_payload| async move {
Err::<serde_json::Value, CordisError>(CordisError::Configuration("gate down".into()))
});
let second = svc.on("another.event".into(), |_p| async move {
Ok(serde_json::Value::Null)
});
second.dispose();
assert_eq!(
svc.listener_count("another.event"),
0,
"erroring veto chain must fail closed"
);
fail_gate.dispose();
allow_gate.dispose();
}
fn p_owned(p: serde_json::Value) -> Pin<Box<dyn Future<Output = Result<serde_json::Value, CordisError>> + Send>> {
Box::pin(async move { Ok(p) })
}
#[tokio::test]
async fn internal_dispatch_observes_non_internal_only() {
let svc = EventsService::new();
let seen = Arc::new(Mutex::new(Vec::<InternalDispatchPayload>::new()));
let s = seen.clone();
svc.on(INTERNAL_DISPATCH_EVENT.into(), move |payload| {
let s = s.clone();
async move {
if let Ok(parsed) =
serde_json::from_value::<InternalDispatchPayload>(payload)
{
s.lock().push(parsed);
}
Ok(serde_json::Value::Null)
}
});
svc.dispatch("observed.a".into(), serde_json::json!({ "n": 1 }), Dispatch::Emit)
.await
.unwrap();
svc.dispatch(
"observed.b".into(),
serde_json::json!({ "n": 2 }),
Dispatch::Waterfall,
)
.await
.unwrap();
svc.dispatch("observed.c".into(), serde_json::json!({}), Dispatch::Bail)
.await
.unwrap();
svc.on(INTERNAL_GET_EVENT.into(), |_p| async move {
Ok(serde_json::json!({ "x": true }))
});
svc.intercept_get("SomeSvc", None).await.unwrap();
for _ in 0..100 {
if seen.lock().len() >= 3 {
break;
}
tokio::time::sleep(std::time::Duration::from_millis(10)).await;
}
let observed = seen.lock().clone();
assert_eq!(
observed.len(),
3,
"meta dispatch must NOT be observed; got {observed:?}"
);
assert_eq!(observed[0].mode, "emit");
assert_eq!(observed[0].name, "observed.a");
assert_eq!(observed[1].mode, "waterfall");
assert_eq!(observed[1].name, "observed.b");
assert_eq!(observed[2].mode, "bail");
assert_eq!(observed[0].args["n"], 1);
}
#[tokio::test]
async fn interceptor_error_fails_fiber_activation() {
use crate::{Context, Fiber};
let ctx = Context::new_root();
let events = Arc::new(EventsService::new());
ctx.provide_arc(events.clone());
events.on(INTERNAL_CONFIG_EVENT.into(), |_raw| async move {
Err::<serde_json::Value, CordisError>(CordisError::Configuration(
"config rejected by policy".into(),
))
});
let fiber = Arc::new(Fiber::new());
fiber.set_reload_context(&ctx);
fiber.set_id(70_101);
fiber.set_raw_config(serde_json::json!({ "model": "base" }));
fiber.set_reload_runner(Box::new(|_| {
panic!("runner must never run when config interception refuses");
}));
fiber.declare_inject::<crate::ReflectService>();
let _prov = ctx.provide(crate::ReflectService::new());
fiber.refresh(&ctx).await;
match fiber.state() {
crate::FiberState::Failed { error } => {
let msg = error.unwrap_or_default();
assert!(
msg.contains("config rejected by policy"),
"failure must carry the interception error, got: {msg}"
);
}
other => panic!("expected Failed activation, got {other:?}"),
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn target_carrying_dispatches_filter_per_dispatch() {
let svc = EventsService::new();
let flat_ran = Arc::new(AtomicUsize::new(0));
let f = flat_ran.clone();
svc.on_with(
"wf.filtered".into(),
EventOptions::default(),
move |_p| {
let f = f.clone();
async move {
f.fetch_add(1, Ordering::SeqCst);
Ok(serde_json::Value::Null)
}
},
);
let wf_ran = Arc::new(AtomicUsize::new(0));
let w = wf_ran.clone();
svc.on_waterfall("wf.filtered".into(), move |_p, next| {
let w = w.clone();
async move {
w.fetch_add(1, Ordering::SeqCst);
next(_p).await
}
});
let out = svc
.bail_from(
"wf.filtered".into(),
serde_json::json!({"n": 1}),
Some(Box::new(|_opts| false)),
)
.await
.unwrap();
assert_eq!(out, serde_json::json!({"n": 1}));
assert_eq!(
flat_ran.load(Ordering::SeqCst),
0,
"rejecting filter must exclude the flat listener"
);
let out = svc
.waterfall_from(
"wf.filtered".into(),
serde_json::json!({"n": 1}),
Some(Box::new(|_opts| false)),
)
.await
.unwrap();
assert_eq!(out, serde_json::json!({"n": 1}));
assert_eq!(
wf_ran.load(Ordering::SeqCst),
0,
"rejecting filter must exclude the waterfall listener"
);
svc.bail_from(
"wf.filtered".into(),
serde_json::json!({"n": 2}),
Some(Box::new(|_opts| true)),
)
.await
.unwrap();
assert_eq!(flat_ran.load(Ordering::SeqCst), 1);
svc.waterfall_async_from(
"wf.filtered".into(),
serde_json::json!({"n": 2}),
Some(Box::new(|_opts| true)),
|mut p| async move {
if let Some(obj) = p.as_object_mut() {
obj.insert("core".into(), serde_json::json!(true));
}
Ok(p)
},
)
.await
.unwrap();
assert_eq!(wf_ran.load(Ordering::SeqCst), 1);
assert_eq!(svc.listener_count("wf.filtered"), 2);
}
}