use std::collections::VecDeque;
use std::sync::Arc;
use std::time::Duration;
use nexo_broker::{handle::BrokerHandle, AnyBroker, Event};
use nexo_config::types::agents::InboundBinding;
use nexo_config::types::event_subscriber::{EventSubscriberBinding, OverflowPolicy, SynthesisMode};
use serde_json::{json, Value};
use thiserror::Error;
use tokio::sync::{Mutex, Notify, Semaphore};
use tokio_util::sync::CancellationToken;
use uuid::Uuid;
pub const EVENT_INBOUND_TOPIC_PREFIX: &str = "plugin.inbound.event";
pub const EVENT_BINDING_CHANNEL: &str = "event";
pub const EVENT_SOURCE_PAYLOAD_FIELD: &str = "_nexo_event_source";
pub fn event_inbound_topic(source_id: &str) -> String {
format!("{EVENT_INBOUND_TOPIC_PREFIX}.{source_id}")
}
pub fn build_synthesised_payload(binding: &EventSubscriberBinding, event: &Event) -> Option<Value> {
if binding.synthesize_inbound == SynthesisMode::Off {
return None;
}
let envelope_id = extract_envelope_id(&event.payload);
let from = format!("event:{}", binding.id);
let timestamp = chrono::Utc::now().timestamp_millis();
let msg_id = envelope_id
.map(|u| u.to_string())
.unwrap_or_else(|| event.id.to_string());
let text = match binding.synthesize_inbound {
SynthesisMode::Tick => format!(
"<event subject=\"{}\" envelope_id=\"{}\"/>",
event.topic,
envelope_id
.map(|u| u.to_string())
.unwrap_or_else(|| "null".to_string())
),
SynthesisMode::Synthesize => match binding.inbound_template.as_deref() {
Some(tmpl) => nexo_tool_meta::render_template(tmpl, &event.payload),
None => serde_json::to_string(&event.payload)
.unwrap_or_else(|_| "<unrenderable payload>".to_string()),
},
SynthesisMode::Off => unreachable!("checked above"),
_ => "<unknown synthesis mode>".to_string(),
};
let event_source = json!({
"subject": event.topic,
"envelope_id": envelope_id,
"synthesis_mode": binding.synthesize_inbound.as_str(),
});
let inbound_kind = match binding.inbound_kind {
nexo_tool_meta::InboundKind::ExternalUser => "external_user",
nexo_tool_meta::InboundKind::InternalSystem => "internal_system",
nexo_tool_meta::InboundKind::InterSession => "inter_session",
_ => "external_user",
};
Some(json!({
"kind": "message",
"from": from,
"chat": from,
"text": text,
"reply_to": Value::Null,
"is_group": false,
"timestamp": timestamp,
"msg_id": msg_id,
"inbound_kind": inbound_kind,
EVENT_SOURCE_PAYLOAD_FIELD: event_source,
}))
}
fn extract_envelope_id(payload: &Value) -> Option<Uuid> {
payload
.as_object()?
.get("envelope_id")?
.as_str()
.and_then(|s| Uuid::parse_str(s).ok())
}
pub fn synthesize_event_inbound_bindings(
declared: &[InboundBinding],
event_subscribers: &[EventSubscriberBinding],
) -> Vec<InboundBinding> {
let mut out = declared.to_vec();
for sub in event_subscribers {
let already = out
.iter()
.any(|b| b.plugin == EVENT_BINDING_CHANNEL && b.instance.as_deref() == Some(&sub.id));
if !already {
let mut binding = InboundBinding::default();
binding.plugin = EVENT_BINDING_CHANNEL.into();
binding.instance = Some(sub.id.clone());
out.push(binding);
}
}
out
}
pub fn extract_nexo_event_source(payload: &Value) -> Option<nexo_tool_meta::EventSourceMeta> {
let raw = payload.as_object()?.get(EVENT_SOURCE_PAYLOAD_FIELD)?;
serde_json::from_value(raw.clone()).ok()
}
pub async fn run_event_subscriber(
sub: Arc<EventSubscriber>,
cancel: CancellationToken,
) -> Result<(), EventSubscriberError> {
if sub.binding.synthesize_inbound == SynthesisMode::Off {
tracing::info!(
agent_id = %sub.agent_id,
binding_id = %sub.binding.id,
"event_subscriber off — task exit"
);
return Ok(());
}
let mut subscription = sub
.broker
.subscribe(&sub.binding.subject_pattern)
.await
.map_err(|e| EventSubscriberError::BrokerSubscribe {
subject: sub.binding.subject_pattern.clone(),
detail: e.to_string(),
})?;
tracing::info!(
agent_id = %sub.agent_id,
binding_id = %sub.binding.id,
subject = %sub.binding.subject_pattern,
"event_subscriber started"
);
let buffer: Arc<Mutex<VecDeque<Event>>> =
Arc::new(Mutex::new(VecDeque::with_capacity(sub.binding.max_buffer)));
let notify = Arc::new(Notify::new());
let consumer_done = Arc::new(Notify::new());
let consumer_sub = Arc::clone(&sub);
let consumer_buffer = Arc::clone(&buffer);
let consumer_notify = Arc::clone(¬ify);
let consumer_done_emit = Arc::clone(&consumer_done);
let consumer_cancel = cancel.clone();
let consumer = tokio::spawn(async move {
loop {
tokio::select! {
biased;
_ = consumer_cancel.cancelled() => break,
_ = consumer_notify.notified() => {}
}
loop {
let event = {
let mut g = consumer_buffer.lock().await;
g.pop_front()
};
let Some(event) = event else { break };
let permit = match &consumer_sub.semaphore {
Some(sem) => match Arc::clone(sem).acquire_owned().await {
Ok(p) => Some(p),
Err(_) => return,
},
None => None,
};
if let Some(payload) = build_synthesised_payload(&consumer_sub.binding, &event) {
let topic = event_inbound_topic(&consumer_sub.binding.id);
let republish = Event::new(
topic.clone(),
format!("event_subscriber:{}", consumer_sub.binding.id),
payload,
);
if let Err(e) = consumer_sub.broker.publish(&topic, republish).await {
tracing::warn!(
agent_id = %consumer_sub.agent_id,
binding_id = %consumer_sub.binding.id,
topic = %topic,
error = %e,
"event_subscriber publish failed — dropping event"
);
}
}
drop(permit);
}
}
consumer_done_emit.notify_one();
});
let mut dropped_oldest: u64 = 0;
let mut dropped_newest: u64 = 0;
loop {
tokio::select! {
biased;
_ = cancel.cancelled() => {
tracing::info!(
agent_id = %sub.agent_id,
binding_id = %sub.binding.id,
"event_subscriber cancel — draining"
);
break;
}
ev = subscription.next() => {
let Some(event) = ev else { break };
let republish_topic = event_inbound_topic(&sub.binding.id);
if event.topic == republish_topic {
tracing::warn!(
agent_id = %sub.agent_id,
binding_id = %sub.binding.id,
topic = %event.topic,
"event_subscriber loop guard fired — dropped self-event"
);
continue;
}
let max_buffer = sub.binding.max_buffer;
let mut buf = buffer.lock().await;
if buf.len() >= max_buffer {
match sub.binding.overflow_policy {
OverflowPolicy::DropOldest => {
buf.pop_front();
buf.push_back(event);
dropped_oldest += 1;
tracing::warn!(
agent_id = %sub.agent_id,
binding_id = %sub.binding.id,
dropped_total = dropped_oldest,
"event_subscriber buffer full — drop-oldest"
);
}
OverflowPolicy::DropNewest => {
dropped_newest += 1;
tracing::warn!(
agent_id = %sub.agent_id,
binding_id = %sub.binding.id,
dropped_total = dropped_newest,
"event_subscriber buffer full — drop-newest"
);
}
_ => {
dropped_newest += 1;
}
}
} else {
buf.push_back(event);
}
drop(buf);
notify.notify_one();
}
}
}
let _ = tokio::time::timeout(Duration::from_secs(1), consumer_done.notified()).await;
consumer.abort();
tracing::info!(
agent_id = %sub.agent_id,
binding_id = %sub.binding.id,
dropped_oldest,
dropped_newest,
"event_subscriber stopped"
);
Ok(())
}
pub struct EventSubscriber {
pub agent_id: String,
pub binding: Arc<EventSubscriberBinding>,
pub broker: AnyBroker,
pub semaphore: Option<Arc<Semaphore>>,
}
impl EventSubscriber {
pub fn new(
agent_id: impl Into<String>,
binding: EventSubscriberBinding,
broker: AnyBroker,
) -> Self {
let semaphore = if binding.max_concurrency == 0 {
None
} else {
Some(Arc::new(Semaphore::new(binding.max_concurrency as usize)))
};
Self {
agent_id: agent_id.into(),
binding: Arc::new(binding),
broker,
semaphore,
}
}
}
#[non_exhaustive]
#[derive(Debug, Error)]
pub enum EventSubscriberError {
#[error("broker subscribe failed on `{subject}`: {detail}")]
BrokerSubscribe {
subject: String,
detail: String,
},
#[error("broker publish failed on `{topic}`: {detail}")]
BrokerPublish {
topic: String,
detail: String,
},
}
#[cfg(test)]
mod skeleton_tests {
use super::*;
use nexo_broker::AnyBroker;
use nexo_config::types::event_subscriber::{OverflowPolicy, SynthesisMode};
fn mk_binding(id: &str, pattern: &str) -> EventSubscriberBinding {
EventSubscriberBinding {
id: id.into(),
subject_pattern: pattern.into(),
synthesize_inbound: SynthesisMode::Synthesize,
inbound_template: None,
max_concurrency: 1,
max_buffer: 64,
overflow_policy: OverflowPolicy::DropOldest,
inbound_kind: nexo_tool_meta::InboundKind::ExternalUser,
}
}
#[test]
fn topic_prefix_lock_down() {
assert_eq!(EVENT_INBOUND_TOPIC_PREFIX, "plugin.inbound.event");
assert_eq!(EVENT_BINDING_CHANNEL, "event");
assert_eq!(EVENT_SOURCE_PAYLOAD_FIELD, "_nexo_event_source");
}
#[test]
fn event_inbound_topic_renders_dotted() {
assert_eq!(
event_inbound_topic("github_main"),
"plugin.inbound.event.github_main"
);
}
#[tokio::test]
async fn new_with_max_concurrency_one_creates_semaphore() {
let broker = AnyBroker::local();
let sub = EventSubscriber::new("ana", mk_binding("a", "x.>"), broker);
assert!(sub.semaphore.is_some());
assert_eq!(sub.semaphore.as_ref().unwrap().available_permits(), 1);
}
#[tokio::test]
async fn new_with_zero_concurrency_skips_semaphore() {
let broker = AnyBroker::local();
let mut binding = mk_binding("a", "x.>");
binding.max_concurrency = 0;
let sub = EventSubscriber::new("ana", binding, broker);
assert!(sub.semaphore.is_none());
}
fn mk_event(topic: &str, payload: serde_json::Value) -> Event {
Event::new(topic, "tester", payload)
}
#[test]
fn build_synthesise_renders_template() {
let mut binding = mk_binding("github", "webhook.>");
binding.inbound_template = Some("got {{event_kind}}".into());
let event = mk_event(
"webhook.github.opened",
serde_json::json!({"event_kind": "opened"}),
);
let payload = build_synthesised_payload(&binding, &event).unwrap();
assert_eq!(payload["kind"], "message");
assert_eq!(payload["from"], "event:github");
assert_eq!(payload["chat"], "event:github");
assert_eq!(payload["text"], "got opened");
assert_eq!(payload["is_group"], false);
let es = &payload[EVENT_SOURCE_PAYLOAD_FIELD];
assert_eq!(es["subject"], "webhook.github.opened");
assert_eq!(es["synthesis_mode"], "synthesize");
}
#[test]
fn build_tick_renders_event_marker() {
let mut binding = mk_binding("alerts", "alert.>");
binding.synthesize_inbound = SynthesisMode::Tick;
let event = mk_event("alert.cpu.high", serde_json::json!({}));
let payload = build_synthesised_payload(&binding, &event).unwrap();
let text = payload["text"].as_str().unwrap();
assert!(text.starts_with("<event subject=\"alert.cpu.high\""));
assert!(text.ends_with("/>"));
assert_eq!(
payload[EVENT_SOURCE_PAYLOAD_FIELD]["synthesis_mode"],
"tick"
);
}
#[test]
fn build_off_returns_none() {
let mut binding = mk_binding("x", "y.>");
binding.synthesize_inbound = SynthesisMode::Off;
let event = mk_event("y.z", serde_json::json!({}));
assert!(build_synthesised_payload(&binding, &event).is_none());
}
#[test]
fn build_synthesise_falls_back_to_json_when_template_missing() {
let binding = mk_binding("x", "y.>");
let event = mk_event("y.z", serde_json::json!({"a": 1, "b": "two"}));
let payload = build_synthesised_payload(&binding, &event).unwrap();
let text = payload["text"].as_str().unwrap();
assert!(text.contains("\"a\":1"));
assert!(text.contains("\"b\":\"two\""));
}
#[test]
fn build_extracts_envelope_id_when_present() {
let binding = mk_binding("github", "webhook.>");
let envelope_id = Uuid::from_u128(42);
let event = mk_event(
"webhook.github.x",
serde_json::json!({"envelope_id": envelope_id.to_string()}),
);
let payload = build_synthesised_payload(&binding, &event).unwrap();
assert_eq!(payload["msg_id"], envelope_id.to_string());
assert_eq!(
payload[EVENT_SOURCE_PAYLOAD_FIELD]["envelope_id"],
envelope_id.to_string()
);
}
#[test]
fn extract_nexo_event_source_happy_path() {
let binding = mk_binding("github", "webhook.>");
let mut binding_t = binding.clone();
binding_t.inbound_template = Some("x".into());
let event = mk_event("webhook.github.x", serde_json::json!({}));
let payload = build_synthesised_payload(&binding_t, &event).unwrap();
let meta = extract_nexo_event_source(&payload).unwrap();
assert_eq!(meta.subject, "webhook.github.x");
assert_eq!(meta.synthesis_mode, "synthesize");
}
#[test]
fn extract_nexo_event_source_returns_none_when_missing() {
let payload = serde_json::json!({"kind": "message", "from": "x"});
assert!(extract_nexo_event_source(&payload).is_none());
}
#[test]
fn extract_nexo_event_source_returns_none_for_malformed_shape() {
let payload = serde_json::json!({
EVENT_SOURCE_PAYLOAD_FIELD: "not-an-object"
});
assert!(extract_nexo_event_source(&payload).is_none());
}
#[test]
fn synthesize_event_inbound_appends_when_absent() {
let declared: Vec<InboundBinding> = vec![];
let subs = vec![mk_binding("github", "x.>"), mk_binding("stripe", "y.>")];
let out = synthesize_event_inbound_bindings(&declared, &subs);
assert_eq!(out.len(), 2);
assert_eq!(out[0].plugin, "event");
assert_eq!(out[0].instance.as_deref(), Some("github"));
assert_eq!(out[1].plugin, "event");
assert_eq!(out[1].instance.as_deref(), Some("stripe"));
}
#[test]
fn synthesize_event_inbound_idempotent_when_manual_present() {
let mut manual = InboundBinding::default();
manual.plugin = "event".into();
manual.instance = Some("github".into());
manual.allowed_tools = Some(vec!["custom_tool".into()]);
let declared = vec![manual.clone()];
let subs = vec![mk_binding("github", "x.>")];
let out = synthesize_event_inbound_bindings(&declared, &subs);
assert_eq!(out.len(), 1, "manual override preserved, no duplicate");
assert_eq!(out[0].allowed_tools.as_ref().unwrap()[0], "custom_tool");
}
#[test]
fn synthesize_event_inbound_partial_overlap() {
let mut manual = InboundBinding::default();
manual.plugin = "event".into();
manual.instance = Some("github".into());
let declared = vec![manual];
let subs = vec![mk_binding("github", "x.>"), mk_binding("stripe", "y.>")];
let out = synthesize_event_inbound_bindings(&declared, &subs);
assert_eq!(out.len(), 2);
assert_eq!(out[0].instance.as_deref(), Some("github"));
assert_eq!(out[1].instance.as_deref(), Some("stripe"));
}
#[test]
fn synthesize_event_inbound_preserves_unrelated_bindings() {
let mut wa = InboundBinding::default();
wa.plugin = "whatsapp".into();
wa.instance = Some("personal".into());
let declared = vec![wa];
let subs = vec![mk_binding("github", "x.>")];
let out = synthesize_event_inbound_bindings(&declared, &subs);
assert_eq!(out.len(), 2);
assert_eq!(out[0].plugin, "whatsapp");
assert_eq!(out[1].plugin, "event");
}
use std::time::Duration;
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn run_loop_synthesise_republishes() {
let broker = AnyBroker::local();
let mut binding = mk_binding("github", "test.>");
binding.inbound_template = Some("got {{event_kind}}".into());
let sub = Arc::new(EventSubscriber::new("ana", binding, broker.clone()));
let cancel = CancellationToken::new();
let mut listener = broker
.subscribe("plugin.inbound.event.github")
.await
.unwrap();
let task = tokio::spawn(run_event_subscriber(Arc::clone(&sub), cancel.clone()));
tokio::time::sleep(Duration::from_millis(50)).await;
broker
.publish(
"test.opened",
Event::new(
"test.opened",
"tester",
serde_json::json!({"event_kind": "opened"}),
),
)
.await
.unwrap();
let received = tokio::time::timeout(Duration::from_secs(2), listener.next())
.await
.expect("timeout waiting for republish")
.expect("no event");
assert_eq!(received.payload["text"], "got opened");
assert_eq!(received.payload["from"], "event:github");
assert_eq!(
received.payload[EVENT_SOURCE_PAYLOAD_FIELD]["synthesis_mode"],
"synthesize"
);
cancel.cancel();
let _ = tokio::time::timeout(Duration::from_secs(2), task).await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn run_loop_tick_republishes_marker() {
let broker = AnyBroker::local();
let mut binding = mk_binding("alerts", "alert.>");
binding.synthesize_inbound = SynthesisMode::Tick;
let sub = Arc::new(EventSubscriber::new("ana", binding, broker.clone()));
let cancel = CancellationToken::new();
let mut listener = broker
.subscribe("plugin.inbound.event.alerts")
.await
.unwrap();
let task = tokio::spawn(run_event_subscriber(Arc::clone(&sub), cancel.clone()));
tokio::time::sleep(Duration::from_millis(50)).await;
broker
.publish(
"alert.cpu.high",
Event::new("alert.cpu.high", "tester", serde_json::json!({})),
)
.await
.unwrap();
let received = tokio::time::timeout(Duration::from_secs(2), listener.next())
.await
.unwrap()
.unwrap();
let text = received.payload["text"].as_str().unwrap();
assert!(text.starts_with("<event subject=\"alert.cpu.high\""));
cancel.cancel();
let _ = tokio::time::timeout(Duration::from_secs(2), task).await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn run_loop_off_mode_exits_immediately() {
let broker = AnyBroker::local();
let mut binding = mk_binding("x", "y.>");
binding.synthesize_inbound = SynthesisMode::Off;
let sub = Arc::new(EventSubscriber::new("ana", binding, broker));
let cancel = CancellationToken::new();
let task = tokio::spawn(run_event_subscriber(Arc::clone(&sub), cancel));
let result = tokio::time::timeout(Duration::from_secs(2), task)
.await
.expect("task should finish without cancel");
assert!(result.unwrap().is_ok());
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn run_loop_cancel_token_stops_within_one_second() {
let broker = AnyBroker::local();
let binding = mk_binding("x", "y.>");
let sub = Arc::new(EventSubscriber::new("ana", binding, broker));
let cancel = CancellationToken::new();
let task = tokio::spawn(run_event_subscriber(Arc::clone(&sub), cancel.clone()));
tokio::time::sleep(Duration::from_millis(50)).await;
cancel.cancel();
let res = tokio::time::timeout(Duration::from_secs(2), task).await;
assert!(res.is_ok(), "task should join within 2s of cancel");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn run_loop_guard_drops_self_events() {
let broker = AnyBroker::local();
let binding = mk_binding("loopy", "plugin.inbound.event.>");
let sub = Arc::new(EventSubscriber::new("ana", binding, broker.clone()));
let cancel = CancellationToken::new();
let mut listener = broker
.subscribe("plugin.inbound.event.loopy")
.await
.unwrap();
let task = tokio::spawn(run_event_subscriber(Arc::clone(&sub), cancel.clone()));
tokio::time::sleep(Duration::from_millis(50)).await;
broker
.publish(
"plugin.inbound.event.loopy",
Event::new(
"plugin.inbound.event.loopy",
"tester",
serde_json::json!({}),
),
)
.await
.unwrap();
let first = tokio::time::timeout(Duration::from_millis(200), listener.next())
.await
.unwrap()
.unwrap();
assert_eq!(first.source, "tester", "first event is the test publish");
let second = tokio::time::timeout(Duration::from_millis(300), listener.next()).await;
assert!(second.is_err(), "no self-republish (loop guard fired)");
cancel.cancel();
let _ = tokio::time::timeout(Duration::from_secs(2), task).await;
}
}