Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 8 additions & 3 deletions crates/core/src/api/llm.rs
Original file line number Diff line number Diff line change
Expand Up @@ -706,6 +706,9 @@ fn emit_optimization_marks_with<F>(
/// This emits an LLM-start event after applying sanitize-request guardrails to
/// the payload recorded for observability.
///
/// If a sanitizer errors or panics, Relay omits the payload and request
/// annotation and does not run remaining sanitizers.
///
/// # Parameters
/// - `name`: Logical provider or model family name recorded on the span.
/// - `request`: Raw [`LlmRequest`] associated with the span.
Expand Down Expand Up @@ -927,12 +930,14 @@ async fn build_llm_end_payload(
/// # Errors
/// Returns an error when the runtime owner check fails or internal state cannot
/// be read safely. Dispatcher submission failures are logged because
/// observability publication is best effort. Sanitizer and response-codec errors
/// discovered during queued publication are also logged and fail open.
/// observability publication is best effort. Sanitizer errors discovered during
/// queued publication are logged and fail closed by omitting the governed payload.
/// Response-codec errors retain their documented fallback behavior.
///
/// # Notes
/// Sanitize-response guardrails affect only the emitted end-event payload, not
/// the caller-owned `response` value.
/// the caller-owned `response` value. If a sanitizer errors or panics, Relay
/// omits the payload and response annotation and does not run remaining sanitizers.
pub fn llm_call_end(params: LlmCallEndParams<'_>) -> Result<()> {
ensure_runtime_owner()?;
let scope_stack = params.handle.captured_scope_stack().clone();
Expand Down
173 changes: 91 additions & 82 deletions crates/core/src/api/runtime/state.rs
Original file line number Diff line number Diff line change
Expand Up @@ -20,8 +20,9 @@ use std::task::{Context, Poll};
use futures_util::{FutureExt, Stream};

use crate::api::event::{
BaseEvent, CategoryProfile, Event, EventCategory, MarkEvent, ScopeCategory, ScopeEvent,
llm_attributes_to_strings, scope_attributes_to_strings, tool_attributes_to_strings,
BaseEvent, CategoryProfile, Event, EventCategory, EventSanitizeFields, MarkEvent,
ScopeCategory, ScopeEvent, llm_attributes_to_strings, scope_attributes_to_strings,
tool_attributes_to_strings,
};
use crate::api::llm::{CreateLlmHandleParams, EndLlmHandleParams};
use crate::api::llm::{LlmHandle, LlmRequest};
Expand Down Expand Up @@ -810,20 +811,28 @@ impl NemoRelayContextState {
event = Arc::try_unwrap(context).unwrap_or_else(|context| (*context).clone());
match outcome {
Ok(Ok(fields)) => event.apply_sanitize_fields(fields),
Ok(Err(error)) => log::error!(
target: "nemo_relay.runtime",
event = "event_sanitizer_failed",
sanitizer = entry.name.as_str(),
event_name = event.name();
"Event sanitizer failed; preserving the last valid event snapshot: {error}"
),
Err(_) => log::error!(
target: "nemo_relay.runtime",
event = "event_sanitizer_panicked",
sanitizer = entry.name.as_str(),
event_name = event.name();
"Event sanitizer panicked; publishing the latest valid event snapshot"
),
Ok(Err(_error)) => {
log::error!(
target: "nemo_relay.runtime",
event = "event_sanitizer_failed",
sanitizer = entry.name.as_str(),
event_name = event.name();
"Event sanitizer failed; clearing observability fields"
);
event.apply_sanitize_fields(EventSanitizeFields::default());
break;
}
Err(_) => {
log::error!(
target: "nemo_relay.runtime",
event = "event_sanitizer_panicked",
sanitizer = entry.name.as_str(),
event_name = event.name();
"Event sanitizer panicked; clearing observability fields"
);
event.apply_sanitize_fields(EventSanitizeFields::default());
break;
}
}
}
event
Expand Down Expand Up @@ -856,36 +865,38 @@ impl NemoRelayContextState {
/// - `entries`: Sanitizer snapshots to evaluate.
///
/// # Returns
/// The sanitized JSON payload after every provided guardrail has run.
/// The sanitized JSON payload after every provided guardrail has run, or
/// `None` when a sanitizer failure omits the payload.
pub(crate) async fn tool_sanitize_request_snapshot_chain(
name: &str,
args: Json,
entries: &[Guardrail<ToolSanitizeFn>],
) -> Json {
let mut value = args;
) -> Option<Json> {
let mut value = Some(args);
for entry in entries {
let callback = Arc::clone(&entry.payload);
let callback_name = name.to_string();
let current = value.clone();
match AssertUnwindSafe(async move { callback(callback_name, current).await })
.catch_unwind()
.await
{
Ok(Ok(next)) => value = next,
Ok(Err(error)) => log::error!(
target: "nemo_relay.runtime",
event = "tool_request_sanitizer_failed",
sanitizer = entry.name.as_str(),
tool_name = name;
"Tool request sanitizer failed; preserving the last valid payload: {error}"
),
Err(_) => log::error!(
target: "nemo_relay.runtime",
event = "tool_request_sanitizer_panicked",
sanitizer = entry.name.as_str(),
tool_name = name;
"Tool request sanitizer panicked; preserving the last valid payload"
),
if let Some(current) = value.take() {
let callback = Arc::clone(&entry.payload);
let callback_name = name.to_string();
match AssertUnwindSafe(async move { callback(callback_name, current).await })
.catch_unwind()
.await
{
Ok(Ok(next)) => value = Some(next),
Ok(Err(_error)) => log::error!(
target: "nemo_relay.runtime",
event = "tool_request_sanitizer_failed",
sanitizer = entry.name.as_str(),
tool_name = name;
"Tool request sanitizer failed; omitting the observability payload"
),
Err(_) => log::error!(
target: "nemo_relay.runtime",
event = "tool_request_sanitizer_panicked",
sanitizer = entry.name.as_str(),
tool_name = name;
"Tool request sanitizer panicked; omitting the observability payload"
),
}
}
}
value
Expand Down Expand Up @@ -918,36 +929,38 @@ impl NemoRelayContextState {
/// - `entries`: Sanitizer snapshots to evaluate.
///
/// # Returns
/// The sanitized JSON payload after every provided guardrail has run.
/// The sanitized JSON payload after every provided guardrail has run, or
/// `None` when a sanitizer failure omits the payload.
pub(crate) async fn tool_sanitize_response_snapshot_chain(
name: &str,
result: Json,
entries: &[Guardrail<ToolSanitizeFn>],
) -> Json {
let mut value = result;
) -> Option<Json> {
let mut value = Some(result);
for entry in entries {
let callback = Arc::clone(&entry.payload);
let callback_name = name.to_string();
let current = value.clone();
match AssertUnwindSafe(async move { callback(callback_name, current).await })
.catch_unwind()
.await
{
Ok(Ok(next)) => value = next,
Ok(Err(error)) => log::error!(
target: "nemo_relay.runtime",
event = "tool_response_sanitizer_failed",
sanitizer = entry.name.as_str(),
tool_name = name;
"Tool response sanitizer failed; preserving the last valid payload: {error}"
),
Err(_) => log::error!(
target: "nemo_relay.runtime",
event = "tool_response_sanitizer_panicked",
sanitizer = entry.name.as_str(),
tool_name = name;
"Tool response sanitizer panicked; preserving the last valid payload"
),
if let Some(current) = value.take() {
let callback = Arc::clone(&entry.payload);
let callback_name = name.to_string();
match AssertUnwindSafe(async move { callback(callback_name, current).await })
.catch_unwind()
.await
{
Ok(Ok(next)) => value = Some(next),
Ok(Err(_error)) => log::error!(
target: "nemo_relay.runtime",
event = "tool_response_sanitizer_failed",
sanitizer = entry.name.as_str(),
tool_name = name;
"Tool response sanitizer failed; omitting the observability payload"
),
Err(_) => log::error!(
target: "nemo_relay.runtime",
event = "tool_response_sanitizer_panicked",
sanitizer = entry.name.as_str(),
tool_name = name;
"Tool response sanitizer panicked; omitting the observability payload"
),
}
}
}
value
Expand Down Expand Up @@ -1226,7 +1239,8 @@ impl NemoRelayContextState {
/// - `entries`: Sanitizer snapshots to evaluate.
///
/// # Returns
/// The sanitized [`LlmRequest`] after every provided guardrail has run.
/// The sanitized [`LlmRequest`] after every provided guardrail has run, or
/// `None` when a sanitizer errors or panics.
pub(crate) async fn llm_sanitize_request_snapshot_chain(
request: LlmRequest,
context: LlmSanitizeRequestContext,
Expand All @@ -1245,24 +1259,21 @@ impl NemoRelayContextState {
.await
{
Ok(Ok(next)) => value = next,
Ok(Err(error)) => {
Ok(Err(_error)) => {
log::error!(
target: "nemo_relay.runtime",
event = "llm_request_sanitizer_failed",
sanitizer = entry.name.as_str(),
preserved_value = "last_valid_request";
"LLM request sanitizer failed; preserving the last valid request: {error}"
sanitizer = entry.name.as_str();
"LLM request sanitizer failed; omitting the observability payload"
);
value = Some(current);
}
Err(_) => {
log::error!(
target: "nemo_relay.runtime",
event = "llm_request_sanitizer_panicked",
sanitizer = entry.name.as_str();
"LLM request sanitizer panicked; preserving the last valid request"
"LLM request sanitizer panicked; omitting the observability payload"
);
value = Some(current);
}
}
}
Expand Down Expand Up @@ -1296,7 +1307,8 @@ impl NemoRelayContextState {
/// - `entries`: Sanitizer snapshots to evaluate.
///
/// # Returns
/// The sanitized response payload after every provided guardrail has run.
/// The sanitized response payload after every provided guardrail has run,
/// or `None` when a sanitizer errors or panics.
pub(crate) async fn llm_sanitize_response_snapshot_chain(
response: Json,
context: LlmSanitizeResponseContext,
Expand All @@ -1315,24 +1327,21 @@ impl NemoRelayContextState {
.await
{
Ok(Ok(next)) => value = next,
Ok(Err(error)) => {
Ok(Err(_error)) => {
log::error!(
target: "nemo_relay.runtime",
event = "llm_response_sanitizer_failed",
sanitizer = entry.name.as_str(),
preserved_value = "last_valid_response";
"LLM response sanitizer failed; preserving the last valid response: {error}"
sanitizer = entry.name.as_str();
"LLM response sanitizer failed; omitting the observability payload"
);
value = Some(current);
}
Err(_) => {
log::error!(
target: "nemo_relay.runtime",
event = "llm_response_sanitizer_panicked",
sanitizer = entry.name.as_str();
"LLM response sanitizer panicked; preserving the last valid response"
"LLM response sanitizer panicked; omitting the observability payload"
);
value = Some(current);
}
}
}
Expand Down
20 changes: 12 additions & 8 deletions crates/core/src/api/runtime/subscriber_dispatcher.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,7 @@

//! Asynchronous subscriber delivery for native targets.

use crate::api::event::Event;
use crate::api::event::{Event, EventSanitizeFields};
use crate::api::registry::Guardrail;
use crate::api::runtime::{
EventSanitizeFn, EventSubscriberFn, NemoRelayContextState, ScopeStackHandle,
Expand Down Expand Up @@ -967,8 +967,8 @@ mod native {

/// Apply a transform and sanitizers on the dispatcher thread. A transform
/// failure drops the event because it may be responsible for inserting the
/// sanitized payload. A sanitizer failure retains the transformed snapshot
/// and continues publication (fail open).
/// sanitized payload. A sanitizer failure clears mutable observability
/// fields before publication.
pub(super) fn sanitize_event_snapshot(
event: Event,
transform: Option<EventTransformFn>,
Expand Down Expand Up @@ -1037,10 +1037,12 @@ mod native {
}
log::error!(
target: "nemo_relay.runtime",
event = "event_sanitizer_fail_open";
"Publishing the transformed event snapshot because event sanitizers could not run"
event = "event_sanitizer_runtime_failed";
"Event sanitizers could not run; clearing observability fields before publication"
);
return (Some(transformed), nested_publications);
let mut cleared = transformed;
cleared.apply_sanitize_fields(EventSanitizeFields::default());
return (Some(cleared), nested_publications);
}
};
let fallback = transformed.clone();
Expand All @@ -1057,9 +1059,11 @@ mod native {
log::error!(
target: "nemo_relay.runtime",
event = "event_sanitizer_panicked";
"Event sanitizer panicked; preserving the last valid event snapshot"
"Event sanitizer panicked; clearing observability fields"
);
Some(fallback)
let mut cleared = fallback;
cleared.apply_sanitize_fields(EventSanitizeFields::default());
Some(cleared)
}
};
(event, nested_publications)
Expand Down
18 changes: 13 additions & 5 deletions crates/core/src/api/shared.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@ use std::sync::Arc;

use uuid::Uuid;

use crate::api::event::{Event, ScopeCategory};
use crate::api::event::{Event, EventSanitizeFields, ScopeCategory};
use crate::api::llm::LlmRequest;
use crate::api::registry::Guardrail;
use crate::api::runtime::global_context;
Expand Down Expand Up @@ -61,9 +61,13 @@ pub(crate) fn snapshot_event_sanitizers(
log::error!(
target: "nemo_relay.runtime",
event = "event_sanitizer_snapshot_failed";
"Event sanitizer snapshot failed open because the scope stack lock is poisoned; publishing without event sanitizers: {error}"
"Event sanitizer snapshot failed; clearing observability fields: {error}"
);
return None;
return Some(vec![Guardrail::new(
"event-sanitizer-snapshot-failure",
i32::MIN,
Arc::new(|_, _| Box::pin(async { Ok(EventSanitizeFields::default()) })),
)]);
}
};
let context = global_context();
Expand All @@ -73,9 +77,13 @@ pub(crate) fn snapshot_event_sanitizers(
log::error!(
target: "nemo_relay.runtime",
event = "event_sanitizer_snapshot_failed";
"Event sanitizer snapshot failed open because the runtime context lock is poisoned; publishing without event sanitizers: {error}"
"Event sanitizer snapshot failed; clearing observability fields: {error}"
);
return None;
return Some(vec![Guardrail::new(
"event-sanitizer-snapshot-failure",
i32::MIN,
Arc::new(|_, _| Box::pin(async { Ok(EventSanitizeFields::default()) })),
)]);
}
};
match &event {
Expand Down
Loading
Loading