1use std::cell::RefCell;
2
3use crate::value::{VmError, VmValue};
4
5#[derive(Clone, Debug)]
6pub struct RuntimeContext {
7 pub task_id: String,
8 pub parent_task_id: Option<String>,
9 pub root_task_id: String,
10 pub task_name: Option<String>,
11 pub task_group_id: Option<String>,
12 pub scope_id: Option<String>,
13 pub values: crate::value::DictMap,
14}
15
16impl RuntimeContext {
17 pub fn root() -> Self {
18 Self {
19 task_id: "task_root".to_string(),
20 parent_task_id: None,
21 root_task_id: "task_root".to_string(),
22 task_name: Some("root".to_string()),
23 task_group_id: None,
24 scope_id: None,
25 values: crate::value::DictMap::new(),
26 }
27 }
28
29 pub fn child_task(
30 &self,
31 task_id: impl Into<String>,
32 task_name: impl Into<String>,
33 task_group_id: Option<String>,
34 ) -> Self {
35 Self {
36 task_id: task_id.into(),
37 parent_task_id: Some(self.task_id.clone()),
38 root_task_id: self.root_task_id.clone(),
39 task_name: Some(task_name.into()),
40 task_group_id,
41 scope_id: self.scope_id.clone(),
42 values: self.values.clone(),
43 }
44 }
45}
46
47impl Default for RuntimeContext {
48 fn default() -> Self {
49 Self::root()
50 }
51}
52
53#[derive(Clone, Debug, Default)]
54pub struct RuntimeContextOverlay {
55 pub workflow_id: Option<String>,
56 pub run_id: Option<String>,
57 pub stage_id: Option<String>,
58 pub worker_id: Option<String>,
59}
60
61thread_local! {
62 static RUNTIME_CONTEXT_OVERLAY_STACK: RefCell<Vec<RuntimeContextOverlay>> =
63 const { RefCell::new(Vec::new()) };
64}
65
66pub struct RuntimeContextOverlayGuard;
67
68pub fn install_runtime_context_overlay(
69 overlay: RuntimeContextOverlay,
70) -> RuntimeContextOverlayGuard {
71 RUNTIME_CONTEXT_OVERLAY_STACK.with(|stack| stack.borrow_mut().push(overlay));
72 RuntimeContextOverlayGuard
73}
74
75pub(crate) fn swap_runtime_context_overlay_stack(
80 next: Vec<RuntimeContextOverlay>,
81) -> Vec<RuntimeContextOverlay> {
82 RUNTIME_CONTEXT_OVERLAY_STACK.with(|stack| std::mem::replace(&mut *stack.borrow_mut(), next))
83}
84
85impl Drop for RuntimeContextOverlayGuard {
86 fn drop(&mut self) {
87 RUNTIME_CONTEXT_OVERLAY_STACK.with(|stack| {
88 stack.borrow_mut().pop();
89 });
90 }
91}
92
93fn current_overlay() -> RuntimeContextOverlay {
94 RUNTIME_CONTEXT_OVERLAY_STACK.with(|stack| {
95 let mut merged = RuntimeContextOverlay::default();
96 for overlay in stack.borrow().iter() {
97 if overlay.workflow_id.is_some() {
98 merged.workflow_id = overlay.workflow_id.clone();
99 }
100 if overlay.run_id.is_some() {
101 merged.run_id = overlay.run_id.clone();
102 }
103 if overlay.stage_id.is_some() {
104 merged.stage_id = overlay.stage_id.clone();
105 }
106 if overlay.worker_id.is_some() {
107 merged.worker_id = overlay.worker_id.clone();
108 }
109 }
110 merged
111 })
112}
113
114pub(crate) fn current_agent_run_ref() -> Option<crate::agent_events::AgentRunRef> {
116 let session_id = crate::agent_sessions::current_session_id()?;
117 let run_id = crate::llm::active_agent_run_id(&session_id)
118 .or_else(|| crate::agent_sessions::active_run_id(&session_id))?;
119 Some(crate::agent_events::AgentRunRef { session_id, run_id })
120}
121
122pub(crate) fn runtime_context_get(
123 vm: &crate::vm::Vm,
124 args: &[VmValue],
125) -> Result<VmValue, VmError> {
126 let key = require_key(args, "HarnessRuntime.context_get")?;
127 Ok(vm
128 .runtime_context
129 .values
130 .get(key.as_str())
131 .cloned()
132 .or_else(|| args.get(1).cloned())
133 .unwrap_or(VmValue::Nil))
134}
135
136pub(crate) fn runtime_context_set(
137 vm: &mut crate::vm::Vm,
138 args: &[VmValue],
139) -> Result<VmValue, VmError> {
140 let key = require_key(args, "HarnessRuntime.context_set")?;
141 let value = args.get(1).cloned().unwrap_or(VmValue::Nil);
142 Ok(vm
143 .runtime_context
144 .values
145 .insert(crate::value::intern_key(&key), value)
146 .unwrap_or(VmValue::Nil))
147}
148
149pub(crate) fn runtime_context_clear(
150 vm: &mut crate::vm::Vm,
151 args: &[VmValue],
152) -> Result<VmValue, VmError> {
153 let key = require_key(args, "HarnessRuntime.context_clear")?;
154 Ok(vm
155 .runtime_context
156 .values
157 .remove(key.as_str())
158 .unwrap_or(VmValue::Nil))
159}
160
161fn require_key(args: &[VmValue], builtin: &str) -> Result<String, VmError> {
162 match args.first() {
163 Some(VmValue::String(value)) => Ok(value.to_string()),
164 _ => Err(VmError::Runtime(format!(
165 "{builtin}: first argument must be a string key"
166 ))),
167 }
168}
169
170pub(crate) fn runtime_context_value(vm: &crate::vm::Vm) -> VmValue {
171 let overlay = current_overlay();
172 let mutation = crate::orchestration::current_mutation_session();
173 let dispatch = crate::triggers::dispatcher::current_dispatch_context();
174 let trace_context = crate::stdlib::tracing::current_trace_context();
175 let agent_session_id = crate::agent_sessions::current_session_id();
176 let agent_ancestry = agent_session_id
177 .as_deref()
178 .and_then(crate::agent_sessions::ancestry);
179 let cancelled = vm
180 .cancel_token
181 .as_ref()
182 .is_some_and(|token| token.load(std::sync::atomic::Ordering::SeqCst));
183
184 let workflow_id = overlay.workflow_id;
185 let run_id = overlay
186 .run_id
187 .or_else(|| mutation.as_ref().and_then(|session| session.run_id.clone()));
188 let stage_id = overlay.stage_id;
189 let worker_id = overlay.worker_id.or_else(|| {
190 mutation
191 .as_ref()
192 .and_then(|session| session.worker_id.clone())
193 });
194
195 let mut values = crate::value::DictMap::new();
196 insert_string(
197 &mut values,
198 "task_id",
199 Some(vm.runtime_context.task_id.clone()),
200 );
201 insert_string(
202 &mut values,
203 "parent_task_id",
204 vm.runtime_context.parent_task_id.clone(),
205 );
206 insert_string(
207 &mut values,
208 "root_task_id",
209 Some(vm.runtime_context.root_task_id.clone()),
210 );
211 insert_string(
212 &mut values,
213 "task_name",
214 vm.runtime_context.task_name.clone(),
215 );
216 insert_string(
217 &mut values,
218 "task_group_id",
219 vm.runtime_context.task_group_id.clone(),
220 );
221 insert_string(&mut values, "scope_id", vm.runtime_context.scope_id.clone());
222 insert_string(&mut values, "workflow_id", workflow_id);
223 insert_string(&mut values, "run_id", run_id);
224 insert_string(&mut values, "stage_id", stage_id);
225 insert_string(&mut values, "worker_id", worker_id);
226 insert_string(&mut values, "agent_session_id", agent_session_id);
227 insert_string(
228 &mut values,
229 "parent_agent_session_id",
230 agent_ancestry
231 .as_ref()
232 .and_then(|ancestry| ancestry.parent_id.clone()),
233 );
234 insert_string(
235 &mut values,
236 "root_agent_session_id",
237 agent_ancestry
238 .as_ref()
239 .map(|ancestry| ancestry.root_id.clone()),
240 );
241 insert_string(&mut values, "agent_name", None);
242
243 if let Some(context) = dispatch {
244 insert_string(&mut values, "trigger_id", Some(context.binding_id.clone()));
245 insert_string(
246 &mut values,
247 "trigger_event_id",
248 Some(context.trigger_event.id.0.clone()),
249 );
250 insert_string(
251 &mut values,
252 "binding_key",
253 Some(format!(
254 "{}@{}",
255 context.binding_id, context.binding_version
256 )),
257 );
258 insert_string(
259 &mut values,
260 "tenant_id",
261 context.trigger_event.tenant_id.map(|tenant| tenant.0),
262 );
263 insert_string(
264 &mut values,
265 "provider",
266 Some(context.trigger_event.provider.0),
267 );
268 insert_string(
269 &mut values,
270 "trace_id",
271 Some(context.trigger_event.trace_id.0),
272 );
273 } else {
274 insert_string(&mut values, "trigger_id", None);
275 insert_string(&mut values, "trigger_event_id", None);
276 insert_string(&mut values, "binding_key", None);
277 insert_string(
283 &mut values,
284 "tenant_id",
285 crate::harness_tenant::current_tenant_id().map(|tenant| tenant.0),
286 );
287 insert_string(&mut values, "provider", None);
288 insert_string(
289 &mut values,
290 "trace_id",
291 trace_context.as_ref().map(|context| context.0.clone()),
292 );
293 }
294
295 insert_string(
296 &mut values,
297 "span_id",
298 trace_context
299 .as_ref()
300 .map(|context| context.1.clone())
301 .or_else(|| crate::tracing::current_span_id().map(|id| id.to_string())),
302 );
303 insert_string(&mut values, "scheduler_key", None);
304 insert_string(&mut values, "runner", None);
305 insert_string(&mut values, "capacity_class", None);
306 values.insert(
307 crate::value::intern_key("context_values"),
308 VmValue::dict(vm.runtime_context.values.clone()),
309 );
310 values.insert(
311 crate::value::intern_key("cancelled"),
312 VmValue::Bool(cancelled),
313 );
314 values.insert(
315 crate::value::intern_key("debug"),
316 debug_context_value(vm, cancelled),
317 );
318 VmValue::dict(values)
319}
320
321fn debug_context_value(vm: &crate::vm::Vm, cancelled: bool) -> VmValue {
322 let mut debug = crate::value::DictMap::new();
323 debug.insert(
324 crate::value::intern_key("cancelled"),
325 VmValue::Bool(cancelled),
326 );
327 debug.insert(crate::value::intern_key("waiting_reason"), VmValue::Nil);
328 debug.insert(
329 crate::value::intern_key("active_task_ids"),
330 VmValue::List(std::sync::Arc::new(
331 vm.spawned_tasks
332 .keys()
333 .map(|id| VmValue::String(arcstr::ArcStr::from(id.as_str())))
334 .collect(),
335 )),
336 );
337 debug.insert(
338 crate::value::intern_key("held_synchronization"),
339 VmValue::List(std::sync::Arc::new(Vec::new())),
340 );
341 debug.insert(
342 crate::value::intern_key("supervisors"),
343 crate::stdlib::supervisor::supervisor_debug_values(),
344 );
345 VmValue::dict(debug)
346}
347
348fn insert_string(values: &mut crate::value::DictMap, key: &str, value: Option<String>) {
349 values.insert(
350 crate::value::intern_key(key),
351 value
352 .map(|value| VmValue::String(arcstr::ArcStr::from(value)))
353 .unwrap_or(VmValue::Nil),
354 );
355}