1use std::fs;
17use std::path::Path;
18
19use serde_json::Value;
20
21use super::canonical::canonical_json;
22use super::decode::{encode_obligation_row, OBLIGATION_COLUMNS};
23use super::folder::{
24 config_record, copy_unmodeled, encode_config, encode_jobs_file, encode_subscriptions_file,
25 write_executions, Flavor, LoadedHome, ProfileIo,
26};
27use super::openclaw::ms_from_iso;
28use super::sqlite::{write_table, Param};
29use crate::ontology::{
30 hermes_source_for_binding, render_hermes_session_key, ArtifactFidelity, Binding, Fidelity,
31};
32use crate::orchestration::Profile;
33
34const HERMES_STATE_V22: &str = include_str!("hermes_state_v22.sql");
40const HERMES_STATE_V22_VERSION: i64 = 22;
41
42const SESSION_COLUMNS: &[&str] = &[
45 "id",
46 "source",
47 "user_id",
48 "session_key",
49 "chat_id",
50 "chat_type",
51 "thread_id",
52 "expiry_finalized",
53 "started_at",
54 "ended_at",
55 "end_reason",
56 "handoff_state",
57 "handoff_platform",
58 "handoff_error",
59 "profile_name",
60];
61
62fn epoch_seconds(iso: Option<&str>) -> Option<f64> {
63 ms_from_iso(iso).map(|ms| ms as f64 / 1000.0)
64}
65
66fn hermes_session_row(profile: &str, slot: &str, b: &Binding) -> Vec<Param> {
68 let text = |v: &Option<String>| v.clone().map(Param::Text).unwrap_or(Param::Null);
69 let hermes_profile = if profile == "default" {
70 "main"
71 } else {
72 profile
73 };
74 let started = epoch_seconds(b.started_at.as_deref())
75 .or_else(|| epoch_seconds(b.last_activity_at.as_deref()))
76 .unwrap_or(0.0);
77 vec![
78 Param::Text(
79 b.worker
80 .session_id
81 .clone()
82 .unwrap_or_else(|| slot.to_string()),
83 ),
84 Param::Text(hermes_source_for_binding(b)),
85 text(&b.key.participant_id),
86 Param::Text(
87 b.key
88 .key
89 .clone()
90 .unwrap_or_else(|| render_hermes_session_key(hermes_profile, &b.key)),
91 ),
92 text(&b.key.chat_id),
93 text(&b.key.kind),
94 text(&b.key.thread_id),
95 Param::Int(i64::from(b.ended_at.is_some())),
96 Param::Real(started),
97 epoch_seconds(b.ended_at.as_deref())
98 .map(Param::Real)
99 .unwrap_or(Param::Null),
100 b.end_reason
101 .map(|r| Param::Text(r.as_str().into()))
102 .unwrap_or(Param::Null),
103 text(&b.handoff.as_ref().map(|h| h.state.clone())),
104 text(&b.handoff.as_ref().and_then(|h| h.to.clone())),
105 text(&b.handoff.as_ref().and_then(|h| h.error.clone())),
106 if profile == "default" {
107 Param::Null
108 } else {
109 Param::Text(profile.into())
110 },
111 ]
112}
113
114fn with_handoffs(profile: &Profile, o: &crate::orchestration::Orchestration) -> Profile {
117 let mut out = profile.clone();
118 for b in out.bindings.values_mut() {
119 let Some(id) = b.worker.session_id.clone() else {
120 continue;
121 };
122 let session = format!("{}:{}", b.worker.harness.as_str(), id);
123 if let Some(record) = o.handoffs.iter().rev().find(|h| h.session == session) {
124 b.handoff = Some(crate::ontology::Handoff {
125 to: o
126 .conversations
127 .get(&record.to)
128 .and_then(|c| c.surface.as_ref())
129 .and_then(|s| s.platform.clone()),
130 state: match record.state {
131 crate::orchestration::HandoffState::Pending => "pending",
132 crate::orchestration::HandoffState::Done => "done",
133 crate::orchestration::HandoffState::Failed => "failed",
134 }
135 .into(),
136 error: record.error.clone(),
137 });
138 }
139 }
140 out
141}
142
143fn write_fresh_store(profiles: &[(&str, &Profile)], path: &Path) -> Result<()> {
148 let placeholders = |n: usize| std::iter::repeat_n("?", n).collect::<Vec<_>>().join(",");
149 write_table(
150 path,
151 HERMES_STATE_V22,
152 "insert into schema_version (version) values (?)",
153 &[vec![Param::Int(HERMES_STATE_V22_VERSION)]],
154 )?;
155 let sessions: Vec<Vec<Param>> = profiles
156 .iter()
157 .flat_map(|(name, profile)| {
158 profile
159 .bindings
160 .iter()
161 .map(move |(slot, b)| hermes_session_row(name, slot, b))
162 })
163 .collect();
164 write_table(
165 path,
166 "",
167 &format!(
168 "insert into sessions ({}) values ({})",
169 SESSION_COLUMNS.join(", "),
170 placeholders(SESSION_COLUMNS.len())
171 ),
172 &sessions,
173 )?;
174 let obligations: Vec<Vec<Param>> = profiles
175 .iter()
176 .flat_map(|(_, profile)| profile.obligations.iter())
177 .map(|o| encode_obligation_row(o).iter().map(Param::from).collect())
178 .collect();
179 write_table(
180 path,
181 "",
182 &format!(
183 "insert into delivery_obligations ({}) values ({})",
184 OBLIGATION_COLUMNS.join(", "),
185 placeholders(OBLIGATION_COLUMNS.len())
186 ),
187 &obligations,
188 )
189}
190use crate::Result;
191
192#[derive(Debug, Clone, PartialEq, Eq)]
194pub struct Refusal {
195 pub file: String,
197 pub reason: String,
199}
200
201#[derive(Debug, Clone, Default)]
203pub struct HermesReport {
204 pub written: Vec<ArtifactFidelity>,
206 pub refused: Vec<Refusal>,
208 pub loss: super::loss::LossReport,
210}
211
212fn write_atomic(path: &Path, text: &str) -> Result<()> {
213 if let Some(parent) = path.parent() {
214 fs::create_dir_all(parent)?;
215 }
216 let tmp = path.with_file_name(format!(
217 "{}.tmp-{}",
218 path.file_name().unwrap().to_string_lossy(),
219 std::process::id()
220 ));
221 fs::write(&tmp, text)?;
222 fs::rename(&tmp, path)?;
223 Ok(())
224}
225
226pub fn from_hermes(home: &Path) -> Result<LoadedHome> {
228 super::folder::load_home(home, Flavor::Hermes)
229}
230
231pub fn to_hermes(loaded: &LoadedHome, dest: &Path, only: Option<&str>) -> Result<HermesReport> {
235 let mut report = HermesReport::default();
236 report.loss = super::loss::loss_report(
237 "hermes",
238 &loaded.orchestration,
239 super::loss::HERMES_FEATURES,
240 );
241 super::loss::write_loss_report(dest, &report.loss)?;
242 let empty = ProfileIo {
243 raw: Default::default(),
244 snapshot: Default::default(),
245 source_dir: None,
246 flavor: Flavor::Orchestrator,
247 jobs_form: None,
248 routes_at_top: false,
249 borrowed_from: None,
250 lenders: Vec::new(),
251 };
252 let mut fresh_rows: Vec<(&str, &Profile)> = Vec::new();
253 for (name, profile) in &loaded.orchestration.profiles {
254 if only.is_some_and(|o| o != name) {
255 continue;
256 }
257 let dir = if name == "default" {
258 dest.to_path_buf()
259 } else {
260 dest.join("profiles").join(name)
261 };
262 fs::create_dir_all(dir.join("cron"))?;
263 let meta = loaded.io.get(name).unwrap_or(&empty);
264 let rel = |p: &str| {
265 if name == "default" {
266 p.to_string()
267 } else {
268 format!("profiles/{name}/{p}")
269 }
270 };
271 let unchanged =
272 |file: &str, record: &Value| meta.snapshot.get(file) == Some(&canonical_json(record));
273 fn emit_file(
274 written: &mut Vec<ArtifactFidelity>,
275 dir: &Path,
276 path: String,
277 file: &str,
278 text: &str,
279 tier: Fidelity,
280 ) -> Result<()> {
281 write_atomic(&dir.join(file), text)?;
282 written.push(ArtifactFidelity {
283 path,
284 fidelity: tier,
285 loss: Vec::new(),
286 });
287 Ok(())
288 }
289 macro_rules! emit {
290 ($file:expr, $text:expr, $tier:expr) => {
291 emit_file(&mut report.written, &dir, rel($file), $file, $text, $tier)?
292 };
293 }
294
295 let cfg_record = config_record(profile);
297 let cfg_empty = profile.routes.is_empty()
298 && profile.channels.is_empty()
299 && profile.residue.config.is_empty()
300 && profile.worker.is_none();
301 let hermes_bytes = meta.flavor == Flavor::Hermes;
305 if hermes_bytes
306 && unchanged("config.yaml", &cfg_record)
307 && meta.raw.contains_key("config.yaml")
308 {
309 emit!(
310 "config.yaml",
311 &meta.raw["config.yaml"],
312 Fidelity::ByteLossless
313 );
314 } else if !cfg_empty || meta.raw.contains_key("config.yaml") {
315 emit!(
316 "config.yaml",
317 &encode_config(
318 profile,
319 Some(meta),
320 Some(&loaded.vault),
321 Flavor::Hermes,
322 &loaded.orchestration.agents
323 ),
324 Fidelity::Semantic
325 );
326 }
327
328 if let Some(persona) = &profile.persona {
330 let src_name = if meta.raw.contains_key("SOUL.md") {
331 "SOUL.md"
332 } else {
333 "AGENTS.md"
334 };
335 let record = serde_json::to_value(&profile.persona).unwrap();
336 if unchanged(src_name, &record) && meta.raw.contains_key(src_name) {
337 emit!(
338 "SOUL.md",
339 &meta.raw[src_name],
340 if src_name == "SOUL.md" {
341 Fidelity::ByteLossless
342 } else {
343 Fidelity::Semantic
344 }
345 );
346 } else {
347 emit!(
348 "SOUL.md",
349 persona.text.as_deref().unwrap_or(""),
350 Fidelity::Semantic
351 );
352 }
353 }
354
355 let jobs_record: Vec<Value> = profile
358 .jobs
359 .values()
360 .map(|j| serde_json::to_value(j).unwrap())
361 .collect();
362 if !profile.jobs.is_empty() || meta.raw.contains_key("cron/jobs.json") {
363 if unchanged("cron/jobs.json", &Value::Array(jobs_record))
364 && meta.raw.contains_key("cron/jobs.json")
365 {
366 emit!(
367 "cron/jobs.json",
368 &meta.raw["cron/jobs.json"],
369 Fidelity::ByteLossless
370 );
371 } else {
372 let mut view: Profile = profile.clone();
373 for job in view.jobs.values_mut() {
374 if let Some(w) = &job.workdir {
375 if !w.starts_with('/') {
376 job.workdir = Some(dir.join(w).display().to_string());
377 }
378 }
379 }
380 emit!(
381 "cron/jobs.json",
382 &encode_jobs_file(&view, Some(meta)),
383 Fidelity::ByteLossless
384 );
385 }
386 }
387
388 let src_exec = meta
390 .source_dir
391 .as_ref()
392 .map(|d| d.join("cron/executions.db"));
393 if !profile.fires.is_empty() || src_exec.as_ref().is_some_and(|p| p.exists()) {
394 let target = dir.join("cron/executions.db");
395 let fires_record = serde_json::to_value(&profile.fires).unwrap();
396 if unchanged("cron/executions.db", &fires_record)
397 && src_exec.as_ref().is_some_and(|p| p.exists())
398 {
399 fs::copy(src_exec.as_ref().unwrap(), &target)?;
400 } else {
401 let tmp =
402 target.with_file_name(format!("executions.db.tmp-{}", std::process::id()));
403 let _ = fs::remove_file(&tmp);
404 write_executions(&tmp, &profile.fires)?;
405 fs::rename(&tmp, &target)?;
406 }
407 report
408 .written
409 .push(ArtifactFidelity::byte(rel("cron/executions.db")));
410 }
411
412 let subs_record: Vec<Value> = profile
414 .subscriptions
415 .values()
416 .map(|s| serde_json::to_value(s).unwrap())
417 .collect();
418 if !profile.subscriptions.is_empty() || meta.raw.contains_key("webhook_subscriptions.json")
419 {
420 if hermes_bytes
421 && unchanged("webhook_subscriptions.json", &Value::Array(subs_record))
422 && meta.raw.contains_key("webhook_subscriptions.json")
423 {
424 emit!(
425 "webhook_subscriptions.json",
426 &meta.raw["webhook_subscriptions.json"],
427 Fidelity::ByteLossless
428 );
429 } else {
430 emit!(
431 "webhook_subscriptions.json",
432 &encode_subscriptions_file(profile, Some(&loaded.vault)),
433 Fidelity::ByteLossless
434 );
435 }
436 }
437
438 if meta.borrowed_from.is_none() {
440 let state_record = serde_json::json!({ "bindings": profile.bindings, "obligations": profile.obligations });
441 let src_state = meta.source_dir.as_ref().map(|d| d.join("state.db"));
442 let lenders_unchanged = meta.lenders.iter().all(|n| {
443 let p = &loaded.orchestration.profiles[n];
444 loaded.io.get(n).and_then(|m| m.snapshot.get("state.db")) == Some(&canonical_json(&serde_json::json!({ "bindings": p.bindings, "obligations": p.obligations })))
445 });
446 let dest_store = dir.join("state.db");
447 if src_state.as_ref().is_some_and(|p| p.exists()) && meta.flavor == Flavor::Hermes {
448 if unchanged("state.db", &state_record) && lenders_unchanged {
449 fs::copy(src_state.as_ref().unwrap(), &dest_store)?;
450 report.written.push(ArtifactFidelity::byte(rel("state.db")));
451 } else {
452 report.refused.push(Refusal { file: rel("state.db"), reason: "bindings/obligations changed since import; writing the change into a Hermes session store is UNI-18 (the first shared-WAL write), not this codec's".into() });
456 }
457 } else if !profile.bindings.is_empty() || !profile.obligations.is_empty() {
458 let _ = dest_store;
461 fresh_rows.push((name.as_str(), profile));
462 }
463 }
464
465 copy_unmodeled(profile, meta, &dir)?;
467 for f in &profile.residue.files {
468 report.written.push(ArtifactFidelity::byte(rel(f)));
469 }
470 }
471 if !fresh_rows.is_empty() {
472 let root_store = dest.join("state.db");
473 if root_store.exists() {
474 report.refused.push(Refusal { file: "state.db".into(), reason: "the destination already holds a Hermes session store; writing into a live store is UNI-18 (the first shared-WAL write), not this codec's".into() });
475 } else {
476 let patched: Vec<(&str, Profile)> = fresh_rows
478 .iter()
479 .map(|(name, p)| (*name, with_handoffs(p, &loaded.orchestration)))
480 .collect();
481 let refs: Vec<(&str, &Profile)> = patched.iter().map(|(n, p)| (*n, p)).collect();
482 write_fresh_store(&refs, &root_store)?;
483 let bindings: usize = fresh_rows.iter().map(|(_, p)| p.bindings.len()).sum();
484 let obligations: usize = fresh_rows.iter().map(|(_, p)| p.obligations.len()).sum();
485 report.written.push(ArtifactFidelity::semantic(
486 "state.db",
487 vec![format!(
488 "a fresh store at schema {HERMES_STATE_V22_VERSION} (Hermes migrates it on open): {bindings} binding(s) as sessions rows across {} profile(s) partitioned by profile_name, {obligations} obligation(s) as delivery rows; the transcripts live in the worker's store and are not carried",
489 fresh_rows.len()
490 )],
491 ));
492 }
493 }
494 Ok(report)
495}