use std::collections::HashMap;
use std::sync::OnceLock;
use std::time::Duration;
use async_trait::async_trait;
use rmpv::Value;
use crate::automation::http_client::{AsyncHttpRequest, ASYNC_HTTP_REQUEST};
use crate::envelope::EventEnvelope;
use crate::function::{AppError, ComposableFunction};
use crate::platform::Platform;
use crate::post_office::PostOffice;
use crate::util::app_config_reader::AppConfigReader;
use crate::util::config_reader::ConfigReader;
use crate::util::multi_level_map::ConfigValue;
use crate::util::w3c_trace;
pub const EVENT_API_SERVICE: &str = "event.api.service";
pub const X_EVENT_API: &str = "x-event-api";
const OCTET_STREAM: &str = "application/octet-stream";
const TEXT_EVENT_STREAM: &str = "text/event-stream";
const STREAM_CALLER_REQUIRED: &str =
"Streaming function requires a caller that accepts text/event-stream";
const X_TTL: &str = "x-ttl";
const X_ASYNC: &str = "x-async";
const X_NO_STREAM: &str = "x-small-payload-as-bytes";
const EVENT_OVER_HTTP_YAML: &str = "yaml.event.over.http";
const DEFAULT_EVENT_OVER_HTTP_YAML: &str = "classpath:/event-over-http.yaml";
const ASYNC_EVENT_HTTP_TIMEOUT: Duration = Duration::from_secs(60);
pub struct EventApiService {
platform: Platform,
}
impl EventApiService {
pub fn new(platform: &Platform) -> Self {
EventApiService {
platform: platform.clone(),
}
}
}
#[async_trait]
impl ComposableFunction for EventApiService {
async fn handle_event(
&self,
_headers: HashMap<String, String>,
input: EventEnvelope,
_instance: usize,
) -> Result<EventEnvelope, AppError> {
let Some(reply_to) = input.reply_to().map(str::to_string) else {
return Ok(EventEnvelope::new());
};
let context_id = input.correlation_id().unwrap_or_default().to_string();
let po = PostOffice::new(&self.platform);
if let Some(response) = self.dispatch(&po, &reply_to, &context_id, input).await? {
let _ = po
.send(response.set_to(&reply_to).set_correlation_id(&context_id))
.await;
}
Ok(EventEnvelope::new())
}
}
impl EventApiService {
async fn dispatch(
&self,
po: &PostOffice,
reply_to: &str,
context_id: &str,
input: EventEnvelope,
) -> Result<Option<EventEnvelope>, AppError> {
let request = AsyncHttpRequest::from_value(input.body());
let timeout_ms = request
.header(X_TTL)
.and_then(|v| v.parse::<u64>().ok())
.unwrap_or(0)
.max(1000);
let is_async = request.header(X_ASYNC) == Some("true");
let accepts_sse = request
.header("accept")
.is_some_and(|accept| accept.contains(TEXT_EVENT_STREAM));
let capable = accepts_sse && !is_async;
let answer = |status: i32, message: &str| -> EventEnvelope {
if capable {
EventEnvelope::new()
.set_status(status)
.set_raw_body(Value::from(message))
} else {
reply(status, error_envelope(status, message))
}
};
let Value::Binary(bytes) = request.body() else {
return Ok(Some(if capable {
answer(500, "Invalid event-over-http data format")
} else {
reply(500, b"Invalid event-over-http data format".to_vec())
}));
};
if is_compact_envelope(bytes) {
return Ok(Some(answer(
400,
"compact format not supported - set event.over.http.format=standard on the sender",
)));
}
let inner = match EventEnvelope::from_bytes(bytes) {
Ok(envelope) => envelope,
Err(e) => return Ok(Some(answer(400, e.message()))),
};
let Some(to) = inner
.to()
.map(|to| crate::platform::bare_route(to).to_string())
else {
return Ok(Some(answer(400, "Missing routing path")));
};
let mut inner = inner;
for (key, value) in request.session() {
inner = inner.set_header(key, value);
}
if !self.platform.has_route(&to) {
return Ok(Some(answer(404, &format!("Route {to} not found"))));
}
if self.platform.is_private(&to) == Some(true) {
return Ok(Some(answer(403, &format!("{to} is private"))));
}
if is_async {
po.send(inner).await?;
let ack = EventEnvelope::new()
.set_status(202)
.set_body(serde_json::json!({
"type": "async",
"delivered": true,
"time": crate::trace::iso8601_utc_now(),
}))?;
Ok(Some(reply(200, ack.to_bytes()?)))
} else if accepts_sse {
let inner = inner.set_reply_to(reply_to).set_correlation_id(context_id);
po.send(inner).await?;
Ok(None)
} else {
match po.request(inner, Duration::from_millis(timeout_ms)).await {
Ok(result) => {
if has_stream_marker(&result) {
Ok(Some(reply(
406,
error_envelope(406, STREAM_CALLER_REQUIRED),
)))
} else {
Ok(Some(reply(200, result.to_bytes()?)))
}
}
Err(e) => Ok(Some(reply(408, error_envelope(408, e.message())))),
}
}
}
}
fn has_stream_marker(event: &EventEnvelope) -> bool {
event
.headers()
.iter()
.any(|(name, _)| name.eq_ignore_ascii_case(crate::event_stream::X_EVENT_STREAM))
}
fn reply(http_status: i32, body: Vec<u8>) -> EventEnvelope {
EventEnvelope::new()
.set_status(http_status)
.set_header("content-type", OCTET_STREAM)
.set_raw_body(Value::Binary(body))
}
fn error_envelope(status: i32, message: &str) -> Vec<u8> {
EventEnvelope::new()
.set_status(status)
.set_raw_body(Value::from(message))
.to_bytes()
.unwrap_or_default()
}
fn is_compact_envelope(bytes: &[u8]) -> bool {
match rmp_serde::from_slice::<Value>(bytes) {
Ok(Value::Map(entries)) if !entries.is_empty() => entries
.iter()
.all(|(k, _)| k.as_str().is_some_and(|s| s.chars().count() == 1)),
_ => false,
}
}
pub async fn event_over_http(
po: &PostOffice,
endpoint: &str,
event: EventEnvelope,
timeout: Duration,
rpc: bool,
) -> Result<EventEnvelope, AppError> {
static NO_HEADERS: OnceLock<HashMap<String, String>> = OnceLock::new();
event_over_http_with_headers(
po,
endpoint,
event,
timeout,
rpc,
NO_HEADERS.get_or_init(HashMap::new),
)
.await
}
pub async fn event_over_http_with_headers(
po: &PostOffice,
endpoint: &str,
event: EventEnvelope,
timeout: Duration,
rpc: bool,
security_headers: &HashMap<String, String>,
) -> Result<EventEnvelope, AppError> {
let (host, path) = split_endpoint(endpoint)?;
let event = crate::post_office::apply_current_trace(event);
let trace_id = event.trace_id().map(str::to_string);
let span_id = event.span_id().map(str::to_string);
let payload = event.to_bytes()?;
let mut http = AsyncHttpRequest::new()
.set_method("POST")
.set_url(&path)
.set_target_host(&host)
.set_header("content-type", OCTET_STREAM)
.set_header(X_NO_STREAM, "true")
.set_header("accept", "*/*")
.set_header(X_TTL, &timeout.as_millis().max(1000).to_string())
.set_header(X_EVENT_API, "true")
.set_body(Value::Binary(payload));
if !rpc {
http = http.set_header(X_ASYNC, "true");
}
for (key, value) in security_headers {
http = http.set_header(key, value);
}
http = stamp_trace_headers(http, trace_id.as_deref(), span_id.as_deref());
let http_event = EventEnvelope::new()
.set_to(ASYNC_HTTP_REQUEST)
.set_raw_body(http.to_value());
let response = po
.request_direct(http_event, timeout + Duration::from_millis(100))
.await?;
match response.body() {
Value::Binary(bytes) => EventEnvelope::from_bytes(bytes),
_ => Ok(response),
}
}
fn stamp_trace_headers(
mut http: AsyncHttpRequest,
trace_id: Option<&str>,
span_id: Option<&str>,
) -> AsyncHttpRequest {
if let Some(trace_id) = trace_id {
http = http.set_header("x-trace-id", trace_id);
if let Some(span_id) = span_id {
if let Some(traceparent) = w3c_trace::format(trace_id, span_id) {
http = http.set_header(w3c_trace::TRACEPARENT, &traceparent);
let custom_traceparent = AppConfigReader::get_instance()
.get_property_or("http.traceparent.header", w3c_trace::TRACEPARENT);
if !custom_traceparent.eq_ignore_ascii_case(w3c_trace::TRACEPARENT) {
http = http.set_header(&custom_traceparent, &traceparent);
}
}
}
}
http
}
fn split_endpoint(endpoint: &str) -> Result<(String, String), AppError> {
let scheme_end = endpoint
.find("://")
.map(|i| i + 3)
.ok_or_else(|| AppError::new(400, format!("Invalid endpoint {endpoint}")))?;
match endpoint[scheme_end..].find('/') {
Some(offset) => {
let split = scheme_end + offset;
Ok((endpoint[..split].to_string(), endpoint[split..].to_string()))
}
None => Ok((endpoint.to_string(), "/api/event".to_string())),
}
}
#[derive(Debug)]
pub struct EventHttpTarget {
pub target: String,
pub headers: HashMap<String, String>,
}
fn event_http_registry() -> &'static HashMap<String, EventHttpTarget> {
static REGISTRY: OnceLock<HashMap<String, EventHttpTarget>> = OnceLock::new();
REGISTRY.get_or_init(load_event_http_routes)
}
pub fn get_event_http_target(route: &str) -> Option<&'static EventHttpTarget> {
let base = match route.find('@') {
Some(at) => &route[..at],
None => route,
};
event_http_registry().get(base)
}
fn load_event_http_routes() -> HashMap<String, EventHttpTarget> {
let mut targets = HashMap::new();
let explicit = AppConfigReader::get_instance().get_property(EVENT_OVER_HTTP_YAML);
let path = explicit
.clone()
.unwrap_or_else(|| DEFAULT_EVENT_OVER_HTTP_YAML.to_string());
let reader = match ConfigReader::load(&path) {
Ok(reader) => reader,
Err(e) => {
if explicit.is_some() {
log::error!("Unable to load event-over-http config - {e}");
}
return targets;
}
};
let Some(ConfigValue::List(entries)) = reader.get("event.http") else {
log::error!(
"Invalid config {path} - the event.http section should be a list of route and target"
);
return targets;
};
for i in 0..entries.len() {
let route = reader
.get_property(&format!("event.http[{i}].route"))
.unwrap_or_default();
let target = reader
.get_property(&format!("event.http[{i}].target"))
.unwrap_or_default();
if route.is_empty() || target.is_empty() {
continue;
}
if crate::platform::validate_route(&route).is_err() {
log::error!("Invalid Event over HTTP config entry - check route {route}");
continue;
}
if split_endpoint(&target).is_err() {
log::error!("Invalid Event over HTTP config entry - check target {target}");
continue;
}
let mut headers = HashMap::new();
if let Some(ConfigValue::Map(map)) = reader.get(&format!("event.http[{i}].headers")) {
for key in map.keys() {
if let Some(value) = reader.get_property(&format!("event.http[{i}].headers.{key}"))
{
headers.insert(key.clone(), value);
}
}
}
log::info!(
"Event-over-HTTP {route} -> {target} with {} header{}",
headers.len(),
if headers.len() == 1 { "" } else { "s" }
);
targets.insert(route, EventHttpTarget { target, headers });
}
log::info!(
"Total {} event-over-http target{} configured",
targets.len(),
if targets.len() == 1 { "" } else { "s" }
);
targets
}
pub(crate) fn send_with_event_http(
platform: &Platform,
event: EventEnvelope,
to: &str,
entry: &'static EventHttpTarget,
) -> Result<(), AppError> {
let callback = event.reply_to().map(str::to_string);
if let Some(callback) = &callback {
if accepts_event_stream_header(&event) {
return relay_event_stream(platform, event, to, entry, callback.clone());
}
}
let event_api_type = if callback.is_some() {
"callback"
} else {
"async"
};
let trace_id = event.trace_id().map(str::to_string);
let trace_path = event.trace_path().map(str::to_string);
let cid = event.correlation_id().map(str::to_string);
let forward = event
.clear_reply_to()
.set_header(X_EVENT_API, event_api_type);
let platform = platform.clone();
let to = to.to_string();
tokio::spawn(async move {
let po = PostOffice::new(&platform);
let outcome = event_over_http_with_headers(
&po,
&entry.target,
forward,
ASYNC_EVENT_HTTP_TIMEOUT,
callback.is_some(),
&entry.headers,
)
.await;
match outcome {
Ok(reply) => {
if let Some(callback) = callback {
let mut response = reply.set_to(&callback).clear_reply_to().set_from(&to);
if let (Some(id), Some(path)) = (&trace_id, &trace_path) {
response = response.set_trace(id, path);
}
if let Some(cid) = &cid {
response = response.set_correlation_id(cid);
}
if let Err(e) = po.send(response).await {
log::error!(
"Error in sending callback event {to} from {} to {callback} - {}",
entry.target,
e.message()
);
}
} else if reply.status() != 202 {
log::error!(
"Error in sending async event {to} to {} - status={}, error={}",
entry.target,
reply.status(),
reply
.body_as::<String>()
.unwrap_or_else(|_| format!("{}", reply.body()))
);
}
}
Err(e) => {
log::error!(
"Error in sending event {to} to {} - {}",
entry.target,
e.message()
);
}
}
});
Ok(())
}
fn accepts_event_stream_header(event: &EventEnvelope) -> bool {
event.headers().iter().any(|(name, value)| {
name.eq_ignore_ascii_case("accept") && value.contains(TEXT_EVENT_STREAM)
})
}
fn relay_event_stream(
platform: &Platform,
event: EventEnvelope,
to: &str,
entry: &'static EventHttpTarget,
callback: String,
) -> Result<(), AppError> {
let ttl_ms = event
.headers()
.iter()
.find(|(name, _)| name.eq_ignore_ascii_case(X_TTL))
.and_then(|(_, value)| value.trim().parse::<u64>().ok())
.map_or(ASYNC_EVENT_HTTP_TIMEOUT.as_millis() as u64, |v| v.max(1000));
let forward = crate::post_office::apply_current_trace(
event
.clear_reply_to()
.set_header(X_EVENT_API, crate::automation::http_client::STREAM_RELAY),
);
let cid = forward.correlation_id().map(str::to_string);
let trace_id = forward.trace_id().map(str::to_string);
let trace_path = forward.trace_path().map(str::to_string);
let span_id = forward.span_id().map(str::to_string);
let platform = platform.clone();
let to = to.to_string();
tokio::spawn(async move {
let (host, path) = match split_endpoint(&entry.target) {
Ok(parts) => parts,
Err(e) => {
log::error!(
"Unable to relay event stream {to} to {} - {}",
entry.target,
e.message()
);
return;
}
};
let payload = match forward.to_bytes() {
Ok(bytes) => bytes,
Err(e) => {
log::error!(
"Unable to relay event stream {to} to {} - {}",
entry.target,
e.message()
);
return;
}
};
let mut http = AsyncHttpRequest::new()
.set_method("POST")
.set_url(&path)
.set_target_host(&host)
.set_header("content-type", OCTET_STREAM)
.set_header(X_NO_STREAM, "true")
.set_header("accept", TEXT_EVENT_STREAM)
.set_header(X_TTL, &ttl_ms.to_string())
.set_header(X_EVENT_API, "true")
.set_body(Value::Binary(payload));
for (key, value) in &entry.headers {
http = http.set_header(key, value);
}
http = stamp_trace_headers(http, trace_id.as_deref(), span_id.as_deref());
let mut http_event = EventEnvelope::new()
.set_to(ASYNC_HTTP_REQUEST)
.set_raw_body(http.to_value())
.set_reply_to(&callback)
.set_from(&to)
.set_header(X_EVENT_API, crate::automation::http_client::STREAM_RELAY);
if let Some(cid) = &cid {
http_event = http_event.set_correlation_id(cid);
}
if let (Some(id), Some(path)) = (&trace_id, &trace_path) {
http_event = http_event.set_trace(id, path);
}
let po = PostOffice::new(&platform);
if let Err(e) = po.send(http_event).await {
log::error!(
"Unable to relay event stream {to} to {} - {}",
entry.target,
e.message()
);
}
});
Ok(())
}