1use std::collections::{BTreeMap, BTreeSet};
16use std::fs;
17use std::path::{Path, PathBuf};
18
19use serde_json::{Map, Value};
20use sha2::{Digest, Sha256};
21
22use super::canonical::canonical_json;
23use super::decode::{
24 decode_channel, decode_fire_row, decode_job, decode_obligation_row, decode_route,
25 decode_subscription, decode_worker, encode_channel, encode_fire_row, encode_job,
26 encode_obligation_folder_extras, encode_obligation_row, encode_route, encode_subscription,
27 load_error, surface_key_string, EXECUTION_COLUMNS, FOLDER_OBLIGATION_EXTRA_COLUMNS,
28 OBLIGATION_COLUMNS,
29};
30use super::dotenv::{parse_dotenv, render_dotenv};
31use super::sqlite::{read_rows, replace_rows, table_exists, write_table, Param};
32use crate::ontology::{
33 hermes_cron_job_id, parse_hermes_session_key, Binding, EndReason, Handoff, HarnessId,
34 HermesSessionRow, Recurrence, Residue, SurfaceKey, Trigger, Worker,
35};
36use crate::orchestration::{
37 Access, AccessPolicy, ChannelConfig, Orchestration, PersonaRef, Profile, ProfileResidue,
38};
39use crate::Result;
40
41pub const OWNED_FILES: &[&str] = &[
43 "config.yaml",
44 "AGENTS.md",
45 "CLAUDE.md",
46 ".env",
47 "cron/jobs.json",
48 "cron/executions.db",
49 "webhook_subscriptions.json",
50 "state.db",
51 "agents.json",
52];
53
54const CONFIG_O_KEYS: &[&str] = &["worker", "budgets"];
55
56#[derive(Debug, Default, serde::Serialize, serde::Deserialize)]
58#[serde(deny_unknown_fields)]
59struct AgentLayerFile {
60 #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
61 agents: BTreeMap<String, crate::orchestration::AgentDecl>,
62 #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
63 conversations: BTreeMap<String, crate::orchestration::Conversation>,
64 #[serde(default, skip_serializing_if = "Vec::is_empty")]
65 transfers: Vec<crate::orchestration::Transfer>,
66 #[serde(default, skip_serializing_if = "Vec::is_empty")]
67 handoffs: Vec<crate::orchestration::agent::Handoff>,
68 #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
69 usage: BTreeMap<String, crate::orchestration::SessionUsage>,
70}
71
72fn agent_layer_record(o: &Orchestration) -> Value {
73 serde_json::to_value(AgentLayerFile {
74 agents: o.agents.clone(),
75 conversations: o.conversations.clone(),
76 transfers: o.transfers.clone(),
77 handoffs: o.handoffs.clone(),
78 usage: o.usage.clone(),
79 })
80 .unwrap()
81}
82
83#[derive(Debug, Clone, Copy, PartialEq, Eq)]
85pub enum Flavor {
86 Orchestrator,
88 Hermes,
90}
91
92#[derive(Debug, Clone, PartialEq)]
94pub struct JobsForm {
95 pub object: bool,
97 pub extras: Map<String, Value>,
99}
100
101#[derive(Debug, Clone)]
104pub struct ProfileIo {
105 pub raw: BTreeMap<String, String>,
107 pub snapshot: BTreeMap<String, String>,
109 pub source_dir: Option<PathBuf>,
111 pub flavor: Flavor,
113 pub jobs_form: Option<JobsForm>,
115 pub routes_at_top: bool,
117 pub borrowed_from: Option<PathBuf>,
119 pub lenders: Vec<String>,
121}
122
123impl ProfileIo {
124 fn new(flavor: Flavor) -> Self {
125 Self {
126 raw: BTreeMap::new(),
127 snapshot: BTreeMap::new(),
128 source_dir: None,
129 flavor,
130 jobs_form: None,
131 routes_at_top: false,
132 borrowed_from: None,
133 lenders: Vec::new(),
134 }
135 }
136}
137
138#[derive(Debug, Clone)]
140pub struct LoadedHome {
141 pub orchestration: Orchestration,
143 pub vault: BTreeMap<String, String>,
145 pub io: BTreeMap<String, ProfileIo>,
147}
148
149fn sha256_hex(text: &str) -> String {
150 let mut h = Sha256::new();
151 h.update(text.as_bytes());
152 h.finalize().iter().map(|b| format!("{b:02x}")).collect()
153}
154
155pub fn persona_ref(text: &str) -> PersonaRef {
157 PersonaRef {
158 path: "AGENTS.md".into(),
159 text: Some(text.to_string()),
160 sha256: sha256_hex(text),
161 }
162}
163
164fn read_text(dir: &Path, rel: &str) -> Result<Option<String>> {
165 let p = dir.join(rel);
166 if !p.is_file() {
167 return Ok(None);
168 }
169 Ok(Some(fs::read_to_string(&p)?))
170}
171
172fn yaml_to_json(file: &str, text: &str) -> Result<Value> {
173 let value: serde_yaml::Value =
174 serde_yaml::from_str(text).map_err(|e| load_error(file, "", format!("YAML: {e}")))?;
175 let json: Value =
176 serde_json::to_value(value).map_err(|e| load_error(file, "", format!("YAML: {e}")))?;
177 Ok(if json.is_null() {
178 Value::Object(Map::new())
179 } else {
180 json
181 })
182}
183
184const YAML11_BOOLS: [&str; 18] = [
187 "yes", "Yes", "YES", "no", "No", "NO", "true", "True", "TRUE", "false", "False", "FALSE", "on",
188 "On", "ON", "off", "Off", "OFF",
189];
190const YAML11_MARK: &str = "__supercode_yaml11_string__";
191
192fn mark_yaml11_strings(value: &Value) -> Value {
193 match value {
194 Value::String(s) if YAML11_BOOLS.contains(&s.as_str()) => {
195 Value::String(format!("{YAML11_MARK}{s}"))
196 }
197 Value::Array(items) => Value::Array(items.iter().map(mark_yaml11_strings).collect()),
198 Value::Object(map) => Value::Object(
199 map.iter()
200 .map(|(k, v)| (k.clone(), mark_yaml11_strings(v)))
201 .collect(),
202 ),
203 other => other.clone(),
204 }
205}
206
207fn json_to_yaml(value: &Value) -> String {
208 if value.as_object().is_some_and(|m| m.is_empty()) {
209 return String::new();
210 }
211 let y: serde_yaml::Value =
213 serde_json::from_value(mark_yaml11_strings(value)).unwrap_or(serde_yaml::Value::Null);
214 let mut out = serde_yaml::to_string(&y).unwrap_or_default();
215 for word in YAML11_BOOLS {
216 out = out.replace(&format!("{YAML11_MARK}{word}"), &format!("'{word}'"));
217 }
218 out
219}
220
221pub fn empty_profile(name: &str, dir: &Path) -> Profile {
223 Profile {
224 name: name.to_string(),
225 dir: dir.to_path_buf(),
226 worker: None,
227 persona: None,
228 channels: BTreeMap::new(),
229 routes: Vec::new(),
230 jobs: BTreeMap::new(),
231 subscriptions: BTreeMap::new(),
232 access: Default::default(),
233 bindings: BTreeMap::new(),
234 fires: Vec::new(),
235 obligations: Vec::new(),
236 budgets: Vec::new(),
237 residue: ProfileResidue::default(),
238 }
239}
240
241pub fn config_record(profile: &Profile) -> Value {
243 serde_json::json!({
244 "worker": profile.worker, "budgets": profile.budgets,
245 "routes": profile.routes, "channels": profile.channels, "residue": profile.residue.config,
246 })
247}
248
249fn state_record(profile: &Profile) -> Value {
250 serde_json::json!({ "bindings": profile.bindings, "obligations": profile.obligations })
251}
252
253pub fn load_home(root: &Path, flavor: Flavor) -> Result<LoadedHome> {
255 if !root.is_dir() {
256 return Err(load_error(
257 &root.display().to_string(),
258 "",
259 "not a directory",
260 ));
261 }
262 let mut vault = BTreeMap::new();
263 let mut io = BTreeMap::new();
264 let mut profiles = BTreeMap::new();
265 let (default, default_io) = load_profile_dir("default", root, flavor, &mut vault)?;
266 profiles.insert("default".to_string(), default);
267 io.insert("default".to_string(), default_io);
268 let profiles_dir = root.join("profiles");
269 if profiles_dir.is_dir() {
270 let mut names: Vec<String> = fs::read_dir(&profiles_dir)?
271 .flatten()
272 .filter(|e| e.path().is_dir())
273 .filter_map(|e| e.file_name().into_string().ok())
274 .filter(|n| n != "node_modules" && !n.starts_with('.'))
275 .collect();
276 names.sort();
277 for name in names {
278 let dir = profiles_dir.join(&name);
279 if name == "default" {
280 return Err(load_error(
281 &dir.display().to_string(),
282 "",
283 "\"default\" is the root folder, not a named profile",
284 ));
285 }
286 let (profile, meta) = load_profile_dir(&name, &dir, flavor, &mut vault)?;
287 profiles.insert(name.clone(), profile);
288 io.insert(name, meta);
289 }
290 }
291 let workflow = match read_text(root, "workflow.yaml")? {
293 Some(text) => {
294 let file = root.join("workflow.yaml").display().to_string();
295 let value = yaml_to_json(&file, &text)?;
296 let workflow: crate::orchestration::Workflow =
297 serde_json::from_value(value).map_err(|e| load_error(&file, "", e.to_string()))?;
298 workflow.check().map_err(|e| load_error(&file, "", e))?;
299 if let Some(meta) = io.get_mut("default") {
300 meta.raw.insert("workflow.yaml".into(), text);
301 meta.snapshot.insert(
302 "workflow.yaml".into(),
303 canonical_json(&serde_json::to_value(&workflow).unwrap()),
304 );
305 }
306 workflow
307 }
308 None => crate::orchestration::hermes_instance(),
309 };
310 let layer = match (flavor, read_text(root, "agents.json")?) {
312 (Flavor::Orchestrator, Some(text)) => {
313 let file = root.join("agents.json").display().to_string();
314 let layer: AgentLayerFile =
315 serde_json::from_str(&text).map_err(|e| load_error(&file, "", e.to_string()))?;
316 if let Some(meta) = io.get_mut("default") {
317 meta.raw.insert("agents.json".into(), text);
318 meta.snapshot.insert(
319 "agents.json".into(),
320 canonical_json(&serde_json::to_value(&layer).unwrap()),
321 );
322 }
323 layer
324 }
325 _ => AgentLayerFile::default(),
326 };
327 let mut loaded = LoadedHome {
328 orchestration: Orchestration {
329 root: root.to_path_buf(),
330 profiles,
331 workflow,
332 agents: layer.agents,
333 conversations: layer.conversations,
334 transfers: layer.transfers,
335 handoffs: layer.handoffs,
336 usage: layer.usage,
337 },
338 vault,
339 io,
340 };
341 if flavor == Flavor::Hermes {
342 partition_shared_store(&mut loaded, root)?;
343 }
344 link_fires(&mut loaded, root)?;
345 Ok(loaded)
346}
347
348fn link_fires(loaded: &mut LoadedHome, root: &Path) -> Result<()> {
356 let names: Vec<String> = loaded.orchestration.profiles.keys().cloned().collect();
357 let root_store = root.join("state.db");
358 for name in &names {
359 let own = loaded.orchestration.profiles[name].dir.join("state.db");
360 let store = if own.is_file() {
361 own
362 } else {
363 root_store.clone()
364 };
365 let sessions: Vec<Map<String, Value>> = if table_exists(&store, "sessions") {
366 read_rows(
367 &store,
368 "select id, started_at, end_reason, parent_session_id, session_key from sessions",
369 &[],
370 )?
371 .unwrap_or_default()
372 } else {
373 Vec::new()
374 };
375 let holders: Vec<String> = {
378 let io = &loaded.io[name];
379 if io.borrowed_from.is_some() || !io.lenders.is_empty() {
380 let mut v = vec!["default".to_string()];
381 v.extend(loaded.io["default"].lenders.iter().cloned());
382 v
383 } else {
384 vec![name.clone()]
385 }
386 };
387 let mut by_fire: BTreeMap<String, (String, String)> = BTreeMap::new();
390 for (holder, profile) in &loaded.orchestration.profiles {
391 for o in &profile.obligations {
392 if let crate::orchestration::ObligationSource::Fire { fire_id } = &o.source {
393 by_fire
394 .entry(fire_id.clone())
395 .or_insert_with(|| (holder.clone(), o.id.clone()));
396 }
397 }
398 }
399 let mut links: Vec<(usize, Option<String>, Option<(String, String)>)> = Vec::new();
400 let profile = &loaded.orchestration.profiles[name];
401 for (index, fire) in profile.fires.iter().enumerate() {
402 if let Some(known) = by_fire.get(&fire.id) {
403 links.push((index, fire.session_id.clone(), Some(known.clone())));
405 continue;
406 }
407 let session_id = fire_session(&sessions, fire).or_else(|| fire.session_id.clone());
408 let session_key = session_id.as_deref().and_then(|id| {
409 sessions
410 .iter()
411 .find(|row| row.get("id").and_then(Value::as_str) == Some(id))
412 .and_then(|row| row.get("session_key"))
413 .and_then(Value::as_str)
414 .filter(|k| !k.is_empty())
415 .map(str::to_string)
416 });
417 let surface = profile.jobs.get(&fire.job_id).and_then(job_surface);
418 let (Some(from), to) = (
419 iso_epoch(&fire.claimed_at),
420 fire.finished_at
421 .as_deref()
422 .and_then(iso_epoch)
423 .unwrap_or(f64::MAX),
424 ) else {
425 links.push((index, session_id, None));
426 continue;
427 };
428 let in_window = |o: &&crate::orchestration::Obligation| {
429 o.created_at
430 .parse::<f64>()
431 .is_ok_and(|at| at >= from && at <= to)
432 };
433 let latest = |mut found: Vec<(&String, &crate::orchestration::Obligation)>| {
434 found.sort_by(|a, b| {
435 let at = |o: &crate::orchestration::Obligation| {
436 o.created_at.parse::<f64>().unwrap_or(0.0)
437 };
438 at(b.1)
439 .partial_cmp(&at(a.1))
440 .unwrap_or(std::cmp::Ordering::Equal)
441 });
442 found
443 .first()
444 .map(|(holder, o)| ((*holder).clone(), o.id.clone()))
445 };
446 let candidates = |pick: &dyn Fn(&crate::orchestration::Obligation) -> bool| {
447 holders
448 .iter()
449 .flat_map(|h| {
450 loaded.orchestration.profiles[h]
451 .obligations
452 .iter()
453 .filter(in_window)
454 .filter(|o| pick(o))
455 .map(move |o| (h, o))
456 })
457 .collect::<Vec<_>>()
458 };
459 let obligation = match &session_key {
460 Some(key) => latest(candidates(&|o| o.session_key.as_deref() == Some(key))),
461 None => None,
462 }
463 .or_else(|| {
464 let (platform, chat_id) = surface.as_ref()?;
465 latest(candidates(&|o| {
466 o.target.platform.as_deref() == Some(platform)
467 && o.target.chat_id.as_deref() == Some(chat_id)
468 }))
469 });
470 links.push((index, session_id, obligation));
471 }
472 for (index, session_id, obligation) in links {
473 let fire_id = {
474 let fire = &mut loaded.orchestration.profiles.get_mut(name).unwrap().fires[index];
475 fire.session_id = session_id;
476 fire.obligation_id = obligation.as_ref().map(|(_, id)| id.clone());
477 fire.id.clone()
478 };
479 if let Some((holder, obligation_id)) = obligation {
480 if let Some(o) = loaded
481 .orchestration
482 .profiles
483 .get_mut(&holder)
484 .and_then(|p| p.obligations.iter_mut().find(|o| o.id == obligation_id))
485 {
486 o.source = crate::orchestration::ObligationSource::Fire { fire_id };
487 }
488 }
489 }
490 }
491 for name in &names {
493 let profile = &loaded.orchestration.profiles[name];
494 let io = loaded.io.get_mut(name).unwrap();
495 io.snapshot.insert(
496 "cron/executions.db".into(),
497 canonical_json(&serde_json::to_value(&profile.fires).unwrap()),
498 );
499 io.snapshot
500 .insert("state.db".into(), canonical_json(&state_record(profile)));
501 }
502 Ok(())
503}
504
505fn job_surface(job: &crate::orchestration::Job) -> Option<(String, String)> {
510 use crate::orchestration::Target;
511 match &job.deliver {
512 Target::Origin => {
513 let origin = job.origin.as_ref()?;
514 Some((origin.platform.clone(), origin.chat_id.clone()?))
515 }
516 Target::Explicit {
517 platform, chat_id, ..
518 } => Some((platform.clone(), chat_id.clone()?)),
519 Target::Home | Target::Local => None,
520 }
521}
522
523fn fire_session(
527 sessions: &[Map<String, Value>],
528 fire: &crate::orchestration::Fire,
529) -> Option<String> {
530 let claimed = instant_key(&fire.claimed_at)?;
531 let finished = fire.finished_at.as_deref().and_then(instant_key);
532 let candidates = sessions.iter().filter_map(|row| {
533 let id = row.get("id").and_then(Value::as_str)?;
534 let key = cron_session_instant(id, &fire.job_id)?;
535 (key >= claimed && finished.is_none_or(|f| key <= f)).then(|| (key, id.to_string()))
536 });
537 let chosen = match finished {
538 Some(_) => candidates.max_by_key(|(key, _)| *key),
539 None => candidates.min_by_key(|(key, _)| *key),
540 }?;
541 let mut current = chosen.1;
542 for _ in 0..32 {
543 let row = sessions
544 .iter()
545 .find(|row| row.get("id").and_then(Value::as_str) == Some(current.as_str()));
546 let compressed = row
547 .and_then(|row| row.get("end_reason"))
548 .and_then(Value::as_str)
549 == Some("compression");
550 if !compressed {
551 return Some(current);
552 }
553 let next = sessions
554 .iter()
555 .filter(|row| {
556 row.get("parent_session_id").and_then(Value::as_str) == Some(current.as_str())
557 })
558 .max_by(|a, b| {
559 let at = |r: &Map<String, Value>| {
560 r.get("started_at").and_then(Value::as_f64).unwrap_or(0.0)
561 };
562 at(a)
563 .partial_cmp(&at(b))
564 .unwrap_or(std::cmp::Ordering::Equal)
565 .then_with(|| {
566 let id = |r: &Map<String, Value>| {
567 r.get("id")
568 .and_then(Value::as_str)
569 .unwrap_or("")
570 .to_string()
571 };
572 id(a).cmp(&id(b))
573 })
574 })
575 .and_then(|row| row.get("id").and_then(Value::as_str).map(str::to_string));
576 match next {
577 None => return Some(current),
578 Some(next) => current = next,
579 }
580 }
581 Some(current)
582}
583
584fn instant_key(iso: &str) -> Option<u64> {
586 let digits: String = iso
587 .chars()
588 .take_while(|c| *c != '+' && *c != 'Z')
589 .filter(char::is_ascii_digit)
590 .collect();
591 (digits.len() >= 14).then(|| digits[..14].parse().ok())?
592}
593
594fn cron_session_instant(session_id: &str, job_id: &str) -> Option<u64> {
596 if hermes_cron_job_id(session_id).as_deref() != Some(job_id) {
597 return None;
598 }
599 let stamp = session_id.rsplit_once('_')?;
600 let date = stamp.0.rsplit_once('_')?.1;
601 format!("{date}{}", stamp.1).parse().ok()
602}
603
604pub(crate) fn iso_epoch(iso: &str) -> Option<f64> {
606 let (instant, offset) = if let Some(instant) = iso.strip_suffix('Z') {
607 (instant, 0.0)
608 } else {
609 let time_at = iso.find('T')?;
610 let sign_at = iso[time_at..].find(['+', '-']).map(|i| i + time_at)?;
611 let (instant, offset) = iso.split_at(sign_at);
612 let (hours, minutes) = offset[1..].split_once(':')?;
613 let seconds = hours.parse::<f64>().ok()? * 3_600.0 + minutes.parse::<f64>().ok()? * 60.0;
614 (
615 instant,
616 if offset.starts_with('-') {
617 -seconds
618 } else {
619 seconds
620 },
621 )
622 };
623 let (date, time) = instant.split_once('T')?;
624 let mut date = date.splitn(3, '-');
625 let year: i64 = date.next()?.parse().ok()?;
626 let month: i64 = date.next()?.parse().ok()?;
627 let day: i64 = date.next()?.parse().ok()?;
628 let mut clock = time.splitn(3, ':');
629 let hour: i64 = clock.next()?.parse().ok()?;
630 let minute: i64 = clock.next()?.parse().ok()?;
631 let seconds: f64 = clock.next()?.parse().ok()?;
632 let year = year - i64::from(month <= 2);
633 let era = year.div_euclid(400);
634 let yoe = year - era * 400;
635 let doy = (153 * (if month > 2 { month - 3 } else { month + 9 }) + 2) / 5 + day - 1;
636 let doe = yoe * 365 + yoe / 4 - yoe / 100 + doy;
637 let days = era * 146_097 + doe - 719_468;
638 Some((days * 86_400 + hour * 3_600 + minute * 60) as f64 + seconds - offset)
639}
640
641fn partition_shared_store(loaded: &mut LoadedHome, root: &Path) -> Result<()> {
646 let root_path = root.join("state.db");
647 if !root_path.is_file() {
648 return Ok(());
649 }
650 let names: Vec<String> = loaded
651 .orchestration
652 .profiles
653 .keys()
654 .filter(|n| *n != "default")
655 .cloned()
656 .collect();
657 for name in names {
658 let has_own = loaded.orchestration.profiles[&name]
659 .dir
660 .join("state.db")
661 .is_file();
662 if has_own {
663 continue;
664 }
665 loaded.io.get_mut(&name).unwrap().borrowed_from = Some(root_path.clone());
666 loaded
667 .io
668 .get_mut("default")
669 .unwrap()
670 .lenders
671 .push(name.clone());
672 if table_exists(&root_path, "sessions") {
673 let rows = read_rows(
674 &root_path,
675 "select * from sessions where profile_name = ?1 order by started_at, id",
676 &[&name],
677 )?
678 .unwrap_or_default();
679 for row in rows {
680 if let Some(b) = binding_from_hermes_session(&root_path, &row, &name) {
681 loaded
682 .orchestration
683 .profiles
684 .get_mut(&name)
685 .unwrap()
686 .bindings
687 .insert(surface_key_string(&b.key), b);
688 }
689 }
690 }
691 let root_profile = loaded.orchestration.profiles.get_mut("default").unwrap();
693 let (mine, rest): (Vec<_>, Vec<_>) = root_profile.obligations.drain(..).partition(|o| {
694 o.session_key
695 .as_deref()
696 .and_then(parse_hermes_session_key)
697 .and_then(|(_, p)| p)
698 .as_deref()
699 == Some(name.as_str())
700 });
701 root_profile.obligations = rest;
702 loaded
703 .orchestration
704 .profiles
705 .get_mut(&name)
706 .unwrap()
707 .obligations = mine;
708 let snap = canonical_json(&state_record(&loaded.orchestration.profiles[&name]));
709 loaded
710 .io
711 .get_mut(&name)
712 .unwrap()
713 .snapshot
714 .insert("state.db".into(), snap);
715 }
716 let snap = canonical_json(&state_record(&loaded.orchestration.profiles["default"]));
717 loaded
718 .io
719 .get_mut("default")
720 .unwrap()
721 .snapshot
722 .insert("state.db".into(), snap);
723 Ok(())
724}
725
726const HERMES_SESSION_MAPPED: &[&str] = &[
727 "id",
728 "source",
729 "session_key",
730 "chat_id",
731 "chat_type",
732 "thread_id",
733 "user_id",
734 "profile_name",
735 "handoff_state",
736 "handoff_platform",
737 "handoff_error",
738 "started_at",
739 "ended_at",
740 "end_reason",
741];
742
743fn routing_scope(dir: &Path) -> String {
745 let sessions = dir.join("sessions");
746 fs::canonicalize(&sessions)
747 .or_else(|_| fs::canonicalize(dir).map(|d| d.join("sessions")))
748 .unwrap_or(sessions)
749 .display()
750 .to_string()
751}
752
753fn routing_entry(profile_name: &str, slot: &str, b: &Binding) -> (String, Value) {
756 let ns = if profile_name == "default" {
757 ""
758 } else {
759 profile_name
760 };
761 let key = b
762 .key
763 .key
764 .clone()
765 .filter(|k| k.starts_with("agent:"))
766 .unwrap_or_else(|| crate::ontology::render_hermes_session_key(ns, &b.key));
767 let entry = serde_json::json!({
768 "session_key": key,
769 "session_id": b.worker.session_id.clone().unwrap_or_else(|| slot.to_string()),
770 "created_at": b.started_at,
771 "updated_at": b.last_activity_at.clone().or_else(|| b.started_at.clone()),
772 "display_name": Value::Null,
773 "platform": b.key.platform,
774 "chat_type": b.key.kind.clone().unwrap_or_else(|| "dm".into()),
775 "origin": {
776 "platform": b.key.platform,
777 "chat_id": b.key.chat_id,
778 "chat_type": b.key.kind,
779 "thread_id": b.key.thread_id,
780 "user_id": b.key.participant_id,
781 },
782 "metadata": { "supercode": { "slot": slot, "binding": b } },
783 "input_tokens": 0, "output_tokens": 0, "cache_read_tokens": 0, "cache_write_tokens": 0,
784 "total_tokens": 0, "last_prompt_tokens": 0, "estimated_cost_usd": 0.0, "cost_status": "unknown",
785 "was_auto_reset": false, "auto_reset_reason": Value::Null, "reset_had_activity": false,
786 });
787 (key, entry)
788}
789
790fn binding_from_routing(file: &Path, row: &Map<String, Value>) -> Option<(String, Binding)> {
793 let entry: Value = serde_json::from_str(row.get("entry_json")?.as_str()?).ok()?;
794 if let Some(sc) = entry.pointer("/metadata/supercode") {
795 let b: Binding = serde_json::from_value(sc.get("binding")?.clone()).ok()?;
796 let slot = sc
797 .get("slot")
798 .and_then(Value::as_str)
799 .map(str::to_string)
800 .unwrap_or_else(|| surface_key_string(&b.key));
801 return Some((slot, b));
802 }
803 let session_key = entry.get("session_key").and_then(Value::as_str)?;
804 let (mut key, _) = parse_hermes_session_key(session_key)?;
805 key.key = None;
806 let text = |k: &str| entry.get(k).and_then(Value::as_str).map(str::to_string);
807 let b = Binding {
808 key: key.clone(),
809 profile: None,
810 worker: Worker {
811 harness: HarnessId::new(HarnessId::HERMES),
812 session_id: text("session_id"),
813 locator: Some(file.display().to_string()),
814 },
815 trigger: Trigger::Unknown,
816 recurrence: None,
817 conversation: None,
818 agent: None,
819 handoff: None,
820 started_at: text("created_at"),
821 last_activity_at: text("updated_at"),
822 ended_at: None,
823 end_reason: None,
824 residue: Residue::default(),
825 };
826 Some((surface_key_string(&key), b))
827}
828
829pub fn binding_from_hermes_session(
831 file: &Path,
832 row: &Map<String, Value>,
833 profile_name: &str,
834) -> Option<Binding> {
835 let text = |k: &str| {
836 row.get(k)
837 .and_then(|v| match v {
838 Value::String(s) => Some(s.clone()),
839 Value::Number(n) => Some(n.to_string()),
840 _ => None,
841 })
842 .filter(|s| !s.is_empty())
843 };
844 let num = |k: &str| row.get(k).and_then(Value::as_f64);
845 let id = text("id")?;
846 let source = text("source");
847 let parsed = text("session_key").and_then(|k| parse_hermes_session_key(&k));
850 let chat_type = text("chat_type").or_else(|| parsed.as_ref().and_then(|(k, _)| k.kind.clone()));
851 let key = if text("session_key").is_some()
852 && chat_type
853 .as_deref()
854 .is_some_and(|c| super::decode::CHAT_TYPES.contains(&c))
855 {
856 SurfaceKey {
857 key: None,
858 platform: source
859 .clone()
860 .or_else(|| parsed.as_ref().and_then(|(k, _)| k.platform.clone())),
861 kind: chat_type,
862 chat_id: text("chat_id")
863 .or_else(|| parsed.as_ref().and_then(|(k, _)| k.chat_id.clone())),
864 thread_id: text("thread_id")
865 .or_else(|| parsed.as_ref().and_then(|(k, _)| k.thread_id.clone())),
866 participant_id: parsed.as_ref().and_then(|(k, _)| k.participant_id.clone()),
867 }
868 } else if source.as_deref() == Some("cron") {
869 let job = crate::ontology::hermes_cron_job_id(&id).unwrap_or_else(|| id.clone());
870 SurfaceKey {
871 key: None,
872 platform: Some("cron".into()),
873 kind: Some("dm".into()),
874 chat_id: Some(job),
875 thread_id: None,
876 participant_id: None,
877 }
878 } else {
879 return None;
880 };
881 if let Some(p) = text("profile_name") {
882 if p != profile_name && !(profile_name == "default" && p == "main") {
883 return None; }
885 }
886 let iso = |v: Option<f64>| v.map(epoch_iso);
887 let end_word = text("end_reason");
888 let end_reason = end_word.as_deref().and_then(EndReason::parse);
889 let mut residue = Residue::default();
890 for (k, v) in row {
891 if !HERMES_SESSION_MAPPED.contains(&k.as_str()) && !v.is_null() {
892 residue.keep(k.clone(), v.clone());
893 }
894 }
895 if let (Some(word), None) = (&end_word, end_reason) {
896 residue.keep("end_reason", Value::String(word.clone()));
897 }
898 if let Some(u) = text("user_id") {
899 residue.keep("user_id", Value::String(u));
900 }
901 let recurrence = if key.platform.as_deref() == Some("cron") {
902 key.chat_id.clone().map(|job_id| Recurrence {
903 job_id,
904 kind: "cron".into(),
905 })
906 } else {
907 None
908 };
909 Some(Binding {
910 trigger: match (recurrence.is_some(), source.as_deref()) {
911 (true, _) => Trigger::Cron,
912 (_, Some(s)) => crate::ontology::hermes_trigger_for_source(s),
913 _ => Trigger::Unknown,
914 },
915 key,
916 profile: None,
917 worker: Worker {
918 harness: HarnessId::new(HarnessId::HERMES),
919 session_id: Some(id),
920 locator: Some(file.display().to_string()),
921 },
922 recurrence,
923 conversation: None,
924 agent: None,
925 handoff: text("handoff_state").map(|state| Handoff {
926 to: text("handoff_platform"),
927 state,
928 error: text("handoff_error"),
929 }),
930 started_at: iso(num("started_at")),
931 last_activity_at: iso(num("ended_at").or_else(|| num("started_at"))),
932 ended_at: iso(num("ended_at")),
933 end_reason,
934 residue,
935 })
936}
937
938fn epoch_iso(seconds: f64) -> String {
940 let row = HermesSessionRow {
941 started_at: Some(seconds),
942 ..Default::default()
943 };
944 Binding::from_hermes_row(&row, None)
945 .started_at
946 .unwrap_or_default()
947}
948
949fn load_profile_dir(
950 name: &str,
951 dir: &Path,
952 flavor: Flavor,
953 vault: &mut BTreeMap<String, String>,
954) -> Result<(Profile, ProfileIo)> {
955 let mut profile = empty_profile(name, dir);
956 let mut config_normalized = false;
957 let mut subs_normalized = false;
958 let mut meta = ProfileIo::new(flavor);
959 meta.source_dir = Some(dir.to_path_buf());
960 let remember = |meta: &mut ProfileIo, rel: &str, raw: Option<String>, record: &Value| {
961 if let Some(raw) = raw {
962 meta.raw.insert(rel.to_string(), raw);
963 }
964 meta.snapshot
965 .insert(rel.to_string(), canonical_json(record));
966 };
967
968 if let Some(env) = read_text(dir, ".env")? {
970 for (k, v) in parse_dotenv(&env) {
971 vault.insert(k, v);
972 }
973 meta.raw.insert(".env".into(), env);
974 }
975
976 let cfg_file = dir.join("config.yaml").display().to_string();
978 let cfg_text = read_text(dir, "config.yaml")?;
979 let cfg = match &cfg_text {
980 Some(text) => yaml_to_json(&cfg_file, text)?,
981 None => Value::Object(Map::new()),
982 };
983 let cfg_map = cfg
984 .as_object()
985 .ok_or_else(|| load_error(&cfg_file, "", "expected a mapping"))?;
986 profile.worker = decode_worker(&cfg_file, cfg_map.get("worker"))?;
987 if let Some(budgets) = cfg_map.get("budgets") {
988 profile.budgets = serde_json::from_value(budgets.clone())
989 .map_err(|e| load_error(&cfg_file, "budgets", e.to_string()))?;
990 }
991 let gateway = cfg_map.get("gateway").and_then(Value::as_object);
992 let routes_raw: Vec<Value> = match cfg_map.get("profile_routes").and_then(Value::as_array) {
993 Some(a) => {
994 meta.routes_at_top = true;
995 a.clone()
996 }
997 None => gateway
998 .and_then(|g| g.get("profile_routes"))
999 .and_then(Value::as_array)
1000 .cloned()
1001 .unwrap_or_default(),
1002 };
1003 for (i, r) in routes_raw.iter().enumerate() {
1004 profile.routes.push(decode_route(&cfg_file, i, r)?);
1005 }
1006 if let Some(platforms) = cfg_map.get("platforms") {
1007 let map = platforms
1008 .as_object()
1009 .ok_or_else(|| load_error(&cfg_file, "platforms", "expected a map"))?;
1010 let known: BTreeSet<String> = vault.keys().cloned().collect();
1011 for (platform, raw) in map {
1012 profile.channels.insert(
1013 platform.clone(),
1014 decode_channel(&cfg_file, platform, raw, vault)?,
1015 );
1016 }
1017 config_normalized =
1018 flavor == Flavor::Orchestrator && vault.keys().any(|k| !known.contains(k));
1019 }
1020 {
1026 let own: BTreeMap<String, String> = meta
1027 .raw
1028 .get(".env")
1029 .map(|text| parse_dotenv(text).into_iter().collect())
1030 .unwrap_or_default();
1031 let set = |k: &str| own.get(k).filter(|v| !v.trim().is_empty()).cloned();
1032 for (platform, token, home) in [
1033 ("discord", "DISCORD_BOT_TOKEN", "DISCORD_HOME_CHANNEL"),
1034 ("telegram", "TELEGRAM_BOT_TOKEN", "TELEGRAM_HOME_CHANNEL"),
1035 ("slack", "SLACK_BOT_TOKEN", "SLACK_HOME_CHANNEL"),
1036 ] {
1037 if profile.channels.contains_key(platform) || set(token).is_none() {
1038 continue;
1039 }
1040 let mut extra = BTreeMap::new();
1041 extra.insert("from_env".to_string(), Value::Bool(true));
1042 if let Some(chat) = set(home) {
1043 let mut hc = Map::new();
1044 hc.insert("platform".into(), Value::String(platform.into()));
1045 hc.insert("chat_id".into(), Value::String(chat));
1046 hc.insert(
1047 "name".into(),
1048 Value::String(set(&format!("{home}_NAME")).unwrap_or_else(|| "Home".into())),
1049 );
1050 if let Some(thread) = set(&format!("{home}_THREAD_ID")) {
1051 hc.insert("thread_id".into(), Value::String(thread));
1052 }
1053 extra.insert("home_channel".to_string(), Value::Object(hc));
1054 }
1055 profile.channels.insert(
1057 platform.to_string(),
1058 ChannelConfig {
1059 platform: platform.to_string(),
1060 enabled: true,
1061 credentials: BTreeMap::new(),
1062 extra,
1063 },
1064 );
1065 }
1066 }
1067 for (k, v) in cfg_map {
1069 if CONFIG_O_KEYS.contains(&k.as_str()) || k == "platforms" || k == "profile_routes" {
1070 continue;
1071 }
1072 if k == "gateway" {
1073 let mut g = v.as_object().cloned().unwrap_or_default();
1074 g.remove("profile_routes");
1075 if !g.is_empty() {
1076 profile
1077 .residue
1078 .config
1079 .insert("gateway".into(), Value::Object(g));
1080 }
1081 continue;
1082 }
1083 profile.residue.config.insert(k.clone(), v.clone());
1084 }
1085 remember(&mut meta, "config.yaml", cfg_text, &config_record(&profile));
1086 if config_normalized {
1087 meta.raw.remove("config.yaml");
1092 meta.snapshot.remove("config.yaml");
1093 }
1094
1095 let persona_file = if flavor == Flavor::Hermes {
1097 "SOUL.md"
1098 } else {
1099 "AGENTS.md"
1100 };
1101 let persona_text = read_text(dir, persona_file)?;
1102 profile.persona = persona_text.as_ref().map(|t| PersonaRef {
1103 path: "AGENTS.md".into(),
1104 text: Some(t.clone()),
1105 sha256: sha256_hex(t),
1106 });
1107 remember(
1108 &mut meta,
1109 persona_file,
1110 persona_text,
1111 &serde_json::to_value(&profile.persona).unwrap(),
1112 );
1113
1114 let jobs_file = dir.join("cron/jobs.json").display().to_string();
1116 let jobs_text = read_text(dir, "cron/jobs.json")?;
1117 if let Some(text) = &jobs_text {
1118 let parsed: Value = serde_json::from_str(text)
1119 .map_err(|e| load_error(&jobs_file, "", format!("JSON: {e}")))?;
1120 let arr: Vec<Value> = match &parsed {
1121 Value::Array(a) => {
1122 meta.jobs_form = Some(JobsForm {
1123 object: false,
1124 extras: Map::new(),
1125 });
1126 a.clone()
1127 }
1128 Value::Object(o) => match o.get("jobs") {
1129 Some(Value::Array(a)) => {
1130 let mut extras = o.clone();
1131 extras.remove("jobs");
1132 meta.jobs_form = Some(JobsForm {
1133 object: true,
1134 extras,
1135 });
1136 a.clone()
1137 }
1138 Some(Value::Object(m)) => {
1139 let mut extras = o.clone();
1140 extras.remove("jobs");
1141 meta.jobs_form = Some(JobsForm {
1142 object: true,
1143 extras,
1144 });
1145 m.iter()
1146 .map(|(id, j)| {
1147 let mut j = j.as_object().cloned().unwrap_or_default();
1148 j.insert("id".into(), Value::String(id.clone()));
1149 Value::Object(j)
1150 })
1151 .collect()
1152 }
1153 _ => {
1154 return Err(load_error(
1155 &jobs_file,
1156 "",
1157 "expected an array of jobs or {\"jobs\": [...]}",
1158 ))
1159 }
1160 },
1161 _ => {
1162 return Err(load_error(
1163 &jobs_file,
1164 "",
1165 "expected an array of jobs or {\"jobs\": [...]}",
1166 ))
1167 }
1168 };
1169 for raw in &arr {
1170 let job = decode_job(&jobs_file, raw)?;
1171 if profile.jobs.contains_key(&job.id) {
1172 return Err(load_error(&jobs_file, &job.id, "duplicate job id"));
1173 }
1174 profile.jobs.insert(job.id.clone(), job);
1175 }
1176 }
1177 let jobs_record: Vec<Value> = profile
1178 .jobs
1179 .values()
1180 .map(|j| serde_json::to_value(j).unwrap())
1181 .collect();
1182 remember(
1183 &mut meta,
1184 "cron/jobs.json",
1185 jobs_text,
1186 &Value::Array(jobs_record),
1187 );
1188
1189 let exec_path = dir.join("cron/executions.db");
1191 if table_exists(&exec_path, "executions") {
1192 for row in read_rows(
1193 &exec_path,
1194 "select * from executions order by claimed_at, id",
1195 &[],
1196 )?
1197 .unwrap_or_default()
1198 {
1199 profile
1200 .fires
1201 .push(decode_fire_row(&exec_path.display().to_string(), &row)?);
1202 }
1203 }
1204 remember(
1205 &mut meta,
1206 "cron/executions.db",
1207 None,
1208 &serde_json::to_value(&profile.fires).unwrap(),
1209 );
1210
1211 let subs_file = dir.join("webhook_subscriptions.json").display().to_string();
1213 let subs_text = read_text(dir, "webhook_subscriptions.json")?;
1214 if let Some(text) = &subs_text {
1215 let parsed: Value = serde_json::from_str(text)
1216 .map_err(|e| load_error(&subs_file, "", format!("JSON: {e}")))?;
1217 let map = parsed
1218 .as_object()
1219 .ok_or_else(|| load_error(&subs_file, "", "expected a map"))?;
1220 let known: BTreeSet<String> = vault.keys().cloned().collect();
1221 for (n, raw) in map {
1222 profile
1223 .subscriptions
1224 .insert(n.clone(), decode_subscription(&subs_file, n, raw, vault)?);
1225 }
1226 subs_normalized =
1227 flavor == Flavor::Orchestrator && vault.keys().any(|k| !known.contains(k));
1228 }
1229 let subs_record: Vec<Value> = profile
1230 .subscriptions
1231 .values()
1232 .map(|s| serde_json::to_value(s).unwrap())
1233 .collect();
1234 remember(
1235 &mut meta,
1236 "webhook_subscriptions.json",
1237 subs_text,
1238 &Value::Array(subs_record),
1239 );
1240 if subs_normalized {
1241 meta.raw.remove("webhook_subscriptions.json");
1242 meta.snapshot.remove("webhook_subscriptions.json");
1243 }
1244
1245 let own_env = meta
1247 .raw
1248 .get(".env")
1249 .map(|t| parse_dotenv(t))
1250 .unwrap_or_default();
1251 profile.access = hermes_access(dir, &own_env, cfg_map.get("platforms"));
1252
1253 let state_path = dir.join("state.db");
1255 if table_exists(&state_path, "delivery_obligations") {
1256 for row in read_rows(
1257 &state_path,
1258 "select * from delivery_obligations order by created_at, obligation_id",
1259 &[],
1260 )?
1261 .unwrap_or_default()
1262 {
1263 profile.obligations.push(decode_obligation_row(
1264 &state_path.display().to_string(),
1265 &row,
1266 )?);
1267 }
1268 }
1269 if flavor == Flavor::Orchestrator {
1270 if table_exists(&state_path, "gateway_routing") {
1273 let scope = routing_scope(dir);
1274 for row in read_rows(
1275 &state_path,
1276 "select session_key, entry_json from gateway_routing where scope = ? order by updated_at",
1277 &[&scope],
1278 )?
1279 .unwrap_or_default()
1280 {
1281 if let Some((slot, b)) = binding_from_routing(&state_path, &row) {
1282 profile.bindings.insert(slot, b);
1283 }
1284 }
1285 }
1286 } else if table_exists(&state_path, "sessions") {
1287 for row in read_rows(
1288 &state_path,
1289 "select * from sessions order by started_at, id",
1290 &[],
1291 )?
1292 .unwrap_or_default()
1293 {
1294 if let Some(b) = binding_from_hermes_session(&state_path, &row, name) {
1295 profile.bindings.insert(surface_key_string(&b.key), b);
1296 }
1297 }
1298 }
1299 remember(&mut meta, "state.db", None, &state_record(&profile));
1300
1301 profile.residue.files = list_unmodeled(dir, flavor)?;
1303 Ok((profile, meta))
1304}
1305
1306fn list_unmodeled(dir: &Path, flavor: Flavor) -> Result<Vec<String>> {
1307 let mut owned: Vec<&str> = OWNED_FILES.to_vec();
1308 if flavor == Flavor::Hermes {
1309 owned.push("SOUL.md");
1310 owned.retain(|f| !["AGENTS.md", "CLAUDE.md"].contains(f));
1311 }
1312 let runtime_artifacts = ["orchestrator.lock", "orchestrator.sock", "service"];
1313 let mut out = Vec::new();
1314 fn walk(
1315 base: &Path,
1316 d: &Path,
1317 owned: &[&str],
1318 runtime: &[&str],
1319 out: &mut Vec<String>,
1320 ) -> Result<()> {
1321 let entries = match fs::read_dir(d) {
1322 Ok(entries) => entries,
1323 Err(error) if d != base && error.kind() == std::io::ErrorKind::NotFound => {
1325 return Ok(());
1326 }
1327 Err(error) => return Err(error.into()),
1328 };
1329 let mut entries = entries.collect::<std::io::Result<Vec<_>>>()?;
1330 entries.sort_by_key(|e| e.file_name());
1331 for entry in entries {
1332 let p = entry.path();
1333 let rel = p
1334 .strip_prefix(base)
1335 .unwrap_or(&p)
1336 .to_string_lossy()
1337 .replace('\\', "/");
1338 let name = entry.file_name().to_string_lossy().into_owned();
1339 if rel == "profiles"
1340 || name == "node_modules"
1341 || name == ".git"
1342 || rel.starts_with("state.db")
1343 || rel.starts_with("cron/executions.db")
1344 {
1345 continue;
1346 }
1347 if runtime.contains(&rel.as_str()) || regex_tmp(&name) {
1348 continue;
1349 }
1350 let st = match fs::symlink_metadata(&p) {
1351 Ok(metadata) => metadata,
1352 Err(error) if error.kind() == std::io::ErrorKind::NotFound => continue,
1354 Err(error) => return Err(error.into()),
1355 };
1356 if st.is_dir() {
1357 walk(base, &p, owned, runtime, out)?;
1358 continue;
1359 }
1360 if !st.is_file() {
1361 continue;
1362 }
1363 if owned.contains(&rel.as_str()) {
1364 continue;
1365 }
1366 out.push(rel);
1367 }
1368 Ok(())
1369 }
1370 walk(dir, dir, &owned, &runtime_artifacts, &mut out)?;
1371 Ok(out)
1372}
1373
1374fn regex_tmp(name: &str) -> bool {
1375 name.rsplit_once(".tmp-")
1377 .is_some_and(|(_, pid)| !pid.is_empty() && pid.chars().all(|c| c.is_ascii_digit()))
1378}
1379
1380fn write_atomic(path: &Path, text: &str) -> Result<()> {
1383 if let Some(parent) = path.parent() {
1384 fs::create_dir_all(parent)?;
1385 }
1386 let tmp = path.with_file_name(format!(
1387 "{}.tmp-{}",
1388 path.file_name().unwrap().to_string_lossy(),
1389 std::process::id()
1390 ));
1391 write_with_mode(&tmp, text, target_mode(path))?;
1392 fs::rename(&tmp, path)?;
1393 Ok(())
1394}
1395
1396#[cfg(unix)]
1400fn target_mode(path: &Path) -> Option<u32> {
1401 use std::os::unix::fs::PermissionsExt;
1402 let current = fs::metadata(path)
1403 .ok()
1404 .map(|m| m.permissions().mode() & 0o777);
1405 if path.file_name().is_some_and(|name| name == ".env") {
1406 return Some(current.unwrap_or(0o600) & 0o600);
1407 }
1408 current
1409}
1410
1411#[cfg(not(unix))]
1412fn target_mode(_path: &Path) -> Option<u32> {
1413 None
1414}
1415
1416#[cfg(unix)]
1417fn write_with_mode(path: &Path, text: &str, mode: Option<u32>) -> Result<()> {
1418 use std::io::Write;
1419 use std::os::unix::fs::{OpenOptionsExt, PermissionsExt};
1420 let mut options = fs::OpenOptions::new();
1421 options.write(true).create(true).truncate(true);
1422 if let Some(mode) = mode {
1423 options.mode(mode);
1424 }
1425 let mut file = options.open(path)?;
1426 if let Some(mode) = mode {
1428 file.set_permissions(fs::Permissions::from_mode(mode))?;
1429 }
1430 file.write_all(text.as_bytes())?;
1431 Ok(())
1432}
1433
1434#[cfg(not(unix))]
1435fn write_with_mode(path: &Path, text: &str, _mode: Option<u32>) -> Result<()> {
1436 fs::write(path, text)?;
1437 Ok(())
1438}
1439
1440pub fn encode_config(
1442 profile: &Profile,
1443 meta: Option<&ProfileIo>,
1444 vault: Option<&BTreeMap<String, String>>,
1445 flavor: Flavor,
1446 agents: &BTreeMap<String, crate::orchestration::AgentDecl>,
1447) -> String {
1448 let mut out = Map::new();
1449 for (k, v) in &profile.residue.config {
1450 if k != "gateway" {
1451 out.insert(k.clone(), v.clone());
1452 }
1453 }
1454 if let Some(w) = &profile.worker {
1455 let mut wm = Map::new();
1456 wm.insert("harness".into(), Value::String(w.harness.as_str().into()));
1457 if let Some(m) = &w.model {
1458 wm.insert("model".into(), Value::String(m.clone()));
1459 }
1460 if let Some(p) = &w.preset {
1461 wm.insert("preset".into(), Value::String(p.clone()));
1462 }
1463 if w.cwd != "." {
1464 wm.insert("cwd".into(), Value::String(w.cwd.clone()));
1465 }
1466 if w.home == crate::orchestration::WorkerHome::User {
1467 wm.insert("home".into(), Value::String("user".into()));
1468 }
1469 if w.surface == crate::orchestration::WorkerSurface::Pane {
1470 wm.insert("surface".into(), Value::String("pane".into()));
1471 }
1472 if let Some(c) = w.capacity {
1473 wm.insert("capacity".into(), Value::from(c));
1474 }
1475 if !w.env.is_empty() {
1476 let env: Map<String, Value> = w
1477 .env
1478 .iter()
1479 .map(|(k, v)| {
1480 let value = match v {
1481 crate::orchestration::EnvValue::Literal(s) => Value::String(s.clone()),
1482 crate::orchestration::EnvValue::Secret(r) => super::decode::render_ref(r),
1483 };
1484 (k.clone(), value)
1485 })
1486 .collect();
1487 wm.insert("env".into(), Value::Object(env));
1488 }
1489 if w.permission.timeout_seconds != 300
1490 || w.permission.default != crate::orchestration::PermissionDefault::Deny
1491 || w.permission.unattended != crate::orchestration::PermissionUnattended::Deny
1492 {
1493 wm.insert(
1494 "permission".into(),
1495 serde_json::to_value(&w.permission).unwrap(),
1496 );
1497 }
1498 out.insert("worker".into(), Value::Object(wm));
1499 }
1500 if flavor == Flavor::Orchestrator && !profile.budgets.is_empty() {
1501 out.insert(
1502 "budgets".into(),
1503 serde_json::to_value(&profile.budgets).unwrap(),
1504 );
1505 }
1506 let mut gateway = profile
1507 .residue
1508 .config
1509 .get("gateway")
1510 .and_then(Value::as_object)
1511 .cloned()
1512 .unwrap_or_default();
1513 let routes: Vec<Value> = profile
1514 .routes
1515 .iter()
1516 .map(|r| encode_route(r, agents, flavor == Flavor::Hermes))
1517 .collect();
1518 if meta.is_some_and(|m| m.routes_at_top) {
1519 if !routes.is_empty() {
1520 out.insert("profile_routes".into(), Value::Array(routes));
1521 }
1522 } else if !routes.is_empty() {
1523 gateway.insert("profile_routes".into(), Value::Array(routes));
1524 }
1525 if !gateway.is_empty() {
1526 out.insert("gateway".into(), Value::Object(gateway));
1527 }
1528 let mut platforms = Map::new();
1529 for (p, ch) in &profile.channels {
1530 if ch.extra.get("from_env") == Some(&Value::Bool(true)) {
1532 continue;
1533 }
1534 platforms.insert(
1535 p.clone(),
1536 encode_channel(
1537 ch,
1538 if flavor == Flavor::Hermes {
1539 vault
1540 } else {
1541 None
1542 },
1543 ),
1544 );
1545 }
1546 if !platforms.is_empty() {
1547 out.insert("platforms".into(), Value::Object(platforms));
1548 }
1549 json_to_yaml(&Value::Object(out))
1550}
1551
1552pub fn encode_jobs_file(profile: &Profile, meta: Option<&ProfileIo>) -> String {
1554 let jobs: Vec<Value> = profile
1555 .jobs
1556 .values()
1557 .map(|j| ordered_object(encode_job(j)))
1558 .collect();
1559 let form = meta.and_then(|m| m.jobs_form.clone()).unwrap_or(JobsForm {
1560 object: true,
1561 extras: Map::new(),
1562 });
1563 let body = if form.object {
1564 let mut pairs = vec![("jobs".to_string(), Value::Array(jobs))];
1565 pairs.extend(form.extras.iter().map(|(k, v)| (k.clone(), v.clone())));
1566 ordered_object(pairs)
1567 } else {
1568 Value::Array(jobs)
1569 };
1570 format!("{}\n", pretty_ordered(&body, 0))
1571}
1572
1573pub(crate) fn ordered_object(pairs: Vec<(String, Value)>) -> Value {
1576 Value::Array(vec![
1579 Value::String("__ordered__".into()),
1580 Value::Array(
1581 pairs
1582 .into_iter()
1583 .map(|(k, v)| serde_json::json!({"__k": k, "__v": v}))
1584 .collect(),
1585 ),
1586 ])
1587}
1588
1589fn is_ordered(value: &Value) -> Option<&Vec<Value>> {
1590 let arr = value.as_array()?;
1591 if arr.len() == 2 && arr[0].as_str() == Some("__ordered__") {
1592 arr[1].as_array()
1593 } else {
1594 None
1595 }
1596}
1597
1598pub(crate) fn pretty_ordered(value: &Value, depth: usize) -> String {
1600 let pad = |d: usize| " ".repeat(d);
1601 if let Some(pairs) = is_ordered(value) {
1602 if pairs.is_empty() {
1603 return "{}".into();
1604 }
1605 let inner: Vec<String> = pairs
1606 .iter()
1607 .map(|p| {
1608 format!(
1609 "{}{}: {}",
1610 pad(depth + 1),
1611 serde_json::to_string(p["__k"].as_str().unwrap_or("")).unwrap(),
1612 pretty_ordered(&p["__v"], depth + 1)
1613 )
1614 })
1615 .collect();
1616 return format!("{{\n{}\n{}}}", inner.join(",\n"), pad(depth));
1617 }
1618 match value {
1619 Value::Array(items) if items.is_empty() => "[]".into(),
1620 Value::Array(items) => {
1621 let inner: Vec<String> = items
1622 .iter()
1623 .map(|v| format!("{}{}", pad(depth + 1), pretty_ordered(v, depth + 1)))
1624 .collect();
1625 format!("[\n{}\n{}]", inner.join(",\n"), pad(depth))
1626 }
1627 Value::Object(o) if o.is_empty() => "{}".into(),
1628 Value::Object(o) => {
1629 let inner: Vec<String> = o
1630 .iter()
1631 .map(|(k, v)| {
1632 format!(
1633 "{}{}: {}",
1634 pad(depth + 1),
1635 serde_json::to_string(k).unwrap(),
1636 pretty_ordered(v, depth + 1)
1637 )
1638 })
1639 .collect();
1640 format!("{{\n{}\n{}}}", inner.join(",\n"), pad(depth))
1641 }
1642 Value::Number(n) => {
1643 if let Some(f) = n.as_f64() {
1644 if n.is_f64() && f.fract() == 0.0 && f.abs() < 1e21 {
1645 return format!("{}", f as i64);
1646 }
1647 }
1648 n.to_string()
1649 }
1650 other => serde_json::to_string(other).unwrap(),
1651 }
1652}
1653
1654pub fn encode_subscriptions_file(
1656 profile: &Profile,
1657 vault: Option<&BTreeMap<String, String>>,
1658) -> String {
1659 let mut out = Map::new();
1660 for (n, s) in &profile.subscriptions {
1661 out.insert(n.clone(), encode_subscription(s, vault));
1662 }
1663 format!("{}\n", pretty_ordered(&Value::Object(out), 0))
1664}
1665
1666fn hermes_access(dir: &Path, env: &BTreeMap<String, String>, platforms: Option<&Value>) -> Access {
1672 let mut access = Access::default();
1673 let truthy = |v: &str| {
1674 matches!(
1675 v.trim().to_ascii_lowercase().as_str(),
1676 "1" | "true" | "yes" | "on"
1677 )
1678 };
1679 let split = |v: &str| {
1680 v.split(',')
1681 .map(|s| s.trim().to_string())
1682 .filter(|s| !s.is_empty())
1683 .collect::<Vec<_>>()
1684 };
1685 let global_open = env
1686 .get("GATEWAY_ALLOW_ALL_USERS")
1687 .is_some_and(|v| truthy(v));
1688 for (k, v) in env {
1689 if let Some(p) = k.strip_suffix("_ALLOWED_USERS") {
1690 access
1691 .allowlist
1692 .entry(p.to_ascii_lowercase())
1693 .or_default()
1694 .extend(split(v));
1695 } else if let Some(p) = k.strip_suffix("_ALLOW_ALL_USERS") {
1696 if truthy(v) {
1697 access
1698 .policy
1699 .insert(p.to_ascii_lowercase(), AccessPolicy::Open);
1700 }
1701 }
1702 }
1703 for sub in ["platforms/pairing", "pairing"] {
1704 let Ok(entries) = fs::read_dir(dir.join(sub)) else {
1705 continue;
1706 };
1707 for entry in entries.flatten() {
1708 let name = entry.file_name().to_string_lossy().into_owned();
1709 let Some(platform) = name.strip_suffix("-approved.json") else {
1710 continue;
1711 };
1712 let Ok(text) = fs::read_to_string(entry.path()) else {
1713 continue;
1714 };
1715 if let Ok(Value::Object(map)) = serde_json::from_str::<Value>(&text) {
1716 access
1717 .allowlist
1718 .entry(platform.to_string())
1719 .or_default()
1720 .extend(map.keys().cloned());
1721 }
1722 }
1723 }
1724 if let Some(Value::Object(map)) = platforms {
1725 for (platform, block) in map {
1726 if global_open {
1727 access.policy.insert(platform.clone(), AccessPolicy::Open);
1728 }
1729 for key in ["allow_admin_from", "group_allow_admin_from"] {
1730 let found = block
1731 .get("extra")
1732 .and_then(|e| e.get(key))
1733 .or_else(|| block.get(key));
1734 let ids = match found {
1735 Some(Value::String(s)) => split(s),
1736 Some(Value::Array(items)) => items
1737 .iter()
1738 .filter_map(|v| {
1739 v.as_str()
1740 .map(str::to_string)
1741 .or_else(|| v.as_i64().map(|n| n.to_string()))
1742 })
1743 .collect(),
1744 _ => Vec::new(),
1745 };
1746 access
1747 .admins
1748 .entry(platform.clone())
1749 .or_default()
1750 .extend(ids);
1751 }
1752 }
1753 }
1754 for list in access
1755 .allowlist
1756 .values_mut()
1757 .chain(access.admins.values_mut())
1758 {
1759 list.sort();
1760 list.dedup();
1761 }
1762 access.allowlist.retain(|_, v| !v.is_empty());
1763 access.admins.retain(|_, v| !v.is_empty());
1764 access
1765}
1766
1767const EXECUTIONS_DDL: &str = "CREATE TABLE IF NOT EXISTS executions (
1768 id TEXT PRIMARY KEY, job_id TEXT NOT NULL, source TEXT NOT NULL, process_id TEXT NOT NULL, pid INTEGER NOT NULL,
1769 process_started_at INTEGER, status TEXT NOT NULL CHECK(status IN ('claimed','running','completed','failed','unknown')),
1770 claimed_at TEXT NOT NULL, started_at TEXT, finished_at TEXT, error TEXT);
1771CREATE INDEX IF NOT EXISTS idx_executions_job_claimed ON executions(job_id, claimed_at DESC, id DESC);
1772CREATE INDEX IF NOT EXISTS idx_executions_status_claimed ON executions(status, claimed_at DESC, id DESC);";
1773const FOLDER_EXECUTIONS_DDL: &str = "CREATE TABLE IF NOT EXISTS executions (
1778 id TEXT PRIMARY KEY, job_id TEXT NOT NULL, source TEXT NOT NULL, process_id TEXT NOT NULL, pid INTEGER NOT NULL,
1779 process_started_at INTEGER, status TEXT NOT NULL CHECK(status IN ('claimed','running','completed','failed','unknown')),
1780 claimed_at TEXT NOT NULL, started_at TEXT, finished_at TEXT, error TEXT, residue_json TEXT,
1781 session_id TEXT, obligation_id TEXT);
1782CREATE INDEX IF NOT EXISTS idx_executions_job_claimed ON executions(job_id, claimed_at DESC, id DESC);
1783CREATE INDEX IF NOT EXISTS idx_executions_status_claimed ON executions(status, claimed_at DESC, id DESC);";
1784
1785const OBLIGATIONS_DDL: &str = "CREATE TABLE IF NOT EXISTS delivery_obligations (
1786 obligation_id TEXT PRIMARY KEY, session_key TEXT NOT NULL, platform TEXT NOT NULL, chat_id TEXT NOT NULL, thread_id TEXT,
1787 content TEXT NOT NULL, state TEXT NOT NULL, attempts INTEGER NOT NULL DEFAULT 0, created_at REAL NOT NULL, updated_at REAL NOT NULL,
1788 owner_pid INTEGER, owner_started_at INTEGER, last_error TEXT, adapter_profile TEXT,
1789 posted_message_id TEXT, source_json TEXT);";
1790
1791pub fn write_executions(path: &Path, fires: &[crate::orchestration::Fire]) -> Result<()> {
1793 write_executions_shaped(path, fires, false)
1794}
1795
1796pub fn write_executions_shaped(
1799 path: &Path,
1800 fires: &[crate::orchestration::Fire],
1801 ours: bool,
1802) -> Result<()> {
1803 let mut cols: Vec<&str> = EXECUTION_COLUMNS.to_vec();
1804 if ours {
1805 cols.extend(["residue_json", "session_id", "obligation_id"]);
1808 }
1809 let insert = format!(
1810 "insert into executions ({}) values ({})",
1811 cols.join(", "),
1812 cols.iter().map(|_| "?").collect::<Vec<_>>().join(",")
1813 );
1814 let rows: Vec<Vec<Param>> = fires
1815 .iter()
1816 .map(|f| {
1817 let mut row: Vec<Param> = encode_fire_row(f).iter().map(Param::from).collect();
1818 if ours {
1819 let rest: serde_json::Map<String, Value> = f
1821 .residue
1822 .0
1823 .iter()
1824 .filter(|(k, _)| {
1825 !["source", "process_id", "pid", "process_started_at"].contains(&k.as_str())
1826 })
1827 .map(|(k, v)| (k.clone(), v.clone()))
1828 .collect();
1829 row.push(if rest.is_empty() {
1830 Param::Null
1831 } else {
1832 Param::Text(serde_json::to_string(&rest).unwrap())
1833 });
1834 let opt = |v: &Option<String>| v.clone().map(Param::Text).unwrap_or(Param::Null);
1835 row.push(opt(&f.session_id));
1836 row.push(opt(&f.obligation_id));
1837 }
1838 row
1839 })
1840 .collect();
1841 write_table(
1842 path,
1843 if ours {
1844 FOLDER_EXECUTIONS_DDL
1845 } else {
1846 EXECUTIONS_DDL
1847 },
1848 &insert,
1849 &rows,
1850 )
1851}
1852
1853const ROUTING_DDL: &str = "CREATE TABLE IF NOT EXISTS gateway_routing (
1855 scope TEXT NOT NULL DEFAULT '', session_key TEXT NOT NULL, entry_json TEXT NOT NULL, updated_at REAL NOT NULL,
1856 PRIMARY KEY (scope, session_key));";
1857
1858fn write_state(dir: &Path, path: &Path, profile: &Profile) -> Result<()> {
1861 fs::create_dir_all(dir.join("sessions"))?;
1862 let scope = routing_scope(dir);
1863 let now = std::time::SystemTime::now()
1864 .duration_since(std::time::UNIX_EPOCH)
1865 .map(|d| d.as_secs_f64())
1866 .unwrap_or(0.0);
1867 let rank = |b: &Binding| {
1870 (
1871 b.ended_at.is_none(),
1872 b.last_activity_at
1873 .clone()
1874 .or_else(|| b.started_at.clone())
1875 .unwrap_or_default(),
1876 )
1877 };
1878 let mut chosen: BTreeMap<String, (&Binding, Value)> = BTreeMap::new();
1879 for (slot, b) in &profile.bindings {
1880 let (key, entry) = routing_entry(&profile.name, slot, b);
1881 if chosen
1882 .get(&key)
1883 .is_none_or(|(held, _)| rank(b) > rank(held))
1884 {
1885 chosen.insert(key, (b, entry));
1886 }
1887 }
1888 let mut mirror = Map::new();
1889 let rows: Vec<Vec<Param>> = chosen
1890 .into_iter()
1891 .map(|(key, (_, entry))| {
1892 mirror.insert(key.clone(), entry.clone());
1893 vec![
1894 Param::Text(scope.clone()),
1895 Param::Text(key),
1896 Param::Text(serde_json::to_string(&entry).unwrap()),
1897 Param::Real(now),
1898 ]
1899 })
1900 .collect();
1901 replace_rows(
1902 path,
1903 ROUTING_DDL,
1904 "gateway_routing",
1905 &[],
1906 ("delete from gateway_routing where scope = ?", &[Param::Text(scope)]),
1907 "insert into gateway_routing (scope, session_key, entry_json, updated_at) values (?, ?, ?, ?)",
1908 &rows,
1909 )?;
1910 let mirror_file = dir.join("sessions/sessions.json");
1911 let tmp = mirror_file.with_file_name(format!("sessions.json.tmp-{}", std::process::id()));
1912 fs::write(
1913 &tmp,
1914 serde_json::to_string_pretty(&Value::Object(mirror)).unwrap(),
1915 )?;
1916 fs::rename(&tmp, &mirror_file)?;
1917 let cols: Vec<&str> = OBLIGATION_COLUMNS
1919 .iter()
1920 .chain(FOLDER_OBLIGATION_EXTRA_COLUMNS.iter())
1921 .copied()
1922 .collect();
1923 let extras: Vec<(&str, &str)> = FOLDER_OBLIGATION_EXTRA_COLUMNS
1924 .iter()
1925 .map(|c| (*c, "TEXT"))
1926 .collect();
1927 let insert = format!(
1928 "insert into delivery_obligations ({}) values ({})",
1929 cols.join(", "),
1930 cols.iter().map(|_| "?").collect::<Vec<_>>().join(",")
1931 );
1932 let rows: Vec<Vec<Param>> = profile
1933 .obligations
1934 .iter()
1935 .map(|o| {
1936 encode_obligation_row(o)
1937 .iter()
1938 .chain(encode_obligation_folder_extras(o).iter())
1939 .map(Param::from)
1940 .collect()
1941 })
1942 .collect();
1943 replace_rows(
1944 path,
1945 OBLIGATIONS_DDL,
1946 "delivery_obligations",
1947 &extras,
1948 ("delete from delivery_obligations", &[]),
1949 &insert,
1950 &rows,
1951 )
1952}
1953
1954fn write_if_changed(
1955 meta: &mut ProfileIo,
1956 dir: &Path,
1957 rel: &str,
1958 record: &Value,
1959 render: impl FnOnce() -> Option<String>,
1960) -> Result<bool> {
1961 let snap = canonical_json(record);
1962 let path = dir.join(rel);
1963 let reuse = meta.flavor == Flavor::Orchestrator && meta.snapshot.get(rel) == Some(&snap);
1966 if reuse && path.exists() {
1967 return Ok(false);
1968 }
1969 if reuse {
1970 if let Some(raw) = meta.raw.get(rel).cloned() {
1971 write_atomic(&path, &raw)?;
1972 return Ok(true);
1973 }
1974 }
1975 let Some(text) = render() else {
1976 return Ok(false);
1977 };
1978 write_atomic(&path, &text)?;
1979 meta.raw.insert(rel.into(), text);
1980 meta.snapshot.insert(rel.into(), snap);
1981 meta.flavor = Flavor::Orchestrator; Ok(true)
1983}
1984
1985pub fn save_home(loaded: &mut LoadedHome, root: Option<&Path>) -> Result<()> {
1987 let root = root
1988 .map(Path::to_path_buf)
1989 .unwrap_or_else(|| loaded.orchestration.root.clone());
1990 fs::create_dir_all(&root)?;
1991 let names: Vec<String> = loaded.orchestration.profiles.keys().cloned().collect();
1992 for name in names {
1993 let dir = if name == "default" {
1994 root.clone()
1995 } else {
1996 root.join("profiles").join(&name)
1997 };
1998 fs::create_dir_all(dir.join("cron"))?;
1999 let profile = loaded.orchestration.profiles[&name].clone();
2000 let meta = loaded
2001 .io
2002 .entry(name.clone())
2003 .or_insert_with(|| ProfileIo::new(Flavor::Orchestrator));
2004 save_profile_dir(
2005 &profile,
2006 meta,
2007 &dir,
2008 &loaded.vault,
2009 &loaded.orchestration.agents,
2010 )?;
2011 }
2012 let workflow = &loaded.orchestration.workflow;
2014 let meta = loaded
2015 .io
2016 .entry("default".into())
2017 .or_insert_with(|| ProfileIo::new(Flavor::Orchestrator));
2018 if meta.raw.contains_key("workflow.yaml")
2019 || *workflow != crate::orchestration::hermes_instance()
2020 {
2021 let record = serde_json::to_value(workflow).unwrap();
2022 write_if_changed(meta, &root, "workflow.yaml", &record, || {
2023 Some(json_to_yaml(&record))
2024 })?;
2025 }
2026 let layer = agent_layer_record(&loaded.orchestration);
2028 let meta = loaded
2029 .io
2030 .entry("default".into())
2031 .or_insert_with(|| ProfileIo::new(Flavor::Orchestrator));
2032 if meta.raw.contains_key("agents.json") || layer.as_object().is_some_and(|m| !m.is_empty()) {
2033 write_if_changed(meta, &root, "agents.json", &layer, || {
2034 Some(serde_json::to_string_pretty(&layer).unwrap() + "\n")
2035 })?;
2036 }
2037 Ok(())
2038}
2039
2040fn save_profile_dir(
2041 profile: &Profile,
2042 meta: &mut ProfileIo,
2043 dir: &Path,
2044 vault: &BTreeMap<String, String>,
2045 agents: &BTreeMap<String, crate::orchestration::AgentDecl>,
2046) -> Result<()> {
2047 let cfg_record = config_record(profile);
2048 write_if_changed(meta, dir, "config.yaml", &cfg_record, || {
2049 Some(encode_config(
2050 profile,
2051 None,
2052 Some(vault),
2053 Flavor::Orchestrator,
2054 agents,
2055 ))
2056 })?;
2057 if let Some(persona) = &profile.persona {
2058 if !dir.join("AGENTS.md").exists() {
2062 let text = persona.text.clone().unwrap_or_default();
2063 write_if_changed(
2064 meta,
2065 dir,
2066 "AGENTS.md",
2067 &serde_json::to_value(&profile.persona).unwrap(),
2068 || Some(text),
2069 )?;
2070 }
2071 if !dir.join("CLAUDE.md").exists() {
2072 write_atomic(&dir.join("CLAUDE.md"), "@AGENTS.md\n")?;
2073 }
2074 }
2075 let jobs_record: Vec<Value> = profile
2076 .jobs
2077 .values()
2078 .map(|j| serde_json::to_value(j).unwrap())
2079 .collect();
2080 let had_jobs = meta.raw.contains_key("cron/jobs.json");
2081 let form = meta.jobs_form.clone();
2082 write_if_changed(
2083 meta,
2084 dir,
2085 "cron/jobs.json",
2086 &Value::Array(jobs_record),
2087 || {
2088 if profile.jobs.is_empty() && !had_jobs {
2089 None
2090 } else {
2091 let stub = ProfileIo {
2092 jobs_form: form.clone(),
2093 ..ProfileIo::new(Flavor::Orchestrator)
2094 };
2095 Some(encode_jobs_file(profile, Some(&stub)))
2096 }
2097 },
2098 )?;
2099 let subs_record: Vec<Value> = profile
2100 .subscriptions
2101 .values()
2102 .map(|s| serde_json::to_value(s).unwrap())
2103 .collect();
2104 let had_subs = meta.raw.contains_key("webhook_subscriptions.json");
2105 write_if_changed(
2106 meta,
2107 dir,
2108 "webhook_subscriptions.json",
2109 &Value::Array(subs_record),
2110 || {
2111 if profile.subscriptions.is_empty() && !had_subs {
2112 None
2113 } else {
2114 Some(encode_subscriptions_file(profile, None))
2115 }
2116 },
2117 )?;
2118 let fires_snap = canonical_json(&serde_json::to_value(&profile.fires).unwrap());
2120 let exec_path = dir.join("cron/executions.db");
2121 if (meta.snapshot.get("cron/executions.db") != Some(&fires_snap) || !exec_path.exists())
2122 && (!profile.fires.is_empty() || exec_path.exists())
2123 {
2124 let tmp = exec_path.with_file_name(format!("executions.db.tmp-{}", std::process::id()));
2125 let _ = fs::remove_file(&tmp);
2126 write_executions_shaped(&tmp, &profile.fires, true)?;
2127 fs::rename(&tmp, &exec_path)?;
2128 meta.snapshot
2129 .insert("cron/executions.db".into(), fires_snap);
2130 }
2131 let state_snap = canonical_json(&state_record(profile));
2132 let state_path = dir.join("state.db");
2133 if (meta.snapshot.get("state.db") != Some(&state_snap) || !state_path.exists())
2134 && (!profile.bindings.is_empty() || !profile.obligations.is_empty() || state_path.exists())
2135 {
2136 write_state(dir, &state_path, profile)?;
2138 meta.snapshot.insert("state.db".into(), state_snap);
2139 }
2140
2141 let mut refs: Vec<String> = Vec::new();
2143 for ch in profile.channels.values() {
2144 for r in ch.credentials.values() {
2145 if let crate::ontology::SecretRef::Dotenv(n) = r {
2146 refs.push(n.clone());
2147 }
2148 }
2149 }
2150 for s in profile.subscriptions.values() {
2151 if let Some(crate::ontology::SecretRef::Dotenv(n)) = &s.secret {
2152 refs.push(n.clone());
2153 }
2154 }
2155 if let Some(w) = &profile.worker {
2156 for v in w.env.values() {
2157 if let crate::orchestration::EnvValue::Secret(crate::ontology::SecretRef::Dotenv(n)) = v
2158 {
2159 refs.push(n.clone());
2160 }
2161 }
2162 }
2163 let mut entries: BTreeMap<String, String> = BTreeMap::new();
2164 for r in refs {
2165 if let Some(v) = vault.get(&r) {
2166 entries.insert(r, v.clone());
2167 }
2168 }
2169 if !entries.is_empty() {
2170 let existing = meta.raw.get(".env").cloned();
2171 let mut merged = existing.as_deref().map(parse_dotenv).unwrap_or_default();
2172 for (k, v) in entries {
2173 merged.insert(k, v);
2174 }
2175 let text = render_dotenv(&merged);
2176 if existing.as_deref() != Some(text.as_str()) {
2177 write_atomic(&dir.join(".env"), &text)?;
2178 meta.raw.insert(".env".into(), text);
2179 }
2180 }
2181 Ok(())
2182}
2183
2184pub fn copy_unmodeled(profile: &Profile, meta: &ProfileIo, dest: &Path) -> Result<()> {
2186 let Some(src) = &meta.source_dir else {
2187 return Ok(());
2188 };
2189 carry_unmodeled(&profile.residue.files, src, dest)?;
2190 Ok(())
2191}
2192
2193pub fn carry_unmodeled(files: &[String], src: &Path, dest: &Path) -> Result<Vec<String>> {
2197 let mut carried = Vec::new();
2198 for rel in files {
2199 let from = src.join(rel);
2200 if !from.is_file() {
2201 continue;
2202 }
2203 let to = dest.join(rel);
2204 if let Some(parent) = to.parent() {
2205 fs::create_dir_all(parent)?;
2206 }
2207 fs::copy(&from, &to)?;
2208 carried.push(rel.clone());
2209 }
2210 Ok(carried)
2211}