Skip to main content

ferrin_core/telemetry/
dispatcher.rs

1//! Fans telemetry events out to the configured integrations and to
2//! `tracing`.
3
4use 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/// Dispatches events to every integration of a [`TelemetryOptions`].
32#[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    /// Filters the telemetry copy without changing application hooks or results.
141    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    /// Wraps `call` with every integration's `execute_language_model_call`;
179    /// the first integration becomes the outermost wrapper.
180    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    /// Wraps `call` with every integration's `execute_tool`.
198    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    /// Convenience for tests: returns `true` when `self` is the same object.
231    #[must_use]
232    pub fn ptr_eq(this: &Arc<Self>, other: &Arc<Self>) -> bool {
233        Arc::ptr_eq(this, other)
234    }
235}