diff --git a/src/app/endpoints/a2a.py b/src/app/endpoints/a2a.py index e3cc39998..0b2c77adf 100644 --- a/src/app/endpoints/a2a.py +++ b/src/app/endpoints/a2a.py @@ -203,7 +203,7 @@ def _record_execution_span( output_text = run_result.response.text if output_text: - span.set_attribute(SpanAttributes.OUTPUT, anonymize_value(output_text)) + span.set_attribute(SpanAttributes.OUTPUT, output_text) async def _persist_compacted_a2a_turn( @@ -420,7 +420,7 @@ async def _process_task_streaming( # pylint: disable=too-many-locals ) return - span.set_attribute(SpanAttributes.INPUT, anonymize_value(user_input)) + span.set_attribute(SpanAttributes.INPUT, user_input) preview = user_input[:200] + ("..." if len(user_input) > 200 else "") logger.info("Processing A2A request: %s", preview) @@ -1094,9 +1094,7 @@ async def _handle_a2a_jsonrpc( # pylint: disable=too-many-locals,too-many-state span, { SpanAttributes.A2A_RPC_METHOD: rpc_method, - SpanAttributes.A2A_REQUEST_ID: ( - anonymize_value(rpc_request_id) if rpc_request_id else "" - ), + SpanAttributes.A2A_REQUEST_ID: rpc_request_id if rpc_request_id else "", SpanAttributes.USER_ID: anonymize_value(auth[0]) if auth[0] else "", }, ) diff --git a/src/app/endpoints/feedback.py b/src/app/endpoints/feedback.py index 77162e4ea..29e4a979d 100644 --- a/src/app/endpoints/feedback.py +++ b/src/app/endpoints/feedback.py @@ -107,8 +107,9 @@ def _record_feedback_request_attributes( """Set high-level feedback attributes on the root span. User-generated free text (the question, LLM response, and comment) is - anonymized before being recorded. Low-cardinality signals (rating and - categories) are recorded as-is. + recorded as raw content; anonymization is handled downstream at ingestion. + Only the user identifier is anonymized at emission. Low-cardinality signals + (rating and categories) are recorded as-is. Parameters: span: The root feedback span to annotate. @@ -121,8 +122,8 @@ def _record_feedback_request_attributes( SpanAttributes.FEEDBACK_OPERATION: "submit", SpanAttributes.USER_ID: anonymize_value(user_id) if user_id else "", SpanAttributes.FEEDBACK_CONVERSATION: feedback_request.conversation_id, - SpanAttributes.INPUT: anonymize_value(feedback_request.user_question), - SpanAttributes.OUTPUT: anonymize_value(feedback_request.llm_response), + SpanAttributes.INPUT: feedback_request.user_question, + SpanAttributes.OUTPUT: feedback_request.llm_response, }, ) if feedback_request.sentiment is not None: @@ -130,7 +131,7 @@ def _record_feedback_request_attributes( if feedback_request.user_feedback: span.set_attribute( SpanAttributes.FEEDBACK_COMMENT, - anonymize_value(feedback_request.user_feedback), + feedback_request.user_feedback, ) if feedback_request.categories: span.set_attribute( diff --git a/src/app/endpoints/query.py b/src/app/endpoints/query.py index 33865574a..64e7b56e3 100644 --- a/src/app/endpoints/query.py +++ b/src/app/endpoints/query.py @@ -161,7 +161,7 @@ async def _handle_query_with_tracing( root_span, { SpanAttributes.USER_ID: anonymize_value(user_id), - SpanAttributes.INPUT: anonymize_value(query_request.query), + SpanAttributes.INPUT: query_request.query, SpanAttributes.REQUEST_ATTACHMENTS_COUNT: ( len(query_request.attachments) if query_request.attachments else 0 ), @@ -342,7 +342,7 @@ async def _handle_query_with_tracing( SpanAttributes.SESSION_ID: conversation_id, SpanAttributes.LLM_USAGE_INPUT_TOKENS: turn_summary.token_usage.input_tokens, SpanAttributes.LLM_USAGE_OUTPUT_TOKENS: turn_summary.token_usage.output_tokens, - SpanAttributes.OUTPUT: anonymize_value(turn_summary.llm_response), + SpanAttributes.OUTPUT: turn_summary.llm_response, }, ) diff --git a/src/app/endpoints/responses.py b/src/app/endpoints/responses.py index 29340d251..20440d14c 100644 --- a/src/app/endpoints/responses.py +++ b/src/app/endpoints/responses.py @@ -189,7 +189,7 @@ def _finalize_responses_root_span( SpanAttributes.LLM_USAGE_OUTPUT_TOKENS: ( turn_summary.token_usage.output_tokens ), - SpanAttributes.OUTPUT: anonymize_value(turn_summary.llm_response), + SpanAttributes.OUTPUT: turn_summary.llm_response, }, ) add_span_event(root_span, SpanEvents.LLM_RESPONSE_COMPLETED) @@ -577,7 +577,7 @@ async def handle_responses_with_tracing( # pylint: disable=too-many-locals span_attributes: dict[str, Any] = { SpanAttributes.USER_ID: anonymize_value(user_id), - SpanAttributes.INPUT: anonymize_value(input_text), + SpanAttributes.INPUT: input_text, SpanAttributes.REQUEST_ATTACHMENTS_COUNT: attachments_count, } # safety_identifier is a caller-supplied, non-PII identifier, so it is diff --git a/src/app/endpoints/rlsapi_v1.py b/src/app/endpoints/rlsapi_v1.py index b5bc5a6f1..497abfb98 100644 --- a/src/app/endpoints/rlsapi_v1.py +++ b/src/app/endpoints/rlsapi_v1.py @@ -52,7 +52,6 @@ SpanAttributes, SpanEvents, add_span_event, - anonymize_value, set_span_attributes, ) from utils.query import ( @@ -689,9 +688,7 @@ async def infer_endpoint( # pylint: disable=R0914,R0915 endpoint_path = ENDPOINT_PATH_INFER request_id = get_suid() - span.set_attribute( - SpanAttributes.INPUT, anonymize_value(infer_request.question) - ) + span.set_attribute(SpanAttributes.INPUT, infer_request.question) logger.info("Processing rlsapi v1 /infer request %s", request_id) @@ -799,7 +796,7 @@ async def infer_endpoint( # pylint: disable=R0914,R0915 { SpanAttributes.LLM_USAGE_INPUT_TOKENS: token_usage.input_tokens, SpanAttributes.LLM_USAGE_OUTPUT_TOKENS: token_usage.output_tokens, - SpanAttributes.OUTPUT: anonymize_value(response_text), + SpanAttributes.OUTPUT: response_text, }, ) diff --git a/src/app/endpoints/streaming_query.py b/src/app/endpoints/streaming_query.py index 483639e25..bad26d314 100644 --- a/src/app/endpoints/streaming_query.py +++ b/src/app/endpoints/streaming_query.py @@ -206,7 +206,7 @@ async def _handle_streaming_query_with_tracing( # pylint: disable=too-many-loca root_span, { SpanAttributes.USER_ID: anonymize_value(user_id), - SpanAttributes.INPUT: anonymize_value(query_request.query), + SpanAttributes.INPUT: query_request.query, SpanAttributes.REQUEST_ATTACHMENTS_COUNT: ( len(query_request.attachments) if query_request.attachments else 0 ), diff --git a/src/utils/agents/streaming.py b/src/utils/agents/streaming.py index 8e9a7256a..8c80f6530 100644 --- a/src/utils/agents/streaming.py +++ b/src/utils/agents/streaming.py @@ -70,7 +70,6 @@ SpanAttributes, SpanEvents, add_span_event, - anonymize_value, set_span_attributes, ) from utils.pydantic_ai_helpers import build_agent, captured_output_items @@ -400,7 +399,7 @@ async def generate_agent_response( # pylint: disable=too-many-statements SpanAttributes.LLM_USAGE_OUTPUT_TOKENS: ( turn_summary.token_usage.output_tokens ), - SpanAttributes.OUTPUT: anonymize_value(turn_summary.llm_response), + SpanAttributes.OUTPUT: turn_summary.llm_response, }, ) add_span_event(root_span, SpanEvents.LLM_RESPONSE_COMPLETED) diff --git a/src/utils/otel_tracing.py b/src/utils/otel_tracing.py index ddbd5948f..96f83df2f 100644 --- a/src/utils/otel_tracing.py +++ b/src/utils/otel_tracing.py @@ -23,10 +23,10 @@ class SpanAttributes(StrEnum): """OpenTelemetry span attribute keys for LCS instrumentation.""" SESSION_ID = "session.id" - USER_ID = "user.id" # anonymized + USER_ID = "user.id" # anonymized (identity field) SAFETY_IDENTIFIER = "request.safety_identifier" # caller-supplied identifier - INPUT = "request.input" # anonymized - OUTPUT = "response.output" # anonymized + INPUT = "request.input" # raw content (anonymization handled downstream) + OUTPUT = "response.output" # raw content (anonymization handled downstream) RESPONSE_ERROR = "response.error" RESPONSE_CAUSE = "response.cause" REQUEST_ATTACHMENTS_COUNT = "request.attachments.count" @@ -49,7 +49,9 @@ class SpanAttributes(StrEnum): FEEDBACK_OPERATION = "feedback.operation" FEEDBACK_CONVERSATION = "feedback.conversation" FEEDBACK_RATING = "feedback.rating" - FEEDBACK_COMMENT = "feedback.comment" # anonymized + FEEDBACK_COMMENT = ( + "feedback.comment" # raw content (anonymization handled downstream) + ) FEEDBACK_CATEGORIES = "feedback.categories" FEEDBACK_STATUS_CODE = "feedback.status.code" FEEDBACK_STORAGE_OUTCOME = "feedback.storage.outcome" diff --git a/src/utils/vector_search.py b/src/utils/vector_search.py index 805585826..8be710b0f 100644 --- a/src/utils/vector_search.py +++ b/src/utils/vector_search.py @@ -28,7 +28,6 @@ SpanAttributes, SpanEvents, add_span_event, - anonymize_value, set_span_attributes, ) from utils.reranker import apply_byok_rerank_boost, rerank_chunks_with_cross_encoder @@ -673,7 +672,7 @@ async def build_rag_context( # pylint: disable=too-many-locals,too-many-branche """ with tracer.start_as_current_span("rag.retrieve") as span: # Set RAG input attribute - span.set_attribute(SpanAttributes.RAG_INPUT, anonymize_value(query)) + span.set_attribute(SpanAttributes.RAG_INPUT, query) if moderation_decision == "blocked": span.set_attribute(SpanAttributes.RAG_SOURCES_COUNT, 0) diff --git a/tests/unit/app/endpoints/responses_otel_helpers.py b/tests/unit/app/endpoints/responses_otel_helpers.py index 706ea2729..d2dcef224 100644 --- a/tests/unit/app/endpoints/responses_otel_helpers.py +++ b/tests/unit/app/endpoints/responses_otel_helpers.py @@ -262,7 +262,7 @@ def assert_root_setup_attributes( assert root.name == ROOT_SPAN_NAME assert root.attributes is not None assert root.attributes[SpanAttributes.USER_ID] == f"[anon:{MOCK_AUTH[0]}]" - assert root.attributes[SpanAttributes.INPUT] == f"[anon:{input_text}]" + assert root.attributes[SpanAttributes.INPUT] == input_text assert ( root.attributes[SpanAttributes.REQUEST_ATTACHMENTS_COUNT] == attachments_count ) diff --git a/tests/unit/app/endpoints/test_a2a.py b/tests/unit/app/endpoints/test_a2a.py index be7e19d27..ac0bf936b 100644 --- a/tests/unit/app/endpoints/test_a2a.py +++ b/tests/unit/app/endpoints/test_a2a.py @@ -1581,7 +1581,7 @@ async def _mock_asgi(_scope: Any, _receive: Any, send: Any) -> None: attrs = dict(span.attributes or {}) assert attrs["a2a.rpc.method"] == "message/send" - assert attrs["a2a.request.id"].startswith("[hash:") + assert attrs["a2a.request.id"] == "req-1" assert "user.id" in attrs event_names = [e.name for e in span.events] @@ -1649,7 +1649,7 @@ async def _mock_asgi(_scope: Any, _receive: Any, send: Any) -> None: attrs = dict(span.attributes or {}) assert attrs["a2a.rpc.method"] == "message/stream" - assert attrs["a2a.request.id"].startswith("[hash:") + assert attrs["a2a.request.id"] == "req-2" event_names = [e.name for e in span.events] assert "a2a.dispatch.start" in event_names diff --git a/tests/unit/app/endpoints/test_feedback.py b/tests/unit/app/endpoints/test_feedback.py index f7ae7ca31..3b0b18c6f 100644 --- a/tests/unit/app/endpoints/test_feedback.py +++ b/tests/unit/app/endpoints/test_feedback.py @@ -614,11 +614,12 @@ async def test_submit_emits_root_and_storage_spans( assert attrs["feedback.rating"] == -1 assert attrs["feedback.categories"] == "incorrect,incomplete" assert attrs["feedback.status.code"] == status.HTTP_200_OK - # Free-text fields are anonymized. + # The user id is anonymized; free-text content is recorded raw + # (anonymization is handled downstream at ingestion). assert str(attrs["user.id"]).startswith("[hash:") - assert str(attrs["request.input"]).startswith("[hash:") - assert str(attrs["response.output"]).startswith("[hash:") - assert str(attrs["feedback.comment"]).startswith("[hash:") + assert attrs["request.input"] == VALID_BASE["user_question"] + assert attrs["response.output"] == VALID_BASE["llm_response"] + assert attrs["feedback.comment"] == "The answer was too vague." storage_attrs = dict(storage.attributes or {}) assert storage_attrs["feedback.storage.outcome"] == "success" diff --git a/tests/unit/app/endpoints/test_query_otel.py b/tests/unit/app/endpoints/test_query_otel.py index 9b76e6f14..84e325230 100644 --- a/tests/unit/app/endpoints/test_query_otel.py +++ b/tests/unit/app/endpoints/test_query_otel.py @@ -120,7 +120,8 @@ async def test_query_root_span_attributes_and_events( The mocked success path validates the request, persists the turn, and completes the LLM response, so the validation/turn-persisted/LLM-response - events are all recorded, and the anonymized user/input attributes are set. + events are all recorded, the user id is anonymized, and the raw + input/output content is set. """ tracer, exporter = otel mocker.patch(f"{MODULE}.configuration", minimal_config) @@ -142,9 +143,12 @@ async def test_query_root_span_attributes_and_events( root = next(s for s in exporter.get_finished_spans() if s.name == QUERY_SPAN_NAME) attrs = dict(root.attributes or {}) assert attrs[SpanAttributes.USER_ID] == f"[anon:{MOCK_AUTH[0]}]" - assert attrs[SpanAttributes.INPUT] == f"[anon:{QUERY_TEXT}]" + assert attrs[SpanAttributes.INPUT] == QUERY_TEXT assert attrs[SpanAttributes.REQUEST_ATTACHMENTS_COUNT] == 0 - assert SpanAttributes.OUTPUT in attrs + assert ( + attrs[SpanAttributes.OUTPUT] + == "Kubernetes is a container orchestration platform" + ) assert SpanAttributes.SESSION_ID in attrs event_names = {event.name for event in root.events} @@ -197,7 +201,7 @@ def _raise_quota_exceeded(*_args: object, **_kwargs: object) -> None: root = next(s for s in exporter.get_finished_spans() if s.name == QUERY_SPAN_NAME) attrs = dict(root.attributes or {}) assert attrs[SpanAttributes.USER_ID] == f"[anon:{MOCK_AUTH[0]}]" - assert attrs[SpanAttributes.INPUT] == f"[anon:{QUERY_TEXT}]" + assert attrs[SpanAttributes.INPUT] == QUERY_TEXT assert attrs[SpanAttributes.REQUEST_ATTACHMENTS_COUNT] == 0 event_names = {event.name for event in root.events} diff --git a/tests/unit/app/endpoints/test_responses_otel.py b/tests/unit/app/endpoints/test_responses_otel.py index ec1be32b8..6f51c1534 100644 --- a/tests/unit/app/endpoints/test_responses_otel.py +++ b/tests/unit/app/endpoints/test_responses_otel.py @@ -84,7 +84,6 @@ class TestFinalizeResponsesRootSpanOtel: # pylint: disable=too-few-public-metho ) def test_finalize_tool_attrs_and_events( self, - mocker: MockerFixture, otel: tuple[Any, InMemorySpanExporter], tool_names: list[str], expect_tool_event: bool, @@ -97,11 +96,6 @@ def test_finalize_tool_attrs_and_events( if tool_names else make_turn_summary_without_tools() ) - mocker.patch( - f"{MODULE}.anonymize_value", - side_effect=lambda value: f"[anon:{value}]", - ) - _finalize_responses_root_span(root_span, turn_summary) root_span.end() @@ -116,7 +110,7 @@ def test_finalize_tool_attrs_and_events( ) assert span.attributes[SpanAttributes.LLM_USAGE_INPUT_TOKENS] == 10 assert span.attributes[SpanAttributes.LLM_USAGE_OUTPUT_TOKENS] == 5 - assert span.attributes[SpanAttributes.OUTPUT] == "[anon:The answer is 42]" + assert span.attributes[SpanAttributes.OUTPUT] == "The answer is 42" event_names = [event.name for event in span.events] assert SpanEvents.LLM_RESPONSE_COMPLETED in event_names diff --git a/tests/unit/app/endpoints/test_rlsapi_v1.py b/tests/unit/app/endpoints/test_rlsapi_v1.py index 97249269d..030012f56 100644 --- a/tests/unit/app/endpoints/test_rlsapi_v1.py +++ b/tests/unit/app/endpoints/test_rlsapi_v1.py @@ -1826,8 +1826,8 @@ async def test_infer_span_success_attributes( # pylint: disable=too-many-locals assert attrs["shield.result"] == "passed" assert "request.input" in attrs assert "response.output" in attrs - assert str(attrs["request.input"]).startswith("[hash:") - assert str(attrs["response.output"]).startswith("[hash:") + assert attrs["request.input"] == "How do I list files?" + assert attrs["response.output"] == "This is a test LLM response." @pytest.mark.asyncio async def test_infer_span_events( diff --git a/tests/unit/app/endpoints/test_streaming_query.py b/tests/unit/app/endpoints/test_streaming_query.py index 4335edb60..ac8985908 100644 --- a/tests/unit/app/endpoints/test_streaming_query.py +++ b/tests/unit/app/endpoints/test_streaming_query.py @@ -725,7 +725,7 @@ async def test_sets_initial_span_attributes( assert root.attributes[SpanAttributes.USER_ID] == ( "[anon:00000001-0001-0001-0001-000000000001]" ) - assert root.attributes[SpanAttributes.INPUT] == "[anon:What is Kubernetes?]" + assert root.attributes[SpanAttributes.INPUT] == "What is Kubernetes?" assert root.attributes[SpanAttributes.REQUEST_ATTACHMENTS_COUNT] == 0 @pytest.mark.asyncio diff --git a/tests/unit/utils/agents/test_streaming.py b/tests/unit/utils/agents/test_streaming.py index 3f1edf434..66ecbd0e4 100644 --- a/tests/unit/utils/agents/test_streaming.py +++ b/tests/unit/utils/agents/test_streaming.py @@ -916,11 +916,6 @@ async def inner() -> AsyncIterator[str]: mock_config.quota_limiters = [] mocker.patch("utils.agents.streaming.configuration", mock_config) - mocker.patch( - "utils.agents.streaming.anonymize_value", - side_effect=lambda v: f"[anon:{v}]", - ) - [ event async for event in generate_agent_response( @@ -940,7 +935,7 @@ async def inner() -> AsyncIterator[str]: assert span.attributes[SpanAttributes.SESSION_ID] == context.conversation_id assert span.attributes[SpanAttributes.LLM_USAGE_INPUT_TOKENS] == 10 assert span.attributes[SpanAttributes.LLM_USAGE_OUTPUT_TOKENS] == 5 - assert span.attributes[SpanAttributes.OUTPUT] == "[anon:The answer is 42]" + assert span.attributes[SpanAttributes.OUTPUT] == "The answer is 42" event_names = [e.name for e in span.events] assert SpanEvents.TURN_PERSISTED in event_names assert SpanEvents.LLM_RESPONSE_COMPLETED in event_names @@ -985,10 +980,6 @@ async def inner() -> AsyncIterator[str]: mock_config = mocker.Mock() mock_config.quota_limiters = [] mocker.patch("utils.agents.streaming.configuration", mock_config) - mocker.patch( - "utils.agents.streaming.anonymize_value", - side_effect=lambda v: f"[anon:{v}]", - ) [ event diff --git a/tests/unit/utils/test_vector_search.py b/tests/unit/utils/test_vector_search.py index a42388763..3970f2072 100644 --- a/tests/unit/utils/test_vector_search.py +++ b/tests/unit/utils/test_vector_search.py @@ -1705,10 +1705,6 @@ async def test_blocked_moderation_sets_zero_sources_without_completed_event( """Blocked moderation skips retrieval and does not emit completed event.""" tracer, exporter = otel mocker.patch("utils.vector_search.tracer", tracer) - mocker.patch( - "utils.vector_search.anonymize_value", - side_effect=lambda value: f"[anon:{value}]", - ) self._patch_rag_config(mocker) client = mocker.AsyncMock() @@ -1720,7 +1716,7 @@ async def test_blocked_moderation_sets_zero_sources_without_completed_event( if span.name == "rag.retrieve" ) assert span.attributes is not None - assert span.attributes[SpanAttributes.RAG_INPUT] == "[anon:test query]" + assert span.attributes[SpanAttributes.RAG_INPUT] == "test query" assert span.attributes[SpanAttributes.RAG_SOURCES_COUNT] == 0 event_names = [event.name for event in span.events] assert SpanEvents.RAG_RETRIEVAL_COMPLETED not in event_names @@ -1734,10 +1730,6 @@ async def test_passed_with_no_chunks_emits_zero_count_event( """Passed moderation with no chunks emits retrieval completed with count 0.""" tracer, exporter = otel mocker.patch("utils.vector_search.tracer", tracer) - mocker.patch( - "utils.vector_search.anonymize_value", - side_effect=lambda value: f"[anon:{value}]", - ) self._patch_rag_config(mocker) mocker.patch( "utils.vector_search._fetch_byok_rag", @@ -1776,10 +1768,6 @@ async def test_passed_with_chunks_sets_sources_and_chunk_count( """Passed moderation with chunks sets source attrs and chunk count event.""" tracer, exporter = otel mocker.patch("utils.vector_search.tracer", tracer) - mocker.patch( - "utils.vector_search.anonymize_value", - side_effect=lambda value: f"[anon:{value}]", - ) self._patch_rag_config(mocker) chunk = RAGChunk( content="chunk text",