diff --git a/crates/core/src/api/runtime/state.rs b/crates/core/src/api/runtime/state.rs index bb0aa6a16..bfc6c5479 100644 --- a/crates/core/src/api/runtime/state.rs +++ b/crates/core/src/api/runtime/state.rs @@ -1788,8 +1788,7 @@ fn validate_event_metadata_attributes(attributes: &BTreeMap) -> Re #[derive(Clone, Copy, PartialEq, Eq)] enum OtelAttributePrimitiveKind { Boolean, - Integer, - Double, + Number, String, } @@ -1799,12 +1798,12 @@ fn otel_compatible_attribute_number_kind( if let Some(value) = value.as_u64() { return i64::try_from(value) .is_ok() - .then_some(OtelAttributePrimitiveKind::Integer); + .then_some(OtelAttributePrimitiveKind::Number); } if value.as_i64().is_some() { - return Some(OtelAttributePrimitiveKind::Integer); + return Some(OtelAttributePrimitiveKind::Number); } - value.as_f64().map(|_| OtelAttributePrimitiveKind::Double) + value.as_f64().map(|_| OtelAttributePrimitiveKind::Number) } fn is_otel_compatible_attribute_number(value: &serde_json::Number) -> bool { diff --git a/crates/core/tests/unit/runtime_state_tests.rs b/crates/core/tests/unit/runtime_state_tests.rs index 08403222a..31714d844 100644 --- a/crates/core/tests/unit/runtime_state_tests.rs +++ b/crates/core/tests/unit/runtime_state_tests.rs @@ -41,6 +41,7 @@ async fn event_metadata_injection_accepts_flat_otel_values_and_empty_output() { ("nv.test.booleans".into(), json!([true, false])), ("nv.test.integers".into(), json!([1, 2])), ("nv.test.doubles".into(), json!([1.0, 2.5])), + ("nv.test.numbers".into(), json!([1, 2.5])), ("nv.test.empty".into(), json!([])), ])) }) @@ -71,6 +72,7 @@ async fn event_metadata_injection_accepts_flat_otel_values_and_empty_output() { assert_eq!(metadata["nv.test.booleans"], json!([true, false])); assert_eq!(metadata["nv.test.integers"], json!([1, 2])); assert_eq!(metadata["nv.test.doubles"], json!([1.0, 2.5])); + assert_eq!(metadata["nv.test.numbers"], json!([1, 2.5])); assert_eq!(metadata["nv.test.empty"], json!([])); } @@ -94,7 +96,6 @@ async fn event_metadata_injection_rejects_invalid_output_atomically() { BTreeMap::from([("nv.test.object".into(), json!({"nested": true}))]), BTreeMap::from([("nv.test.nested_list".into(), json!([[1]]))]), BTreeMap::from([("nv.test.mixed_list".into(), json!([1, "two"]))]), - BTreeMap::from([("nv.test.mixed_numbers".into(), json!([1, 2.5]))]), BTreeMap::from([("nv.test.oversized_number".into(), json!(u64::MAX))]), BTreeMap::from([("nv.test.oversized_list".into(), json!([u64::MAX]))]), ]; diff --git a/crates/ffi/nemo_relay.h b/crates/ffi/nemo_relay.h index 0d2bc97ef..c38b3cba7 100644 --- a/crates/ffi/nemo_relay.h +++ b/crates/ffi/nemo_relay.h @@ -252,12 +252,13 @@ typedef struct Option_NemoRelayFinalizerCb Option_NemoRelayFinalizerCb; typedef struct Option_NemoRelayPluginValidateCb Option_NemoRelayPluginValidateCb; /** - * Callback for mark and scope event sanitizers. - * The returned JSON string transfers to Relay and is freed exactly once. + * Callback for event metadata injection. + * + * The returned string must contain a JSON object whose properties are proposed + * metadata additions. It transfers to Relay and is freed exactly once. Return + * null after setting the last error message to report a callback failure. */ -typedef char *(*NemoRelayEventSanitizeCb)(void *user_data, - const struct FfiEvent *event, - const char *fields_json); +typedef char *(*NemoRelayEventMetadataInjectorCb)(void *user_data, const struct FfiEvent *event); /** * Optional destructor for user data passed to callbacks. @@ -269,6 +270,14 @@ typedef char *(*NemoRelayEventSanitizeCb)(void *user_data, */ typedef void (*NemoRelayFreeFn)(void *user_data); +/** + * Callback for mark and scope event sanitizers. + * The returned JSON string transfers to Relay and is freed exactly once. + */ +typedef char *(*NemoRelayEventSanitizeCb)(void *user_data, + const struct FfiEvent *event, + const char *fields_json); + /** * Callback for LLM execution (default callable). Receives a native JSON C string, * returns the response as a JSON C string. @@ -871,6 +880,50 @@ NemoRelayStatus nemo_relay_adaptive_build_cache_telemetry_event(const char *opti */ NemoRelayStatus nemo_relay_adaptive_set_latency_sensitivity(uint32_t value); +/** + * Register a global event metadata injector. + * + * # Safety + * Pointers must remain valid for the documented call lifetime. The callback + * and user data remain owned by Relay until deregistration. + */ +NemoRelayStatus nemo_relay_register_event_metadata_injector(const char *name, + int32_t priority, + NemoRelayEventMetadataInjectorCb cb, + void *user_data, + NemoRelayFreeFn free_fn); + +/** + * Deregister a global event metadata injector. + * + * # Safety + * `name` must be a valid C string. + */ +NemoRelayStatus nemo_relay_deregister_event_metadata_injector(const char *name); + +/** + * Register an event metadata injector owned by an active scope. + * + * # Safety + * Pointers must remain valid for the documented call lifetime. The callback + * and user data remain owned by Relay until deregistration or scope cleanup. + */ +NemoRelayStatus nemo_relay_scope_register_event_metadata_injector(const char *scope_uuid, + const char *name, + int32_t priority, + NemoRelayEventMetadataInjectorCb cb, + void *user_data, + NemoRelayFreeFn free_fn); + +/** + * Deregister an event metadata injector owned by an active scope. + * + * # Safety + * String pointers must be valid C strings. + */ +NemoRelayStatus nemo_relay_scope_deregister_event_metadata_injector(const char *scope_uuid, + const char *name); + /** * Register a global mark event sanitizer. * # Safety @@ -2057,6 +2110,20 @@ NemoRelayStatus nemo_relay_plugin_context_register_subscriber(struct FfiPluginCo void *user_data, NemoRelayFreeFn free_fn); +/** + * Register an event metadata injector into a plugin context. + * + * # Safety + * Pointers must remain valid for the documented call lifetime. The callback + * and user data remain owned by the plugin registration until rollback. + */ +NemoRelayStatus nemo_relay_plugin_context_register_event_metadata_injector(struct FfiPluginContext *ctx, + const char *name, + int32_t priority, + NemoRelayEventMetadataInjectorCb cb, + void *user_data, + NemoRelayFreeFn free_fn); + /** * Register a mark event sanitizer into a plugin context. * # Safety diff --git a/crates/ffi/src/api/event_registry.rs b/crates/ffi/src/api/event_registry.rs index 13af2b8d8..e3416aa50 100644 --- a/crates/ffi/src/api/event_registry.rs +++ b/crates/ffi/src/api/event_registry.rs @@ -2,8 +2,9 @@ // SPDX-License-Identifier: Apache-2.0 use super::{ - NemoRelayEventSanitizeCb, NemoRelayFreeFn, NemoRelayStatus, c_char, c_str_to_string, - clear_last_error, core_registry_api, set_last_error, status_from_error, wrap_event_sanitize_fn, + NemoRelayEventMetadataInjectorCb, NemoRelayEventSanitizeCb, NemoRelayFreeFn, NemoRelayStatus, + c_char, c_str_to_string, clear_last_error, core_registry_api, set_last_error, + status_from_error, wrap_event_metadata_injector_fn, wrap_event_sanitize_fn, }; #[derive(Clone, Copy)] @@ -67,6 +68,108 @@ fn parse_scope_uuid(value: *const c_char) -> Result }) } +/// Register a global event metadata injector. +/// +/// # Safety +/// Pointers must remain valid for the documented call lifetime. The callback +/// and user data remain owned by Relay until deregistration. +#[unsafe(no_mangle)] +pub unsafe extern "C" fn nemo_relay_register_event_metadata_injector( + name: *const c_char, + priority: i32, + cb: NemoRelayEventMetadataInjectorCb, + user_data: *mut libc::c_void, + free_fn: NemoRelayFreeFn, +) -> NemoRelayStatus { + clear_last_error(); + let Some(cb) = cb else { + set_last_error("event metadata injector callback is null"); + return NemoRelayStatus::NullPointer; + }; + let callback = wrap_event_metadata_injector_fn(cb, user_data, free_fn); + let name = match c_str_to_string(name) { + Ok(value) => value, + Err(status) => return status, + }; + core_registry_api::register_event_metadata_injector(&name, priority, callback) + .map(|()| NemoRelayStatus::Ok) + .unwrap_or_else(|error| status_from_error(&error)) +} + +/// Deregister a global event metadata injector. +/// +/// # Safety +/// `name` must be a valid C string. +#[unsafe(no_mangle)] +pub unsafe extern "C" fn nemo_relay_deregister_event_metadata_injector( + name: *const c_char, +) -> NemoRelayStatus { + clear_last_error(); + let name = match c_str_to_string(name) { + Ok(value) => value, + Err(status) => return status, + }; + core_registry_api::deregister_event_metadata_injector(&name) + .map(|_| NemoRelayStatus::Ok) + .unwrap_or_else(|error| status_from_error(&error)) +} + +/// Register an event metadata injector owned by an active scope. +/// +/// # Safety +/// Pointers must remain valid for the documented call lifetime. The callback +/// and user data remain owned by Relay until deregistration or scope cleanup. +#[unsafe(no_mangle)] +pub unsafe extern "C" fn nemo_relay_scope_register_event_metadata_injector( + scope_uuid: *const c_char, + name: *const c_char, + priority: i32, + cb: NemoRelayEventMetadataInjectorCb, + user_data: *mut libc::c_void, + free_fn: NemoRelayFreeFn, +) -> NemoRelayStatus { + clear_last_error(); + let Some(cb) = cb else { + set_last_error("event metadata injector callback is null"); + return NemoRelayStatus::NullPointer; + }; + let callback = wrap_event_metadata_injector_fn(cb, user_data, free_fn); + let uuid = match parse_scope_uuid(scope_uuid) { + Ok(value) => value, + Err(status) => return status, + }; + let name = match c_str_to_string(name) { + Ok(value) => value, + Err(status) => return status, + }; + core_registry_api::scope_register_event_metadata_injector(&uuid, &name, priority, callback) + .map(|()| NemoRelayStatus::Ok) + .unwrap_or_else(|error| status_from_error(&error)) +} + +/// Deregister an event metadata injector owned by an active scope. +/// +/// # Safety +/// String pointers must be valid C strings. +#[unsafe(no_mangle)] +pub unsafe extern "C" fn nemo_relay_scope_deregister_event_metadata_injector( + scope_uuid: *const c_char, + name: *const c_char, +) -> NemoRelayStatus { + clear_last_error(); + let uuid = match parse_scope_uuid(scope_uuid) { + Ok(value) => value, + Err(status) => return status, + }; + let name = match c_str_to_string(name) { + Ok(value) => value, + Err(status) => return status, + }; + core_registry_api::scope_deregister_event_metadata_injector(&uuid, &name) + .map(|_| NemoRelayStatus::Ok) + .unwrap_or_else(|error| status_from_error(&error)) +} + unsafe fn register_scope( scope_uuid: *const c_char, name: *const c_char, diff --git a/crates/ffi/src/api/mod.rs b/crates/ffi/src/api/mod.rs index d201a7475..e09e82d09 100644 --- a/crates/ffi/src/api/mod.rs +++ b/crates/ffi/src/api/mod.rs @@ -14,12 +14,13 @@ use std::sync::{Arc, OnceLock}; use std::time::Duration; use crate::callable::{ - NemoRelayCodecDecodeFn, NemoRelayCodecEncodeFn, NemoRelayCollectorCb, NemoRelayEventSanitizeCb, - NemoRelayEventSubscriberCb, NemoRelayFinalizerCb, NemoRelayFreeFn, NemoRelayLlmConditionalCb, - NemoRelayLlmExecCb, NemoRelayLlmExecInterceptCb, NemoRelayLlmRequestInterceptCb, - NemoRelayLlmSanitizeRequestCb, NemoRelayLlmSanitizeResponseCb, NemoRelayPluginRegisterCb, - NemoRelayPluginValidateCb, NemoRelayToolConditionalCb, NemoRelayToolExecCb, - NemoRelayToolExecInterceptCb, NemoRelayToolSanitizeCb, wrap_codec_fn, wrap_collector_fn, + NemoRelayCodecDecodeFn, NemoRelayCodecEncodeFn, NemoRelayCollectorCb, + NemoRelayEventMetadataInjectorCb, NemoRelayEventSanitizeCb, NemoRelayEventSubscriberCb, + NemoRelayFinalizerCb, NemoRelayFreeFn, NemoRelayLlmConditionalCb, NemoRelayLlmExecCb, + NemoRelayLlmExecInterceptCb, NemoRelayLlmRequestInterceptCb, NemoRelayLlmSanitizeRequestCb, + NemoRelayLlmSanitizeResponseCb, NemoRelayPluginRegisterCb, NemoRelayPluginValidateCb, + NemoRelayToolConditionalCb, NemoRelayToolExecCb, NemoRelayToolExecInterceptCb, + NemoRelayToolSanitizeCb, wrap_codec_fn, wrap_collector_fn, wrap_event_metadata_injector_fn, wrap_event_sanitize_fn, wrap_event_subscriber, wrap_finalizer_fn, wrap_llm_conditional_fn, wrap_llm_exec_fn, wrap_llm_exec_intercept_fn, wrap_llm_request_intercept_fn, wrap_llm_sanitize_request_fn, wrap_llm_sanitize_response_fn, wrap_llm_stream_exec_fn, diff --git a/crates/ffi/src/api/plugin.rs b/crates/ffi/src/api/plugin.rs index 09cd151f8..6fb7d8f41 100644 --- a/crates/ffi/src/api/plugin.rs +++ b/crates/ffi/src/api/plugin.rs @@ -3,20 +3,21 @@ use super::{ Arc, CStr, ConfigDiagnostic, DiagnosticLevel, DynamicPluginActivationSpec, FfiPluginActivation, - FfiPluginContext, Future, NemoRelayEventSanitizeCb, NemoRelayEventSubscriberCb, - NemoRelayFreeFn, NemoRelayLlmConditionalCb, NemoRelayLlmExecInterceptCb, - NemoRelayLlmRequestInterceptCb, NemoRelayLlmSanitizeRequestCb, NemoRelayLlmSanitizeResponseCb, - NemoRelayPluginRegisterCb, NemoRelayPluginValidateCb, NemoRelayStatus, - NemoRelayToolConditionalCb, NemoRelayToolExecInterceptCb, NemoRelayToolSanitizeCb, Pin, Plugin, - PluginConfig, PluginError, PluginHostActivation, PluginRegistrationContext, - active_plugin_report, c_char, c_str_to_json, c_str_to_string, clear_last_error, - clear_plugin_configuration, deregister_plugin, initialize_plugins, json_to_c_string, - last_error_message, list_plugin_kinds, nemo_relay_string_free, register_adaptive_component, - register_plugin, set_last_error, status_from_plugin_error, tokio_runtime, - validate_plugin_config, wrap_event_sanitize_fn, wrap_event_subscriber, wrap_llm_conditional_fn, - wrap_llm_exec_intercept_fn, wrap_llm_request_intercept_fn, wrap_llm_sanitize_request_fn, - wrap_llm_sanitize_response_fn, wrap_llm_stream_exec_intercept_fn, wrap_tool_conditional_fn, - wrap_tool_exec_intercept_fn, wrap_tool_request_intercept_fn, wrap_tool_sanitize_fn, + FfiPluginContext, Future, NemoRelayEventMetadataInjectorCb, NemoRelayEventSanitizeCb, + NemoRelayEventSubscriberCb, NemoRelayFreeFn, NemoRelayLlmConditionalCb, + NemoRelayLlmExecInterceptCb, NemoRelayLlmRequestInterceptCb, NemoRelayLlmSanitizeRequestCb, + NemoRelayLlmSanitizeResponseCb, NemoRelayPluginRegisterCb, NemoRelayPluginValidateCb, + NemoRelayStatus, NemoRelayToolConditionalCb, NemoRelayToolExecInterceptCb, + NemoRelayToolSanitizeCb, Pin, Plugin, PluginConfig, PluginError, PluginHostActivation, + PluginRegistrationContext, active_plugin_report, c_char, c_str_to_json, c_str_to_string, + clear_last_error, clear_plugin_configuration, deregister_plugin, initialize_plugins, + json_to_c_string, last_error_message, list_plugin_kinds, nemo_relay_string_free, + register_adaptive_component, register_plugin, set_last_error, status_from_plugin_error, + tokio_runtime, validate_plugin_config, wrap_event_metadata_injector_fn, wrap_event_sanitize_fn, + wrap_event_subscriber, wrap_llm_conditional_fn, wrap_llm_exec_intercept_fn, + wrap_llm_request_intercept_fn, wrap_llm_sanitize_request_fn, wrap_llm_sanitize_response_fn, + wrap_llm_stream_exec_intercept_fn, wrap_tool_conditional_fn, wrap_tool_exec_intercept_fn, + wrap_tool_request_intercept_fn, wrap_tool_sanitize_fn, }; use crate::api::event_registry::Surface; use nemo_relay_pii_redaction::component::register_pii_redaction_component; @@ -542,6 +543,40 @@ pub unsafe extern "C" fn nemo_relay_plugin_context_register_subscriber( } } +/// Register an event metadata injector into a plugin context. +/// +/// # Safety +/// Pointers must remain valid for the documented call lifetime. The callback +/// and user data remain owned by the plugin registration until rollback. +#[unsafe(no_mangle)] +pub unsafe extern "C" fn nemo_relay_plugin_context_register_event_metadata_injector( + ctx: *mut FfiPluginContext, + name: *const c_char, + priority: i32, + cb: NemoRelayEventMetadataInjectorCb, + user_data: *mut libc::c_void, + free_fn: NemoRelayFreeFn, +) -> NemoRelayStatus { + clear_last_error(); + if ctx.is_null() { + set_last_error("plugin context is null"); + return NemoRelayStatus::NullPointer; + } + let Some(cb) = cb else { + set_last_error("event metadata injector callback is null"); + return NemoRelayStatus::NullPointer; + }; + let callback = wrap_event_metadata_injector_fn(cb, user_data, free_fn); + let name = match c_str_to_string(name) { + Ok(value) => value, + Err(status) => return status, + }; + match unsafe { &mut *((*ctx).0) }.register_event_metadata_injector(&name, priority, callback) { + Ok(()) => NemoRelayStatus::Ok, + Err(error) => status_from_plugin_error(&error), + } +} + unsafe fn plugin_register_event_sanitizer( ctx: *mut FfiPluginContext, name: *const c_char, diff --git a/crates/ffi/src/callable.rs b/crates/ffi/src/callable.rs index 280b9dfb1..906873cee 100644 --- a/crates/ffi/src/callable.rs +++ b/crates/ffi/src/callable.rs @@ -16,6 +16,7 @@ //! free function in an `Arc` so the closure is `Send + Sync` and the //! free function is called exactly once when all references are dropped. +use std::collections::BTreeMap; use std::ffi::{CStr, CString}; use std::future::Future; use std::pin::Pin; @@ -23,10 +24,11 @@ use std::sync::Arc; use libc::c_char; use nemo_relay::api::runtime::{ - EventSanitizeFn, EventSubscriberFn, LlmCodecIdentity, LlmConditionalFn, LlmExecutionNextFn, - LlmJsonStream, LlmRequestInterceptFn, LlmSanitizeRequestContext, LlmSanitizeRequestFn, - LlmSanitizeResponseContext, LlmSanitizeResponseFn, LlmStreamExecutionNextFn, ToolConditionalFn, - ToolExecutionFn, ToolExecutionNextFn, ToolInterceptFn, ToolSanitizeFn, + EventMetadataInjectorFn, EventSanitizeFn, EventSubscriberFn, LlmCodecIdentity, + LlmConditionalFn, LlmExecutionNextFn, LlmJsonStream, LlmRequestInterceptFn, + LlmSanitizeRequestContext, LlmSanitizeRequestFn, LlmSanitizeResponseContext, + LlmSanitizeResponseFn, LlmStreamExecutionNextFn, ToolConditionalFn, ToolExecutionFn, + ToolExecutionNextFn, ToolInterceptFn, ToolSanitizeFn, }; use serde_json::Value as Json; use tokio_stream::StreamExt; @@ -200,6 +202,18 @@ pub type NemoRelayLlmExecInterceptCb = unsafe extern "C" fn( pub type NemoRelayEventSubscriberCb = unsafe extern "C" fn(user_data: *mut libc::c_void, event: *const FfiEvent); +/// Callback for event metadata injection. +/// +/// The returned string must contain a JSON object whose properties are proposed +/// metadata additions. It transfers to Relay and is freed exactly once. Return +/// null after setting the last error message to report a callback failure. +pub type NemoRelayEventMetadataInjectorCb = Option< + unsafe extern "C" fn(user_data: *mut libc::c_void, event: *const FfiEvent) -> *mut c_char, +>; + +type FfiEventMetadataInjectorFn = + unsafe extern "C" fn(user_data: *mut libc::c_void, event: *const FfiEvent) -> *mut c_char; + /// Callback for mark and scope event sanitizers. /// The returned JSON string transfers to Relay and is freed exactly once. pub type NemoRelayEventSanitizeCb = unsafe extern "C" fn( @@ -990,6 +1004,34 @@ pub fn wrap_event_subscriber( }) } +/// Wrap a C event metadata injector callback into a Rust closure. +pub fn wrap_event_metadata_injector_fn( + cb: FfiEventMetadataInjectorFn, + user_data: *mut libc::c_void, + free_fn: NemoRelayFreeFn, +) -> EventMetadataInjectorFn { + let ud = make_user_data(user_data, free_fn); + Arc::new(move |event: Arc| { + let ud = ud.clone(); + Box::pin(async move { + clear_last_error(); + let ffi_event = FfiEvent((*event).clone()); + let result_ptr = unsafe { cb(ud.ptr, &ffi_event) }; + let result = + json_result_from_ptr(result_ptr, "event metadata injector callback returned null") + .and_then(|value| { + serde_json::from_value::>(value).map_err(|error| { + FlowError::Internal(format!( + "invalid event metadata injector result: {error}" + )) + }) + }); + unsafe { nemo_relay_string_free_internal(result_ptr) }; + result + }) + }) +} + /// Wrap a C event sanitizer callback into a Rust closure. pub fn wrap_event_sanitize_fn( cb: NemoRelayEventSanitizeCb, diff --git a/crates/ffi/tests/unit/api/plugin_tests.rs b/crates/ffi/tests/unit/api/plugin_tests.rs index 1dfa4850e..3e445838c 100644 --- a/crates/ffi/tests/unit/api/plugin_tests.rs +++ b/crates/ffi/tests/unit/api/plugin_tests.rs @@ -4,6 +4,17 @@ //! Unit tests for plugin in the NeMo Relay FFI crate. use super::*; +use nemo_relay::plugin::rollback_registrations; + +unsafe extern "C" fn plugin_event_metadata_injector_cb( + _user_data: *mut libc::c_void, + event: *const FfiEvent, +) -> *mut c_char { + let name = unsafe { take_string(nemo_relay_event_name(event)) }.unwrap_or_default(); + CString::new(json!({"ffi.injected": name}).to_string()) + .unwrap() + .into_raw() +} #[test] fn test_ffi_dynamic_plugin_activation_rejects_empty_specs_without_outputs() { @@ -255,6 +266,115 @@ fn test_ffi_plugin_registration_validation_and_cleanup() { assert_eq!(*lock_unpoisoned(plugin_frees()), 1); } +#[test] +fn test_ffi_plugin_context_event_metadata_injector_is_rolled_back() { + let _guard = TEST_MUTEX.lock().unwrap_or_else(|error| error.into_inner()); + reset_globals(); + + let subscriber_name = cstring(&unique_name("ffi_plugin_metadata_subscriber")); + let injector_name = cstring("metadata"); + let mut registrations = PluginRegistrationContext::with_namespace("ffi_plugin::"); + let mut ctx = FfiPluginContext(&mut registrations as *mut _); + + unsafe { + assert_status!( + nemo_relay_register_subscriber( + subscriber_name.as_ptr(), + subscriber_cb, + ptr::null_mut(), + None, + ), + NemoRelayStatus::Ok + ); + assert_status!( + nemo_relay_plugin_context_register_event_metadata_injector( + &mut ctx, + injector_name.as_ptr(), + 10, + Some(plugin_event_metadata_injector_cb), + ptr::null_mut(), + None, + ), + NemoRelayStatus::Ok + ); + + let active_name = cstring("ffi-plugin-metadata-active"); + assert_status!( + nemo_relay_event(active_name.as_ptr(), ptr::null(), ptr::null(), ptr::null()), + NemoRelayStatus::Ok + ); + assert_status!(nemo_relay_flush_subscribers(), NemoRelayStatus::Ok); + + let mut registrations = registrations.into_registrations(); + rollback_registrations(&mut registrations); + let cleanup_name = cstring("ffi-plugin-metadata-cleanup"); + assert_status!( + nemo_relay_event(cleanup_name.as_ptr(), ptr::null(), ptr::null(), ptr::null(),), + NemoRelayStatus::Ok + ); + assert_status!(nemo_relay_flush_subscribers(), NemoRelayStatus::Ok); + + let events = lock_unpoisoned(event_log()); + let active = events + .iter() + .find(|event| event["name"] == "ffi-plugin-metadata-active") + .expect("plugin-active mark should be delivered"); + let cleanup = events + .iter() + .find(|event| event["name"] == "ffi-plugin-metadata-cleanup") + .expect("plugin-cleanup mark should be delivered"); + assert_eq!( + active["metadata"]["ffi.injected"], + json!("ffi-plugin-metadata-active") + ); + assert!(cleanup["metadata"].is_null()); + drop(events); + + assert_status!( + nemo_relay_deregister_subscriber(subscriber_name.as_ptr()), + NemoRelayStatus::Ok + ); + } +} + +#[test] +fn test_ffi_plugin_context_event_metadata_injector_rejects_null_callback() { + let _guard = TEST_MUTEX.lock().unwrap_or_else(|error| error.into_inner()); + reset_globals(); + + let injector_name = cstring("metadata"); + let mut registrations = PluginRegistrationContext::with_namespace("ffi_plugin_null::"); + let mut ctx = FfiPluginContext(&mut registrations as *mut _); + + unsafe { + assert_status!( + nemo_relay_plugin_context_register_event_metadata_injector( + &mut ctx, + injector_name.as_ptr(), + 10, + None, + ptr::null_mut(), + None, + ), + NemoRelayStatus::NullPointer + ); + assert_status!( + nemo_relay_plugin_context_register_event_metadata_injector( + &mut ctx, + injector_name.as_ptr(), + 10, + Some(plugin_event_metadata_injector_cb), + ptr::null_mut(), + None, + ), + NemoRelayStatus::Ok + ); + } + + let mut registrations = registrations.into_registrations(); + rollback_registrations(&mut registrations); +} + #[test] fn test_ffi_plugin_validation_failure_modes_are_reported() { let _guard = TEST_MUTEX.lock().unwrap(); diff --git a/crates/ffi/tests/unit/api/registry_tests.rs b/crates/ffi/tests/unit/api/registry_tests.rs index 6a7e8467a..99204fd93 100644 --- a/crates/ffi/tests/unit/api/registry_tests.rs +++ b/crates/ffi/tests/unit/api/registry_tests.rs @@ -126,6 +126,321 @@ unsafe extern "C" fn invalid_event_sanitize_cb( CString::new("not-json").unwrap().into_raw() } +unsafe extern "C" fn event_metadata_injector_cb( + _user_data: *mut libc::c_void, + event: *const FfiEvent, +) -> *mut c_char { + let name = unsafe { take_string(nemo_relay_event_name(event)) }.unwrap_or_default(); + CString::new( + json!({ + "ffi.injected": name, + "ffi.integers": [1, 2], + "ffi.doubles": [1.25, 2.5], + "ffi.numbers": [1, 2.5], + }) + .to_string(), + ) + .unwrap() + .into_raw() +} + +unsafe extern "C" fn event_metadata_local_injector_cb( + _user_data: *mut libc::c_void, + _event: *const FfiEvent, +) -> *mut c_char { + CString::new(json!({"ffi.local": true}).to_string()) + .unwrap() + .into_raw() +} + +unsafe extern "C" fn event_metadata_injector_fail_cb( + _user_data: *mut libc::c_void, + _event: *const FfiEvent, +) -> *mut c_char { + crate::error::set_last_error("event metadata injector callback failed"); + ptr::null_mut() +} + +unsafe extern "C" fn event_metadata_injector_invalid_cb( + _user_data: *mut libc::c_void, + _event: *const FfiEvent, +) -> *mut c_char { + CString::new("[]").unwrap().into_raw() +} + +unsafe extern "C" fn event_metadata_injector_mixed_values_cb( + _user_data: *mut libc::c_void, + _event: *const FfiEvent, +) -> *mut c_char { + CString::new( + json!({ + "ffi.invalid.mixed_values": [1, "two"], + "ffi.invalid.sentinel": "must-be-omitted", + }) + .to_string(), + ) + .unwrap() + .into_raw() +} + +#[test] +fn test_ffi_event_metadata_injector_registries_and_failure_paths() { + let _lock = TEST_MUTEX.lock().unwrap_or_else(|error| error.into_inner()); + reset_globals(); + + unsafe { + let stack = fresh_scope_stack(); + let subscriber_name = cstring(&unique_name("ffi_event_metadata_subscriber")); + assert_status!( + nemo_relay_register_subscriber( + subscriber_name.as_ptr(), + subscriber_cb, + ptr::null_mut(), + None, + ), + NemoRelayStatus::Ok + ); + + let global_name = cstring(&unique_name("ffi_event_metadata_global")); + let failure_name = cstring(&unique_name("ffi_event_metadata_failure")); + let invalid_name = cstring(&unique_name("ffi_event_metadata_invalid")); + let mixed_values_name = cstring(&unique_name("ffi_event_metadata_mixed_values")); + assert_status!( + nemo_relay_register_event_metadata_injector( + global_name.as_ptr(), + 10, + Some(event_metadata_injector_cb), + ptr::null_mut(), + None, + ), + NemoRelayStatus::Ok + ); + assert_status!( + nemo_relay_register_event_metadata_injector( + failure_name.as_ptr(), + 20, + Some(event_metadata_injector_fail_cb), + ptr::null_mut(), + None, + ), + NemoRelayStatus::Ok + ); + assert_status!( + nemo_relay_register_event_metadata_injector( + invalid_name.as_ptr(), + 30, + Some(event_metadata_injector_invalid_cb), + ptr::null_mut(), + None, + ), + NemoRelayStatus::Ok + ); + assert_status!( + nemo_relay_register_event_metadata_injector( + mixed_values_name.as_ptr(), + 40, + Some(event_metadata_injector_mixed_values_cb), + ptr::null_mut(), + None, + ), + NemoRelayStatus::Ok + ); + + let scope_name = cstring("ffi-event-metadata-scope"); + let mut scope = ptr::null_mut(); + assert_status!( + nemo_relay_push_scope( + scope_name.as_ptr(), + NemoRelayScopeType::Custom, + ptr::null(), + 0, + ptr::null(), + ptr::null(), + ptr::null(), + &mut scope, + ), + NemoRelayStatus::Ok + ); + + let scope_uuid = cstring(&take_string(nemo_relay_scope_handle_uuid(scope)).unwrap()); + let local_name = cstring(&unique_name("ffi_event_metadata_local")); + assert_status!( + nemo_relay_scope_register_event_metadata_injector( + scope_uuid.as_ptr(), + local_name.as_ptr(), + 5, + Some(event_metadata_local_injector_cb), + ptr::null_mut(), + None, + ), + NemoRelayStatus::Ok + ); + + let mark_name = cstring("ffi-event-metadata-mark"); + assert_status!( + nemo_relay_event(mark_name.as_ptr(), scope, ptr::null(), ptr::null()), + NemoRelayStatus::Ok + ); + assert_status!( + nemo_relay_pop_scope(scope, ptr::null()), + NemoRelayStatus::Ok + ); + nemo_relay_scope_handle_free(scope); + + for name in [ + &global_name, + &failure_name, + &invalid_name, + &mixed_values_name, + ] { + assert_status!( + nemo_relay_deregister_event_metadata_injector(name.as_ptr()), + NemoRelayStatus::Ok + ); + } + + let cleanup_name = cstring("ffi-event-metadata-cleanup"); + assert_status!( + nemo_relay_event(cleanup_name.as_ptr(), ptr::null(), ptr::null(), ptr::null(),), + NemoRelayStatus::Ok + ); + assert_status!(nemo_relay_flush_subscribers(), NemoRelayStatus::Ok); + + let events = lock_unpoisoned(event_log()); + let scope_start = events + .iter() + .find(|event| { + event["name"] == "ffi-event-metadata-scope" + && event["json"]["scope_category"] == "start" + }) + .expect("scope start should be delivered"); + let mark = events + .iter() + .find(|event| event["name"] == "ffi-event-metadata-mark") + .expect("mark should be delivered"); + let scope_end = events + .iter() + .find(|event| { + event["name"] == "ffi-event-metadata-scope" + && event["json"]["scope_category"] == "end" + }) + .expect("scope end should be delivered"); + let cleanup = events + .iter() + .find(|event| event["name"] == "ffi-event-metadata-cleanup") + .expect("cleanup mark should be delivered"); + assert_eq!( + scope_start["metadata"]["ffi.injected"], + json!("ffi-event-metadata-scope") + ); + assert!(scope_start["metadata"].get("ffi.local").is_none()); + assert_eq!( + mark["metadata"]["ffi.injected"], + json!("ffi-event-metadata-mark") + ); + assert_eq!(mark["metadata"]["ffi.integers"], json!([1, 2])); + assert_eq!(mark["metadata"]["ffi.doubles"], json!([1.25, 2.5])); + assert_eq!(mark["metadata"]["ffi.numbers"], json!([1, 2.5])); + assert!(mark["metadata"].get("ffi.invalid.mixed_values").is_none()); + assert!(mark["metadata"].get("ffi.invalid.sentinel").is_none()); + assert_eq!(mark["metadata"]["ffi.local"], json!(true)); + assert_eq!( + scope_end["metadata"]["ffi.injected"], + json!("ffi-event-metadata-scope") + ); + assert_eq!(scope_end["metadata"]["ffi.local"], json!(true)); + assert!(cleanup["metadata"].is_null()); + drop(events); + + assert_status!( + nemo_relay_deregister_subscriber(subscriber_name.as_ptr()), + NemoRelayStatus::Ok + ); + nemo_relay_scope_stack_free(stack); + } +} + +#[test] +fn test_ffi_event_metadata_injector_rejects_null_callbacks() { + let _lock = TEST_MUTEX.lock().unwrap_or_else(|error| error.into_inner()); + reset_globals(); + + unsafe { + let stack = fresh_scope_stack(); + let global_name = cstring(&unique_name("ffi_event_metadata_null_global")); + assert_status!( + nemo_relay_register_event_metadata_injector( + global_name.as_ptr(), + 10, + None, + ptr::null_mut(), + None, + ), + NemoRelayStatus::NullPointer + ); + assert_status!( + nemo_relay_register_event_metadata_injector( + global_name.as_ptr(), + 10, + Some(event_metadata_injector_cb), + ptr::null_mut(), + None, + ), + NemoRelayStatus::Ok + ); + assert_status!( + nemo_relay_deregister_event_metadata_injector(global_name.as_ptr()), + NemoRelayStatus::Ok + ); + + let scope_name = cstring("ffi-event-metadata-null-scope"); + let mut scope = ptr::null_mut(); + assert_status!( + nemo_relay_push_scope( + scope_name.as_ptr(), + NemoRelayScopeType::Custom, + ptr::null(), + 0, + ptr::null(), + ptr::null(), + ptr::null(), + &mut scope, + ), + NemoRelayStatus::Ok + ); + let scope_uuid = cstring(&take_string(nemo_relay_scope_handle_uuid(scope)).unwrap()); + let local_name = cstring(&unique_name("ffi_event_metadata_null_local")); + assert_status!( + nemo_relay_scope_register_event_metadata_injector( + scope_uuid.as_ptr(), + local_name.as_ptr(), + 10, + None, + ptr::null_mut(), + None, + ), + NemoRelayStatus::NullPointer + ); + assert_status!( + nemo_relay_scope_register_event_metadata_injector( + scope_uuid.as_ptr(), + local_name.as_ptr(), + 10, + Some(event_metadata_local_injector_cb), + ptr::null_mut(), + None, + ), + NemoRelayStatus::Ok + ); + assert_status!( + nemo_relay_pop_scope(scope, ptr::null()), + NemoRelayStatus::Ok + ); + nemo_relay_scope_handle_free(scope); + nemo_relay_scope_stack_free(stack); + } +} + #[test] fn test_ffi_event_sanitizer_registries_and_error_paths() { let _lock = TEST_MUTEX.lock().unwrap_or_else(|e| e.into_inner()); diff --git a/crates/node/tests/event_metadata_injection_tests.mjs b/crates/node/tests/event_metadata_injection_tests.mjs index 45132a877..06955688f 100644 --- a/crates/node/tests/event_metadata_injection_tests.mjs +++ b/crates/node/tests/event_metadata_injection_tests.mjs @@ -115,7 +115,7 @@ describe('event metadata injector bindings', () => { assert.equal(marks['node-event-metadata-after-deregister'].metadata, null); }); - it('requires homogeneous numeric arrays at runtime', async () => { + it('accepts numeric arrays with integer and fractional values at runtime', async () => { const events = capture('node-event-metadata-numeric-sub'); lib.registerEventMetadataInjector('node-event-metadata-integers', 10, () => ({ 'node.injector.integers': [1, 2], @@ -140,6 +140,7 @@ describe('event metadata injector bindings', () => { assert.deepEqual(events.at(-1).metadata, { 'node.injector.doubles': [1.25, 2.5], 'node.injector.integers': [1, 2], + 'node.injector.mixed_numbers': [1, 2.5], }); }); diff --git a/go/nemo_relay/callbacks.go b/go/nemo_relay/callbacks.go index 4cdb8690a..60f9bff78 100644 --- a/go/nemo_relay/callbacks.go +++ b/go/nemo_relay/callbacks.go @@ -48,6 +48,7 @@ typedef char* (*NemoRelayLlmConditionalCb)(void* user_data, const FfiLLMRequest* typedef char* (*NemoRelayLlmExecFn)(void* user_data, const char* native_json); typedef char* (*NemoRelayLlmSanitizeResponseCb)(void* user_data, const char* response_json, NemoRelayLlmSanitizeResponseContext context); typedef void (*NemoRelayEventSubscriberFn)(void* user_data, const FfiEvent* event); +typedef char* (*NemoRelayEventMetadataInjectorFn)(void* user_data, const FfiEvent* event); typedef char* (*NemoRelayEventSanitizeFn)(void* user_data, const FfiEvent* event, const char* fields_json); typedef struct FfiPluginContext FfiPluginContext; @@ -100,6 +101,8 @@ import ( // The ID is passed as void* user_data to C callbacks. // --------------------------------------------------------------------------- +var errEventMetadataInjectorCallbackNil = errors.New("event metadata injector callback is nil") + var ( closureRegistryMu sync.Mutex closureRegistry = make(map[uintptr]interface{}) @@ -318,6 +321,16 @@ type FinalizerFunc func() string // the callback, so it is safe to retain the event after the callback returns. type EventSubscriberFunc func(event Event) +// EventMetadata contains flat metadata additions proposed by an injector. +// Relay validates keys and values and never overwrites metadata already present +// on the event. +type EventMetadata map[string]any + +// EventMetadataInjectorFunc inspects an event and proposes metadata additions. +// Return nil metadata with a nil error for a valid no-op. Returning an error +// rejects this callback's additions without preventing event delivery. +type EventMetadataInjectorFunc func(event Event) (EventMetadata, error) + // EventSanitizeFields contains the observability fields an event sanitizer may replace. // // A sanitizer result replaces all three fields; it does not patch the input fields. @@ -663,6 +676,25 @@ func goEventSubscriberTrampoline(userData unsafe.Pointer, event *C.FfiEvent) { fn(goEvent) } +//export goEventMetadataInjectorTrampoline +func goEventMetadataInjectorTrampoline(userData unsafe.Pointer, event *C.FfiEvent) *C.char { + fn := lookupClosure(userData).(EventMetadataInjectorFunc) + metadata, err := fn(newEvent(event)) + if err != nil { + setLastErrorMessage(err.Error()) + return nil + } + if metadata == nil { + metadata = EventMetadata{} + } + result, err := json.Marshal(metadata) + if err != nil { + setLastErrorMessage(err.Error()) + return nil + } + return C.CString(string(result)) +} + //export goEventSanitizeTrampoline func goEventSanitizeTrampoline(userData unsafe.Pointer, event *C.FfiEvent, fieldsJSON *C.char) *C.char { fn := lookupClosure(userData).(EventSanitizeFunc) diff --git a/go/nemo_relay/event_metadata_injectors_test.go b/go/nemo_relay/event_metadata_injectors_test.go new file mode 100644 index 000000000..0fe58ca5d --- /dev/null +++ b/go/nemo_relay/event_metadata_injectors_test.go @@ -0,0 +1,288 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +package nemo_relay + +import ( + "encoding/json" + "errors" + "reflect" + "sync" + "testing" +) + +func TestEventMetadataInjectorGlobalScopeLocalAndFailureBehavior(t *testing.T) { + runTestInIsolatedWorkingDirectory(t, func(t *testing.T) { + runTestWithScopeStack(t, testEventMetadataInjectorGlobalScopeLocalAndFailureBehavior) + }) +} + +func testEventMetadataInjectorGlobalScopeLocalAndFailureBehavior(t *testing.T) { + var mu sync.Mutex + var events []Event + if err := RegisterSubscriber("go-event-metadata-subscriber", func(event Event) { + mu.Lock() + events = append(events, event) + mu.Unlock() + }); err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = DeregisterSubscriber("go-event-metadata-subscriber") }) + + if err := RegisterEventMetadataInjector("go-event-metadata-first", 10, func(event Event) (EventMetadata, error) { + return EventMetadata{ + "go.injected.global": event.Kind(), + "go.injected.collision": "first", + "go.existing": "replacement", + "go.injected.integers": []int64{1, 2}, + "go.injected.doubles": []float64{1, 2.5}, + }, nil + }); err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = DeregisterEventMetadataInjector("go-event-metadata-first") }) + + if err := RegisterEventMetadataInjector("go-event-metadata-later", 20, func(Event) (EventMetadata, error) { + return EventMetadata{"go.injected.collision": "later"}, nil + }); err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = DeregisterEventMetadataInjector("go-event-metadata-later") }) + + if err := RegisterEventMetadataInjector("go-event-metadata-failure", 30, func(Event) (EventMetadata, error) { + return nil, errors.New("expected Go injector failure") + }); err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = DeregisterEventMetadataInjector("go-event-metadata-failure") }) + + if err := RegisterEventMetadataInjector("go-event-metadata-mixed-values", 40, func(Event) (EventMetadata, error) { + return EventMetadata{ + "go.invalid.mixed_values": []any{1, "two"}, + "go.invalid.sentinel": "must-be-omitted", + }, nil + }); err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = DeregisterEventMetadataInjector("go-event-metadata-mixed-values") }) + + scope, err := PushScope("go-event-metadata-scope", ScopeTypeCustom) + if err != nil { + t.Fatal(err) + } + if err := ScopeRegisterEventMetadataInjector(scope.UUID(), "go-event-metadata-local", 5, func(Event) (EventMetadata, error) { + return EventMetadata{"go.injected.local": true}, nil + }); err != nil { + t.Fatal(err) + } + if err := EmitEvent( + "go-event-metadata-mark", + WithEventMetadata(json.RawMessage(`{"go.existing":"original"}`)), + ); err != nil { + t.Fatal(err) + } + if err := PopScope(scope); err != nil { + t.Fatal(err) + } + + for _, name := range []string{ + "go-event-metadata-mixed-values", + "go-event-metadata-failure", + "go-event-metadata-later", + "go-event-metadata-first", + } { + if err := DeregisterEventMetadataInjector(name); err != nil { + t.Fatal(err) + } + } + if err := EmitEvent("go-event-metadata-cleanup"); err != nil { + t.Fatal(err) + } + if err := FlushSubscribers(); err != nil { + t.Fatal(err) + } + + mu.Lock() + defer mu.Unlock() + if len(events) != 4 { + t.Fatalf("expected four delivered events, got %d", len(events)) + } + for index, event := range events[:3] { + metadata := decodeEventMetadata(t, event) + if metadata["go.injected.collision"] != "first" { + t.Fatalf("event %d did not preserve first-injector precedence: %#v", index, metadata) + } + if _, ok := metadata["go.injected.global"]; !ok { + t.Fatalf("event %d is missing global metadata: %#v", index, metadata) + } + } + if metadata := decodeEventMetadata(t, events[0]); metadata["go.injected.local"] != nil { + t.Fatalf("scope start unexpectedly contains scope-local metadata: %#v", metadata) + } + for _, event := range events[1:3] { + if metadata := decodeEventMetadata(t, event); metadata["go.injected.local"] != true { + t.Fatalf("event %s is missing scope-local metadata: %#v", event.Name(), metadata) + } + } + if metadata := decodeEventMetadata(t, events[1]); metadata["go.existing"] != "original" { + t.Fatalf("existing metadata was overwritten: %#v", metadata) + } + metadata := decodeEventMetadata(t, events[1]) + if got, want := metadata["go.injected.integers"], []any{float64(1), float64(2)}; !reflect.DeepEqual(got, want) { + t.Fatalf("homogeneous integer metadata = %#v, want %#v", got, want) + } + if got, want := metadata["go.injected.doubles"], []any{float64(1), 2.5}; !reflect.DeepEqual(got, want) { + t.Fatalf("homogeneous double metadata = %#v, want %#v", got, want) + } + if _, ok := metadata["go.invalid.mixed_values"]; ok { + t.Fatalf("mixed primitive metadata was accepted: %#v", metadata) + } + if _, ok := metadata["go.invalid.sentinel"]; ok { + t.Fatalf("invalid callback output was partially applied: %#v", metadata) + } + if metadata := decodeEventMetadata(t, events[3]); len(metadata) != 0 { + t.Fatalf("cleanup event retained injector metadata: %#v", metadata) + } +} + +func TestPluginContextEventMetadataInjectorLifecycle(t *testing.T) { + runTestInIsolatedWorkingDirectory(t, func(t *testing.T) { + runTestWithScopeStack(t, func(t *testing.T) { + const kind = "go.event.metadata.plugin" + var mu sync.Mutex + var events []Event + + if err := RegisterSubscriber("go-plugin-metadata-subscriber", func(event Event) { + mu.Lock() + events = append(events, event) + mu.Unlock() + }); err != nil { + t.Fatal(err) + } + defer DeregisterSubscriber("go-plugin-metadata-subscriber") + + if err := RegisterPlugin(kind, PluginFuncs{RegisterFunc: func(_ map[string]any, ctx *PluginContext) error { + return ctx.RegisterEventMetadataInjector("configured", 10, func(Event) (EventMetadata, error) { + return EventMetadata{"go.injected.plugin": true}, nil + }) + }}); err != nil { + t.Fatal(err) + } + defer DeregisterPlugin(kind) + + if _, err := InitializePlugins(PluginConfig{ + Version: 1, + Components: []PluginComponentSpec{{ + Kind: kind, + Enabled: true, + }}, + }); err != nil { + t.Fatal(err) + } + if err := EmitEvent("go-plugin-metadata-active"); err != nil { + t.Fatal(err) + } + if err := ClearPluginConfiguration(); err != nil { + t.Fatal(err) + } + if err := EmitEvent("go-plugin-metadata-cleanup"); err != nil { + t.Fatal(err) + } + if err := FlushSubscribers(); err != nil { + t.Fatal(err) + } + + mu.Lock() + defer mu.Unlock() + if len(events) != 2 { + t.Fatalf("expected two delivered events, got %d", len(events)) + } + if metadata := decodeEventMetadata(t, events[0]); metadata["go.injected.plugin"] != true { + t.Fatalf("plugin metadata was not injected: %#v", metadata) + } + if metadata := decodeEventMetadata(t, events[1]); len(metadata) != 0 { + t.Fatalf("plugin metadata remained after cleanup: %#v", metadata) + } + }) + }) +} + +func TestEventMetadataInjectorRegistrationErrorsReleaseCallbacks(t *testing.T) { + baseline := registeredClosureCount() + callback := func(Event) (EventMetadata, error) { return nil, nil } + if err := RegisterEventMetadataInjector("go-event-metadata-duplicate", 0, callback); err != nil { + t.Fatal(err) + } + if err := RegisterEventMetadataInjector("go-event-metadata-duplicate", 0, callback); err == nil { + t.Fatal("expected duplicate event metadata injector registration to fail") + } + if current := registeredClosureCount(); current != baseline+1 { + t.Fatalf("duplicate registration leaked callback: baseline=%d current=%d", baseline, current) + } + if err := DeregisterEventMetadataInjector("go-event-metadata-duplicate"); err != nil { + t.Fatal(err) + } +} + +func TestEventMetadataInjectorNilCallbacksDoNotRegister(t *testing.T) { + runTestInIsolatedWorkingDirectory(t, func(t *testing.T) { + runTestWithScopeStack(t, func(t *testing.T) { + baseline := registeredClosureCount() + + if err := RegisterEventMetadataInjector("go-event-metadata-nil-global", 0, nil); !errors.Is(err, errEventMetadataInjectorCallbackNil) { + t.Fatalf("RegisterEventMetadataInjector() error = %v, want %v", err, errEventMetadataInjectorCallbackNil) + } + if current := registeredClosureCount(); current != baseline { + t.Fatalf("nil global callback changed registry size: baseline=%d current=%d", baseline, current) + } + + scope, err := PushScope("go-event-metadata-nil-scope", ScopeTypeCustom) + if err != nil { + t.Fatal(err) + } + if err := ScopeRegisterEventMetadataInjector(scope.UUID(), "go-event-metadata-nil-local", 0, nil); !errors.Is(err, errEventMetadataInjectorCallbackNil) { + t.Fatalf("ScopeRegisterEventMetadataInjector() error = %v, want %v", err, errEventMetadataInjectorCallbackNil) + } + if current := registeredClosureCount(); current != baseline { + t.Fatalf("nil scope callback changed registry size: baseline=%d current=%d", baseline, current) + } + if err := PopScope(scope); err != nil { + t.Fatal(err) + } + + const kind = "go.event.metadata.nil.plugin" + if err := RegisterPlugin(kind, PluginFuncs{RegisterFunc: func(_ map[string]any, ctx *PluginContext) error { + return ctx.RegisterEventMetadataInjector("nil", 0, nil) + }}); err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = DeregisterPlugin(kind) }) + + pluginBaseline := registeredClosureCount() + if _, err := InitializePlugins(PluginConfig{ + Version: 1, + Components: []PluginComponentSpec{{ + Kind: kind, + Enabled: true, + }}, + }); err == nil { + t.Fatal("expected nil plugin callback registration to fail") + } + if current := registeredClosureCount(); current != pluginBaseline { + t.Fatalf("nil plugin callback changed registry size: baseline=%d current=%d", pluginBaseline, current) + } + }) + }) +} + +func decodeEventMetadata(t *testing.T, event Event) map[string]any { + t.Helper() + metadata := map[string]any{} + if len(event.Metadata()) == 0 { + return metadata + } + if err := json.Unmarshal(event.Metadata(), &metadata); err != nil { + t.Fatalf("decode event metadata: %v", err) + } + return metadata +} diff --git a/go/nemo_relay/nemo_relay.go b/go/nemo_relay/nemo_relay.go index 72da4d715..487803775 100644 --- a/go/nemo_relay/nemo_relay.go +++ b/go/nemo_relay/nemo_relay.go @@ -176,10 +176,13 @@ extern int32_t nemo_relay_deregister_llm_stream_execution_intercept(const char* // Subscribers typedef void (*NemoRelayEventSubscriberFn)(void* user_data, const FfiEvent* event); +typedef char* (*NemoRelayEventMetadataInjectorFn)(void* user_data, const FfiEvent* event); typedef char* (*NemoRelayEventSanitizeFn)(void* user_data, const FfiEvent* event, const char* fields_json); extern int32_t nemo_relay_register_subscriber(const char* name, NemoRelayEventSubscriberFn cb, void* user_data, NemoRelayFreeFn free_fn); extern int32_t nemo_relay_deregister_subscriber(const char* name); extern int32_t nemo_relay_flush_subscribers(void); +extern int32_t nemo_relay_register_event_metadata_injector(const char* name, int32_t priority, NemoRelayEventMetadataInjectorFn cb, void* user_data, NemoRelayFreeFn free_fn); +extern int32_t nemo_relay_deregister_event_metadata_injector(const char* name); extern int32_t nemo_relay_register_mark_sanitize_guardrail(const char* name, int32_t priority, NemoRelayEventSanitizeFn cb, void* user_data, NemoRelayFreeFn free_fn); extern int32_t nemo_relay_deregister_mark_sanitize_guardrail(const char* name); extern int32_t nemo_relay_register_scope_sanitize_start_guardrail(const char* name, int32_t priority, NemoRelayEventSanitizeFn cb, void* user_data, NemoRelayFreeFn free_fn); @@ -188,6 +191,8 @@ extern int32_t nemo_relay_register_scope_sanitize_end_guardrail(const char* name extern int32_t nemo_relay_deregister_scope_sanitize_end_guardrail(const char* name); // Scope-local tool guardrails +extern int32_t nemo_relay_scope_register_event_metadata_injector(const char* scope_uuid, const char* name, int32_t priority, NemoRelayEventMetadataInjectorFn cb, void* user_data, NemoRelayFreeFn free_fn); +extern int32_t nemo_relay_scope_deregister_event_metadata_injector(const char* scope_uuid, const char* name); extern int32_t nemo_relay_scope_register_mark_sanitize_guardrail(const char* scope_uuid, const char* name, int32_t priority, NemoRelayEventSanitizeFn cb, void* user_data, NemoRelayFreeFn free_fn); extern int32_t nemo_relay_scope_deregister_mark_sanitize_guardrail(const char* scope_uuid, const char* name); extern int32_t nemo_relay_scope_register_scope_sanitize_start_guardrail(const char* scope_uuid, const char* name, int32_t priority, NemoRelayEventSanitizeFn cb, void* user_data, NemoRelayFreeFn free_fn); @@ -296,6 +301,7 @@ extern void nemo_relay_otel_metric_subscriber_free(void*); // Go trampoline forward declarations (defined via //export in callbacks.go) extern char* goToolSanitizeTrampoline(void*, const char*, const char*); +extern char* goEventMetadataInjectorTrampoline(void*, const FfiEvent*); extern char* goEventSanitizeTrampoline(void*, const FfiEvent*, const char*); extern char* goToolConditionalTrampoline(void*, const char*, const char*); extern char* goToolExecTrampoline(void*, const char*); @@ -1426,6 +1432,32 @@ func LlmStreamCallExecute(name string, request interface{}, fn LLMExecutionFunc, // Guardrail/Intercept registration (Tool) // --------------------------------------------------------------------------- +// RegisterEventMetadataInjector registers a global event metadata injector. +// Injectors run in ascending priority order and may only add metadata keys that +// are not already present on the event. +func RegisterEventMetadataInjector(name string, priority int32, fn EventMetadataInjectorFunc) error { + if fn == nil { + return errEventMetadataInjectorCallbackNil + } + id := registerClosure(fn) + cName := C.CString(name) + defer C.free(unsafe.Pointer(cName)) + return checkStatus(C.nemo_relay_register_event_metadata_injector( + cName, + C.int32_t(priority), + C.NemoRelayEventMetadataInjectorFn(C.goEventMetadataInjectorTrampoline), + id, + C.NemoRelayFreeFn(C.goFreeTrampoline), + )) +} + +// DeregisterEventMetadataInjector removes a global event metadata injector. +func DeregisterEventMetadataInjector(name string) error { + cName := C.CString(name) + defer C.free(unsafe.Pointer(cName)) + return checkStatus(C.nemo_relay_deregister_event_metadata_injector(cName)) +} + func registerEventSanitizer(name string, priority int32, fn EventSanitizeFunc, kind int) error { id := registerClosure(fn) cName := C.CString(name) @@ -2880,6 +2912,40 @@ func (s *OpenTelemetryMetricSubscriber) Close() { // Scope-local guardrail/intercept registration (Tool) // --------------------------------------------------------------------------- +// ScopeRegisterEventMetadataInjector registers an event metadata injector +// owned by an active scope. +func ScopeRegisterEventMetadataInjector(scopeUUID, name string, priority int32, fn EventMetadataInjectorFunc) error { + if fn == nil { + return errEventMetadataInjectorCallbackNil + } + id := registerClosure(fn) + cScopeUUID := C.CString(scopeUUID) + defer C.free(unsafe.Pointer(cScopeUUID)) + cName := C.CString(name) + defer C.free(unsafe.Pointer(cName)) + return checkStatus(C.nemo_relay_scope_register_event_metadata_injector( + cScopeUUID, + cName, + C.int32_t(priority), + C.NemoRelayEventMetadataInjectorFn(C.goEventMetadataInjectorTrampoline), + id, + C.NemoRelayFreeFn(C.goFreeTrampoline), + )) +} + +// ScopeDeregisterEventMetadataInjector removes an event metadata injector +// owned by an active scope. +func ScopeDeregisterEventMetadataInjector(scopeUUID, name string) error { + cScopeUUID := C.CString(scopeUUID) + defer C.free(unsafe.Pointer(cScopeUUID)) + cName := C.CString(name) + defer C.free(unsafe.Pointer(cName)) + return checkStatus(C.nemo_relay_scope_deregister_event_metadata_injector( + cScopeUUID, + cName, + )) +} + func registerScopeEventSanitizer(scopeUUID, name string, priority int32, fn EventSanitizeFunc, kind int) error { id := registerClosure(fn) cScopeUUID := C.CString(scopeUUID) diff --git a/go/nemo_relay/plugin.go b/go/nemo_relay/plugin.go index 1de3de731..71359bd3a 100644 --- a/go/nemo_relay/plugin.go +++ b/go/nemo_relay/plugin.go @@ -19,6 +19,7 @@ typedef void (*NemoRelayFreeFn)(void* user_data); typedef char* (*NemoRelayPluginValidateCb)(void* user_data, const char* plugin_config_json); typedef int32_t (*NemoRelayPluginRegisterCb)(void* user_data, const char* plugin_config_json, FfiPluginContext* ctx); typedef void (*NemoRelayEventSubscriberFn)(void* user_data, const void* event); +typedef char* (*NemoRelayEventMetadataInjectorFn)(void* user_data, const void* event); typedef char* (*NemoRelayEventSanitizeFn)(void* user_data, const void* event, const char* fields_json); typedef char* (*NemoRelayToolSanitizeFn)(void* user_data, const char* name, const char* args_json); typedef char* (*NemoRelayToolConditionalFn)(void* user_data, const char* name, const char* args_json); @@ -44,6 +45,7 @@ extern int32_t nemo_relay_deregister_plugin(const char* plugin_kind); extern void nemo_relay_string_free(char* ptr); extern int32_t nemo_relay_plugin_context_register_subscriber(FfiPluginContext* ctx, const char* name, NemoRelayEventSubscriberFn cb, void* user_data, NemoRelayFreeFn free_fn); +extern int32_t nemo_relay_plugin_context_register_event_metadata_injector(FfiPluginContext* ctx, const char* name, int32_t priority, NemoRelayEventMetadataInjectorFn cb, void* user_data, NemoRelayFreeFn free_fn); extern int32_t nemo_relay_plugin_context_register_mark_sanitize_guardrail(FfiPluginContext* ctx, const char* name, int32_t priority, NemoRelayEventSanitizeFn cb, void* user_data, NemoRelayFreeFn free_fn); extern int32_t nemo_relay_plugin_context_register_scope_sanitize_start_guardrail(FfiPluginContext* ctx, const char* name, int32_t priority, NemoRelayEventSanitizeFn cb, void* user_data, NemoRelayFreeFn free_fn); extern int32_t nemo_relay_plugin_context_register_scope_sanitize_end_guardrail(FfiPluginContext* ctx, const char* name, int32_t priority, NemoRelayEventSanitizeFn cb, void* user_data, NemoRelayFreeFn free_fn); @@ -62,6 +64,7 @@ extern int32_t nemo_relay_plugin_context_register_tool_execution_intercept(FfiPl extern char* goPluginValidateTrampoline(void*, const char*); extern int32_t goPluginRegisterTrampoline(void*, const char*, FfiPluginContext*); extern void goEventSubscriberTrampoline(void*, const void*); +extern char* goEventMetadataInjectorTrampoline(void*, const void*); extern char* goEventSanitizeTrampoline(void*, const void*, const char*); extern void goFreeTrampoline(void*); extern char* goToolSanitizeTrampoline(void*, const char*, const char*); @@ -583,6 +586,29 @@ func (ctx *PluginContext) RegisterSubscriber(name string, fn EventSubscriberFunc )) } +// RegisterEventMetadataInjector registers an event metadata injector for this +// component. Relay qualifies its name and removes it when plugin configuration +// is cleared or registration rolls back. +func (ctx *PluginContext) RegisterEventMetadataInjector(name string, priority int32, fn EventMetadataInjectorFunc) error { + if ctx == nil || ctx.ptr == nil { + return errors.New(errPluginContextClosed) + } + if fn == nil { + return errEventMetadataInjectorCallbackNil + } + cName := C.CString(name) + defer C.free(unsafe.Pointer(cName)) + userData := registerClosure(fn) + return checkStatus(C.nemo_relay_plugin_context_register_event_metadata_injector( + ctx.ptr, + cName, + C.int32_t(priority), + (C.NemoRelayEventMetadataInjectorFn)(C.goEventMetadataInjectorTrampoline), + userData, + (C.NemoRelayFreeFn)(C.goFreeTrampoline), + )) +} + func (ctx *PluginContext) registerEventSanitizer(name string, priority int32, fn EventSanitizeFunc, surface int) error { if ctx == nil || ctx.ptr == nil { return errors.New(errPluginContextClosed) diff --git a/go/nemo_relay/plugin_gap_test.go b/go/nemo_relay/plugin_gap_test.go index 856b6080f..734c5bf07 100644 --- a/go/nemo_relay/plugin_gap_test.go +++ b/go/nemo_relay/plugin_gap_test.go @@ -49,6 +49,9 @@ func TestClosedPluginContextRejectsEveryRegistrationSurface(t *testing.T) { call func() error }{ {name: "subscriber", call: func() error { return ctx.RegisterSubscriber("closed_subscriber", nil) }}, + {name: "event metadata injector", call: func() error { + return ctx.RegisterEventMetadataInjector("closed_event_metadata", 0, nil) + }}, {name: "mark sanitizer", call: func() error { return ctx.RegisterMarkSanitizeGuardrail("closed_mark", 0, nil) }}, {name: "scope-start sanitizer", call: func() error { return ctx.RegisterScopeSanitizeStartGuardrail("closed_scope_start", 0, nil) }}, {name: "scope-end sanitizer", call: func() error { return ctx.RegisterScopeSanitizeEndGuardrail("closed_scope_end", 0, nil) }}, diff --git a/python/tests/test_event_metadata_injection.py b/python/tests/test_event_metadata_injection.py index 86927bbfb..37d7ec17e 100644 --- a/python/tests/test_event_metadata_injection.py +++ b/python/tests/test_event_metadata_injection.py @@ -68,7 +68,9 @@ def fail(_event: nemo_relay.Event) -> nemo_relay.EventMetadata: } -async def test_python_injectors_require_homogeneous_numeric_lists(subscribed_events): +async def test_python_injectors_accept_numeric_lists_with_integer_and_fractional_values( + subscribed_events, +): integers_name = f"python-integers-{uuid4()}" doubles_name = f"python-doubles-{uuid4()}" mixed_name = f"python-mixed-numbers-{uuid4()}" @@ -100,6 +102,7 @@ async def test_python_injectors_require_homogeneous_numeric_lists(subscribed_eve assert event.metadata == { "python.injector.doubles": [1.0, 2.5], "python.injector.integers": [1, 2], + "python.injector.mixed_numbers": [1, 2.5], }