use std::collections::BTreeMap;
use std::fmt;
use std::sync::Arc;
use ferrin_spec::BoxFuture;
use ferrin_spec::GenerateResult;
use ferrin_spec::JsonValue;
use ferrin_spec::StreamResult;
use ferrin_spec::ToolCallId;
use ferrin_spec::ToolName;
use ferrin_tool::ToolError;
use crate::error::Error;
pub(crate) mod dispatcher;
mod events;
mod redact;
mod redact_provider;
pub(crate) mod spans;
pub(crate) use dispatcher::TelemetryDispatcher;
pub use events::AbortEvent;
pub use events::EmbedEndEvent;
pub use events::EmbedStartEvent;
pub use events::EndEvent;
pub use events::ErrorEvent;
pub use events::ErrorPhase;
pub use events::ModelCallEndEvent;
pub use events::ModelCallStartEvent;
pub use events::ModelIdentity;
pub use events::RecordedInputs;
pub use events::RerankEndEvent;
pub use events::RerankStartEvent;
pub use events::StartEvent;
pub use events::StepEndEvent;
pub use events::StepStartEvent;
pub use events::ToolExecutionEndEvent;
pub use events::ToolExecutionStartEvent;
pub use events::ToolOutcome;
pub use spans::WARNINGS_TARGET;
#[derive(Debug, Clone)]
pub struct ModelCallContext {
pub call_id: String,
pub step_number: u32,
pub model: ModelIdentity,
pub function_id: Option<String>,
}
#[derive(Debug, Clone)]
pub struct ToolExecutionContext {
pub call_id: String,
pub tool_call_id: ToolCallId,
pub tool_name: ToolName,
pub input: Option<JsonValue>,
}
#[derive(Debug)]
#[non_exhaustive]
pub enum ModelCallOutcome {
Generate(Box<GenerateResult>),
Stream(Box<StreamResult>),
}
pub trait Telemetry: Send + Sync + 'static {
fn on_start(&self, _event: &StartEvent) {}
fn on_step_start(&self, _event: &StepStartEvent) {}
fn on_language_model_call_start(&self, _event: &ModelCallStartEvent) {}
fn on_language_model_call_end(&self, _event: &ModelCallEndEvent) {}
fn on_tool_execution_start(&self, _event: &ToolExecutionStartEvent) {}
fn on_tool_execution_end(&self, _event: &ToolExecutionEndEvent) {}
fn on_step_end(&self, _event: &StepEndEvent) {}
fn on_embed_start(&self, _event: &EmbedStartEvent) {}
fn on_embed_end(&self, _event: &EmbedEndEvent) {}
fn on_rerank_start(&self, _event: &RerankStartEvent) {}
fn on_rerank_end(&self, _event: &RerankEndEvent) {}
fn on_end(&self, _event: &EndEvent) {}
fn on_abort(&self, _event: &AbortEvent) {}
fn on_error(&self, _event: &ErrorEvent<'_>) {}
fn execute_language_model_call<'a>(
&'a self,
_ctx: &'a ModelCallContext,
call: BoxFuture<'a, Result<ModelCallOutcome, Error>>,
) -> BoxFuture<'a, Result<ModelCallOutcome, Error>> {
call
}
fn execute_tool<'a>(
&'a self,
_ctx: &'a ToolExecutionContext,
call: BoxFuture<'a, Result<ToolOutcome, ToolError>>,
) -> BoxFuture<'a, Result<ToolOutcome, ToolError>> {
call
}
}
#[derive(Clone, Default)]
pub struct TelemetryOptions {
pub enabled: bool,
pub record_inputs: bool,
pub record_outputs: bool,
pub function_id: Option<String>,
pub metadata: BTreeMap<String, JsonValue>,
pub include_tools_context: bool,
pub integrations: Vec<Arc<dyn Telemetry>>,
}
impl TelemetryOptions {
#[must_use]
pub fn enabled() -> Self {
Self {
enabled: true,
..Self::default()
}
}
#[must_use]
pub fn with_integration(mut self, integration: Arc<dyn Telemetry>) -> Self {
self.integrations.push(integration);
self
}
}
impl fmt::Debug for TelemetryOptions {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("TelemetryOptions")
.field("enabled", &self.enabled)
.field("record_inputs", &self.record_inputs)
.field("record_outputs", &self.record_outputs)
.field("function_id", &self.function_id)
.field("metadata", &self.metadata)
.field("include_tools_context", &self.include_tools_context)
.field("integrations", &self.integrations.len())
.finish()
}
}