Skip to main content

harn_vm/
runtime_context.rs

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
75/// Per-task ambient-scope swap of the runtime-context overlay stack. See
76/// `orchestration::ambient_scope`: the overlay attributes events (worker_id,
77/// run_id) to the running task, so it must follow that task across `.await`
78/// rather than leak to cooperatively-scheduled siblings.
79pub(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        // Outside a trigger dispatch, fall back to the ambient
270        // `enter_tenant` scope hosts install (today: harn-serve binds
271        // it from the authenticated principal). Keeps
272        // `runtime_context.tenant_id` consistent across trigger and
273        // non-trigger code paths.
274        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}