-
Notifications
You must be signed in to change notification settings - Fork 97
Expand file tree
/
Copy pathresponses_telemetry.py
More file actions
164 lines (146 loc) · 6.07 KB
/
Copy pathresponses_telemetry.py
File metadata and controls
164 lines (146 loc) · 6.07 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
"""Splunk telemetry helpers for the Responses API endpoint.
Extracted from responses.py to reduce module size while keeping telemetry
functions co-located with the endpoint they serve.
"""
from datetime import UTC, datetime
from typing import Optional
from fastapi import BackgroundTasks
from lightspeed_stack.log import get_logger
from lightspeed_stack.models.common.responses.contexts import ResponsesContext
from lightspeed_stack.models.common.responses.responses_api_params import (
ResponsesApiParams,
)
from lightspeed_stack.models.common.turn_summary import TurnSummary
from lightspeed_stack.observability import ResponsesEventData, build_responses_event
from lightspeed_stack.observability.splunk import dispatch_splunk_event
from lightspeed_stack.utils.suid import normalize_conversation_id
logger = get_logger(__name__)
def queue_responses_splunk_event( # pylint: disable=too-many-arguments,too-many-positional-arguments
background_tasks: Optional[BackgroundTasks],
input_text: str,
response_text: str,
conversation_id: str,
model: str,
rh_identity_context: tuple[str, str],
inference_time: float,
sourcetype: str,
input_tokens: int = 0,
output_tokens: int = 0,
fire_and_forget: bool = False,
user_agent: Optional[str] = None,
) -> None:
"""Build and queue a Splunk telemetry event for the responses endpoint.
No-op when background_tasks is None and fire_and_forget is False
(Splunk telemetry disabled).
Args:
background_tasks: FastAPI background task manager, or None if disabled.
input_text: User input text.
response_text: Response text from LLM or shield.
conversation_id: Conversation identifier.
model: Model name used for inference.
rh_identity_context: Tuple of (org_id, system_id) from RH identity.
inference_time: Request processing duration in seconds.
sourcetype: Splunk sourcetype for the event.
input_tokens: Number of prompt tokens consumed.
output_tokens: Number of completion tokens produced.
fire_and_forget: When True, dispatch via asyncio.create_task() instead
of background_tasks. Use for error paths where an HTTPException
follows, since FastAPI discards BackgroundTasks on non-2xx responses.
user_agent: Sanitized User-Agent string from the request header, or None.
"""
if not fire_and_forget and background_tasks is None:
return
event_data = ResponsesEventData(
input_text=input_text,
response_text=response_text,
conversation_id=conversation_id,
model=model,
org_id=rh_identity_context[0],
system_id=rh_identity_context[1],
inference_time=inference_time,
input_tokens=input_tokens,
output_tokens=output_tokens,
user_agent=user_agent,
)
event = build_responses_event(event_data)
dispatch_splunk_event(
event,
sourcetype,
background_tasks=background_tasks,
fire_and_forget=fire_and_forget,
)
def queue_responses_error_event(
error: Exception,
api_params: ResponsesApiParams,
context: ResponsesContext,
) -> None:
"""Queue fire-and-forget Splunk telemetry for a Responses API error.
Args:
error: The backend exception being converted into an HTTP error.
api_params: Responses API parameters for the failed request.
context: Request-scoped Responses API context.
"""
queue_responses_splunk_event(
background_tasks=context.background_tasks,
input_text=context.input_text,
response_text=type(error).__name__,
conversation_id=normalize_conversation_id(api_params.conversation),
model=api_params.model,
rh_identity_context=context.rh_identity_context,
inference_time=(datetime.now(UTC) - context.started_at).total_seconds(),
sourcetype="responses_error",
fire_and_forget=True,
user_agent=context.user_agent,
)
def queue_blocked_response_event(
api_params: ResponsesApiParams,
context: ResponsesContext,
response_text: str,
) -> None:
"""Queue Splunk telemetry for a shield-blocked Responses API request.
Args:
api_params: Responses API parameters for the blocked request.
context: Request-scoped Responses API context.
response_text: Refusal text sent to the client.
"""
queue_responses_splunk_event(
background_tasks=context.background_tasks,
input_text=context.input_text,
response_text=response_text,
conversation_id=normalize_conversation_id(api_params.conversation),
model=api_params.model,
rh_identity_context=context.rh_identity_context,
inference_time=(datetime.now(UTC) - context.started_at).total_seconds(),
sourcetype="responses_shield_blocked",
user_agent=context.user_agent,
)
def queue_completed_response_event(
api_params: ResponsesApiParams,
context: ResponsesContext,
turn_summary: TurnSummary,
completed_at: datetime,
response_text: str,
) -> None:
"""Queue Splunk telemetry for a completed Responses API request.
Args:
api_params: Responses API parameters for the completed request.
context: Request-scoped Responses API context.
turn_summary: Summary containing token usage for telemetry.
completed_at: Time when response handling completed.
response_text: Final text sent to the client.
"""
if context.moderation_result.decision != "passed":
return
queue_responses_splunk_event(
background_tasks=context.background_tasks,
input_text=context.input_text,
response_text=response_text,
conversation_id=normalize_conversation_id(api_params.conversation),
model=api_params.model,
rh_identity_context=context.rh_identity_context,
inference_time=(completed_at - context.started_at).total_seconds(),
sourcetype="responses_completed",
input_tokens=turn_summary.token_usage.input_tokens,
output_tokens=turn_summary.token_usage.output_tokens,
user_agent=context.user_agent,
)