use serde_json::{Value, json};
use std::sync::Arc;
use turbomcp_core::{JsonRpcNotification, LogLevel};
use turbomcp_protocol::methods;
use crate::subscriptions::request_writer;
#[derive(Clone, Debug)]
pub struct LogSender {
inner: Option<Arc<Inner>>,
}
#[derive(Debug)]
struct Inner {
min: LogLevel,
connection: String,
session: String,
}
impl LogSender {
#[must_use]
pub(crate) fn disabled() -> Self {
Self { inner: None }
}
#[must_use]
pub(crate) fn new(min: LogLevel, connection: String, session: String) -> Self {
Self {
inner: Some(Arc::new(Inner {
min,
connection,
session,
})),
}
}
#[must_use]
pub fn would_log(&self, level: LogLevel) -> bool {
self.inner.as_ref().is_some_and(|i| level >= i.min)
}
pub async fn log_with(&self, level: LogLevel, logger: Option<&str>, data: Value) {
let Some(inner) = &self.inner else { return };
if level < inner.min {
return;
}
let Some(writer) = request_writer(&inner.connection, &inner.session) else {
tracing::debug!("no stream for log notification; dropped");
return;
};
let mut params = json!({ "level": level, "data": data });
if let (Some(obj), Some(logger)) = (params.as_object_mut(), logger) {
obj.insert("logger".to_owned(), json!(logger));
}
if writer
.send(JsonRpcNotification::new(methods::notification::MESSAGE, Some(params)).into())
.await
.is_err()
{
tracing::debug!("request stream closed before log notification; dropped");
}
}
pub async fn log(&self, level: LogLevel, data: Value) {
self.log_with(level, None, data).await;
}
pub async fn debug(&self, data: Value) {
self.log(LogLevel::Debug, data).await;
}
pub async fn info(&self, data: Value) {
self.log(LogLevel::Info, data).await;
}
pub async fn warning(&self, data: Value) {
self.log(LogLevel::Warning, data).await;
}
pub async fn error(&self, data: Value) {
self.log(LogLevel::Error, data).await;
}
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
async fn disabled_sender_is_inert() {
let sender = LogSender::disabled();
assert!(!sender.would_log(LogLevel::Emergency));
sender.error(json!("nope")).await; }
#[tokio::test]
async fn severity_filter_and_shape() {
let (tx, mut rx) = tokio::sync::mpsc::channel(8);
let _guard = turbomcp_service::outbound::register("log-test-conn", tx);
let sender = LogSender::new(LogLevel::Info, "log-test-conn".into(), String::new());
assert!(!sender.would_log(LogLevel::Debug));
assert!(sender.would_log(LogLevel::Error));
sender.debug(json!("filtered")).await;
sender
.log_with(LogLevel::Error, Some("db"), json!({ "oops": true }))
.await;
let only = rx.try_recv().expect("error message sent");
assert!(rx.try_recv().is_err(), "debug was filtered");
let turbomcp_core::JsonRpcMessage::Notification(n) = only else {
panic!("expected a notification");
};
assert_eq!(n.method, "notifications/message");
let params = n.params.unwrap();
assert_eq!(params["level"], "error");
assert_eq!(params["logger"], "db");
assert_eq!(params["data"]["oops"], true);
}
}