hyperactor/mailbox/
headers.rs1use std::any::type_name;
15use std::time::SystemTime;
16
17use hyperactor_config::Flattrs;
18use hyperactor_config::attrs::OPERATION_CONTEXT_HEADER;
19use hyperactor_config::attrs::declare_attrs;
20use hyperactor_config::global;
21
22use crate::ActorAddr;
23use crate::PortAddr;
24use crate::metrics::MESSAGE_LATENCY_MICROS;
25use crate::ordering::SeqInfo;
26
27declare_attrs! {
28 pub attr SEND_TIMESTAMP: SystemTime;
30
31 pub attr RUST_MESSAGE_TYPE: String;
33
34 pub attr SENDER_ACTOR_ID_HASH: u64;
36
37 pub attr SENDER_ACTOR_ID: ActorAddr;
52
53 pub attr TELEMETRY_MESSAGE_ID: u64;
55
56 pub attr TELEMETRY_PORT_INDEX: u64;
58
59 @meta(OPERATION_CONTEXT_HEADER = true)
78 pub attr OPERATION_ENDPOINT: String;
79
80 @meta(OPERATION_CONTEXT_HEADER = true)
84 pub attr OPERATION_ADVERB: String;
85}
86
87pub fn set_send_timestamp(headers: &mut Flattrs) {
89 if !headers.contains_key(SEND_TIMESTAMP) {
90 let time = std::time::SystemTime::now();
91 headers.set(SEND_TIMESTAMP, time);
92 }
93}
94
95pub fn set_rust_message_type<M>(headers: &mut Flattrs) {
97 headers.set(RUST_MESSAGE_TYPE, type_name::<M>().to_string());
98}
99
100#[doc(hidden)]
111pub fn stamp_sender_actor_id(
112 headers: &mut Flattrs,
113 seq_info: &SeqInfo,
114 dest: &PortAddr,
115 owner: &ActorAddr,
116) {
117 if let SeqInfo::Session { seq, .. } = seq_info {
118 let early_session = *seq <= 4;
119 let caller_set_stale = headers.contains_key(SENDER_ACTOR_ID);
120 if (early_session || caller_set_stale) && dest.is_handler_port() {
121 headers.set(SENDER_ACTOR_ID, owner.clone());
122 }
123 }
124}
125
126pub(crate) fn stamp_sender_actor_id_fresh(
130 headers: &mut Flattrs,
131 seq: u64,
132 dest: &PortAddr,
133 owner: &ActorAddr,
134) {
135 if seq <= 4 && dest.is_handler_port() {
136 headers.set(SENDER_ACTOR_ID, owner.clone());
137 }
138}
139
140pub fn log_message_latency_if_sampling(headers: &Flattrs, actor_id: String) {
144 if fastrand::f32() > global::get(crate::config::MESSAGE_LATENCY_SAMPLING_RATE) {
145 return;
146 }
147
148 if !headers.contains_key(SEND_TIMESTAMP) {
149 tracing::debug!(
150 actor_id = actor_id,
151 "SEND_TIMESTAMP missing from message headers, cannot measure latency"
152 );
153 return;
154 }
155
156 let metric_pairs = hyperactor_telemetry::kv_pairs!(
157 "actor_id" => actor_id
158 );
159 let Some(send_timestamp) = headers.get(SEND_TIMESTAMP) else {
160 return;
161 };
162 let now = std::time::SystemTime::now();
163 let latency = now.duration_since(send_timestamp).unwrap_or_default();
164 MESSAGE_LATENCY_MICROS.record(latency.as_micros() as f64, metric_pairs);
165}
166
167#[cfg(test)]
168mod tests {
169 use uuid::Uuid;
170
171 use super::*;
172 use crate::port::ControlPort;
173 use crate::port::Port;
174 use crate::testing::ids::test_actor_id;
175
176 fn session(seq: u64) -> SeqInfo {
177 SeqInfo::Session {
178 session_id: Uuid::now_v7(),
179 seq,
180 }
181 }
182
183 fn handler_port(actor_name: &str) -> (ActorAddr, PortAddr) {
184 let addr: ActorAddr = test_actor_id(actor_name, "worker");
185 let port = addr.port_addr(Port::handler::<TestHandlerMsg>());
186 (addr, port)
187 }
188
189 fn non_handler_port(actor_name: &str) -> (ActorAddr, PortAddr) {
190 let addr: ActorAddr = test_actor_id(actor_name, "worker");
191 let port = addr.port_addr(Port::from(1));
195 (addr, port)
196 }
197
198 fn control_port(actor_name: &str) -> (ActorAddr, PortAddr) {
199 let addr: ActorAddr = test_actor_id(actor_name, "worker");
200 let port = addr.port_addr(Port::control(ControlPort::Introspect));
201 (addr, port)
202 }
203
204 #[derive(typeuri::Named)]
207 struct TestHandlerMsg;
208
209 #[test]
210 fn test_stamp_helper_sets_sender_on_seq_1() {
211 let (owner, dest) = handler_port("test_0");
212 let mut headers = Flattrs::new();
213 stamp_sender_actor_id(&mut headers, &session(1), &dest, &owner);
214 assert_eq!(headers.get(SENDER_ACTOR_ID), Some(owner));
215 }
216
217 #[test]
218 fn test_stamp_helper_sets_sender_on_seq_4() {
219 let (owner, dest) = handler_port("test_0");
220 let mut headers = Flattrs::new();
221 stamp_sender_actor_id(&mut headers, &session(4), &dest, &owner);
222 assert_eq!(headers.get(SENDER_ACTOR_ID), Some(owner));
223 }
224
225 #[test]
226 fn test_stamp_helper_skips_seq_5_no_stale() {
227 let (owner, dest) = handler_port("test_0");
228 let mut headers = Flattrs::new();
229 stamp_sender_actor_id(&mut headers, &session(5), &dest, &owner);
230 assert_eq!(headers.get(SENDER_ACTOR_ID), None);
231 }
232
233 #[test]
234 fn test_stamp_helper_overwrites_stale_at_seq_5() {
235 let (owner, dest) = handler_port("test_0");
236 let fake_owner: ActorAddr = test_actor_id("fake_0", "imposter");
237 let mut headers = Flattrs::new();
238 headers.set(SENDER_ACTOR_ID, fake_owner.clone());
239 stamp_sender_actor_id(&mut headers, &session(5), &dest, &owner);
240 assert_eq!(headers.get(SENDER_ACTOR_ID), Some(owner));
241 }
242
243 #[test]
244 fn test_stamp_helper_skips_non_handler_port() {
245 let (owner, dest) = non_handler_port("test_0");
246 let mut headers = Flattrs::new();
247 stamp_sender_actor_id(&mut headers, &session(1), &dest, &owner);
248 assert_eq!(headers.get(SENDER_ACTOR_ID), None);
249 }
250
251 #[test]
252 fn test_stamp_helper_skips_control_port() {
253 let (owner, dest) = control_port("test_0");
254 let mut headers = Flattrs::new();
255 stamp_sender_actor_id(&mut headers, &session(1), &dest, &owner);
256 assert_eq!(headers.get(SENDER_ACTOR_ID), None);
257 }
258
259 #[test]
260 fn test_stamp_helper_skips_seq_info_direct() {
261 let (owner, dest) = handler_port("test_0");
262 let mut headers = Flattrs::new();
263 stamp_sender_actor_id(&mut headers, &SeqInfo::Direct, &dest, &owner);
264 assert_eq!(headers.get(SENDER_ACTOR_ID), None);
265 }
266
267 #[test]
268 fn test_stamp_fresh_helper_sets_on_seq_4() {
269 let (owner, dest) = handler_port("test_0");
270 let mut headers = Flattrs::new();
271 stamp_sender_actor_id_fresh(&mut headers, 4, &dest, &owner);
272 assert_eq!(headers.get(SENDER_ACTOR_ID), Some(owner));
273 }
274
275 #[test]
276 fn test_stamp_fresh_helper_skips_on_seq_5() {
277 let (owner, dest) = handler_port("test_0");
278 let mut headers = Flattrs::new();
279 stamp_sender_actor_id_fresh(&mut headers, 5, &dest, &owner);
280 assert_eq!(headers.get(SENDER_ACTOR_ID), None);
281 }
282}