use std::{error::Error, fmt, time::Instant};
use saddle_core::{CallContext, ErrorKind, SaddleError};
use serde_json::{Value, json};
use crate::{EventLevel, Observer, OutputStage, logger::LogRecord};
const MAX_IDENTITY_BYTES: usize = 256;
#[derive(Clone, Debug, Eq, PartialEq)]
pub enum ChainFieldError {
Empty,
TooLong,
ControlCharacter,
UnsafeAuthority,
ZeroAttempt,
}
impl fmt::Display for ChainFieldError {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.write_str(match self {
Self::Empty => "observability identity must not be empty",
Self::TooLong => "observability identity exceeds 256 bytes",
Self::ControlCharacter => "observability identity contains a control character",
Self::UnsafeAuthority => "outbound authority must be a credential-free target label",
Self::ZeroAttempt => "attempt must be greater than zero",
})
}
}
impl Error for ChainFieldError {}
macro_rules! safe_identity {
($name:ident) => {
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct $name(String);
impl $name {
pub fn new(value: impl Into<String>) -> Result<Self, ChainFieldError> {
validate_identity(value.into()).map(Self)
}
pub fn as_str(&self) -> &str {
&self.0
}
}
};
}
safe_identity!(RequestIdentity);
safe_identity!(RouteIdentity);
safe_identity!(BottleneckIdentity);
safe_identity!(RejectReason);
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct OutboundAuthority(String);
impl OutboundAuthority {
pub fn new(value: impl Into<String>) -> Result<Self, ChainFieldError> {
let value = validate_identity(value.into())?;
if value.contains(['@', '?', '#', '/', '\\']) || value.contains("://") {
return Err(ChainFieldError::UnsafeAuthority);
}
Ok(Self(value))
}
}
fn validate_identity(value: String) -> Result<String, ChainFieldError> {
if value.is_empty() {
return Err(ChainFieldError::Empty);
}
if value.len() > MAX_IDENTITY_BYTES {
return Err(ChainFieldError::TooLong);
}
if value.chars().any(char::is_control) {
return Err(ChainFieldError::ControlCharacter);
}
Ok(value)
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum Stage {
Ingress,
Admission,
Handler,
Database,
ProfuseContract,
Response,
ResourceFinalization,
}
impl Stage {
const fn as_str(self) -> &'static str {
match self {
Self::Ingress => "ingress",
Self::Admission => "admission",
Self::Handler => "handler",
Self::Database => "database",
Self::ProfuseContract => "profuse_contract",
Self::Response => "response",
Self::ResourceFinalization => "resource_finalization",
}
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum StageOutcome {
Success,
Failure,
Rejected,
Cancelled,
}
impl StageOutcome {
const fn as_str(self) -> &'static str {
match self {
Self::Success => "success",
Self::Failure => "failure",
Self::Rejected => "rejected",
Self::Cancelled => "cancelled",
}
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct EventContext {
request: RequestIdentity,
route: RouteIdentity,
attempt: u32,
}
impl EventContext {
pub fn new(
request: RequestIdentity,
route: RouteIdentity,
attempt: u32,
) -> Result<Self, ChainFieldError> {
if attempt == 0 {
return Err(ChainFieldError::ZeroAttempt);
}
Ok(Self {
request,
route,
attempt,
})
}
}
pub struct ActiveStage {
observer: Observer,
context: CallContext,
parent_rpc_id: String,
stage: Stage,
request: RequestIdentity,
route: RouteIdentity,
attempt: u32,
started_at: Instant,
finished: bool,
}
impl Observer {
pub fn start_stage(
&self,
parent: &CallContext,
stage: Stage,
event_context: EventContext,
) -> ActiveStage {
let context = CallContext::new(
parent.application().clone(),
parent.module().clone(),
parent.service().clone(),
parent.operation().clone(),
parent.trace_id(),
self.new_span_id(),
)
.with_trace_correlation_id(parent.trace_correlation_id().clone());
let active = ActiveStage {
observer: self.clone(),
context,
parent_rpc_id: parent.span_id().to_string(),
stage,
request: event_context.request,
route: event_context.route,
attempt: event_context.attempt,
started_at: Instant::now(),
finished: false,
};
active.emit_started();
active
}
pub fn record_capacity(
&self,
context: &CallContext,
event_context: &EventContext,
value: CapacityObservation,
) {
let mut record = base_record(context, EventLevel::Info, "framework.capacity", "admission");
add_event_context(&mut record, event_context);
if let Some(dimension) = value.dimension {
record.data.insert(
"capacity_dimension".into(),
Value::String(dimension.as_str().into()),
);
}
record.data.insert("budget".into(), json!(value.budget));
record.data.insert("limit".into(), json!(value.limit));
record.data.insert("used".into(), json!(value.used));
record
.data
.insert("elapsed_ms".into(), json!(value.elapsed_ms));
record
.data
.insert("bottleneck".into(), Value::String(value.bottleneck.0));
record.data.insert(
"outcome".into(),
Value::String(
if value.reject_reason.is_some() {
"rejected"
} else {
"accepted"
}
.into(),
),
);
if let Some(reason) = value.reject_reason {
record
.data
.insert("reject_reason".into(), Value::String(reason.0));
}
self.emit(record);
}
pub fn record_database_disposition(
&self,
context: &CallContext,
event_context: &EventContext,
disposition: DatabaseDisposition,
) {
let mut record = base_record(
context,
EventLevel::Info,
"framework.database.disposition",
"database",
);
add_event_context(&mut record, event_context);
record.data.insert(
"db_disposition".into(),
Value::String(disposition.as_str().into()),
);
record
.data
.insert("outcome".into(), Value::String(disposition.as_str().into()));
self.emit(record);
}
pub fn record_resource_finalization(
&self,
context: &CallContext,
event_context: &EventContext,
disposition: DatabaseDisposition,
elapsed_ms: u64,
) {
let mut record = base_record(
context,
EventLevel::Info,
"framework.resource.finalized",
"resource_finalization",
);
add_event_context(&mut record, event_context);
record.data.insert("elapsed_ms".into(), json!(elapsed_ms));
record
.data
.insert("credit".into(), Value::String("released".into()));
record.data.insert(
"db_disposition".into(),
Value::String(disposition.as_str().into()),
);
record
.data
.insert("outcome".into(), Value::String("success".into()));
self.emit(record);
}
pub fn record_outbound(
&self,
context: &CallContext,
event_context: &EventContext,
observation: OutboundObservation,
) {
let mut record = base_record(
context,
EventLevel::Info,
"framework.outbound",
"profuse_contract",
);
add_event_context(&mut record, event_context);
record
.data
.insert("zone".into(), Value::String(observation.zone.0));
record
.data
.insert("authority".into(), Value::String(observation.authority.0));
record.data.insert(
"outbound_result".into(),
Value::String(observation.result.as_str().into()),
);
record.data.insert(
"outcome".into(),
Value::String(observation.result.as_str().into()),
);
self.emit(record);
}
pub fn record_lifecycle(
&self,
context: &CallContext,
event_context: &EventContext,
state: LifecycleState,
health: Health,
) {
let mut record = base_record(
context,
if health == Health::Healthy {
EventLevel::Info
} else {
EventLevel::Error
},
"framework.lifecycle",
"resource_finalization",
);
add_event_context(&mut record, event_context);
record
.data
.insert("lifecycle".into(), Value::String(state.as_str().into()));
record
.data
.insert("health".into(), Value::String(health.as_str().into()));
record
.data
.insert("outcome".into(), Value::String(health.as_str().into()));
self.emit(record);
}
pub fn record_logger_health(
&self,
context: &CallContext,
event_context: &EventContext,
dropped: u64,
failure: Option<OutputStage>,
) {
let mut record = base_record(
context,
if failure.is_some() || dropped > 0 {
EventLevel::Error
} else {
EventLevel::Info
},
"framework.logger.health",
"logger",
);
add_event_context(&mut record, event_context);
record.data.insert("logger_dropped".into(), json!(dropped));
record.data.insert(
"logger_health".into(),
Value::String(
if failure.is_some() {
"output_failed"
} else {
"healthy"
}
.into(),
),
);
record.data.insert(
"outcome".into(),
Value::String(
if failure.is_some() {
"failure"
} else {
"success"
}
.into(),
),
);
if let Some(stage) = failure {
record.data.insert(
"output_failed".into(),
Value::String(format!("{stage:?}").to_ascii_lowercase()),
);
}
self.emit(record);
}
}
impl ActiveStage {
pub fn context(&self) -> &CallContext {
&self.context
}
pub fn succeed(mut self) {
self.finish(StageOutcome::Success, None);
}
pub fn reject(mut self, error: &SaddleError) {
self.finish(StageOutcome::Rejected, Some(error));
}
pub fn fail(mut self, error: &SaddleError) {
self.finish(StageOutcome::Failure, Some(error));
}
fn emit_started(&self) {
let mut record = self.record(EventLevel::Info, "framework.stage.started");
record
.data
.insert("outcome".into(), Value::String("started".into()));
self.observer.emit(record);
}
fn finish(&mut self, outcome: StageOutcome, error: Option<&SaddleError>) {
if self.finished {
return;
}
let mut record = self.record(
if outcome == StageOutcome::Success {
EventLevel::Info
} else {
EventLevel::Error
},
"framework.stage.finished",
);
record
.data
.insert("outcome".into(), Value::String(outcome.as_str().into()));
record.data.insert(
"elapsed_ms".into(),
json!(u64::try_from(self.started_at.elapsed().as_millis()).unwrap_or(u64::MAX)),
);
if let Some(error) = error {
record
.data
.insert("error_code".into(), Value::String(error.code().to_owned()));
record.data.insert(
"error_kind".into(),
Value::String(error_kind(error.kind()).into()),
);
}
self.observer.emit(record);
self.finished = true;
}
fn record(&self, level: EventLevel, event: &'static str) -> LogRecord {
let mut record = base_record(&self.context, level, event, self.stage.as_str());
record.parent = Some(self.parent_rpc_id.clone());
record.parent_span_id = Some(self.parent_rpc_id.clone());
record.data.insert(
"request_identity".into(),
Value::String(self.request.0.clone()),
);
record
.data
.insert("route".into(), Value::String(self.route.0.clone()));
record.data.insert("attempt".into(), json!(self.attempt));
record.data.insert("elapsed_ms".into(), json!(0));
record
.data
.insert("error_code".into(), Value::String("none".into()));
record
}
}
impl Drop for ActiveStage {
fn drop(&mut self) {
self.finish(StageOutcome::Cancelled, None);
}
}
pub struct CapacityObservation {
dimension: Option<CapacityDimension>,
budget: u64,
limit: u64,
used: u64,
bottleneck: BottleneckIdentity,
reject_reason: Option<RejectReason>,
elapsed_ms: u64,
}
impl CapacityObservation {
pub fn accepted(budget: u64, limit: u64, used: u64, bottleneck: BottleneckIdentity) -> Self {
Self {
dimension: None,
budget,
limit,
used,
bottleneck,
reject_reason: None,
elapsed_ms: 0,
}
}
pub fn rejected(
budget: u64,
limit: u64,
used: u64,
bottleneck: BottleneckIdentity,
reason: RejectReason,
) -> Self {
Self {
dimension: None,
budget,
limit,
used,
bottleneck,
reject_reason: Some(reason),
elapsed_ms: 0,
}
}
pub fn accepted_dimension(
dimension: CapacityDimension,
budget: u64,
limit: u64,
used: u64,
bottleneck: BottleneckIdentity,
elapsed_ms: u64,
) -> Self {
Self {
dimension: Some(dimension),
budget,
limit,
used,
bottleneck,
reject_reason: None,
elapsed_ms,
}
}
pub fn rejected_dimension(
dimension: CapacityDimension,
budget: u64,
limit: u64,
used: u64,
bottleneck: BottleneckIdentity,
reason: RejectReason,
elapsed_ms: u64,
) -> Self {
Self {
dimension: Some(dimension),
budget,
limit,
used,
bottleneck,
reject_reason: Some(reason),
elapsed_ms,
}
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum CapacityDimension {
Cpu,
Memory,
Database,
ProfuseContract,
}
impl CapacityDimension {
const fn as_str(self) -> &'static str {
match self {
Self::Cpu => "cpu",
Self::Memory => "memory",
Self::Database => "database",
Self::ProfuseContract => "profuse_contract",
}
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum DatabaseDisposition {
NotUsed,
Returned,
Discarded,
}
impl DatabaseDisposition {
const fn as_str(self) -> &'static str {
match self {
Self::NotUsed => "not_used",
Self::Returned => "returned",
Self::Discarded => "discarded",
}
}
}
pub struct OutboundObservation {
zone: RouteIdentity,
authority: OutboundAuthority,
result: OutboundResult,
}
impl OutboundObservation {
pub fn new(zone: RouteIdentity, authority: OutboundAuthority, result: OutboundResult) -> Self {
Self {
zone,
authority,
result,
}
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum OutboundResult {
Success,
Failure,
Rejected,
Timeout,
}
impl OutboundResult {
const fn as_str(self) -> &'static str {
match self {
Self::Success => "success",
Self::Failure => "failure",
Self::Rejected => "rejected",
Self::Timeout => "timeout",
}
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum LifecycleState {
Starting,
Running,
Draining,
Stopped,
}
impl LifecycleState {
const fn as_str(self) -> &'static str {
match self {
Self::Starting => "starting",
Self::Running => "running",
Self::Draining => "draining",
Self::Stopped => "stopped",
}
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum Health {
Healthy,
Degraded,
Failed,
}
impl Health {
const fn as_str(self) -> &'static str {
match self {
Self::Healthy => "healthy",
Self::Degraded => "degraded",
Self::Failed => "failed",
}
}
}
fn base_record(
context: &CallContext,
level: EventLevel,
event: &'static str,
stage: &'static str,
) -> LogRecord {
let mut record = LogRecord::new(level, event);
record.trace_id = Some(context.trace_correlation_id().to_string());
record.span = Some(stage.into());
record.span_id = Some(context.span_id().to_string());
record
.data
.insert("timestamp".into(), json!(record.timestamp_unix_ms));
record
.data
.insert("stage".into(), Value::String(stage.into()));
record.data.insert(
"rpc_id".into(),
Value::String(context.span_id().to_string()),
);
record.data.insert(
"request".into(),
Value::String(context.operation().to_string()),
);
record
}
fn add_event_context(record: &mut LogRecord, context: &EventContext) {
record.data.insert(
"request_identity".into(),
Value::String(context.request.0.clone()),
);
record
.data
.insert("route".into(), Value::String(context.route.0.clone()));
record.data.insert("attempt".into(), json!(context.attempt));
record.data.insert("elapsed_ms".into(), json!(0));
record
.data
.insert("error_code".into(), Value::String("none".into()));
}
const fn error_kind(kind: ErrorKind) -> &'static str {
match kind {
ErrorKind::InvalidArgument => "invalid_argument",
ErrorKind::NotFound => "not_found",
ErrorKind::Conflict => "conflict",
ErrorKind::Business => "business",
ErrorKind::Unavailable => "unavailable",
ErrorKind::Infrastructure => "infrastructure",
ErrorKind::Internal => "internal",
_ => "unknown",
}
}
#[cfg(test)]
mod tests {
use std::{
future::Future,
io,
sync::{Arc, Mutex},
task::{Context, Poll, Wake, Waker},
thread,
};
use super::*;
use crate::ObserverConfig;
#[derive(Clone, Default)]
struct Capture(Arc<Mutex<Vec<u8>>>);
impl io::Write for Capture {
fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
self.0.lock().unwrap().extend_from_slice(bytes);
Ok(bytes.len())
}
fn flush(&mut self) -> io::Result<()> {
Ok(())
}
}
struct ThreadWaker(thread::Thread);
impl Wake for ThreadWaker {
fn wake(self: Arc<Self>) {
self.0.unpark();
}
}
fn block_on<T>(future: impl Future<Output = T>) -> T {
let waker = Waker::from(Arc::new(ThreadWaker(thread::current())));
let mut context = Context::from_waker(&waker);
let mut future = std::pin::pin!(future);
loop {
match future.as_mut().poll(&mut context) {
Poll::Ready(output) => return output,
Poll::Pending => thread::park(),
}
}
}
fn event_context(attempt: u32) -> EventContext {
EventContext::new(
RequestIdentity::new("request-7").unwrap(),
RouteIdentity::new("orders.create").unwrap(),
attempt,
)
.unwrap()
}
fn records(capture: &Capture) -> Vec<Value> {
String::from_utf8(capture.0.lock().unwrap().clone())
.unwrap()
.lines()
.map(|line| serde_json::from_str(line).unwrap())
.collect()
}
#[test]
fn stage_chain_is_correlated_ordered_and_has_one_terminal() {
let capture = Capture::default();
let observer = Observer::with_writer(ObserverConfig::default(), capture.clone()).unwrap();
let (root, _) = observer.start_external_call(
"shop",
"entry",
"orders",
"create",
Some("00112233445566778899aabbccddeeff"),
);
let stages = [
Stage::Ingress,
Stage::Admission,
Stage::Handler,
Stage::Database,
Stage::ProfuseContract,
Stage::Response,
Stage::ResourceFinalization,
];
for (index, stage) in stages.into_iter().enumerate() {
observer
.start_stage(root.context(), stage, event_context((index + 1) as u32))
.succeed();
}
drop(observer.start_stage(root.context(), Stage::Handler, event_context(8)));
root.succeed();
block_on(observer.flush()).unwrap();
let records = records(&capture);
let stage_records: Vec<_> = records
.iter()
.filter(|value| {
value["event"]
.as_str()
.unwrap()
.starts_with("framework.stage.")
})
.collect();
assert_eq!(stage_records.len(), 16);
for pair in stage_records.chunks_exact(2) {
assert_eq!(pair[0]["event"], "framework.stage.started");
assert_eq!(pair[1]["event"], "framework.stage.finished");
assert_eq!(pair[0]["rpc_id"], pair[1]["rpc_id"]);
assert_eq!(pair[0]["trace_id"], "00112233445566778899aabbccddeeff");
for field in [
"timestamp_unix_ms",
"timestamp",
"level",
"event",
"stage",
"trace_id",
"rpc_id",
"request_identity",
"route",
"attempt",
"outcome",
"error_code",
] {
assert!(pair[0].get(field).is_some(), "missing {field}");
}
assert!(pair[1].get("elapsed_ms").is_some());
}
assert_eq!(stage_records.last().unwrap()["outcome"], "cancelled");
}
#[test]
fn closed_observations_have_common_fields_and_no_sensitive_payload() {
let capture = Capture::default();
let observer = Observer::with_writer(ObserverConfig::default(), capture.clone()).unwrap();
let (root, _) = observer.start_external_call("shop", "entry", "orders", "create", None);
let common = event_context(1);
observer.record_capacity(
root.context(),
&common,
CapacityObservation::rejected(
100,
80,
80,
BottleneckIdentity::new("request_slots").unwrap(),
RejectReason::new("limit_reached").unwrap(),
),
);
observer.record_database_disposition(
root.context(),
&common,
DatabaseDisposition::Returned,
);
observer.record_outbound(
root.context(),
&common,
OutboundObservation::new(
RouteIdentity::new("cn-hz-a").unwrap(),
OutboundAuthority::new("inventory-service").unwrap(),
OutboundResult::Success,
),
);
observer.record_lifecycle(
root.context(),
&common,
LifecycleState::Draining,
Health::Healthy,
);
observer.record_logger_health(root.context(), &common, 3, Some(OutputStage::Record));
root.succeed();
let _ = block_on(observer.flush());
let records: Vec<_> = records(&capture)
.into_iter()
.filter(|value| {
matches!(
value["event"].as_str(),
Some(
"framework.capacity"
| "framework.database.disposition"
| "framework.outbound"
| "framework.lifecycle"
| "framework.logger.health"
)
)
})
.collect();
assert_eq!(records.len(), 5);
for record in records {
for field in [
"timestamp_unix_ms",
"timestamp",
"level",
"event",
"stage",
"trace_id",
"rpc_id",
"request_identity",
"route",
"attempt",
"elapsed_ms",
"outcome",
"error_code",
] {
assert!(record.get(field).is_some(), "missing {field}");
}
let encoded = serde_json::to_string(&record).unwrap();
for forbidden in [
"request_body",
"response_body",
"cookie",
"session",
"token",
"password",
"connection_string",
"db_value",
] {
assert!(!encoded.contains(forbidden));
}
}
}
#[test]
fn typed_capacity_and_resource_terminal_preserve_closed_schema() {
let capture = Capture::default();
let observer = Observer::with_writer(ObserverConfig::default(), capture.clone()).unwrap();
let (root, _) = observer.start_external_call("shop", "entry", "orders", "create", None);
let common = event_context(1);
observer.record_capacity(
root.context(),
&common,
CapacityObservation::rejected_dimension(
CapacityDimension::Database,
4,
2,
2,
BottleneckIdentity::new("database").unwrap(),
RejectReason::new("at_limit").unwrap(),
17,
),
);
observer.record_resource_finalization(
root.context(),
&common,
DatabaseDisposition::NotUsed,
23,
);
root.succeed();
block_on(observer.flush()).unwrap();
let records = records(&capture);
let capacity = records
.iter()
.find(|record| record["event"] == "framework.capacity")
.unwrap();
assert_eq!(capacity["capacity_dimension"], "database");
assert_eq!(capacity["budget"], 4);
assert_eq!(capacity["limit"], 2);
assert_eq!(capacity["used"], 2);
assert_eq!(capacity["bottleneck"], "database");
assert_eq!(capacity["reject_reason"], "at_limit");
assert_eq!(capacity["elapsed_ms"], 17);
let terminal = records
.iter()
.find(|record| record["event"] == "framework.resource.finalized")
.unwrap();
assert_eq!(terminal["stage"], "resource_finalization");
assert_eq!(terminal["credit"], "released");
assert_eq!(terminal["db_disposition"], "not_used");
assert_eq!(terminal["elapsed_ms"], 23);
assert_eq!(terminal["outcome"], "success");
assert_eq!(capacity["trace_id"], terminal["trace_id"]);
assert_eq!(capacity["rpc_id"], terminal["rpc_id"]);
assert_eq!(capacity["request_identity"], terminal["request_identity"]);
assert_eq!(capacity["route"], terminal["route"]);
assert_eq!(capacity["attempt"], terminal["attempt"]);
}
#[test]
fn unsafe_or_unbounded_identifiers_and_zero_attempt_are_rejected() {
assert_eq!(RequestIdentity::new(""), Err(ChainFieldError::Empty));
assert_eq!(
RouteIdentity::new("x".repeat(257)),
Err(ChainFieldError::TooLong)
);
assert_eq!(
RejectReason::new("bad\nreason"),
Err(ChainFieldError::ControlCharacter)
);
for unsafe_value in ["https://user:secret@host/path", "host/path", "host?token=x"] {
assert_eq!(
OutboundAuthority::new(unsafe_value),
Err(ChainFieldError::UnsafeAuthority)
);
}
assert_eq!(
EventContext::new(
RequestIdentity::new("r").unwrap(),
RouteIdentity::new("route").unwrap(),
0
),
Err(ChainFieldError::ZeroAttempt),
);
}
}