ferrin_core/telemetry/
dispatcher.rs1use std::sync::Arc;
5
6use ferrin_spec::BoxFuture;
7use ferrin_tool::ToolError;
8
9use super::AbortEvent;
10use super::EmbedEndEvent;
11use super::EmbedStartEvent;
12use super::EndEvent;
13use super::ErrorEvent;
14use super::ModelCallContext;
15use super::ModelCallEndEvent;
16use super::ModelCallOutcome;
17use super::ModelCallStartEvent;
18use super::RerankEndEvent;
19use super::RerankStartEvent;
20use super::StartEvent;
21use super::StepEndEvent;
22use super::StepStartEvent;
23use super::Telemetry;
24use super::TelemetryOptions;
25use super::ToolExecutionContext;
26use super::ToolExecutionEndEvent;
27use super::ToolExecutionStartEvent;
28use super::ToolOutcome;
29use crate::error::Error;
30
31#[derive(Clone, Debug)]
33pub(crate) struct TelemetryDispatcher {
34 options: Arc<TelemetryOptions>,
35}
36
37macro_rules! dispatch {
38 ($name:ident, $event:ty) => {
39 pub(crate) fn $name(&self, event: &$event) {
40 if !self.options.enabled {
41 return;
42 }
43 for integration in &self.options.integrations {
44 integration.$name(event);
45 }
46 }
47 };
48}
49
50impl TelemetryDispatcher {
51 pub(crate) fn new(options: TelemetryOptions) -> Self {
52 Self {
53 options: Arc::new(options),
54 }
55 }
56
57 pub(crate) fn record_inputs(&self) -> bool {
58 self.options.enabled && self.options.record_inputs
59 }
60
61 pub(crate) fn record_outputs(&self) -> bool {
62 self.options.enabled && self.options.record_outputs
63 }
64
65 dispatch!(on_start, StartEvent);
66 dispatch!(on_step_start, StepStartEvent);
67 dispatch!(on_language_model_call_start, ModelCallStartEvent);
68
69 dispatch!(on_tool_execution_start, ToolExecutionStartEvent);
70 pub(crate) fn on_tool_execution_end(&self, event: &ToolExecutionEndEvent) {
71 if !self.options.enabled {
72 return;
73 }
74 let mut recorded = event.clone();
75 if !self.record_outputs() {
76 recorded.output = None;
77 recorded.error = recorded
78 .error
79 .map(|_| crate::generate_text::ToolErrorInfo::text(super::redact::REDACTED));
80 }
81 for integration in &self.options.integrations {
82 integration.on_tool_execution_end(&recorded);
83 }
84 }
85
86 dispatch!(on_abort, AbortEvent);
87
88 pub(crate) fn on_language_model_call_end(&self, event: &ModelCallEndEvent) {
89 if !self.options.enabled {
90 return;
91 }
92 let mut recorded = event.clone();
93 if !(self.record_inputs() && self.record_outputs()) {
94 recorded.warnings = super::redact::warnings(&recorded.warnings);
95 }
96 if !self.record_outputs() {
97 recorded.content = None;
98 recorded.response.body = None;
99 }
100 for integration in &self.options.integrations {
101 integration.on_language_model_call_end(&recorded);
102 }
103 }
104
105 pub(crate) fn on_step_end(&self, event: &StepEndEvent) {
106 if !self.options.enabled {
107 return;
108 }
109 let recorded = StepEndEvent {
110 call_id: event.call_id.clone(),
111 step: Arc::new(self.recorded_step(&event.step)),
112 };
113 for integration in &self.options.integrations {
114 integration.on_step_end(&recorded);
115 }
116 }
117
118 pub(crate) fn on_end(&self, event: &EndEvent) {
119 if !self.options.enabled {
120 return;
121 }
122 let recorded = EndEvent {
123 call_id: event.call_id.clone(),
124 steps: event
125 .steps
126 .iter()
127 .map(|step| self.recorded_step(step))
128 .collect(),
129 total_usage: event.total_usage.clone(),
130 output_recorded: self
131 .record_outputs()
132 .then(|| event.output_recorded.clone())
133 .flatten(),
134 };
135 for integration in &self.options.integrations {
136 integration.on_end(&recorded);
137 }
138 }
139
140 fn recorded_step(
142 &self,
143 step: &crate::generate_text::StepResult,
144 ) -> crate::generate_text::StepResult {
145 let mut recorded = step.clone();
146 if !(self.record_inputs() && self.record_outputs()) {
147 recorded.warnings = super::redact::warnings(&recorded.warnings);
148 }
149 if !self.record_inputs() {
150 recorded.request.body = None;
151 recorded.request.messages = None;
152 }
153 if !self.record_outputs() {
154 recorded.content.clear();
155 recorded.response.body = None;
156 recorded.response.messages.clear();
157 recorded.provider_metadata = None;
158 }
159 recorded
160 }
161
162 pub(crate) fn on_error(&self, event: &ErrorEvent<'_>) {
163 if !self.options.enabled {
164 return;
165 }
166 let error = (!(self.record_inputs() && self.record_outputs()))
167 .then(|| super::redact::redact_error(event.error, &self.options));
168 let recorded = ErrorEvent {
169 call_id: event.call_id,
170 error: error.as_ref().unwrap_or(event.error),
171 phase: event.phase,
172 };
173 for integration in &self.options.integrations {
174 integration.on_error(&recorded);
175 }
176 }
177
178 pub(crate) fn execute_language_model_call<'a>(
181 &'a self,
182 ctx: &'a ModelCallContext,
183 call: BoxFuture<'a, Result<ModelCallOutcome, Error>>,
184 ) -> BoxFuture<'a, Result<ModelCallOutcome, Error>> {
185 if !self.options.enabled {
186 return call;
187 }
188 self.options
189 .integrations
190 .iter()
191 .rev()
192 .fold(call, |inner, integration| {
193 integration.execute_language_model_call(ctx, inner)
194 })
195 }
196
197 pub(crate) fn execute_tool<'a>(
199 &'a self,
200 ctx: &'a ToolExecutionContext,
201 call: BoxFuture<'a, Result<ToolOutcome, ToolError>>,
202 ) -> BoxFuture<'a, Result<ToolOutcome, ToolError>> {
203 if !self.options.enabled {
204 return call;
205 }
206 self.options
207 .integrations
208 .iter()
209 .rev()
210 .fold(call, |inner, integration| {
211 integration.execute_tool(ctx, inner)
212 })
213 }
214}
215
216impl TelemetryDispatcher {
217 dispatch!(on_embed_start, EmbedStartEvent);
218 dispatch!(on_embed_end, EmbedEndEvent);
219 dispatch!(on_rerank_start, RerankStartEvent);
220 dispatch!(on_rerank_end, RerankEndEvent);
221}
222
223impl Default for TelemetryDispatcher {
224 fn default() -> Self {
225 Self::new(TelemetryOptions::default())
226 }
227}
228
229impl dyn Telemetry {
230 #[must_use]
232 pub fn ptr_eq(this: &Arc<Self>, other: &Arc<Self>) -> bool {
233 Arc::ptr_eq(this, other)
234 }
235}