harn_vm/orchestration/
mod.rs1use std::path::PathBuf;
2use std::{cell::RefCell, thread_local};
3
4use serde::{Deserialize, Serialize};
5
6use crate::llm::vm_value_to_json;
7use crate::value::{VmError, VmValue};
8
9pub(crate) fn now_rfc3339() -> String {
10 use std::time::{SystemTime, UNIX_EPOCH};
11 let ts = SystemTime::now()
12 .duration_since(UNIX_EPOCH)
13 .unwrap_or_default()
14 .as_secs();
15 format!("{ts}")
16}
17
18pub(crate) fn new_id(prefix: &str) -> String {
19 format!("{prefix}_{}", uuid::Uuid::now_v7())
20}
21
22pub(crate) fn default_run_dir() -> PathBuf {
23 let base = std::env::current_dir().unwrap_or_else(|_| PathBuf::from("."));
24 crate::runtime_paths::run_root(&base)
25}
26
27mod hooks;
28pub use hooks::*;
29#[cfg(test)]
30mod tests_lazy_hooks;
31
32mod pipeline_lifecycle;
33pub use pipeline_lifecycle::*;
34
35mod settlement_agent;
36pub use settlement_agent::*;
37
38mod lifecycle_receipts;
39pub use lifecycle_receipts::*;
40
41mod command_policy;
42pub use command_policy::*;
43
44mod compaction;
45pub use compaction::*;
46
47mod repair_ledger;
48
49mod compact_lifecycle;
50pub use compact_lifecycle::*;
51
52mod compaction_policy_registry;
53pub use compaction_policy_registry::*;
54
55pub mod agent_inbox;
56
57mod artifacts;
58pub use artifacts::*;
59
60mod assemble;
61pub use assemble::*;
62
63mod handoffs;
64pub use handoffs::*;
65
66mod friction;
67pub use friction::*;
68
69mod crystallize;
70pub use crystallize::*;
71
72mod release_fixture;
73pub use release_fixture::*;
74
75mod replay_oracle;
76pub use replay_oracle::*;
77
78mod replay_bench;
79pub use replay_bench::*;
80
81mod policy;
82#[cfg(test)]
83pub(crate) use policy::swap_execution_policy_stack;
84pub use policy::*;
85
86mod ambient_scope;
87pub(crate) use ambient_scope::{scope_ambient, AmbientExecutionScope};
88pub use ambient_scope::{
89 scope_llm_runtime_overrides, scope_llm_runtime_overrides_with_provider_endpoints,
90};
91
92mod stage_options;
93pub use stage_options::*;
94
95mod workflow;
96pub use workflow::*;
97
98mod workflow_bundle;
99pub use workflow_bundle::*;
100
101mod workflow_patch;
102pub use workflow_patch::*;
103
104mod safe_function_tools;
105pub use safe_function_tools::*;
106
107mod nested_invocation;
108pub use nested_invocation::*;
109
110#[cfg(test)]
111mod workflow_test_fixtures;
112
113mod records;
114pub use records::*;
115
116mod context_eval;
117pub use context_eval::*;
118
119mod skill_gate;
120pub use skill_gate::*;
121
122mod merge_captain_audit;
123pub use merge_captain_audit::*;
124
125mod merge_captain_driver;
126pub use merge_captain_driver::*;
127
128mod merge_captain_ladder;
129pub use merge_captain_ladder::*;
130
131mod merge_captain_iteration;
132pub use merge_captain_iteration::*;
133
134pub mod playground;
135
136thread_local! {
137 static CURRENT_MUTATION_SESSION: RefCell<Option<MutationSessionRecord>> = const { RefCell::new(None) };
138 static CURRENT_WORKFLOW_SKILL_CONTEXT: RefCell<Option<WorkflowSkillContext>> = const { RefCell::new(None) };
144}
145
146#[derive(Clone, Default)]
151pub struct WorkflowSkillContext {
152 pub registry: Option<VmValue>,
153 pub match_config: Option<VmValue>,
154}
155
156pub fn install_workflow_skill_context(context: Option<WorkflowSkillContext>) {
157 CURRENT_WORKFLOW_SKILL_CONTEXT.with(|slot| {
158 *slot.borrow_mut() = context;
159 });
160}
161
162pub fn current_workflow_skill_context() -> Option<WorkflowSkillContext> {
163 CURRENT_WORKFLOW_SKILL_CONTEXT.with(|slot| slot.borrow().clone())
164}
165
166pub struct WorkflowSkillContextGuard;
170
171impl Drop for WorkflowSkillContextGuard {
172 fn drop(&mut self) {
173 install_workflow_skill_context(None);
174 }
175}
176
177#[derive(Clone, Debug, Default, Serialize, Deserialize, PartialEq, Eq)]
178#[serde(default)]
179pub struct MutationSessionRecord {
180 pub session_id: String,
181 pub parent_session_id: Option<String>,
182 pub run_id: Option<String>,
183 pub worker_id: Option<String>,
184 pub execution_kind: Option<String>,
185 pub mutation_scope: String,
186 pub approval_policy: Option<ToolApprovalPolicy>,
190}
191
192impl MutationSessionRecord {
193 pub fn normalize(mut self) -> Self {
194 if self.session_id.is_empty() {
195 self.session_id = new_id("session");
196 }
197 if self.mutation_scope.is_empty() {
198 self.mutation_scope = "read_only".to_string();
199 }
200 self
201 }
202}
203
204pub fn install_current_mutation_session(session: Option<MutationSessionRecord>) {
205 CURRENT_MUTATION_SESSION.with(|slot| {
206 *slot.borrow_mut() = session.map(MutationSessionRecord::normalize);
207 });
208}
209
210pub fn current_mutation_session() -> Option<MutationSessionRecord> {
211 CURRENT_MUTATION_SESSION.with(|slot| slot.borrow().clone())
212}
213
214pub(crate) fn swap_mutation_session(
222 next: Option<MutationSessionRecord>,
223) -> Option<MutationSessionRecord> {
224 CURRENT_MUTATION_SESSION.with(|slot| std::mem::replace(&mut *slot.borrow_mut(), next))
225}
226pub(crate) fn parse_json_payload<T: for<'de> Deserialize<'de>>(
227 json: serde_json::Value,
228 label: &str,
229) -> Result<T, VmError> {
230 let payload = json.to_string();
231 let mut deserializer = serde_json::Deserializer::from_str(&payload);
232 let mut tracker = serde_path_to_error::Track::new();
233 let path_deserializer = serde_path_to_error::Deserializer::new(&mut deserializer, &mut tracker);
234 T::deserialize(path_deserializer).map_err(|error| {
235 let snippet = if payload.len() > 600 {
236 format!("{}...", &payload[..600])
237 } else {
238 payload.clone()
239 };
240 VmError::Runtime(format!(
241 "{label} parse error at {}: {} | payload={}",
242 tracker.path(),
243 error,
244 snippet
245 ))
246 })
247}
248
249pub(crate) fn parse_json_value<T: for<'de> Deserialize<'de>>(
250 value: &VmValue,
251) -> Result<T, VmError> {
252 parse_json_payload(vm_value_to_json(value), "orchestration")
253}
254
255#[cfg(test)]
256mod tests;
257
258#[cfg(test)]
259mod policy_restriction_tests;
260
261#[cfg(test)]
262mod typed_options_parity;