use std::io;
use std::time::{Duration, Instant};
use super::protocol::DaemonResponse;
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum DispatchOutcome {
Success,
Error,
Degraded,
}
impl DispatchOutcome {
#[must_use]
pub fn from_response(response: &DaemonResponse) -> Self {
if response.error.is_some() {
Self::Error
} else if response.degraded_codes.is_empty() {
Self::Success
} else {
Self::Degraded
}
}
#[must_use]
pub fn label(self) -> &'static str {
match self {
Self::Success => "success",
Self::Error => "error",
Self::Degraded => "degraded",
}
}
}
pub trait DaemonMetricsCollector: Send + Sync {
fn record_dispatch(&self, method: &str, outcome: DispatchOutcome, elapsed: Duration);
fn record_accept_loop_terminated(&self, _kind: io::ErrorKind) {}
fn record_worker_spawn_failure(&self, _kind: io::ErrorKind) {}
fn record_stream_clone_failure(&self, _kind: io::ErrorKind) {}
fn record_handler_panic(&self, _method: &str) {}
}
#[derive(Clone, Copy, Debug, Default)]
pub struct NoopMetricsCollector;
impl DaemonMetricsCollector for NoopMetricsCollector {
#[inline]
fn record_dispatch(&self, _method: &str, _outcome: DispatchOutcome, _elapsed: Duration) {}
}
pub fn instrument_dispatch<F>(
method: &str,
collector: &dyn DaemonMetricsCollector,
dispatch_fn: F,
) -> DaemonResponse
where
F: FnOnce() -> DaemonResponse,
{
let start = Instant::now();
let response = dispatch_fn();
let elapsed = start.elapsed();
collector.record_dispatch(method, DispatchOutcome::from_response(&response), elapsed);
response
}
#[cfg(test)]
mod tests {
use std::io;
use std::sync::Mutex;
use std::time::Duration;
use super::super::protocol::{DaemonResponse, DaemonResponseError};
use super::{
DaemonMetricsCollector, DispatchOutcome, NoopMetricsCollector, instrument_dispatch,
};
#[derive(Default)]
struct CapturingCollector {
samples: Mutex<Vec<(String, DispatchOutcome)>>,
events: Mutex<Vec<CapturedEvent>>,
}
#[derive(Clone, Debug, Eq, PartialEq)]
enum CapturedEvent {
AcceptLoopTerminated(io::ErrorKind),
WorkerSpawnFailure(io::ErrorKind),
StreamCloneFailure(io::ErrorKind),
HandlerPanic(String),
}
impl DaemonMetricsCollector for CapturingCollector {
fn record_dispatch(&self, method: &str, outcome: DispatchOutcome, _elapsed: Duration) {
self.samples
.lock()
.expect("capturing collector mutex must not be poisoned")
.push((method.to_owned(), outcome));
}
fn record_accept_loop_terminated(&self, kind: io::ErrorKind) {
self.events
.lock()
.expect("capturing collector mutex must not be poisoned")
.push(CapturedEvent::AcceptLoopTerminated(kind));
}
fn record_worker_spawn_failure(&self, kind: io::ErrorKind) {
self.events
.lock()
.expect("capturing collector mutex must not be poisoned")
.push(CapturedEvent::WorkerSpawnFailure(kind));
}
fn record_stream_clone_failure(&self, kind: io::ErrorKind) {
self.events
.lock()
.expect("capturing collector mutex must not be poisoned")
.push(CapturedEvent::StreamCloneFailure(kind));
}
fn record_handler_panic(&self, method: &str) {
self.events
.lock()
.expect("capturing collector mutex must not be poisoned")
.push(CapturedEvent::HandlerPanic(method.to_owned()));
}
}
fn ok_response() -> DaemonResponse {
DaemonResponse::ok(
"req-1",
"agent-metrics-test",
None,
serde_json::json!({"ok": true}),
)
}
fn error_response() -> DaemonResponse {
DaemonResponse {
schema: super::super::DAEMON_RESPONSE_SCHEMA_V1.to_owned(),
request_id: "req-2".to_owned(),
agent_id: "agent-metrics-test".to_owned(),
workspace_id: None,
result: None,
error: Some(DaemonResponseError {
code: "daemon_unknown_method".to_owned(),
message: "nope".to_owned(),
}),
degraded_codes: Vec::new(),
delivery: None,
}
}
fn degraded_response() -> DaemonResponse {
DaemonResponse::ok("req-3", "agent-metrics-test", None, serde_json::Value::Null)
.with_degraded("daemon_overloaded")
}
#[test]
fn outcome_classifies_success_error_and_degraded() {
assert_eq!(
DispatchOutcome::from_response(&ok_response()),
DispatchOutcome::Success
);
assert_eq!(
DispatchOutcome::from_response(&error_response()),
DispatchOutcome::Error
);
assert_eq!(
DispatchOutcome::from_response(°raded_response()),
DispatchOutcome::Degraded
);
}
#[test]
fn error_dominates_even_when_degraded_codes_present() {
let mut response = error_response();
response.degraded_codes.push("daemon_overloaded".to_owned());
assert_eq!(
DispatchOutcome::from_response(&response),
DispatchOutcome::Error
);
}
#[test]
fn instrument_dispatch_records_method_and_outcome_and_returns_response() {
let collector = CapturingCollector::default();
let response = instrument_dispatch("ee.daemon.echo", &collector, ok_response);
assert_eq!(response.request_id, "req-1");
assert!(response.error.is_none());
let samples = collector.samples.lock().unwrap();
assert_eq!(samples.len(), 1);
assert_eq!(samples[0].0, "ee.daemon.echo");
assert_eq!(samples[0].1, DispatchOutcome::Success);
}
#[test]
fn observer_hooks_record_accept_spawn_clone_and_panic_events() {
let collector = CapturingCollector::default();
collector.record_accept_loop_terminated(io::ErrorKind::AddrInUse);
collector.record_worker_spawn_failure(io::ErrorKind::Other);
collector.record_stream_clone_failure(io::ErrorKind::UnexpectedEof);
collector.record_handler_panic("ee.daemon.context");
let events = collector.events.lock().unwrap();
assert_eq!(
events.as_slice(),
&[
CapturedEvent::AcceptLoopTerminated(io::ErrorKind::AddrInUse),
CapturedEvent::WorkerSpawnFailure(io::ErrorKind::Other),
CapturedEvent::StreamCloneFailure(io::ErrorKind::UnexpectedEof),
CapturedEvent::HandlerPanic("ee.daemon.context".to_owned()),
]
);
}
#[test]
fn noop_collector_records_nothing_observable_and_passes_response_through() {
let response = instrument_dispatch(
"ee.daemon.context",
&NoopMetricsCollector,
degraded_response,
);
assert_eq!(response.request_id, "req-3");
assert!(
response
.degraded_codes
.contains(&"daemon_overloaded".to_owned())
);
}
#[test]
fn outcome_labels_are_stable() {
assert_eq!(DispatchOutcome::Success.label(), "success");
assert_eq!(DispatchOutcome::Error.label(), "error");
assert_eq!(DispatchOutcome::Degraded.label(), "degraded");
}
}