Skip to main content

ai_agents_runtime/optimization/
observability.rs

1use std::collections::HashMap;
2use std::future::Future;
3
4use ai_agents_observability::{
5    EventStatus, EventType, ObservabilityManager, ObservationPurpose,
6    with_updated_observation_context,
7};
8use std::sync::Arc;
9
10use super::branch::{RuntimeCommitBehavior, RuntimeOptimizationKind};
11
12/// Runs a branch future with task-local labels that defer observed events until finalization.
13pub async fn with_branch_observation<F, T>(
14    branch_id: &str,
15    optimization: RuntimeOptimizationKind,
16    commit_behavior: RuntimeCommitBehavior,
17    future: F,
18) -> T
19where
20    F: Future<Output = T>,
21{
22    let branch = branch_id.to_string();
23    with_updated_observation_context(
24        move |context| {
25            context
26                .with_tag("runtime.defer_observation", "true")
27                .with_tag("runtime.branch_id", branch)
28                .with_tag("runtime.optimization", optimization.as_label())
29                .with_tag("runtime.commit_behavior", commit_behavior.as_label())
30                .with_tag("runtime.speculative", "true")
31        },
32        future,
33    )
34    .await
35}
36
37/// Builds finalization tags shared by all branch outcomes.
38pub fn branch_finalization_tags(
39    optimization: RuntimeOptimizationKind,
40    commit_behavior: RuntimeCommitBehavior,
41) -> HashMap<String, String> {
42    let mut tags = HashMap::new();
43    tags.insert(
44        "runtime.optimization".to_string(),
45        optimization.as_label().to_string(),
46    );
47    tags.insert(
48        "optimization".to_string(),
49        optimization.as_label().to_string(),
50    );
51    tags.insert(
52        "runtime.commit_behavior".to_string(),
53        commit_behavior.as_label().to_string(),
54    );
55    tags.insert(
56        "commit_behavior".to_string(),
57        commit_behavior.as_label().to_string(),
58    );
59    tags.insert("runtime.speculative".to_string(), "true".to_string());
60    tags.insert("speculative".to_string(), "true".to_string());
61    tags
62}
63
64/// Finalizes a branch if an observability manager is configured.
65pub fn finalize_branch(
66    manager: Option<&Arc<ObservabilityManager>>,
67    branch_id: &str,
68    status: &str,
69    winner: bool,
70    optimization: RuntimeOptimizationKind,
71    commit_behavior: RuntimeCommitBehavior,
72) {
73    if let Some(manager) = manager {
74        let tags = branch_finalization_tags(optimization, commit_behavior);
75        let finalized =
76            manager.finalize_pending_branch(branch_id, status.to_string(), winner, tags.clone());
77        if finalized == 0 {
78            let mut lifecycle_tags = tags;
79            lifecycle_tags.insert("runtime.branch_status".to_string(), status.to_string());
80            lifecycle_tags.insert("branch_status".to_string(), status.to_string());
81            lifecycle_tags.insert("runtime.winner".to_string(), winner.to_string());
82            lifecycle_tags.insert("winner".to_string(), winner.to_string());
83            manager.record_lifecycle_event(
84                EventType::MemoryOperation {
85                    operation: format!("runtime_branch_{}", status),
86                },
87                ObservationPurpose::Other("runtime_branch".to_string()),
88                event_status_for_branch(status),
89                0,
90                lifecycle_tags,
91                None,
92            );
93        }
94    }
95}
96
97fn event_status_for_branch(status: &str) -> EventStatus {
98    match status {
99        "failed" => EventStatus::Error,
100        "cancelled" => EventStatus::Cancelled,
101        "discarded" => EventStatus::Skipped,
102        _ => EventStatus::Success,
103    }
104}