use std::io::Write;
use anyhow::Result;
use rmcp::{
service::{RunningService, ServiceExt},
RoleServer,
};
use tracing_subscriber::EnvFilter;
use super::OmniDevServer;
use crate::utils::env::{EnvSource, SystemEnv};
fn resolve_log_directive(env: &impl EnvSource, settings_log_level: Option<&str>) -> String {
env.var("RUST_LOG")
.filter(|s| !s.is_empty())
.or_else(|| {
settings_log_level
.filter(|s| !s.is_empty())
.map(str::to_string)
})
.unwrap_or_else(|| "warn".to_string())
}
pub fn try_init_tracing(settings_log_level: Option<&str>) -> Result<()> {
let directive = resolve_log_directive(&SystemEnv, settings_log_level);
tracing_subscriber::fmt()
.with_writer(std::io::stderr)
.with_ansi(false)
.with_env_filter(EnvFilter::new(directive))
.try_init()
.map_err(|e| anyhow::anyhow!("tracing subscriber already set: {e}"))?;
Ok(())
}
pub async fn serve_with<T, E, A>(transport: T) -> Result<()>
where
T: rmcp::transport::IntoTransport<RoleServer, E, A>,
E: std::error::Error + Send + Sync + 'static,
{
let service: RunningService<RoleServer, OmniDevServer> =
OmniDevServer::new().serve(transport).await?;
service.waiting().await?;
Ok(())
}
pub fn feature_flags() -> &'static str {
"mcp"
}
pub fn log_startup_event() {
let version = env!("CARGO_PKG_VERSION");
let features = feature_flags();
tracing::info!(version, features, "starting omni-dev MCP server");
}
pub fn write_error_chain<W: Write>(writer: &mut W, err: &anyhow::Error) -> std::io::Result<()> {
writeln!(writer, "Error: {err}")?;
let mut source = err.source();
while let Some(inner) = source {
writeln!(writer, " Caused by: {inner}")?;
source = inner.source();
}
Ok(())
}
#[cfg(test)]
#[allow(clippy::unwrap_used, clippy::expect_used)]
mod tests {
use super::*;
use crate::test_support::env::MapEnv;
use anyhow::{anyhow, Context};
#[test]
fn write_error_chain_single_error() {
let err = anyhow!("only failure");
let mut buf = Vec::new();
write_error_chain(&mut buf, &err).unwrap();
let out = String::from_utf8(buf).unwrap();
assert_eq!(out, "Error: only failure\n");
}
#[test]
fn write_error_chain_preserves_chain() {
let result: Result<(), anyhow::Error> =
Err(anyhow!("root")).context("middle").context("outermost");
let err = result.expect_err("constructed Err");
let mut buf = Vec::new();
write_error_chain(&mut buf, &err).unwrap();
let out = String::from_utf8(buf).unwrap();
assert!(out.starts_with("Error: outermost\n"), "got: {out:?}");
assert!(out.contains(" Caused by: middle\n"));
assert!(out.contains(" Caused by: root\n"));
}
#[tokio::test]
async fn serve_with_handles_peer_disconnect() {
let (server_transport, client_transport) = tokio::io::duplex(4096);
let server_handle = tokio::spawn(async move { serve_with(server_transport).await });
drop(client_transport);
let result = server_handle.await.unwrap();
let _ = result;
}
#[test]
fn feature_flags_includes_mcp() {
let flags = feature_flags();
assert!(
flags.contains("mcp"),
"expected feature flags to include mcp, got {flags:?}"
);
}
#[test]
fn log_startup_event_does_not_panic() {
log_startup_event();
}
#[test]
fn try_init_tracing_is_idempotent_or_errors() {
let _ = try_init_tracing(None);
let second = try_init_tracing(Some("info"));
assert!(second.is_err(), "second init should report already-set");
}
#[test]
fn resolve_log_directive_rust_log_beats_settings() {
let env = MapEnv::new().with("RUST_LOG", "debug");
assert_eq!(resolve_log_directive(&env, Some("info")), "debug");
}
#[test]
fn resolve_log_directive_settings_used_when_rust_log_unset() {
let env = MapEnv::new();
assert_eq!(resolve_log_directive(&env, Some("info")), "info");
}
#[test]
fn resolve_log_directive_defaults_to_warn() {
let env = MapEnv::new();
assert_eq!(resolve_log_directive(&env, None), "warn");
}
#[test]
fn resolve_log_directive_ignores_empty_values() {
let env = MapEnv::new().with("RUST_LOG", "");
assert_eq!(resolve_log_directive(&env, Some("trace")), "trace");
assert_eq!(resolve_log_directive(&env, Some("")), "warn");
}
}