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