use std::any::Any;
use std::cell::{Cell, RefCell};
use std::collections::{BTreeMap, HashMap};
use std::fmt::Write as _;
use std::io;
use std::path::Path;
use std::rc::{Rc, Weak};
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::mpsc::Receiver;
use std::time::{Instant, SystemTime, UNIX_EPOCH};
use sema_core::cycle::GcEdge;
use sema_core::runtime::{IdCounter, ScopeId, TaskContextHandle, TaskLocalValue, Trace};
use sema_core::Value;
use crate::event::WorkflowEvent;
use crate::journal::Journal;
use crate::RUNS_ROOT;
const FIXED_TS_ENV: &str = "SEMA_WORKFLOW_FIXED_TS";
const RUN_ID_ENV: &str = "SEMA_WORKFLOW_RUN_ID";
const RUN_DIR_ENV: &str = "SEMA_WORKFLOW_RUN_DIR";
pub const MEMO_MAX_COUNT: u64 = 4096;
pub const MEMO_FILE_MAX_BYTES: usize = 1 << 20; const DIGEST_MAX_BYTES: usize = 1 << 20;
struct CappedWriter {
buf: String,
cap: usize,
truncated: bool,
}
impl std::fmt::Write for CappedWriter {
fn write_str(&mut self, s: &str) -> std::fmt::Result {
if self.truncated {
return Err(std::fmt::Error);
}
let remaining = self.cap.saturating_sub(self.buf.len());
if s.len() <= remaining {
self.buf.push_str(s);
Ok(())
} else {
let mut end = remaining;
while end > 0 && !s.is_char_boundary(end) {
end -= 1;
}
self.buf.push_str(&s[..end]);
self.truncated = true;
Err(std::fmt::Error)
}
}
}
pub fn compact_capped(v: &Value, cap: usize) -> (String, bool) {
let mut w = CappedWriter {
buf: String::new(),
cap,
truncated: false,
};
let _ = write!(w, "{v}");
(w.buf, w.truncated)
}
thread_local! {
static WORKFLOW: Rc<WorkflowTaskState> = Rc::new(WorkflowTaskState::default());
static HOST_CONFIG: RefCell<Option<WorkflowHostConfig>> = const { RefCell::new(None) };
}
#[derive(Debug, Clone)]
pub struct WorkflowHostConfig {
pub runs_root: String,
pub explicit_run_id: Option<String>,
pub resuming: bool,
pub code_version: String,
pub approval_code_version: String,
pub args_json: String,
pub approval_public_key: String,
pub entry_file: String,
pub workspace_root: String,
}
pub struct WorkflowHostConfigGuard {
previous: Option<WorkflowHostConfig>,
}
impl Drop for WorkflowHostConfigGuard {
fn drop(&mut self) {
HOST_CONFIG.with(|slot| {
*slot.borrow_mut() = self.previous.take();
});
}
}
pub fn install_host_config(config: WorkflowHostConfig) -> WorkflowHostConfigGuard {
let previous = HOST_CONFIG.with(|slot| slot.borrow_mut().replace(config));
WorkflowHostConfigGuard { previous }
}
fn host_config() -> Option<WorkflowHostConfig> {
HOST_CONFIG.with(|slot| slot.borrow().clone())
}
pub fn host_workspace_root() -> Option<std::path::PathBuf> {
host_config().map(|config| config.workspace_root.into())
}
pub struct WorkflowCtx {
pub run_id: String,
workflow_name: RefCell<String>,
journal: Rc<RefCell<Journal>>,
state: Rc<RefCell<BTreeMap<String, Value>>>,
seq: Cell<u64>,
event_counts: RefCell<BTreeMap<&'static str, u64>>,
start: Instant,
cost_limit: Option<f64>,
token_limit: Option<u64>,
cost_spent: Cell<f64>,
tokens_spent: Cell<u64>,
over_budget: Cell<bool>,
approval_failure: RefCell<Option<String>>,
cur_phase: RefCell<Option<(u64, String)>>,
agent_n: RefCell<BTreeMap<String, u64>>,
resuming: Cell<bool>,
code_version: RefCell<String>,
approval_code_version: RefCell<String>,
approval_public_key: RefCell<String>,
resume_memos: RefCell<HashMap<String, Value>>,
key_seen: RefCell<HashMap<String, u32>>,
memo_count: Cell<u64>,
args_json: String,
args_fingerprint: String,
fixed_ts: Option<String>,
mcp_declared: RefCell<Vec<String>>,
mcp_handles: RefCell<BTreeMap<String, Value>>,
}
impl WorkflowCtx {
pub fn new(
run_id: String,
journal: Journal,
budget: BTreeMap<String, Value>,
) -> Rc<WorkflowCtx> {
Self::new_with_args(run_id, journal, budget, String::new())
}
pub fn new_with_args(
run_id: String,
journal: Journal,
budget: BTreeMap<String, Value>,
args_json: String,
) -> Rc<WorkflowCtx> {
let fixed_ts = std::env::var(FIXED_TS_ENV).ok();
let args_fingerprint = canonical_args_fingerprint(&args_json);
let cost_limit = budget
.get("usd")
.and_then(|v| v.as_float().or_else(|| v.as_int().map(|i| i as f64)));
let token_limit = budget
.get("tokens")
.and_then(|v| v.as_int().or_else(|| v.as_float().map(|f| f as i64)))
.map(|i| i as u64);
Rc::new(WorkflowCtx {
run_id,
workflow_name: RefCell::new(String::new()),
journal: Rc::new(RefCell::new(journal)),
state: Rc::new(RefCell::new(BTreeMap::new())),
seq: Cell::new(0),
event_counts: RefCell::new(BTreeMap::new()),
start: Instant::now(),
cost_limit,
token_limit,
cost_spent: Cell::new(0.0),
tokens_spent: Cell::new(0),
over_budget: Cell::new(false),
approval_failure: RefCell::new(None),
cur_phase: RefCell::new(None),
agent_n: RefCell::new(BTreeMap::new()),
resuming: Cell::new(false),
code_version: RefCell::new(String::new()),
approval_code_version: RefCell::new(String::new()),
approval_public_key: RefCell::new(String::new()),
resume_memos: RefCell::new(HashMap::new()),
key_seen: RefCell::new(HashMap::new()),
memo_count: Cell::new(0),
args_json,
args_fingerprint,
fixed_ts,
mcp_declared: RefCell::new(Vec::new()),
mcp_handles: RefCell::new(BTreeMap::new()),
})
}
pub fn args_json(&self) -> &str {
&self.args_json
}
pub fn set_workflow_name(&self, name: impl Into<String>) {
*self.workflow_name.borrow_mut() = name.into();
}
pub fn workflow_name(&self) -> String {
self.workflow_name.borrow().clone()
}
pub fn approval_code_version(&self) -> String {
self.approval_code_version.borrow().clone()
}
pub fn approval_public_key(&self) -> String {
self.approval_public_key.borrow().clone()
}
pub fn approval_args_digest(&self) -> String {
let normalized = if self.args_json.trim().is_empty() {
String::new()
} else {
serde_json::from_str::<serde_json::Value>(&self.args_json)
.ok()
.and_then(|json| serde_json::to_string(&json).ok())
.unwrap_or_else(|| self.args_json.clone())
};
crate::approval::sha256_bytes(normalized.as_bytes())
}
pub fn run_dir(&self) -> std::path::PathBuf {
self.journal.borrow().dir().to_path_buf()
}
pub fn open_phase(&self, start_seq: u64, label: String) {
*self.cur_phase.borrow_mut() = Some((start_seq, label));
}
pub fn take_open_phase(&self) -> Option<(u64, String)> {
self.cur_phase.borrow_mut().take()
}
pub fn phase_seq(&self) -> Option<u64> {
self.cur_phase.borrow().as_ref().map(|(seq, _)| *seq)
}
pub fn next_agent_id(&self, name: &str) -> String {
let mut m = self.agent_n.borrow_mut();
let n = m.entry(name.to_string()).or_insert(0);
*n += 1;
format!("{name}_{n}")
}
pub fn content_key(&self, key: &str, value_digest: &str) -> String {
let h = format!(
"{:x}",
md5::compute(format!("{key}:{value_digest}").as_bytes())
);
format!("ck_{}", &h[..8])
}
pub fn next_seq(&self) -> u64 {
let n = self.seq.get();
self.seq.set(n + 1);
n
}
pub fn ts(&self) -> String {
if let Some(ref fixed) = self.fixed_ts {
return fixed.clone();
}
rfc3339_now()
}
pub fn dur_ms(&self) -> u64 {
if self.fixed_ts.is_some() {
return 0;
}
self.start.elapsed().as_millis() as u64
}
pub fn emit(&self, event: WorkflowEvent) {
let kind = event.kind();
let mut counts = self.event_counts.borrow_mut();
*counts.entry(kind).or_insert(0) += 1;
drop(counts);
self.journal.borrow().write(&event);
}
pub fn has_event(&self, kind: &str) -> bool {
self.event_counts
.borrow()
.get(kind)
.is_some_and(|count| *count > 0)
}
pub fn deterministic(&self) -> bool {
self.fixed_ts.is_some()
}
pub fn run_id(&self) -> String {
self.run_id.clone()
}
pub fn store_checkpoint(&self, key: &str, val: Value) {
self.state.borrow_mut().insert(key.to_string(), val);
}
pub fn read_checkpoint(&self, key: &str) -> Option<Value> {
self.state.borrow().get(key).cloned()
}
pub fn value_digest(&self, v: &Value) -> String {
let (compact, truncated) = compact_capped(v, DIGEST_MAX_BYTES);
if truncated {
return format!("oversized_{:x}", md5::compute(compact.as_bytes()));
}
let json = sema_core::json::value_to_json_lossy(v);
let bytes = serde_json::to_vec(&json).unwrap_or_default();
format!("{:x}", md5::compute(bytes))
}
pub fn write_result(&self, envelope: &Value) {
let json = sema_core::json::value_to_json_lossy(envelope);
self.journal.borrow().write_result(&json);
}
pub fn has_budget(&self) -> bool {
self.cost_limit.is_some() || self.token_limit.is_some()
}
pub fn budget_limit_for_event(&self) -> Option<u64> {
self.token_limit
}
pub fn charge(&self, cost: Option<f64>, tokens: u64) -> bool {
if let Some(c) = cost {
self.cost_spent.set(self.cost_spent.get() + c);
}
self.tokens_spent.set(self.tokens_spent.get() + tokens);
let over = self
.cost_limit
.is_some_and(|lim| self.cost_spent.get() > lim)
|| self
.token_limit
.is_some_and(|lim| self.tokens_spent.get() > lim);
if over {
self.over_budget.set(true);
}
over
}
pub fn over_budget(&self) -> bool {
self.over_budget.get()
}
pub fn fail_approval(&self, message: impl Into<String>) {
let mut failure = self.approval_failure.borrow_mut();
if failure.is_none() {
*failure = Some(message.into());
}
}
pub fn approval_failure(&self) -> Option<String> {
self.approval_failure.borrow().clone()
}
pub fn set_code_version(&self, v: String) {
*self.code_version.borrow_mut() = v;
}
pub fn set_approval_code_version(&self, v: String) {
*self.approval_code_version.borrow_mut() = v;
}
pub fn set_approval_public_key(&self, v: String) {
*self.approval_public_key.borrow_mut() = v;
}
pub fn enter_resume(&self, memos: HashMap<String, Value>) {
self.resuming.set(true);
*self.resume_memos.borrow_mut() = memos;
}
pub fn resuming(&self) -> bool {
self.resuming.get()
}
pub fn cur_phase_label(&self) -> String {
self.cur_phase
.borrow()
.as_ref()
.map(|(_, label)| label.clone())
.unwrap_or_default()
}
fn next_occurrence(&self, base: &str) -> u32 {
let mut m = self.key_seen.borrow_mut();
let n = m.entry(base.to_string()).or_insert(0);
let cur = *n;
*n += 1;
cur
}
pub fn agent_content_key(
&self,
prompt: &str,
schema_repr: &str,
name: &str,
phase: &str,
policy_fingerprint: &str,
) -> String {
let cv = self.code_version.borrow().clone();
let base = hash_fields(&[
"agent",
&cv,
&self.args_fingerprint,
phase,
name,
prompt,
schema_repr,
policy_fingerprint,
]);
format!("{base}_{}", self.next_occurrence(&base))
}
pub fn checkpoint_content_key(&self, key: &str, phase: &str) -> String {
let cv = self.code_version.borrow().clone();
let base = hash_fields(&["checkpoint", &cv, &self.args_fingerprint, phase, key]);
format!("{base}_{}", self.next_occurrence(&base))
}
pub fn approval_occurrence(&self, key: &str, subject_digest: &str, phase: &str) -> u32 {
let cv = self.approval_code_version.borrow().clone();
let base = crate::approval::sha256_fields(&[
"approval",
&cv,
&self.args_fingerprint,
phase,
key,
subject_digest,
]);
self.next_occurrence(&base)
}
pub fn memo_lookup(&self, content_key: &str) -> Option<Value> {
self.resume_memos.borrow().get(content_key).cloned()
}
pub fn memo_store(&self, content_key: &str, v: &Value) {
if self.memo_count.get() >= MEMO_MAX_COUNT {
return;
}
let (_, truncated) = compact_capped(v, MEMO_FILE_MAX_BYTES);
if truncated {
return;
}
let json = sema_core::json::value_to_json_lossy(v);
if sema_core::json::json_to_value(&json) != *v {
return;
}
let serialized = serde_json::to_vec(&json).unwrap_or_default();
if serialized.len() > MEMO_FILE_MAX_BYTES {
return;
}
self.memo_count.set(self.memo_count.get() + 1);
self.journal.borrow().write_memo(content_key, &json);
self.resume_memos
.borrow_mut()
.insert(content_key.to_string(), v.clone());
}
pub fn request_flush(&self) -> Receiver<()> {
self.journal.borrow().request_flush()
}
pub fn flush(&self) {
self.journal.borrow().flush_blocking();
}
pub fn set_mcp_declared(&self, aliases: Vec<String>) {
*self.mcp_declared.borrow_mut() = aliases;
}
pub fn is_mcp_declared(&self, alias: &str) -> bool {
self.mcp_declared.borrow().iter().any(|a| a == alias)
}
pub fn set_mcp_handles(&self, handles: BTreeMap<String, Value>) {
*self.mcp_handles.borrow_mut() = handles;
}
pub fn mcp_handle(&self, alias: &str) -> Option<Value> {
self.mcp_handles.borrow().get(alias).cloned()
}
}
impl Trace for WorkflowCtx {
fn trace(&self, sink: &mut dyn FnMut(GcEdge<'_>)) -> bool {
let (Ok(state), Ok(memos), Ok(handles)) = (
self.state.try_borrow(),
self.resume_memos.try_borrow(),
self.mcp_handles.try_borrow(),
) else {
return false;
};
for value in state.values() {
sink(GcEdge::Value(value));
}
for value in memos.values() {
sink(GcEdge::Value(value));
}
for value in handles.values() {
sink(GcEdge::Value(value));
}
true
}
}
struct WorkflowScope {
token: Option<ScopeId>,
ctx: Rc<WorkflowCtx>,
}
struct WorkflowTaskInner {
tokens: IdCounter<ScopeId>,
scopes: Vec<WorkflowScope>,
cur_agent: Option<String>,
}
pub struct WorkflowTaskState {
inner: RefCell<WorkflowTaskInner>,
}
impl Default for WorkflowTaskState {
fn default() -> Self {
Self {
inner: RefCell::new(WorkflowTaskInner {
tokens: IdCounter::new(),
scopes: Vec::new(),
cur_agent: None,
}),
}
}
}
impl WorkflowTaskState {
fn install(&self, ctx: Rc<WorkflowCtx>) -> ScopeId {
let mut inner = self.inner.borrow_mut();
let token = inner
.tokens
.allocate()
.expect("workflow scope identity space exhausted");
inner.scopes.push(WorkflowScope {
token: Some(token),
ctx,
});
token
}
fn remove(&self, token: ScopeId) -> bool {
let mut inner = self.inner.borrow_mut();
match inner.scopes.iter().position(|s| s.token == Some(token)) {
Some(pos) => {
inner.scopes.remove(pos);
true
}
None => false,
}
}
fn current_ctx(&self) -> Option<Rc<WorkflowCtx>> {
self.inner.borrow().scopes.last().map(|s| Rc::clone(&s.ctx))
}
fn scope_depth(&self) -> usize {
self.inner.borrow().scopes.len()
}
fn current_scope_is_owned(&self) -> bool {
self.inner
.borrow()
.scopes
.last()
.is_some_and(|scope| scope.token.is_some())
}
fn cur_agent(&self) -> Option<String> {
self.inner.borrow().cur_agent.clone()
}
fn set_cur_agent(&self, agent_id: Option<String>) {
self.inner.borrow_mut().cur_agent = agent_id;
}
}
impl Trace for WorkflowTaskState {
fn trace(&self, sink: &mut dyn FnMut(GcEdge<'_>)) -> bool {
let Ok(inner) = self.inner.try_borrow() else {
return false;
};
for scope in &inner.scopes {
if !scope.ctx.trace(sink) {
return false;
}
}
true
}
}
impl TaskLocalValue for WorkflowTaskState {
fn inherit(&self) -> Rc<dyn TaskLocalValue> {
let inner = self.inner.borrow();
let scopes = inner
.scopes
.iter()
.map(|s| WorkflowScope {
token: None,
ctx: Rc::clone(&s.ctx),
})
.collect();
Rc::new(Self {
inner: RefCell::new(WorkflowTaskInner {
tokens: IdCounter::new(),
scopes,
cur_agent: inner.cur_agent.clone(),
}),
})
}
fn as_any(&self) -> &dyn Any {
self
}
fn preflight_error(&self) -> Option<sema_core::SemaError> {
self.current_ctx()
.and_then(|ctx| ctx.approval_failure())
.map(|message| sema_core::SemaError::WorkflowApprovalFailed { message })
}
}
pub struct WorkflowGuard {
state: Weak<WorkflowTaskState>,
token: ScopeId,
}
impl Drop for WorkflowGuard {
fn drop(&mut self) {
if let Some(state) = self.state.upgrade() {
state.remove(self.token);
}
}
}
fn host_state() -> Rc<WorkflowTaskState> {
WORKFLOW.with(Rc::clone)
}
fn resolve_state(task_context: Option<&TaskContextHandle>) -> Rc<WorkflowTaskState> {
if let Some(handle) = task_context {
if let Some(state) = handle.get_rc::<WorkflowTaskState>() {
return state;
}
let state = Rc::new(WorkflowTaskState::default());
handle.borrow_mut().insert(Rc::clone(&state));
return state;
}
host_state()
}
pub fn install_scope(
task_context: Option<&TaskContextHandle>,
ctx: Rc<WorkflowCtx>,
) -> WorkflowGuard {
let state = resolve_state(task_context);
let token = state.install(ctx);
WorkflowGuard {
state: Rc::downgrade(&state),
token,
}
}
pub fn current_for(task_context: Option<&TaskContextHandle>) -> Option<Rc<WorkflowCtx>> {
if let Some(handle) = task_context {
if let Some(state) = handle.get_rc::<WorkflowTaskState>() {
if let Some(ctx) = state.current_ctx() {
return Some(ctx);
}
}
}
if !sema_core::in_runtime_quantum() {
return host_state().current_ctx();
}
None
}
pub fn approval_scope_is_root_owner(task_context: Option<&TaskContextHandle>) -> bool {
let state = if let Some(handle) = task_context {
handle.get_rc::<WorkflowTaskState>()
} else if !sema_core::in_runtime_quantum() {
Some(host_state())
} else {
None
};
state.is_some_and(|state| state.scope_depth() == 1 && state.current_scope_is_owned())
}
pub fn scope_depth_for(task_context: Option<&TaskContextHandle>) -> usize {
if let Some(handle) = task_context {
return handle
.get_rc::<WorkflowTaskState>()
.map_or(0, |state| state.scope_depth());
}
if !sema_core::in_runtime_quantum() {
return host_state().scope_depth();
}
0
}
pub fn cur_agent_for(task_context: Option<&TaskContextHandle>) -> Option<String> {
if let Some(handle) = task_context {
if let Some(state) = handle.get_rc::<WorkflowTaskState>() {
return state.cur_agent();
}
}
if !sema_core::in_runtime_quantum() {
return host_state().cur_agent();
}
None
}
pub fn set_cur_agent_for(task_context: Option<&TaskContextHandle>, agent_id: Option<String>) {
resolve_state(task_context).set_cur_agent(agent_id);
}
fn redact_meta_secrets(mut meta_json: serde_json::Value) -> serde_json::Value {
let Some(mcp) = meta_json.get_mut("mcp").and_then(|v| v.as_object_mut()) else {
return meta_json;
};
for spec in mcp.values_mut() {
let Some(spec_obj) = spec.as_object_mut() else {
continue;
};
for field in ["headers", "env"] {
let Some(values) = spec_obj.get_mut(field).and_then(|v| v.as_object_mut()) else {
continue;
};
for value in values.values_mut() {
*value = serde_json::Value::String("<redacted>".to_string());
}
}
}
meta_json
}
pub fn set_workflow_scope(
name: &str,
doc: &str,
meta: &Value,
task_context: Option<&TaskContextHandle>,
) -> io::Result<WorkflowGuard> {
let host = host_config();
let outermost = scope_depth_for(task_context) == 0;
let runs_root = host
.as_ref()
.map(|config| config.runs_root.clone())
.unwrap_or_else(resolve_runs_root_from_env);
let code_version = host
.as_ref()
.map(|config| config.code_version.clone())
.unwrap_or_else(|| std::env::var(CODE_VERSION_ENV).unwrap_or_default());
let approval_code_version = host
.as_ref()
.map(|config| config.approval_code_version.clone())
.unwrap_or_else(|| {
std::env::var(APPROVAL_CODE_VERSION_ENV).unwrap_or_else(|_| code_version.clone())
});
let approval_public_key = host
.as_ref()
.map(|config| config.approval_public_key.clone())
.unwrap_or_default();
let resuming = outermost
&& host.as_ref().map_or_else(
|| std::env::var(RESUME_ENV).map(|v| v == "1").unwrap_or(false),
|config| config.resuming,
);
let configured_id = if outermost {
host.as_ref().map_or_else(
|| std::env::var(RUN_ID_ENV).ok(),
|config| config.explicit_run_id.clone(),
)
} else {
None
};
let explicit_id = match configured_id {
Some(id) if !id.is_empty() => {
validate_explicit_run_id(&id)?;
Some(id)
}
_ => None,
};
let (run_id, journal) = if resuming {
let id = explicit_id.ok_or_else(|| {
io::Error::new(
io::ErrorKind::InvalidInput,
"workflow resume requires an explicit run id (set SEMA_WORKFLOW_RUN_ID)",
)
})?;
let events = Path::new(&runs_root).join(&id).join("events.jsonl");
if !events.exists() {
return Err(io::Error::new(
io::ErrorKind::NotFound,
format!(
"cannot resume: no prior run journal at {}",
events.display()
),
));
}
let journal = crate::journal::next_resume_segment(&runs_root, &id)?;
(id, journal)
} else if let Some(id) = explicit_id {
let journal = Journal::open(&runs_root, &id).map_err(|e| annotate_fresh_open(e, &id))?;
(id, journal)
} else {
open_fresh_generated(&runs_root)?
};
let metadata = serde_json::json!({
"workflow": name,
"doc": doc,
"run_id": run_id,
"code_version": code_version,
"approval_code_version": approval_code_version,
"approval_authority_public_key": approval_public_key,
"entry_file": host.as_ref().map(|config| config.entry_file.as_str()).unwrap_or(""),
"meta": redact_meta_secrets(sema_core::json::value_to_json_lossy(meta)),
});
journal.write_metadata(&metadata);
let args_json = host
.as_ref()
.map(|config| config.args_json.clone())
.unwrap_or_else(|| std::env::var("SEMA_WORKFLOW_ARGS_JSON").unwrap_or_default());
let ctx = WorkflowCtx::new_with_args(run_id.clone(), journal, parse_budget(meta), args_json);
ctx.set_workflow_name(name);
ctx.set_code_version(code_version);
ctx.set_approval_code_version(approval_code_version);
ctx.set_approval_public_key(approval_public_key);
if resuming {
let memos: HashMap<String, Value> = crate::journal::load_memos(&runs_root, &run_id)
.into_iter()
.map(|(ck, json)| (ck, sema_core::json::json_to_value(&json)))
.collect();
ctx.enter_resume(memos);
}
Ok(install_scope(task_context, ctx))
}
pub fn parse_budget(meta: &Value) -> BTreeMap<String, Value> {
let mut out = BTreeMap::new();
if let Some(m) = meta.as_map_rc() {
if let Some(b) = m.get(&Value::keyword("budget")).and_then(|v| v.as_map_rc()) {
for (k, v) in b.iter() {
if let Some(name) = k.as_keyword() {
out.insert(name, v.clone());
}
}
}
}
out
}
pub fn resolve_runs_root() -> String {
host_config()
.map(|config| config.runs_root)
.unwrap_or_else(resolve_runs_root_from_env)
}
fn resolve_runs_root_from_env() -> String {
std::env::var(RUN_DIR_ENV).unwrap_or_else(|_| RUNS_ROOT.to_string())
}
fn hash_fields(fields: &[&str]) -> String {
let mut buf = Vec::new();
for f in fields {
buf.extend_from_slice(&(f.len() as u64).to_le_bytes());
buf.extend_from_slice(f.as_bytes());
}
let h = format!("{:x}", md5::compute(&buf));
h[..16].to_string()
}
fn canonical_args_fingerprint(args_json: &str) -> String {
let normalized = if args_json.trim().is_empty() {
String::new()
} else {
serde_json::from_str::<serde_json::Value>(args_json)
.ok()
.and_then(|json| serde_json::to_string(&json).ok())
.unwrap_or_else(|| args_json.to_string())
};
hash_fields(&["args", &normalized])
}
const RESUME_ENV: &str = "SEMA_WORKFLOW_RESUME";
const CODE_VERSION_ENV: &str = "SEMA_WORKFLOW_CODE_VERSION";
const APPROVAL_CODE_VERSION_ENV: &str = "SEMA_WORKFLOW_APPROVAL_CODE_VERSION";
static RUN_ID_NONCE: AtomicU64 = AtomicU64::new(0);
const MAX_FRESH_ATTEMPTS: u32 = 8;
fn generate_run_id() -> String {
let now = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default();
let nonce = RUN_ID_NONCE.fetch_add(1, Ordering::Relaxed);
format!(
"wf_{}_{}_{}_{}",
now.as_secs(),
now.subsec_nanos(),
std::process::id(),
nonce
)
}
pub fn validate_explicit_run_id(id: &str) -> io::Result<()> {
let reject = |why: &str| {
io::Error::new(
io::ErrorKind::InvalidInput,
format!("workflow run id {id:?} is not a safe directory name: {why}"),
)
};
if id.is_empty() {
return Err(reject("must not be empty"));
}
if id.contains('/') || id.contains('\\') {
return Err(reject("must not contain a path separator"));
}
if id.contains("..") {
return Err(reject("must not contain '..'"));
}
if id.bytes().all(|b| b == b'.') {
return Err(reject("must not be only '.' characters"));
}
if id.chars().any(|c| c == '\0' || c.is_control()) {
return Err(reject("must not contain NUL or control characters"));
}
Ok(())
}
pub fn resolve_run_id() -> io::Result<String> {
match std::env::var(RUN_ID_ENV) {
Ok(id) if !id.is_empty() => {
validate_explicit_run_id(&id)?;
Ok(id)
}
_ => Ok(generate_run_id()),
}
}
fn open_fresh_generated(runs_root: &str) -> io::Result<(String, Journal)> {
open_fresh_with(runs_root, generate_run_id)
}
fn open_fresh_with(
runs_root: &str,
mut next_id: impl FnMut() -> String,
) -> io::Result<(String, Journal)> {
let mut last_err = None;
for _ in 0..MAX_FRESH_ATTEMPTS {
let id = next_id();
match Journal::open(runs_root, &id) {
Ok(journal) => return Ok((id, journal)),
Err(e) if e.kind() == io::ErrorKind::AlreadyExists => last_err = Some(e),
Err(e) => return Err(e),
}
}
Err(last_err.unwrap_or_else(|| {
io::Error::new(
io::ErrorKind::AlreadyExists,
"could not allocate a unique workflow run directory",
)
}))
}
fn annotate_fresh_open(err: io::Error, run_id: &str) -> io::Error {
if err.kind() == io::ErrorKind::AlreadyExists {
io::Error::new(
io::ErrorKind::AlreadyExists,
format!(
"a workflow run journal for {run_id:?} already exists; \
choose a fresh run id or resume it with --resume"
),
)
} else {
err
}
}
pub(crate) fn rfc3339_now() -> String {
let dur = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default();
let secs = dur.as_secs();
let days = (secs / 86_400) as i64;
let rem = secs % 86_400;
let (hour, min, sec) = (rem / 3600, (rem % 3600) / 60, rem % 60);
let (y, m, d) = civil_from_days(days);
format!("{y:04}-{m:02}-{d:02}T{hour:02}:{min:02}:{sec:02}Z")
}
fn civil_from_days(z: i64) -> (i64, u32, u32) {
let z = z + 719_468;
let era = if z >= 0 { z } else { z - 146_096 } / 146_097;
let doe = (z - era * 146_097) as u64; let yoe = (doe - doe / 1460 + doe / 36524 - doe / 146_096) / 365; let y = yoe as i64 + era * 400;
let doy = doe - (365 * yoe + yoe / 4 - yoe / 100); let mp = (5 * doy + 2) / 153; let d = (doy - (153 * mp + 2) / 5 + 1) as u32; let m = if mp < 10 { mp + 3 } else { mp - 9 } as u32; (if m <= 2 { y + 1 } else { y }, m, d)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn seq_is_monotonic_from_zero() {
let ctx = WorkflowCtx::new("wf_t".into(), Journal::null(), BTreeMap::new());
assert_eq!(ctx.next_seq(), 0);
assert_eq!(ctx.next_seq(), 1);
assert_eq!(ctx.next_seq(), 2);
}
#[test]
fn fixed_ts_freezes_ts_and_dur() {
std::env::set_var(FIXED_TS_ENV, "1970-01-01T00:00:00Z");
let ctx = WorkflowCtx::new("wf_t".into(), Journal::null(), BTreeMap::new());
assert_eq!(ctx.ts(), "1970-01-01T00:00:00Z");
assert_eq!(ctx.dur_ms(), 0);
std::env::remove_var(FIXED_TS_ENV);
}
fn ctx_with_budget(pairs: &[(&str, Value)]) -> Rc<WorkflowCtx> {
let mut b = BTreeMap::new();
for (k, v) in pairs {
b.insert(k.to_string(), v.clone());
}
WorkflowCtx::new_with_args("wf_t".into(), Journal::null(), b, String::new())
}
#[test]
fn charge_trips_usd_cap_and_latches() {
let ctx = ctx_with_budget(&[("usd", Value::float(0.01))]);
assert!(!ctx.charge(Some(0.005), 10), "under cap must not trip");
assert!(!ctx.over_budget());
assert!(ctx.charge(Some(0.02), 100), "crossing cap trips");
assert!(ctx.over_budget(), "latch is sticky");
let _ = ctx.charge(Some(0.0), 0);
assert!(ctx.over_budget());
}
#[test]
fn charge_enforces_tokens_when_cost_unknown() {
let ctx = ctx_with_budget(&[("tokens", Value::int(50))]);
assert!(!ctx.charge(None, 40), "cost None still counts tokens");
assert!(!ctx.over_budget());
assert!(ctx.charge(None, 20), "60 > 50 trips on tokens alone");
assert!(ctx.over_budget());
assert_eq!(ctx.budget_limit_for_event(), Some(50));
}
#[test]
fn no_budget_never_trips() {
let ctx = WorkflowCtx::new("wf_t".into(), Journal::null(), BTreeMap::new());
assert!(!ctx.has_budget());
assert!(!ctx.charge(Some(9999.0), 9_999_999));
assert!(!ctx.over_budget());
assert_eq!(ctx.budget_limit_for_event(), None);
}
#[test]
fn parse_budget_extracts_caps_and_tolerates_absence() {
let mut bm = BTreeMap::new();
bm.insert(Value::keyword("usd"), Value::float(2.5));
bm.insert(Value::keyword("tokens"), Value::int(1000));
let mut meta = BTreeMap::new();
meta.insert(Value::keyword("budget"), Value::map(bm));
let parsed = parse_budget(&Value::map(meta));
assert_eq!(parsed.get("usd").and_then(|v| v.as_float()), Some(2.5));
assert_eq!(parsed.get("tokens").and_then(|v| v.as_int()), Some(1000));
assert!(parse_budget(&Value::map(BTreeMap::new())).is_empty());
}
#[test]
fn content_keys_are_stable_distinct_and_length_prefixed() {
let ctx = WorkflowCtx::new("wf_t".into(), Journal::null(), BTreeMap::new());
ctx.set_code_version("v1".into());
let k_a = ctx.agent_content_key("audit a.php", "[:list :string]", "auditor", "Audit", "");
let k_b = ctx.agent_content_key("audit b.php", "[:list :string]", "auditor", "Audit", "");
assert_ne!(k_a, k_b, "different prompts ⇒ different keys");
let k1 = ctx.agent_content_key("a", "bc", "n", "p", "");
let k2 = ctx.agent_content_key("ab", "c", "n", "p", "");
assert_ne!(
k1, k2,
"length-prefixed fields can't collide via concatenation"
);
let r1 = ctx.checkpoint_content_key("files", "Inventory");
let r2 = ctx.checkpoint_content_key("files", "Inventory");
assert_ne!(
r1, r2,
"repeated identical checkpoint ⇒ distinct occurrence key"
);
}
#[test]
fn code_version_changes_invalidate_keys() {
let ctx1 = WorkflowCtx::new("a".into(), Journal::null(), BTreeMap::new());
ctx1.set_code_version("v1".into());
let ctx2 = WorkflowCtx::new("b".into(), Journal::null(), BTreeMap::new());
ctx2.set_code_version("v2".into());
assert_ne!(
ctx1.agent_content_key("p", "s", "n", "ph", ""),
ctx2.agent_content_key("p", "s", "n", "ph", ""),
"a changed code-version produces different content-keys (auto-invalidation)"
);
}
#[test]
fn args_changes_invalidate_keys() {
let ctx1 = WorkflowCtx::new_with_args(
"a".into(),
Journal::null(),
BTreeMap::new(),
r#"{"batch":1}"#.into(),
);
ctx1.set_code_version("v1".into());
let ctx2 = WorkflowCtx::new_with_args(
"b".into(),
Journal::null(),
BTreeMap::new(),
r#"{"batch":2}"#.into(),
);
ctx2.set_code_version("v1".into());
assert_ne!(
ctx1.checkpoint_content_key("files", "ph"),
ctx2.checkpoint_content_key("files", "ph"),
"changed workflow args produce different content-keys"
);
}
#[test]
fn memo_store_round_trip_guard_skips_unsurvivable_values() {
let ctx = WorkflowCtx::new("wf_t".into(), Journal::null(), BTreeMap::new());
ctx.memo_store("ck_text", &Value::string("hello"));
assert_eq!(ctx.memo_lookup("ck_text"), Some(Value::string("hello")));
let mut m = BTreeMap::new();
m.insert(Value::keyword("body"), Value::string("x"));
let kw_map = Value::map(m);
ctx.memo_store("ck_map", &kw_map);
assert_eq!(ctx.memo_lookup("ck_map"), Some(kw_map));
let mut bad = BTreeMap::new();
bad.insert(Value::int(1), Value::int(2));
ctx.memo_store("ck_bad", &Value::map(bad));
assert_eq!(
ctx.memo_lookup("ck_bad"),
None,
"a non-round-trippable value must be left un-memoized"
);
}
#[test]
fn checkpoint_round_trips() {
let ctx = WorkflowCtx::new("wf_t".into(), Journal::null(), BTreeMap::new());
assert_eq!(ctx.read_checkpoint("files"), None);
ctx.store_checkpoint("files", Value::int(3));
assert_eq!(ctx.read_checkpoint("files"), Some(Value::int(3)));
}
#[test]
fn mcp_handle_registry_starts_empty_and_undeclared() {
let ctx = WorkflowCtx::new("wf_t".into(), Journal::null(), BTreeMap::new());
assert_eq!(ctx.mcp_handle("asana"), None);
assert!(!ctx.is_mcp_declared("asana"));
}
#[test]
fn mcp_declared_tracks_aliases_before_handles_resolve() {
let ctx = WorkflowCtx::new("wf_t".into(), Journal::null(), BTreeMap::new());
ctx.set_mcp_declared(vec!["asana".to_string(), "fs".to_string()]);
assert!(ctx.is_mcp_declared("asana"));
assert_eq!(ctx.mcp_handle("asana"), None);
assert!(!ctx.is_mcp_declared("zebra"));
}
#[test]
fn mcp_handle_returns_resolved_handle_by_alias() {
let ctx = WorkflowCtx::new("wf_t".into(), Journal::null(), BTreeMap::new());
ctx.set_mcp_declared(vec!["asana".to_string(), "fs".to_string()]);
let mut handles = BTreeMap::new();
handles.insert("asana".to_string(), Value::string("mcp-1"));
handles.insert("fs".to_string(), Value::string("mcp-2"));
ctx.set_mcp_handles(handles);
assert_eq!(ctx.mcp_handle("asana"), Some(Value::string("mcp-1")));
assert_eq!(ctx.mcp_handle("fs"), Some(Value::string("mcp-2")));
assert_eq!(ctx.mcp_handle("nope"), None);
}
#[test]
fn workflow_ctx_traces_state_memo_and_mcp_values() {
let ctx = WorkflowCtx::new("t".into(), Journal::null(), BTreeMap::new());
ctx.store_checkpoint("k", Value::int(1));
let mut memos = HashMap::new();
memos.insert("ck".to_string(), Value::int(2));
ctx.enter_resume(memos);
let mut handles = BTreeMap::new();
handles.insert("asana".to_string(), Value::string("handle"));
ctx.set_mcp_handles(handles);
let mut edges = 0;
assert!(ctx.trace(&mut |edge| {
assert!(matches!(edge, GcEdge::Value(_)));
edges += 1;
}));
assert_eq!(
edges, 3,
"state bag + resume memo + MCP handle each trace once"
);
}
#[test]
fn host_scope_restores_previous_on_drop() {
assert!(current_for(None).is_none());
let outer = WorkflowCtx::new("outer".into(), Journal::null(), BTreeMap::new());
let g_outer = install_scope(None, outer);
assert_eq!(
current_for(None).map(|c| c.run_id.clone()).as_deref(),
Some("outer")
);
{
let inner = WorkflowCtx::new("inner".into(), Journal::null(), BTreeMap::new());
let _g_inner = install_scope(None, inner);
assert_eq!(
current_for(None).map(|c| c.run_id.clone()).as_deref(),
Some("inner")
);
}
assert_eq!(
current_for(None).map(|c| c.run_id.clone()).as_deref(),
Some("outer")
);
drop(g_outer);
assert!(current_for(None).is_none());
}
#[test]
fn task_state_removes_the_exact_token_out_of_lifo() {
let state = WorkflowTaskState::default();
let outer = WorkflowCtx::new("outer".into(), Journal::null(), BTreeMap::new());
let inner = WorkflowCtx::new("inner".into(), Journal::null(), BTreeMap::new());
let outer_token = state.install(outer);
let inner_token = state.install(inner);
assert_eq!(
state.current_ctx().map(|c| c.run_id.clone()).as_deref(),
Some("inner")
);
assert!(state.remove(outer_token));
assert!(
!state.remove(outer_token),
"removing the same token twice is idempotent"
);
assert_eq!(
state.current_ctx().map(|c| c.run_id.clone()).as_deref(),
Some("inner"),
"removing the outer token leaves the inner scope live and on top"
);
assert!(state.remove(inner_token));
assert!(state.current_ctx().is_none());
}
#[test]
fn child_inherits_run_and_agent_but_not_removal_authority() {
let state = Rc::new(WorkflowTaskState::default());
let run = WorkflowCtx::new("shared-run".into(), Journal::null(), BTreeMap::new());
let parent_token = state.install(run);
state.set_cur_agent(Some("scout_1".to_string()));
let child = state.inherit();
let child = child
.as_any()
.downcast_ref::<WorkflowTaskState>()
.expect("inherited workflow state");
assert_eq!(
child.current_ctx().map(|c| c.run_id.clone()).as_deref(),
Some("shared-run"),
"child observes the spawner's active run"
);
assert_eq!(child.cur_agent().as_deref(), Some("scout_1"));
assert!(!child.remove(parent_token));
assert_eq!(
child.current_ctx().map(|c| c.run_id.clone()).as_deref(),
Some("shared-run")
);
assert!(state.remove(parent_token));
assert!(state.current_ctx().is_none());
}
#[test]
fn generated_run_id_has_secs_nanos_pid_and_nonce() {
let a = generate_run_id();
let b = generate_run_id();
assert_ne!(a, b, "the process nonce makes back-to-back ids distinct");
for id in [&a, &b] {
assert!(id.starts_with("wf_"), "id keeps the wf_ prefix: {id}");
let parts: Vec<&str> = id.split('_').collect();
assert_eq!(parts.len(), 5, "wf_<secs>_<nanos>_<pid>_<nonce>: {id}");
for field in &parts[1..] {
assert!(
!field.is_empty() && field.bytes().all(|c| c.is_ascii_digit()),
"each generated id field is a non-empty number: {id}"
);
}
}
}
#[test]
fn validate_explicit_run_id_accepts_safe_names_and_rejects_unsafe() {
for ok in ["wf_test_0001", "run-42", "abc.def", "a"] {
assert!(validate_explicit_run_id(ok).is_ok(), "should accept {ok:?}");
}
for bad in [
"", "a/b", "a\\b", "..", "a..b", ".", "...", "a\0b", "a\nb", ] {
let err = validate_explicit_run_id(bad).expect_err(&format!("should reject {bad:?}"));
assert_eq!(err.kind(), io::ErrorKind::InvalidInput, "for {bad:?}");
}
}
#[test]
fn open_fresh_with_retries_past_a_colliding_id() {
let mut root = std::env::temp_dir();
root.push(format!(
"sema-wf-fresh-retry-{}-{}",
std::process::id(),
SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap()
.as_nanos()
));
let root_str = root.to_string_lossy().to_string();
for taken in ["taken_1", "taken_2"] {
std::fs::create_dir_all(root.join(taken)).unwrap();
std::fs::write(root.join(taken).join("events.jsonl"), "{}\n").unwrap();
}
let mut candidates = ["taken_1", "taken_2", "free_3"].into_iter();
let (id, _journal) =
open_fresh_with(&root_str, || candidates.next().unwrap().to_string()).unwrap();
assert_eq!(id, "free_3", "opener retried past the colliding ids");
std::fs::remove_dir_all(&root).ok();
}
#[test]
fn civil_date_epoch() {
assert_eq!(civil_from_days(0), (1970, 1, 1));
assert_eq!(civil_from_days(20_628), (2026, 6, 24));
}
#[test]
fn redacts_mcp_headers_and_env_values() {
let meta = serde_json::json!({
"budget": {"usd": 1.0},
"mcp": {
"asana": {
"url": "https://mcp.asana.com/mcp",
"headers": {"Authorization": "Bearer secret-token"},
"persist": "workflow"
},
"fs": {
"command": "npx",
"env": {"API_TOKEN": "supersecret", "PLAIN": "not-a-secret-name"}
}
}
});
let redacted = redact_meta_secrets(meta);
assert_eq!(
redacted["mcp"]["asana"]["headers"]["Authorization"],
"<redacted>"
);
assert_eq!(redacted["mcp"]["fs"]["env"]["API_TOKEN"], "<redacted>");
assert_eq!(redacted["mcp"]["fs"]["env"]["PLAIN"], "<redacted>");
}
#[test]
fn redaction_keeps_header_and_env_keys_and_sibling_fields() {
let meta = serde_json::json!({
"mcp": {
"asana": {
"url": "https://mcp.asana.com/mcp",
"headers": {"Authorization": "Bearer secret-token", "X-Trace": "abc"},
"tools": ["create_task"],
"persist": "workflow"
}
}
});
let redacted = redact_meta_secrets(meta);
assert!(redacted["mcp"]["asana"]["headers"]
.as_object()
.unwrap()
.contains_key("Authorization"));
assert!(redacted["mcp"]["asana"]["headers"]
.as_object()
.unwrap()
.contains_key("X-Trace"));
assert_eq!(redacted["mcp"]["asana"]["url"], "https://mcp.asana.com/mcp");
assert_eq!(redacted["mcp"]["asana"]["tools"][0], "create_task");
assert_eq!(redacted["mcp"]["asana"]["persist"], "workflow");
}
#[test]
fn meta_without_mcp_passes_through_unchanged() {
let meta = serde_json::json!({
"budget": {"usd": 1.0},
"args": {"repo": "sema-lisp/sema"},
"phases": ["Triage"],
});
let redacted = redact_meta_secrets(meta.clone());
assert_eq!(redacted, meta);
}
#[test]
fn mcp_alias_without_headers_or_env_passes_through_unchanged() {
let meta = serde_json::json!({
"mcp": {"asana": {"url": "https://mcp.asana.com/mcp", "persist": "workflow"}}
});
let redacted = redact_meta_secrets(meta.clone());
assert_eq!(redacted, meta);
}
}