Files
bsodfather 0e8f1ca03a
CI / Rust Checks (push) Failing after 4m28s
CI / UI Checks (push) Has been skipped
CI / Frontend E2E (push) Has been skipped
CI / Community Image Smoke (push) Has been skipped
CI / Deploy (push) Has been skipped
наблюдаемость: завершить базовый контур Community
Добавить структурированные журналы, метрики, трассировку и безопасный канал критических ошибок. Усилить границы рантайма, тесты, проверку зависимостей и сценарии развёртывания.
2026-07-31 01:01:14 +03:00

315 lines
11 KiB
Rust

use std::fmt;
use serde_json::{Map, Number, Value};
use time::{OffsetDateTime, format_description::well_known::Rfc3339};
use tracing::{Event, Subscriber, field::Visit};
use tracing_subscriber::{
EnvFilter, Layer,
filter::filter_fn,
fmt::{FmtContext, FormatEvent, FormatFields, MakeWriter, format::Writer},
layer::SubscriberExt,
registry::LookupSpan,
};
use crate::{
ObservabilityConfig, ObservabilityInitError, RedactionLimits, ServiceIdentity,
propagation::current_trace_id,
redaction::{redact_value, truncate_string},
schema::LogEnvelope,
};
pub fn build_subscriber<W>(
config: ObservabilityConfig,
writer: W,
) -> Result<impl Subscriber + Send + Sync, ObservabilityInitError>
where
W: for<'writer> MakeWriter<'writer> + Send + Sync + 'static,
{
build_subscriber_with_tracer(config, writer, None)
}
pub(crate) fn build_subscriber_with_tracer<W>(
config: ObservabilityConfig,
writer: W,
tracer: Option<opentelemetry_sdk::trace::SdkTracer>,
) -> Result<impl Subscriber + Send + Sync, ObservabilityInitError>
where
W: for<'writer> MakeWriter<'writer> + Send + Sync + 'static,
{
let (identity, filter, limits) = config.into_parts();
limits.validate()?;
let filter = EnvFilter::try_new(filter).map_err(|_| ObservabilityInitError::InvalidFilter)?;
let formatter = JsonEventFormatter::new(identity, limits);
let fmt_layer = tracing_subscriber::fmt::layer()
.with_ansi(false)
.event_format(formatter)
.with_writer(writer)
.with_filter(filter);
let otel_layer = tracer.map(|tracer| {
tracing_opentelemetry::layer()
.with_tracer(tracer)
.with_filter(filter_fn(|metadata| {
metadata.is_span() && metadata.target() == "crank::trace"
}))
});
Ok(tracing_subscriber::registry()
.with(fmt_layer)
.with(otel_layer))
}
#[derive(Clone, Debug)]
struct JsonEventFormatter {
identity: ServiceIdentity,
limits: RedactionLimits,
}
impl JsonEventFormatter {
fn new(identity: ServiceIdentity, limits: RedactionLimits) -> Self {
Self { identity, limits }
}
fn envelope(&self, event: &Event<'_>) -> Result<LogEnvelope, fmt::Error> {
let metadata = event.metadata();
let mut visitor = JsonFieldVisitor::default();
event.record(&mut visitor);
let mut raw_fields = visitor.fields;
let request_id = take_correlation_id(&mut raw_fields, "request_id")
.map(|value| truncate_string(&value, self.limits.max_string_bytes));
let trace_id = take_correlation_id(&mut raw_fields, "trace_id")
.or_else(current_trace_id)
.map(|value| truncate_string(&value, self.limits.max_string_bytes));
let cleaned = redact_value(&Value::Object(raw_fields), self.limits);
let fields = cleaned.as_object().cloned().unwrap_or_default();
let timestamp = OffsetDateTime::now_utc()
.format(&Rfc3339)
.map_err(|_| fmt::Error)?;
Ok(LogEnvelope {
timestamp,
level: metadata.level().as_str().to_owned(),
service: self.identity.service().to_owned(),
version: self.identity.version().to_owned(),
environment: self.identity.environment().to_owned(),
target: truncate_string(metadata.target(), self.limits.max_string_bytes),
event: truncate_string(metadata.name(), self.limits.max_string_bytes),
request_id,
trace_id,
fields,
})
}
fn serialize_bounded(&self, mut envelope: LogEnvelope) -> Result<String, fmt::Error> {
let line_budget = self.limits.max_event_bytes.saturating_sub(1);
let serialized = serde_json::to_string(&envelope).map_err(|_| fmt::Error)?;
if serialized.len() <= line_budget {
return Ok(serialized);
}
envelope.fields = Map::from_iter([("truncated".to_owned(), Value::Bool(true))]);
let fallback_string_limit = self.limits.max_string_bytes.min(64);
envelope.target = truncate_string(&envelope.target, fallback_string_limit);
envelope.event = truncate_string(&envelope.event, fallback_string_limit);
envelope.request_id = envelope
.request_id
.map(|value| truncate_string(&value, fallback_string_limit));
envelope.trace_id = envelope
.trace_id
.map(|value| truncate_string(&value, fallback_string_limit));
let serialized = serde_json::to_string(&envelope).map_err(|_| fmt::Error)?;
if serialized.len() <= line_budget {
return Ok(serialized);
}
envelope.request_id = None;
envelope.trace_id = None;
envelope.target = truncate_string(&envelope.target, 16);
envelope.event = truncate_string(&envelope.event, 16);
let serialized = serde_json::to_string(&envelope).map_err(|_| fmt::Error)?;
(serialized.len() <= line_budget)
.then_some(serialized)
.ok_or(fmt::Error)
}
}
impl<S, N> FormatEvent<S, N> for JsonEventFormatter
where
S: Subscriber + for<'lookup> LookupSpan<'lookup>,
N: for<'writer> FormatFields<'writer> + 'static,
{
fn format_event(
&self,
_ctx: &FmtContext<'_, S, N>,
mut writer: Writer<'_>,
event: &Event<'_>,
) -> fmt::Result {
let serialized = self.serialize_bounded(self.envelope(event)?)?;
writer.write_str(&serialized)?;
writer.write_char('\n')
}
}
#[derive(Default)]
struct JsonFieldVisitor {
fields: Map<String, Value>,
}
impl JsonFieldVisitor {
fn insert(&mut self, field: &tracing::field::Field, value: Value) {
self.fields.insert(field.name().to_owned(), value);
}
}
impl Visit for JsonFieldVisitor {
fn record_i64(&mut self, field: &tracing::field::Field, value: i64) {
self.insert(field, Value::Number(value.into()));
}
fn record_u64(&mut self, field: &tracing::field::Field, value: u64) {
self.insert(field, Value::Number(value.into()));
}
fn record_bool(&mut self, field: &tracing::field::Field, value: bool) {
self.insert(field, Value::Bool(value));
}
fn record_f64(&mut self, field: &tracing::field::Field, value: f64) {
let value = Number::from_f64(value)
.map(Value::Number)
.unwrap_or(Value::Null);
self.insert(field, value);
}
fn record_str(&mut self, field: &tracing::field::Field, value: &str) {
self.insert(field, Value::String(value.to_owned()));
}
fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn fmt::Debug) {
let rendered = format!("{value:?}");
let value = if is_correlation_field(field.name()) {
Value::String(debug_scalar(&rendered))
} else {
match serde_json::from_str(&rendered) {
Ok(value @ (Value::Object(_) | Value::Array(_))) => value,
_ if field.name() == "message" || is_safe_display_scalar(&rendered) => {
Value::String(rendered)
}
_ => Value::String(crate::REDACTED_MARKER.to_owned()),
}
};
self.insert(field, value);
}
}
fn take_correlation_id(fields: &mut Map<String, Value>, name: &str) -> Option<String> {
let value = fields.remove(name)?;
let value = match value {
Value::String(value) => value,
Value::Number(value) => value.to_string(),
Value::Bool(value) => value.to_string(),
Value::Null | Value::Array(_) | Value::Object(_) => return None,
};
(!value.is_empty()).then_some(value)
}
fn is_correlation_field(name: &str) -> bool {
matches!(name, "request_id" | "trace_id" | "correlation_id")
}
fn debug_scalar(rendered: &str) -> String {
serde_json::from_str::<String>(rendered).unwrap_or_else(|_| rendered.to_owned())
}
fn is_safe_display_scalar(rendered: &str) -> bool {
!rendered.is_empty()
&& rendered.bytes().all(|byte| {
byte.is_ascii_alphanumeric()
|| matches!(byte, b'-' | b'_' | b'.' | b':' | b'/' | b'+' | b'@')
})
}
#[cfg(test)]
mod tests {
use std::{
io,
sync::{Arc, Mutex},
};
use opentelemetry::{global, trace::TracerProvider as _};
use opentelemetry_sdk::{
error::OTelSdkResult,
propagation::TraceContextPropagator,
trace::{SdkTracerProvider, SpanData, SpanExporter},
};
use tracing::info;
use super::build_subscriber_with_tracer;
use crate::{
ObservabilityConfig, RedactionLimits, ServiceIdentity, inject_current_trace_context,
};
#[test]
fn trace_spans_ignore_the_log_level_filter() {
global::set_text_map_propagator(TraceContextPropagator::new());
let provider = SdkTracerProvider::builder().build();
let tracer = provider.tracer("trace-filter-test");
let subscriber = build_subscriber_with_tracer(test_config("warn"), io::sink, Some(tracer))
.expect("subscriber must build");
let dispatch = tracing::Dispatch::new(subscriber);
let _dispatch_guard = tracing::dispatcher::set_default(&dispatch);
let span = tracing::info_span!(target: "crank::trace", "http.request");
let _span_guard = span.enter();
let mut headers = axum::http::HeaderMap::new();
assert!(inject_current_trace_context(&mut headers));
assert!(headers.contains_key("traceparent"));
}
#[test]
fn otel_layer_does_not_export_events() {
let exported = Arc::new(Mutex::new(Vec::new()));
let provider = SdkTracerProvider::builder()
.with_simple_exporter(CapturingExporter(Arc::clone(&exported)))
.build();
let tracer = provider.tracer("event-filter-test");
let subscriber = build_subscriber_with_tracer(test_config("info"), io::sink, Some(tracer))
.expect("subscriber must build");
let dispatch = tracing::Dispatch::new(subscriber);
tracing::dispatcher::with_default(&dispatch, || {
let span = tracing::info_span!(target: "crank::trace", "http.request");
let _span_guard = span.enter();
info!(password = "canary-secret", "sensitive event");
});
provider.force_flush().expect("span must be exported");
let spans = exported.lock().expect("capture lock");
assert_eq!(spans.len(), 1);
assert!(spans[0].events.is_empty());
assert!(
!format!("{:?}", spans[0])
.as_bytes()
.windows(b"canary-secret".len())
.any(|window| window == b"canary-secret")
);
}
fn test_config(filter: &str) -> ObservabilityConfig {
ObservabilityConfig::new(
ServiceIdentity::try_new("admin-api", "test", "test").unwrap(),
filter,
RedactionLimits::default(),
)
}
#[derive(Clone, Debug)]
struct CapturingExporter(Arc<Mutex<Vec<SpanData>>>);
impl SpanExporter for CapturingExporter {
async fn export(&self, batch: Vec<SpanData>) -> OTelSdkResult {
self.0.lock().expect("capture lock").extend(batch);
Ok(())
}
}
}