ai_agents_runtime/optimization/
observability.rs1use 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
12pub 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
37pub 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
64pub 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}