use std::net::SocketAddr;
use axum::{routing::post, Json, Router};
use serde_json::Value;
use tokio::sync::mpsc;
use super::traits::IncomingMessage;
pub type WebhookParser = fn(&Value) -> Vec<IncomingMessage>;
pub fn json_route(path: &str, parse: WebhookParser, tx: mpsc::Sender<IncomingMessage>) -> Router {
Router::new().route(
path,
post(move |Json(body): Json<Value>| {
let tx = tx.clone();
async move {
for incoming in parse(&body) {
if tx.send(incoming).await.is_err() {
break;
}
}
axum::http::StatusCode::OK
}
}),
)
}
pub async fn spawn(addr: SocketAddr, app: Router) -> anyhow::Result<()> {
let listener = tokio::net::TcpListener::bind(addr)
.await
.map_err(|e| anyhow::anyhow!("webhook receiver could not bind {addr}: {e}"))?;
tokio::spawn(async move {
if let Err(e) = axum::serve(listener, app).await {
tracing::error!("webhook receiver on {addr} stopped: {e}");
}
});
Ok(())
}
pub async fn serve_json(
addr: SocketAddr,
path: &str,
parse: WebhookParser,
tx: mpsc::Sender<IncomingMessage>,
) -> anyhow::Result<()> {
spawn(addr, json_route(path, parse, tx)).await
}