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 runtime_context_get(
115 vm: &crate::vm::Vm,
116 args: &[VmValue],
117) -> Result<VmValue, VmError> {
118 let key = require_key(args, "HarnessRuntime.context_get")?;
119 Ok(vm
120 .runtime_context
121 .values
122 .get(key.as_str())
123 .cloned()
124 .or_else(|| args.get(1).cloned())
125 .unwrap_or(VmValue::Nil))
126}
127
128pub(crate) fn runtime_context_set(
129 vm: &mut crate::vm::Vm,
130 args: &[VmValue],
131) -> Result<VmValue, VmError> {
132 let key = require_key(args, "HarnessRuntime.context_set")?;
133 let value = args.get(1).cloned().unwrap_or(VmValue::Nil);
134 Ok(vm
135 .runtime_context
136 .values
137 .insert(crate::value::intern_key(&key), value)
138 .unwrap_or(VmValue::Nil))
139}
140
141pub(crate) fn runtime_context_clear(
142 vm: &mut crate::vm::Vm,
143 args: &[VmValue],
144) -> Result<VmValue, VmError> {
145 let key = require_key(args, "HarnessRuntime.context_clear")?;
146 Ok(vm
147 .runtime_context
148 .values
149 .remove(key.as_str())
150 .unwrap_or(VmValue::Nil))
151}
152
153fn require_key(args: &[VmValue], builtin: &str) -> Result<String, VmError> {
154 match args.first() {
155 Some(VmValue::String(value)) => Ok(value.to_string()),
156 _ => Err(VmError::Runtime(format!(
157 "{builtin}: first argument must be a string key"
158 ))),
159 }
160}
161
162pub(crate) fn runtime_context_value(vm: &crate::vm::Vm) -> VmValue {
163 let overlay = current_overlay();
164 let mutation = crate::orchestration::current_mutation_session();
165 let dispatch = crate::triggers::dispatcher::current_dispatch_context();
166 let trace_context = crate::stdlib::tracing::current_trace_context();
167 let agent_session_id = crate::agent_sessions::current_session_id();
168 let agent_ancestry = agent_session_id
169 .as_deref()
170 .and_then(crate::agent_sessions::ancestry);
171 let cancelled = vm
172 .cancel_token
173 .as_ref()
174 .is_some_and(|token| token.load(std::sync::atomic::Ordering::SeqCst));
175
176 let workflow_id = overlay.workflow_id;
177 let run_id = overlay
178 .run_id
179 .or_else(|| mutation.as_ref().and_then(|session| session.run_id.clone()));
180 let stage_id = overlay.stage_id;
181 let worker_id = overlay.worker_id.or_else(|| {
182 mutation
183 .as_ref()
184 .and_then(|session| session.worker_id.clone())
185 });
186
187 let mut values = crate::value::DictMap::new();
188 insert_string(
189 &mut values,
190 "task_id",
191 Some(vm.runtime_context.task_id.clone()),
192 );
193 insert_string(
194 &mut values,
195 "parent_task_id",
196 vm.runtime_context.parent_task_id.clone(),
197 );
198 insert_string(
199 &mut values,
200 "root_task_id",
201 Some(vm.runtime_context.root_task_id.clone()),
202 );
203 insert_string(
204 &mut values,
205 "task_name",
206 vm.runtime_context.task_name.clone(),
207 );
208 insert_string(
209 &mut values,
210 "task_group_id",
211 vm.runtime_context.task_group_id.clone(),
212 );
213 insert_string(&mut values, "scope_id", vm.runtime_context.scope_id.clone());
214 insert_string(&mut values, "workflow_id", workflow_id);
215 insert_string(&mut values, "run_id", run_id);
216 insert_string(&mut values, "stage_id", stage_id);
217 insert_string(&mut values, "worker_id", worker_id);
218 insert_string(&mut values, "agent_session_id", agent_session_id);
219 insert_string(
220 &mut values,
221 "parent_agent_session_id",
222 agent_ancestry
223 .as_ref()
224 .and_then(|ancestry| ancestry.parent_id.clone()),
225 );
226 insert_string(
227 &mut values,
228 "root_agent_session_id",
229 agent_ancestry
230 .as_ref()
231 .map(|ancestry| ancestry.root_id.clone()),
232 );
233 insert_string(&mut values, "agent_name", None);
234
235 if let Some(context) = dispatch {
236 insert_string(&mut values, "trigger_id", Some(context.binding_id.clone()));
237 insert_string(
238 &mut values,
239 "trigger_event_id",
240 Some(context.trigger_event.id.0.clone()),
241 );
242 insert_string(
243 &mut values,
244 "binding_key",
245 Some(format!(
246 "{}@{}",
247 context.binding_id, context.binding_version
248 )),
249 );
250 insert_string(
251 &mut values,
252 "tenant_id",
253 context.trigger_event.tenant_id.map(|tenant| tenant.0),
254 );
255 insert_string(
256 &mut values,
257 "provider",
258 Some(context.trigger_event.provider.0),
259 );
260 insert_string(
261 &mut values,
262 "trace_id",
263 Some(context.trigger_event.trace_id.0),
264 );
265 } else {
266 insert_string(&mut values, "trigger_id", None);
267 insert_string(&mut values, "trigger_event_id", None);
268 insert_string(&mut values, "binding_key", None);
269 insert_string(
275 &mut values,
276 "tenant_id",
277 crate::harness_tenant::current_tenant_id().map(|tenant| tenant.0),
278 );
279 insert_string(&mut values, "provider", None);
280 insert_string(
281 &mut values,
282 "trace_id",
283 trace_context.as_ref().map(|context| context.0.clone()),
284 );
285 }
286
287 insert_string(
288 &mut values,
289 "span_id",
290 trace_context
291 .as_ref()
292 .map(|context| context.1.clone())
293 .or_else(|| crate::tracing::current_span_id().map(|id| id.to_string())),
294 );
295 insert_string(&mut values, "scheduler_key", None);
296 insert_string(&mut values, "runner", None);
297 insert_string(&mut values, "capacity_class", None);
298 values.insert(
299 crate::value::intern_key("context_values"),
300 VmValue::dict(vm.runtime_context.values.clone()),
301 );
302 values.insert(
303 crate::value::intern_key("cancelled"),
304 VmValue::Bool(cancelled),
305 );
306 values.insert(
307 crate::value::intern_key("debug"),
308 debug_context_value(vm, cancelled),
309 );
310 VmValue::dict(values)
311}
312
313fn debug_context_value(vm: &crate::vm::Vm, cancelled: bool) -> VmValue {
314 let mut debug = crate::value::DictMap::new();
315 debug.insert(
316 crate::value::intern_key("cancelled"),
317 VmValue::Bool(cancelled),
318 );
319 debug.insert(crate::value::intern_key("waiting_reason"), VmValue::Nil);
320 debug.insert(
321 crate::value::intern_key("active_task_ids"),
322 VmValue::List(std::sync::Arc::new(
323 vm.spawned_tasks
324 .keys()
325 .map(|id| VmValue::String(arcstr::ArcStr::from(id.as_str())))
326 .collect(),
327 )),
328 );
329 debug.insert(
330 crate::value::intern_key("held_synchronization"),
331 VmValue::List(std::sync::Arc::new(Vec::new())),
332 );
333 debug.insert(
334 crate::value::intern_key("supervisors"),
335 crate::stdlib::supervisor::supervisor_debug_values(),
336 );
337 VmValue::dict(debug)
338}
339
340fn insert_string(values: &mut crate::value::DictMap, key: &str, value: Option<String>) {
341 values.insert(
342 crate::value::intern_key(key),
343 value
344 .map(|value| VmValue::String(arcstr::ArcStr::from(value)))
345 .unwrap_or(VmValue::Nil),
346 );
347}