diff --git a/crates/cli/tests/coverage/server_tests.rs b/crates/cli/tests/coverage/server_tests.rs index 8cd63de0f..5581ce409 100644 --- a/crates/cli/tests/coverage/server_tests.rs +++ b/crates/cli/tests/coverage/server_tests.rs @@ -521,6 +521,152 @@ async fn serve_listener_hermes_api_hooks_write_atof_category_profile_and_fidelit assert_eq!(lossy_start["data"]["content"]["message_count"], json!(2)); } +#[tokio::test] +async fn serve_listener_hermes_api_request_error_writes_lossy_atof_error_event() { + let _guard = PLUGIN_TEST_LOCK.lock().await; + let _ = nemo_relay::plugin::clear_plugin_configuration(); + + let temp = tempfile::tempdir().unwrap(); + let atof_dir = temp.path().join("atof"); + std::fs::create_dir_all(&atof_dir).unwrap(); + let mut config = test_config(); + config.plugin_config = Some(json!({ + "version": 1, + "components": [ + { + "kind": "observability", + "enabled": true, + "config": { + "version": 1, + "atof": { + "enabled": true, + "output_directory": atof_dir, + "filename": "events.jsonl", + "mode": "overwrite" + } + } + } + ] + })); + + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let address = listener.local_addr().unwrap(); + let url = format!("http://{address}"); + let (shutdown_tx, shutdown_rx) = oneshot::channel(); + let handle = + tokio::spawn(async move { serve_listener(listener, config, Some(shutdown_rx)).await }); + + wait_for_gateway(&url).await; + let client = test_http_client(); + + let response = client + .post(format!("{url}/hooks/hermes")) + .json(&json!({ + "hook_event_name": "pre_api_request", + "session_id": "hermes-atof-error", + "extra": { + "task_id": "task-err", + "api_request_id": "turn-1:api:3", + "api_call_count": 3, + "model": "qwen", + "provider": "custom", + "request": { + "method": "POST", + "body": { + "model": "qwen", + "messages": [ + { "role": "user", "content": "hello" } + ] + } + } + } + })) + .send() + .await + .unwrap(); + assert_eq!(response.status(), StatusCode::OK); + + let response = client + .post(format!("{url}/hooks/hermes")) + .json(&json!({ + "hook_event_name": "api_request_error", + "session_id": "hermes-atof-error", + "extra": { + "task_id": "task-err", + "api_request_id": "turn-1:api:3", + "api_call_count": 3, + "model": "qwen", + "provider": "custom", + "status_code": 502, + "retry_count": 1, + "max_retries": 2, + "retryable": true, + "reason": "upstream", + "error": { + "type": "BadGateway", + "message": "gateway upstream error" + } + } + })) + .send() + .await + .unwrap(); + assert_eq!(response.status(), StatusCode::OK); + + shutdown_tx.send(()).unwrap(); + handle.await.unwrap().unwrap(); + + let events = std::fs::read_to_string(temp.path().join("atof/events.jsonl")).unwrap(); + let llm_events = events + .lines() + .map(|line| serde_json::from_str::(line).unwrap()) + .filter(|event| event["category"] == "llm") + .collect::>(); + assert_eq!( + llm_events.len(), + 2, + "expected Hermes error-path LLM exports, got {llm_events:?}" + ); + + let end = llm_events + .iter() + .find(|event| { + event["scope_category"] == "end" + && event["metadata"]["api_call_id"] == json!("turn-1:api:3") + }) + .unwrap(); + let start = llm_events + .iter() + .find(|event| { + event["scope_category"] == "start" + && event["metadata"]["api_call_id"] == json!("turn-1:api:3") + }) + .unwrap(); + assert_eq!(start["metadata"]["provider_payload_exact"], json!(true)); + assert_eq!( + start["metadata"]["fidelity_source"], + json!("hermes_api_hooks_sanitized") + ); + assert_eq!( + start["data"]["content"]["messages"][0]["content"], + json!("hello") + ); + assert_eq!(end["category_profile"]["model_name"], json!("qwen")); + assert_eq!(end["metadata"]["provider_payload_exact"], json!(false)); + assert_eq!( + end["metadata"]["fidelity_source"], + json!("hermes_api_hooks") + ); + assert_eq!(end["data"]["status_code"], json!(502)); + assert_eq!(end["data"]["retry_count"], json!(1)); + assert_eq!(end["data"]["retryable"], json!(true)); + assert_eq!(end["data"]["reason"], json!("upstream")); + assert_eq!( + end["data"]["error"]["message"], + json!("gateway upstream error") + ); +} + #[tokio::test] async fn serve_listener_routed_gateway_wire_formats_write_atof_category_profile_and_usage() { let _guard = PLUGIN_TEST_LOCK.lock().await; diff --git a/crates/cli/tests/coverage/session_tests.rs b/crates/cli/tests/coverage/session_tests.rs index a4824a1da..bfd4981bb 100644 --- a/crates/cli/tests/coverage/session_tests.rs +++ b/crates/cli/tests/coverage/session_tests.rs @@ -1562,6 +1562,128 @@ async fn hermes_exact_api_hooks_write_atif_request_response_and_cost() { })); } +#[tokio::test] +async fn hermes_api_request_error_writes_atif_error_step_and_fidelity() { + let _guard = OBSERVABILITY_PLUGIN_TEST_LOCK.lock().await; + let temp = tempfile::tempdir().unwrap(); + let atif_dir = temp.path().join("atif"); + install_test_atif_plugin(&atif_dir).await; + let config = session_test_config(); + let manager = SessionManager::new(config); + let headers = HeaderMap::new(); + + for payload in [ + json!({ + "hook_event_name": "on_session_start", + "session_id": "hermes-error" + }), + json!({ + "hook_event_name": "pre_api_request", + "session_id": "hermes-error", + "extra": { + "task_id": "task-err", + "api_request_id": "turn-1:api:3", + "api_call_count": 3, + "provider": "custom", + "model": "qwen", + "request": { + "method": "POST", + "body": { + "model": "qwen", + "messages": [ + { "role": "user", "content": "hello" } + ] + } + } + } + }), + json!({ + "hook_event_name": "api_request_error", + "session_id": "hermes-error", + "extra": { + "task_id": "task-err", + "api_request_id": "turn-1:api:3", + "api_call_count": 3, + "provider": "custom", + "model": "qwen", + "status_code": 502, + "retry_count": 1, + "max_retries": 2, + "retryable": true, + "reason": "upstream", + "error": { + "type": "BadGateway", + "message": "gateway upstream error" + } + } + }), + json!({ + "hook_event_name": "on_session_finalize", + "session_id": "hermes-error" + }), + ] { + let outcome = crate::adapters::hermes::adapt(payload, &headers); + manager + .apply_events(&headers, outcome.events) + .await + .unwrap(); + } + + clear_plugin_configuration().unwrap(); + let atif = read_atif_for_session(&atif_dir, "hermes-error"); + let steps = atif["steps"].as_array().unwrap(); + assert_eq!(steps.len(), 2); + assert_eq!(steps[0]["message"], json!("hello")); + assert_eq!(steps[1]["source"], json!("agent")); + assert_eq!(steps[1]["extra"]["llm_response"]["status_code"], json!(502)); + assert_eq!(steps[1]["extra"]["llm_response"]["retry_count"], json!(1)); + assert_eq!(steps[1]["extra"]["llm_response"]["retryable"], json!(true)); + assert_eq!( + steps[1]["extra"]["llm_response"]["reason"], + json!("upstream") + ); + assert_eq!( + steps[1]["extra"]["llm_response"]["error"]["message"], + json!("gateway upstream error") + ); + let observed_events = atif["extra"]["observed_events"].as_array().unwrap(); + assert!( + observed_events.len() >= 4, + "expected Hermes error trajectory to keep observed events, got {}", + serde_json::to_string_pretty(&atif["extra"]["observed_events"]).unwrap() + ); + let error_event = observed_events + .iter() + .find(|event| { + event["scope_category"] == json!("end") + && event["metadata"]["api_call_id"] == json!("turn-1:api:3") + }) + .unwrap(); + let request_event = observed_events + .iter() + .find(|event| { + event["scope_category"] == json!("start") + && event["metadata"]["api_call_id"] == json!("turn-1:api:3") + }) + .unwrap(); + assert_eq!( + request_event["metadata"]["provider_payload_exact"], + json!(true) + ); + assert_eq!( + request_event["metadata"]["fidelity_source"], + json!("hermes_api_hooks_sanitized") + ); + assert_eq!( + error_event["metadata"]["provider_payload_exact"], + json!(false) + ); + assert_eq!( + error_event["metadata"]["fidelity_source"], + json!("hermes_api_hooks") + ); +} + #[tokio::test] async fn hermes_lossy_api_hooks_write_atif_fidelity_markers() { let _guard = OBSERVABILITY_PLUGIN_TEST_LOCK.lock().await; diff --git a/crates/core/src/observability/openinference.rs b/crates/core/src/observability/openinference.rs index deb1b55c1..092ed669c 100644 --- a/crates/core/src/observability/openinference.rs +++ b/crates/core/src/observability/openinference.rs @@ -647,6 +647,11 @@ fn start_attributes(event: &Event) -> Vec { let is_llm = event .category() .is_some_and(|category| category.as_str() == "llm"); + if is_llm { + // Final span metadata should reflect the completed event, especially for mixed-fidelity + // Hermes flows where the request can be exact but the terminal error is lossy. + attributes.retain(|attribute| attribute.key.as_str() != oi::METADATA.as_str()); + } let handle_attributes = event.attributes(); if handle_attributes.is_some_and(|attributes| !attributes.is_empty()) { push_serialized( @@ -700,6 +705,10 @@ fn end_attributes(event: &Event) -> Vec { .category() .is_some_and(|category| category.as_str() == "llm"); + if let Some(metadata) = event.metadata().and_then(to_json_string) { + attributes.push(KeyValue::new(oi::METADATA, metadata)); + } + push_serialized( &mut attributes, "nemo_relay.end.output_json", diff --git a/crates/core/tests/unit/observability/openinference_tests.rs b/crates/core/tests/unit/observability/openinference_tests.rs index da710ec71..4488b4b99 100644 --- a/crates/core/tests/unit/observability/openinference_tests.rs +++ b/crates/core/tests/unit/observability/openinference_tests.rs @@ -2486,6 +2486,116 @@ fn hermes_exact_api_payloads_emit_openinference_text_usage_and_metadata() { ); } +#[test] +fn hermes_api_request_error_emits_openinference_json_output_and_metadata() { + let (provider, exporter) = make_provider(); + let mut processor = + OpenInferenceEventProcessor::new(provider.clone(), "test-scope".to_string()); + let uuid = Uuid::now_v7(); + let start_metadata = json!({ + "provider_payload_exact": true, + "fidelity_source": "hermes_api_hooks_sanitized" + }); + let end_metadata = json!({ + "provider_payload_exact": false, + "fidelity_source": "hermes_api_hooks" + }); + + processor.process(&Event::Scope(ScopeEvent::new( + BaseEvent::builder() + .uuid(uuid) + .name("custom") + .data(json!({ + "model": "qwen", + "messages": [{ "role": "user", "content": "hello" }] + })) + .metadata(start_metadata) + .build(), + ScopeCategory::Start, + Vec::new(), + EventCategory::llm(), + Some(CategoryProfile::builder().model_name("qwen").build()), + ))); + processor.process(&Event::Scope(ScopeEvent::new( + BaseEvent::builder() + .uuid(uuid) + .name("custom") + .data(json!({ + "status_code": 502, + "retry_count": 1, + "max_retries": 2, + "retryable": true, + "reason": "upstream", + "error": { + "type": "BadGateway", + "message": "gateway upstream error" + } + })) + .metadata(end_metadata) + .build(), + ScopeCategory::End, + Vec::new(), + EventCategory::llm(), + Some(CategoryProfile::builder().model_name("qwen").build()), + ))); + + processor.force_flush().unwrap(); + + let spans = exporter.get_finished_spans().unwrap(); + assert_eq!(spans.len(), 1); + let attributes = attr_map(&spans[0].attributes); + assert_eq!( + attributes.get("openinference.span.kind"), + Some(&"LLM".to_string()) + ); + assert_eq!(attributes.get("llm.model_name"), Some(&"qwen".to_string())); + assert_eq!( + attributes.get("input.value"), + Some(&"user: hello".to_string()) + ); + assert_eq!( + serde_json::from_str::(attributes.get("output.value").unwrap()).unwrap(), + json!({ + "status_code": 502, + "retry_count": 1, + "max_retries": 2, + "retryable": true, + "reason": "upstream", + "error": { + "type": "BadGateway", + "message": "gateway upstream error" + } + }) + ); + assert_eq!( + attributes.get("output.mime_type"), + Some(&"application/json".to_string()) + ); + assert_eq!( + serde_json::from_str::( + attributes.get("nemo_relay.end.output_json").unwrap(), + ) + .unwrap(), + json!({ + "status_code": 502, + "retry_count": 1, + "max_retries": 2, + "retryable": true, + "reason": "upstream", + "error": { + "type": "BadGateway", + "message": "gateway upstream error" + } + }) + ); + assert_attr_contains(&attributes, "metadata", "\"provider_payload_exact\":false"); + assert_attr_contains( + &attributes, + "metadata", + "\"fidelity_source\":\"hermes_api_hooks\"", + ); +} + #[test] fn llm_end_with_inconsistent_manual_usage_omits_invalid_total_tokens() { let (provider, exporter) = make_provider();