use pulpo_common::api::EventPushRequest;
use pulpo_common::event::PulpoEvent;
#[cfg(not(coverage))]
use tracing::{debug, info, warn};
#[cfg(not(coverage))]
const MAX_BATCH_SIZE: usize = 10;
#[cfg(not(coverage))]
pub async fn run_event_push_loop(
controller_url: String,
controller_token: String,
_node_name: String,
mut event_rx: tokio::sync::broadcast::Receiver<PulpoEvent>,
mut shutdown_rx: tokio::sync::watch::Receiver<bool>,
) {
let client = reqwest::Client::builder()
.timeout(std::time::Duration::from_secs(10))
.build()
.expect("failed to build reqwest client");
let push_url = format!("{controller_url}/api/v1/events/push");
info!("Node event push loop started (controller={controller_url})");
loop {
tokio::select! {
result = event_rx.recv() => {
match result {
Ok(event) => {
let mut batch = vec![event];
while batch.len() < MAX_BATCH_SIZE {
match event_rx.try_recv() {
Ok(ev) => batch.push(ev),
Err(_) => break,
}
}
let req_body = EventPushRequest { events: batch };
let request = client
.post(&push_url)
.bearer_auth(&controller_token)
.json(&req_body);
match request.send().await {
Ok(resp) if resp.status().is_success() => {
debug!(
events = req_body.events.len(),
"Pushed events to controller"
);
}
Ok(resp) => {
warn!(
status = %resp.status(),
"Controller rejected event push"
);
}
Err(e) => {
warn!(error = %e, "Failed to push events to controller");
}
}
}
Err(tokio::sync::broadcast::error::RecvError::Lagged(n)) => {
debug!(missed = n, "Event push loop lagged, skipping events");
}
Err(tokio::sync::broadcast::error::RecvError::Closed) => {
info!("Event bus closed, stopping event push loop");
break;
}
}
}
_ = shutdown_rx.changed() => {
info!("Event push loop shutting down");
break;
}
}
}
}
#[cfg(coverage)]
pub async fn run_event_push_loop(
_controller_url: String,
_controller_token: String,
_node_name: String,
_event_rx: tokio::sync::broadcast::Receiver<PulpoEvent>,
_shutdown_rx: tokio::sync::watch::Receiver<bool>,
) {
}
#[cfg_attr(coverage, allow(dead_code))]
pub const fn build_push_request(events: Vec<PulpoEvent>) -> EventPushRequest {
EventPushRequest { events }
}
#[cfg(test)]
mod tests {
use super::*;
use pulpo_common::event::SessionEvent;
#[cfg(coverage)]
#[tokio::test]
async fn test_coverage_stub_returns_immediately() {
let (_tx, rx) = tokio::sync::broadcast::channel::<PulpoEvent>(16);
let (_shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(false);
tokio::time::timeout(
std::time::Duration::from_secs(2),
run_event_push_loop(
"http://localhost:9999".into(),
"node-token".into(),
"test-node".into(),
rx,
shutdown_rx,
),
)
.await
.expect("coverage stub should return immediately");
}
#[test]
fn test_build_push_request() {
let events = vec![PulpoEvent::Session(SessionEvent {
session_id: "s1".into(),
session_name: "test".into(),
status: "active".into(),
previous_status: None,
node_name: "node-1".into(),
output_snippet: None,
timestamp: "2026-01-01T00:00:00Z".into(),
..Default::default()
})];
let req = build_push_request(events);
assert_eq!(req.events.len(), 1);
}
#[test]
fn test_build_push_request_empty() {
let req = build_push_request(vec![]);
assert!(req.events.is_empty());
}
#[cfg(not(coverage))]
#[test]
fn test_max_batch_size_constant() {
assert_eq!(MAX_BATCH_SIZE, 10);
}
}