1use std::collections::BTreeMap;
12use std::fs;
13use std::path::{Path, PathBuf};
14
15use serde_json::{Map, Value};
16
17use super::canonical::canonical_json;
18use super::decode::{load_error, placeholder_ref, surface_key_string};
19use super::folder::{empty_profile, iso_epoch, ordered_object, persona_ref, pretty_ordered};
20use super::sqlite::{read_rows, table_exists, write_table, Param};
21use crate::ontology::{
22 ArtifactFidelity, Binding, Fidelity, HarnessId, Recurrence, Residue, SecretRef, SurfaceKey,
23 Trigger, Worker,
24};
25use crate::orchestration::{
26 ChannelConfig, Fire, FireStatus, Job, JobOrigin, Obligation, ObligationSource, ObligationState,
27 Orchestration, OutboundContent, Profile, Route, RouteMatch, Schedule, Target,
28 WebhookSubscription,
29};
30use crate::Result;
31
32pub const OPENCLAW_CONFIG: &str = "openclaw.json";
34pub const OPENCLAW_STATE_DB: &str = "state/openclaw.sqlite";
36pub const OPENCLAW_ACCOUNT_KEYS: &[&str] = &[
38 "accountId",
39 "account_id",
40 "account",
41 "teamId",
42 "appId",
43 "userId",
44];
45
46pub const OPENCLAW_DDL: &str = r#"CREATE TABLE IF NOT EXISTS schema_meta (
48 meta_key TEXT NOT NULL PRIMARY KEY,
49 role TEXT NOT NULL,
50 schema_version INTEGER NOT NULL,
51 agent_id TEXT,
52 app_version TEXT,
53 created_at INTEGER NOT NULL,
54 updated_at INTEGER NOT NULL
55);
56CREATE TABLE IF NOT EXISTS cron_jobs (
57 store_key TEXT NOT NULL,
58 job_id TEXT NOT NULL,
59 declaration_key TEXT,
60 display_name TEXT,
61 owner_agent_id TEXT,
62 owner_session_key TEXT,
63 name TEXT NOT NULL,
64 description TEXT,
65 enabled INTEGER NOT NULL,
66 delete_after_run INTEGER,
67 created_at_ms INTEGER NOT NULL,
68 agent_id TEXT,
69 session_key TEXT,
70 schedule_kind TEXT NOT NULL,
71 schedule_expr TEXT,
72 schedule_tz TEXT,
73 every_ms INTEGER,
74 anchor_ms INTEGER,
75 at TEXT,
76 stagger_ms INTEGER,
77 session_target TEXT NOT NULL,
78 wake_mode TEXT NOT NULL,
79 trigger_script TEXT,
80 trigger_once INTEGER,
81 payload_kind TEXT NOT NULL,
82 payload_message TEXT,
83 payload_model TEXT,
84 payload_fallbacks_json TEXT,
85 payload_thinking TEXT,
86 payload_timeout_seconds INTEGER,
87 payload_allow_unsafe_external_content INTEGER,
88 payload_external_content_source_json TEXT,
89 payload_light_context INTEGER,
90 payload_tools_allow_json TEXT,
91 payload_tools_allow_is_default INTEGER,
92 delivery_mode TEXT,
93 delivery_channel TEXT,
94 delivery_to TEXT,
95 delivery_thread_id TEXT,
96 delivery_thread_id_type TEXT,
97 delivery_account_id TEXT,
98 delivery_best_effort INTEGER,
99 delivery_completion_mode TEXT,
100 delivery_completion_to TEXT,
101 failure_delivery_mode TEXT,
102 failure_delivery_channel TEXT,
103 failure_delivery_to TEXT,
104 failure_delivery_account_id TEXT,
105 failure_alert_disabled INTEGER,
106 failure_alert_after INTEGER,
107 failure_alert_channel TEXT,
108 failure_alert_to TEXT,
109 failure_alert_cooldown_ms INTEGER,
110 failure_alert_include_skipped INTEGER,
111 failure_alert_mode TEXT,
112 failure_alert_account_id TEXT,
113 next_run_at_ms INTEGER,
114 running_at_ms INTEGER,
115 last_run_at_ms INTEGER,
116 last_run_status TEXT,
117 last_error TEXT,
118 last_duration_ms INTEGER,
119 consecutive_errors INTEGER,
120 consecutive_skipped INTEGER,
121 schedule_error_count INTEGER,
122 last_delivery_status TEXT,
123 last_delivery_error TEXT,
124 last_delivered INTEGER,
125 last_failure_alert_at_ms INTEGER,
126 job_json TEXT NOT NULL,
127 state_json TEXT NOT NULL DEFAULT '{}',
128 runtime_updated_at_ms INTEGER,
129 schedule_identity TEXT,
130 sort_order INTEGER NOT NULL DEFAULT 0,
131 updated_at INTEGER NOT NULL,
132 PRIMARY KEY (store_key, job_id)
133);
134CREATE INDEX IF NOT EXISTS idx_cron_jobs_store_updated
135 ON cron_jobs(store_key, sort_order ASC, updated_at DESC, job_id);
136CREATE INDEX IF NOT EXISTS idx_cron_jobs_store_order
137 ON cron_jobs(store_key, sort_order ASC, updated_at ASC, job_id);
138CREATE INDEX IF NOT EXISTS idx_cron_jobs_enabled_next_run
139 ON cron_jobs(store_key, enabled, next_run_at_ms, job_id)
140 WHERE next_run_at_ms IS NOT NULL;
141CREATE INDEX IF NOT EXISTS idx_cron_jobs_agent_session
142 ON cron_jobs(agent_id, session_key, updated_at DESC, job_id)
143 WHERE agent_id IS NOT NULL OR session_key IS NOT NULL;
144CREATE TABLE IF NOT EXISTS cron_run_logs (
145 store_key TEXT NOT NULL,
146 job_id TEXT NOT NULL,
147 seq INTEGER NOT NULL,
148 ts INTEGER NOT NULL,
149 status TEXT,
150 error TEXT,
151 summary TEXT,
152 diagnostics_summary TEXT,
153 delivery_status TEXT,
154 delivery_error TEXT,
155 delivered INTEGER,
156 session_id TEXT,
157 session_key TEXT,
158 run_id TEXT,
159 run_at_ms INTEGER,
160 duration_ms INTEGER,
161 next_run_at_ms INTEGER,
162 model TEXT,
163 provider TEXT,
164 total_tokens INTEGER,
165 entry_json TEXT NOT NULL,
166 created_at INTEGER NOT NULL,
167 PRIMARY KEY (store_key, job_id, seq)
168);
169CREATE INDEX IF NOT EXISTS idx_cron_run_logs_store_ts
170 ON cron_run_logs(store_key, ts DESC, seq DESC);
171CREATE INDEX IF NOT EXISTS idx_cron_run_logs_job_status
172 ON cron_run_logs(store_key, job_id, status, ts DESC, seq DESC);
173CREATE INDEX IF NOT EXISTS idx_cron_run_logs_delivery
174 ON cron_run_logs(store_key, delivery_status, ts DESC, seq DESC)
175 WHERE delivery_status IS NOT NULL;
176CREATE TABLE IF NOT EXISTS delivery_queue_entries (
177 queue_name TEXT NOT NULL,
178 id TEXT NOT NULL,
179 status TEXT NOT NULL,
180 entry_kind TEXT,
181 session_key TEXT,
182 channel TEXT,
183 target TEXT,
184 account_id TEXT,
185 retry_count INTEGER NOT NULL DEFAULT 0,
186 last_attempt_at INTEGER,
187 last_error TEXT,
188 recovery_state TEXT,
189 platform_send_started_at INTEGER,
190 entry_json TEXT NOT NULL,
191 enqueued_at INTEGER NOT NULL,
192 updated_at INTEGER NOT NULL,
193 failed_at INTEGER,
194 PRIMARY KEY (queue_name, id)
195);
196CREATE INDEX IF NOT EXISTS idx_delivery_queue_pending
197 ON delivery_queue_entries(queue_name, status, enqueued_at, id);
198CREATE INDEX IF NOT EXISTS idx_delivery_queue_failed
199 ON delivery_queue_entries(queue_name, status, failed_at, id);
200CREATE INDEX IF NOT EXISTS idx_delivery_queue_session
201 ON delivery_queue_entries(queue_name, status, session_key, enqueued_at, id)
202 WHERE session_key IS NOT NULL;
203CREATE INDEX IF NOT EXISTS idx_delivery_queue_target
204 ON delivery_queue_entries(queue_name, status, channel, target, enqueued_at, id)
205 WHERE channel IS NOT NULL AND target IS NOT NULL;
206"#;
207pub const CRON_JOB_COLUMNS: &[&str] = &[
209 "store_key",
210 "job_id",
211 "declaration_key",
212 "display_name",
213 "owner_agent_id",
214 "owner_session_key",
215 "name",
216 "description",
217 "enabled",
218 "delete_after_run",
219 "created_at_ms",
220 "agent_id",
221 "session_key",
222 "schedule_kind",
223 "schedule_expr",
224 "schedule_tz",
225 "every_ms",
226 "anchor_ms",
227 "at",
228 "stagger_ms",
229 "session_target",
230 "wake_mode",
231 "trigger_script",
232 "trigger_once",
233 "payload_kind",
234 "payload_message",
235 "payload_model",
236 "payload_fallbacks_json",
237 "payload_thinking",
238 "payload_timeout_seconds",
239 "payload_allow_unsafe_external_content",
240 "payload_external_content_source_json",
241 "payload_light_context",
242 "payload_tools_allow_json",
243 "payload_tools_allow_is_default",
244 "delivery_mode",
245 "delivery_channel",
246 "delivery_to",
247 "delivery_thread_id",
248 "delivery_thread_id_type",
249 "delivery_account_id",
250 "delivery_best_effort",
251 "delivery_completion_mode",
252 "delivery_completion_to",
253 "failure_delivery_mode",
254 "failure_delivery_channel",
255 "failure_delivery_to",
256 "failure_delivery_account_id",
257 "failure_alert_disabled",
258 "failure_alert_after",
259 "failure_alert_channel",
260 "failure_alert_to",
261 "failure_alert_cooldown_ms",
262 "failure_alert_include_skipped",
263 "failure_alert_mode",
264 "failure_alert_account_id",
265 "next_run_at_ms",
266 "running_at_ms",
267 "last_run_at_ms",
268 "last_run_status",
269 "last_error",
270 "last_duration_ms",
271 "consecutive_errors",
272 "consecutive_skipped",
273 "schedule_error_count",
274 "last_delivery_status",
275 "last_delivery_error",
276 "last_delivered",
277 "last_failure_alert_at_ms",
278 "job_json",
279 "state_json",
280 "runtime_updated_at_ms",
281 "schedule_identity",
282 "sort_order",
283 "updated_at",
284];
285pub const CRON_RUN_LOG_COLUMNS: &[&str] = &[
287 "store_key",
288 "job_id",
289 "seq",
290 "ts",
291 "status",
292 "error",
293 "summary",
294 "diagnostics_summary",
295 "delivery_status",
296 "delivery_error",
297 "delivered",
298 "session_id",
299 "session_key",
300 "run_id",
301 "run_at_ms",
302 "duration_ms",
303 "next_run_at_ms",
304 "model",
305 "provider",
306 "total_tokens",
307 "entry_json",
308 "created_at",
309];
310pub const DELIVERY_QUEUE_COLUMNS: &[&str] = &[
312 "queue_name",
313 "id",
314 "status",
315 "entry_kind",
316 "session_key",
317 "channel",
318 "target",
319 "account_id",
320 "retry_count",
321 "last_attempt_at",
322 "last_error",
323 "recovery_state",
324 "platform_send_started_at",
325 "entry_json",
326 "enqueued_at",
327 "updated_at",
328 "failed_at",
329];
330pub const SCHEMA_META_COLUMNS: &[&str] = &[
332 "meta_key",
333 "role",
334 "schema_version",
335 "agent_id",
336 "app_version",
337 "created_at",
338 "updated_at",
339];
340
341pub fn strip_json5(text: &str) -> String {
347 let mut out = String::with_capacity(text.len());
348 let mut chars = text.chars().peekable();
349 let mut in_string = false;
350 let mut escaped = false;
351 while let Some(ch) = chars.next() {
352 if in_string {
353 out.push(ch);
354 if escaped {
355 escaped = false;
356 } else if ch == '\\' {
357 escaped = true;
358 } else if ch == '"' {
359 in_string = false;
360 }
361 continue;
362 }
363 match ch {
364 '"' => {
365 in_string = true;
366 out.push(ch);
367 }
368 '/' if chars.peek() == Some(&'/') => {
369 for next in chars.by_ref() {
370 if next == '\n' {
371 out.push('\n');
372 break;
373 }
374 }
375 }
376 '/' if chars.peek() == Some(&'*') => {
377 chars.next();
378 let mut previous = '\0';
379 for next in chars.by_ref() {
380 if previous == '*' && next == '/' {
381 break;
382 }
383 previous = next;
384 }
385 out.push(' ');
386 }
387 _ => out.push(ch),
388 }
389 }
390 let bytes: Vec<char> = out.chars().collect();
391 let mut cleaned = String::with_capacity(out.len());
392 let mut index = 0usize;
393 let mut in_string = false;
394 let mut escaped = false;
395 while index < bytes.len() {
396 let ch = bytes[index];
397 if in_string {
398 cleaned.push(ch);
399 if escaped {
400 escaped = false;
401 } else if ch == '\\' {
402 escaped = true;
403 } else if ch == '"' {
404 in_string = false;
405 }
406 index += 1;
407 continue;
408 }
409 if ch == '"' {
410 in_string = true;
411 cleaned.push(ch);
412 index += 1;
413 continue;
414 }
415 if ch == ',' {
416 let mut lookahead = index + 1;
417 while lookahead < bytes.len() && bytes[lookahead].is_whitespace() {
418 lookahead += 1;
419 }
420 if lookahead < bytes.len() && (bytes[lookahead] == '}' || bytes[lookahead] == ']') {
421 index += 1;
422 continue;
423 }
424 }
425 cleaned.push(ch);
426 index += 1;
427 }
428 cleaned
429}
430
431pub fn parse_json5(file: &str, text: &str) -> Result<Value> {
433 serde_json::from_str(&strip_json5(text))
434 .map_err(|e| load_error(file, "", format!("JSON5: {e}")))
435}
436
437pub const NON_SECRET_KEYS: &[&str] = &[
441 "sessionkey",
442 "session_key",
443 "storekey",
444 "store_key",
445 "metakey",
446 "meta_key",
447 "bindingkey",
448 "binding_key",
449 "declarationkey",
450 "declaration_key",
451 "idempotencykey",
452 "idempotency_key",
453];
454
455pub fn is_credential_key(key: &str) -> bool {
457 let lower = key.to_ascii_lowercase();
458 if NON_SECRET_KEYS.contains(&lower.as_str()) {
459 return false;
460 }
461 ["token", "key", "secret", "password", "credential"]
462 .iter()
463 .any(|m| lower.ends_with(m))
464}
465
466pub fn credential_ref_name(scope: &str, key: &str) -> String {
468 format!("OPENCLAW_{scope}_{key}")
469 .chars()
470 .map(|c| {
471 if c.is_ascii_alphanumeric() {
472 c.to_ascii_uppercase()
473 } else {
474 '_'
475 }
476 })
477 .collect()
478}
479
480pub fn redact_secrets(value: &Value, scope: &str, vault: &mut BTreeMap<String, String>) -> Value {
482 match value {
483 Value::Array(items) => Value::Array(
484 items
485 .iter()
486 .enumerate()
487 .map(|(i, v)| redact_secrets(v, &format!("{scope}_{i}"), vault))
488 .collect(),
489 ),
490 Value::Object(map) => {
491 let mut out = Map::new();
492 for (k, v) in map {
493 match v {
494 Value::String(s) if is_credential_key(k) => {
495 let r = credential_ref_name(scope, k);
496 vault.insert(r.clone(), s.clone());
497 out.insert(k.clone(), Value::String(format!("${{{r}}}")));
498 }
499 Value::Object(_) | Value::Array(_) => {
500 out.insert(k.clone(), redact_secrets(v, &format!("{scope}_{k}"), vault));
501 }
502 other => {
503 out.insert(k.clone(), other.clone());
504 }
505 }
506 }
507 Value::Object(out)
508 }
509 other => other.clone(),
510 }
511}
512
513fn resolve_ref(name: &str, vault: &BTreeMap<String, String>, depth: usize) -> Option<String> {
514 let value = vault.get(name)?;
515 if depth > 4 {
516 return Some(value.clone());
517 }
518 match placeholder_ref(value) {
519 Some(next) => resolve_ref(next, vault, depth + 1).or_else(|| Some(value.clone())),
520 None => Some(value.clone()),
521 }
522}
523
524pub fn inline_secrets(value: &Value, vault: &BTreeMap<String, String>) -> Value {
526 match value {
527 Value::String(s) => match placeholder_ref(s).and_then(|n| resolve_ref(n, vault, 0)) {
528 Some(v) => Value::String(v),
529 None => value.clone(),
530 },
531 Value::Array(items) => {
532 Value::Array(items.iter().map(|v| inline_secrets(v, vault)).collect())
533 }
534 Value::Object(map) => Value::Object(
535 map.iter()
536 .map(|(k, v)| (k.clone(), inline_secrets(v, vault)))
537 .collect(),
538 ),
539 other => other.clone(),
540 }
541}
542
543fn unresolved_refs(value: &Value, path: &str) -> Vec<(String, String)> {
545 match value {
546 Value::String(s) => placeholder_ref(s)
547 .map(|n| vec![(path.to_string(), n.to_string())])
548 .unwrap_or_default(),
549 Value::Array(items) => items
550 .iter()
551 .enumerate()
552 .flat_map(|(i, v)| unresolved_refs(v, &format!("{path}[{i}]")))
553 .collect(),
554 Value::Object(map) => map
555 .iter()
556 .flat_map(|(k, v)| {
557 unresolved_refs(
558 v,
559 &if path.is_empty() {
560 k.clone()
561 } else {
562 format!("{path}.{k}")
563 },
564 )
565 })
566 .collect(),
567 _ => Vec::new(),
568 }
569}
570
571fn civil(ms: i64) -> (i64, u32, u32, u32, u32, u32, u32) {
574 let secs = ms.div_euclid(1000);
575 let sub = ms.rem_euclid(1000) as u32;
576 let days = secs.div_euclid(86_400);
577 let sod = secs.rem_euclid(86_400);
578 let z = days + 719_468;
579 let era = z.div_euclid(146_097);
580 let doe = z - era * 146_097;
581 let yoe = (doe - doe / 1460 + doe / 36_524 - doe / 146_096) / 365;
582 let y = yoe + era * 400;
583 let doy = doe - (365 * yoe + yoe / 4 - yoe / 100);
584 let mp = (5 * doy + 2) / 153;
585 let d = (doy - (153 * mp + 2) / 5 + 1) as u32;
586 let m = if mp < 10 { mp + 3 } else { mp - 9 } as u32;
587 let y = if m <= 2 { y + 1 } else { y };
588 (
589 y,
590 m,
591 d,
592 (sod / 3600) as u32,
593 ((sod % 3600) / 60) as u32,
594 (sod % 60) as u32,
595 sub,
596 )
597}
598
599pub fn iso_seconds(ms: Option<i64>) -> Option<String> {
601 ms.map(|ms| {
602 let (y, mo, d, h, mi, s, _) = civil(ms);
603 format!("{y:04}-{mo:02}-{d:02}T{h:02}:{mi:02}:{s:02}Z")
604 })
605}
606
607pub fn iso_millis(ms: Option<i64>) -> Option<String> {
609 ms.map(|ms| {
610 let (y, mo, d, h, mi, s, sub) = civil(ms);
611 format!("{y:04}-{mo:02}-{d:02}T{h:02}:{mi:02}:{s:02}.{sub:03}Z")
612 })
613}
614
615pub fn ms_from_iso(iso: Option<&str>) -> Option<i64> {
617 let s = iso?.trim();
618 let (date, rest) = s.split_once('T')?;
619 let mut dp = date.split('-');
620 let (y, mo, d): (i64, i64, i64) = (
621 dp.next()?.parse().ok()?,
622 dp.next()?.parse().ok()?,
623 dp.next()?.parse().ok()?,
624 );
625 let (time, offset) = if let Some(t) = rest.strip_suffix('Z') {
626 (t, 0i64)
627 } else if let Some(idx) = rest.rfind(['+', '-']) {
628 let (t, off) = rest.split_at(idx);
629 let sign = if off.starts_with('-') { -1 } else { 1 };
630 let mut op = off[1..].split(':');
631 let (oh, om): (i64, i64) = (
632 op.next()?.parse().ok()?,
633 op.next().unwrap_or("0").parse().ok()?,
634 );
635 (t, sign * (oh * 3600 + om * 60))
636 } else {
637 (rest, 0)
638 };
639 let (hms, frac) = match time.split_once('.') {
640 Some((a, b)) => (a, b),
641 None => (time, ""),
642 };
643 let mut tp = hms.split(':');
644 let (h, mi, sec): (i64, i64, i64) = (
645 tp.next()?.parse().ok()?,
646 tp.next()?.parse().ok()?,
647 tp.next().unwrap_or("0").parse().ok()?,
648 );
649 let millis: i64 = if frac.is_empty() {
650 0
651 } else {
652 format!("{:0<3}", &frac[..frac.len().min(3)]).parse().ok()?
653 };
654 let (yy, mm) = if mo <= 2 {
656 (y - 1, mo + 9)
657 } else {
658 (y, mo - 3)
659 };
660 let era = yy.div_euclid(400);
661 let yoe = yy - era * 400;
662 let doy = (153 * mm + 2) / 5 + d - 1;
663 let doe = yoe * 365 + yoe / 4 - yoe / 100 + doy;
664 let days = era * 146_097 + doe - 719_468;
665 Some(((days * 86_400 + h * 3600 + mi * 60 + sec) - offset) * 1000 + millis)
666}
667
668fn ms_of(v: Option<&Value>) -> Option<i64> {
669 match v? {
670 Value::Number(n) => n.as_i64().or_else(|| n.as_f64().map(|f| f as i64)),
671 Value::String(s) => s.parse::<f64>().ok().map(|f| f as i64),
672 _ => None,
673 }
674}
675
676pub fn agent_entries(config: &Value) -> Vec<(String, Value)> {
680 let agents = config.get("agents");
681 if let Some(list) = agents.and_then(|a| a.get("list")).and_then(Value::as_array) {
682 return list
683 .iter()
684 .filter_map(|e| {
685 e.get("id")
686 .or_else(|| e.get("agentId"))
687 .and_then(Value::as_str)
688 .filter(|s| !s.is_empty())
689 .map(|id| (id.to_string(), e.clone()))
690 })
691 .collect();
692 }
693 if let Some(entries) = agents
694 .and_then(|a| a.get("entries"))
695 .and_then(Value::as_object)
696 {
697 return entries
698 .iter()
699 .map(|(k, v)| (k.clone(), v.clone()))
700 .collect();
701 }
702 Vec::new()
703}
704
705pub fn agents_form(config: &Value) -> Option<&'static str> {
707 let agents = config.get("agents")?;
708 if agents.get("list").is_some_and(Value::is_array) {
709 return Some("list");
710 }
711 if agents.get("entries").is_some_and(Value::is_object) {
712 return Some("entries");
713 }
714 None
715}
716
717pub fn default_agent_id(config: &Value) -> String {
719 let entries = agent_entries(config);
720 entries
721 .iter()
722 .find(|(_, e)| e.get("default") == Some(&Value::Bool(true)))
723 .or_else(|| entries.first())
724 .map(|(id, _)| id.clone())
725 .unwrap_or_else(|| "main".into())
726}
727
728#[derive(Debug, Clone, PartialEq)]
732pub struct ParsedKey {
733 pub agent: Option<String>,
735 pub key: SurfaceKey,
737 pub residue: Map<String, Value>,
739 pub recurrence: Option<Recurrence>,
741}
742
743const OPENCLAW_CHAT_KINDS: &[&str] = &["dm", "group", "channel", "thread"];
744
745fn key_of(platform: &str, kind: &str, chat_id: &str, thread_id: Option<String>) -> SurfaceKey {
746 SurfaceKey {
747 key: None,
748 platform: Some(platform.into()),
749 kind: Some(kind.into()),
750 chat_id: Some(chat_id.into()),
751 thread_id,
752 participant_id: None,
753 }
754}
755
756pub fn parse_openclaw_session_key(key: &str) -> Option<ParsedKey> {
761 let parts: Vec<&str> = key.split(':').collect();
762 let mut residue = Map::new();
763 residue.insert("session_key".into(), Value::String(key.into()));
764 match parts.first().copied() {
765 Some("agent") if parts.len() >= 3 => {
766 let agent = Some(parts[1].to_string());
767 if parts[2] == "main" {
768 residue.insert("dm_collapse".into(), Value::Bool(true));
769 return Some(ParsedKey {
770 agent,
771 key: key_of("main", "dm", "main", None),
772 residue,
773 recurrence: None,
774 });
775 }
776 if parts.len() < 5 {
777 return None;
778 }
779 let kind = if OPENCLAW_CHAT_KINDS.contains(&parts[3]) {
780 parts[3]
781 } else {
782 residue.insert("chat_kind".into(), Value::String(parts[3].into()));
783 "dm"
784 };
785 let mut thread_id = None;
786 if (parts.get(5) == Some(&"thread") || parts.get(5) == Some(&"topic"))
787 && parts.get(6).is_some_and(|t| !t.is_empty())
788 {
789 thread_id = Some(parts[6].to_string());
790 if parts[5] == "topic" {
791 residue.insert("thread_word".into(), Value::String("topic".into()));
792 }
793 }
794 let k = key_of(parts[2], kind, parts[4], thread_id);
795 let recurrence = if parts[2] == "cron" {
796 Some(Recurrence {
797 job_id: parts[4].into(),
798 kind: "cron".into(),
799 })
800 } else {
801 None
802 };
803 Some(ParsedKey {
804 agent,
805 key: k,
806 residue,
807 recurrence,
808 })
809 }
810 Some("cron") if parts.len() >= 2 => {
811 let job_id = parts[1..].join(":");
812 Some(ParsedKey {
813 agent: None,
814 key: key_of("cron", "dm", &job_id, None),
815 residue,
816 recurrence: Some(Recurrence {
817 job_id,
818 kind: "cron".into(),
819 }),
820 })
821 }
822 Some("hook") if parts.len() >= 2 => Some(ParsedKey {
823 agent: None,
824 key: key_of("webhook", "dm", parts[1], None),
825 residue,
826 recurrence: None,
827 }),
828 Some("acp-bridge") if parts.len() >= 2 => Some(ParsedKey {
829 agent: None,
830 key: key_of("acp", "dm", &parts[1..].join(":"), None),
831 residue,
832 recurrence: None,
833 }),
834 _ => None,
835 }
836}
837
838#[derive(Debug, Clone, Default)]
842pub struct OpenclawRootIo {
843 pub state_dir: PathBuf,
845 pub config_raw: String,
847 pub config_present: bool,
849 pub config_snapshot: String,
851 pub cron_jobs: BTreeMap<String, Map<String, Value>>,
853 pub cron_run_logs: BTreeMap<String, Map<String, Value>>,
855 pub delivery_queue_entries: BTreeMap<String, Map<String, Value>>,
857 pub schema_meta: Vec<Map<String, Value>>,
859 pub store_key: String,
861 pub db_present: bool,
863 pub default_agent: String,
865 pub legacy_jobs_raw: Option<String>,
868 pub legacy_jobs: BTreeMap<String, Value>,
870 pub legacy_jobs_object_form: bool,
872}
873
874#[derive(Debug, Clone, Default)]
876pub struct OpenclawProfileIo {
877 pub agent_id: String,
879 pub source_dir: PathBuf,
881 pub store_snapshot: String,
883 pub bindings_snapshot: String,
885}
886
887#[derive(Debug, Clone)]
889pub struct OpenclawLoaded {
890 pub orchestration: Orchestration,
892 pub vault: BTreeMap<String, String>,
894 pub root: OpenclawRootIo,
896 pub profiles: BTreeMap<String, OpenclawProfileIo>,
898}
899
900impl OpenclawLoaded {
901 pub fn from_orchestration(
904 orchestration: Orchestration,
905 vault: BTreeMap<String, String>,
906 ) -> Self {
907 let default_agent = orchestration.profiles["default"]
908 .residue
909 .config
910 .get("openclaw")
911 .and_then(|o| o.get("default_agent"))
912 .and_then(Value::as_str)
913 .unwrap_or("main")
914 .to_string();
915 let profiles = orchestration
916 .profiles
917 .keys()
918 .map(|n| {
919 (
920 n.clone(),
921 OpenclawProfileIo {
922 agent_id: if n == "default" {
923 default_agent.clone()
924 } else {
925 n.clone()
926 },
927 ..Default::default()
928 },
929 )
930 })
931 .collect();
932 Self {
933 root: OpenclawRootIo {
934 state_dir: orchestration.root.clone(),
935 default_agent,
936 ..Default::default()
937 },
938 orchestration,
939 vault,
940 profiles,
941 }
942 }
943}
944
945pub fn account_entries(entry: &Value) -> Vec<(String, Value)> {
949 let mut out: Vec<(String, Value)> = match entry.get("accounts") {
950 Some(Value::Object(m)) => m.iter().map(|(k, v)| (k.clone(), v.clone())).collect(),
951 Some(Value::Array(a)) => a
952 .iter()
953 .filter_map(|x| {
954 x.get("id")
955 .or_else(|| x.get("accountId"))
956 .and_then(Value::as_str)
957 .filter(|s| !s.is_empty())
958 .map(|id| (id.to_string(), x.clone()))
959 })
960 .collect(),
961 _ => Vec::new(),
962 };
963 out.sort_by(|a, b| a.0.cmp(&b.0));
964 out
965}
966
967fn credentials_of(
968 block: &Value,
969 scope: &str,
970 vault: &mut BTreeMap<String, String>,
971) -> BTreeMap<String, SecretRef> {
972 let mut creds = BTreeMap::new();
973 let Some(map) = block.as_object() else {
974 return creds;
975 };
976 for (k, v) in map {
977 match v {
978 Value::String(s) if placeholder_ref(s).is_some() => {
979 creds.insert(
980 k.clone(),
981 SecretRef::Dotenv(placeholder_ref(s).unwrap().to_string()),
982 );
983 }
984 Value::String(s) if is_credential_key(k) => {
985 let r = credential_ref_name(scope, k);
986 vault.insert(r.clone(), s.clone());
987 creds.insert(k.clone(), SecretRef::Dotenv(r));
988 }
989 _ => {}
990 }
991 }
992 creds
993}
994
995fn decode_channels(
996 config: &Value,
997 vault: &mut BTreeMap<String, String>,
998) -> BTreeMap<String, ChannelConfig> {
999 let mut channels = BTreeMap::new();
1000 let Some(map) = config.get("channels").and_then(Value::as_object) else {
1001 return channels;
1002 };
1003 for (kind, entry) in map {
1004 let Some(entry_map) = entry.as_object() else {
1005 continue;
1006 };
1007 let mut shared = entry_map.clone();
1008 shared.remove("accounts");
1009 let mut shared_block = redact_secrets(&Value::Object(shared), kind, vault);
1010 let channel_enabled = entry_map.get("enabled").and_then(Value::as_bool);
1011 if let Some(m) = shared_block.as_object_mut() {
1012 m.remove("enabled");
1013 }
1014 let accounts = account_entries(entry);
1015 let inline_account = OPENCLAW_ACCOUNT_KEYS
1016 .iter()
1017 .find_map(|k| entry_map.get(*k).and_then(Value::as_str))
1018 .map(str::to_string);
1019 if accounts.is_empty() {
1020 let mut extra = BTreeMap::new();
1021 extra.insert("kind".into(), Value::String(kind.clone()));
1022 extra.insert(
1023 "accountId".into(),
1024 inline_account
1025 .clone()
1026 .map(Value::String)
1027 .unwrap_or(Value::Null),
1028 );
1029 extra.insert(
1030 "account_source".into(),
1031 if inline_account.is_some() {
1032 Value::String("entry".into())
1033 } else {
1034 Value::Null
1035 },
1036 );
1037 extra.insert("accounts_form".into(), Value::Null);
1038 extra.insert(
1039 "enabled_on".into(),
1040 if channel_enabled.is_some() {
1041 Value::String("channel".into())
1042 } else {
1043 Value::Null
1044 },
1045 );
1046 extra.insert("channel_enabled".into(), Value::Null);
1047 extra.insert("channel_block".into(), shared_block.clone());
1048 channels.insert(
1049 kind.clone(),
1050 ChannelConfig {
1051 platform: kind.clone(),
1052 enabled: channel_enabled.unwrap_or(true),
1053 credentials: credentials_of(entry, kind, vault),
1054 extra,
1055 },
1056 );
1057 continue;
1058 }
1059 let form = if entry_map.get("accounts").is_some_and(Value::is_object) {
1060 "object"
1061 } else {
1062 "array"
1063 };
1064 for (id, account) in accounts {
1065 let name = format!("{kind}/{id}");
1066 let mut block = redact_secrets(&account, &name, vault);
1067 let account_enabled = account.get("enabled").and_then(Value::as_bool);
1068 if let Some(m) = block.as_object_mut() {
1069 m.remove("enabled");
1070 }
1071 let enabled = account_enabled.or(channel_enabled).unwrap_or(true);
1072 let mut credentials = credentials_of(&account, &name, vault);
1073 for (k, v) in credentials_of(entry, kind, vault) {
1074 credentials.insert(k, v);
1075 }
1076 let mut extra = BTreeMap::new();
1077 extra.insert("kind".into(), Value::String(kind.clone()));
1078 extra.insert("accountId".into(), Value::String(id.clone()));
1079 extra.insert("account_source".into(), Value::String("accounts".into()));
1080 extra.insert("accounts_form".into(), Value::String(form.into()));
1081 extra.insert(
1082 "enabled_on".into(),
1083 if account_enabled.is_some() {
1084 Value::String("account".into())
1085 } else if channel_enabled.is_some() {
1086 Value::String("channel".into())
1087 } else {
1088 Value::Null
1089 },
1090 );
1091 extra.insert(
1092 "channel_enabled".into(),
1093 channel_enabled.map(Value::Bool).unwrap_or(Value::Null),
1094 );
1095 extra.insert("channel_block".into(), shared_block.clone());
1096 extra.insert("account_block".into(), block);
1097 channels.insert(
1098 name.clone(),
1099 ChannelConfig {
1100 platform: name,
1101 enabled,
1102 credentials,
1103 extra,
1104 },
1105 );
1106 }
1107 }
1108 channels
1109}
1110
1111const CHANNEL_BOOKKEEPING: &[&str] = &[
1112 "kind",
1113 "accountId",
1114 "account_source",
1115 "accounts_form",
1116 "enabled_on",
1117 "channel_enabled",
1118 "channel_block",
1119 "account_block",
1120];
1121
1122fn block_from_model(ch: &ChannelConfig) -> Value {
1123 let mut out = Map::new();
1124 for (k, v) in &ch.extra {
1125 if CHANNEL_BOOKKEEPING.contains(&k.as_str()) {
1126 continue;
1127 }
1128 out.insert(k.strip_prefix("extra.").unwrap_or(k).to_string(), v.clone());
1129 }
1130 for (k, r) in &ch.credentials {
1131 out.insert(k.clone(), serde_json::to_value(r).unwrap());
1132 }
1133 Value::Object(out)
1134}
1135
1136fn ordered_channel_block(block: Value) -> Vec<(String, Value)> {
1137 let Some(map) = block.as_object() else {
1138 return Vec::new();
1139 };
1140 let mut out: Vec<(String, Value)> = Vec::new();
1141 if let Some(e) = map.get("enabled") {
1142 out.push(("enabled".into(), e.clone()));
1143 }
1144 for (k, v) in map {
1145 if k != "enabled" {
1146 out.push((k.clone(), v.clone()));
1147 }
1148 }
1149 out
1150}
1151
1152fn encode_channels(profile: &Profile, vault: &BTreeMap<String, String>) -> Vec<(String, Value)> {
1153 let mut by_kind: Vec<(String, Vec<&ChannelConfig>)> = Vec::new();
1154 for ch in profile.channels.values() {
1155 let kind = ch
1156 .extra
1157 .get("kind")
1158 .and_then(Value::as_str)
1159 .unwrap_or(&ch.platform)
1160 .to_string();
1161 match by_kind.iter_mut().find(|(k, _)| *k == kind) {
1162 Some((_, rows)) => rows.push(ch),
1163 None => by_kind.push((kind, vec![ch])),
1164 }
1165 }
1166 let mut out = Vec::new();
1167 for (kind, rows) in by_kind {
1168 let first = rows[0];
1169 let channel_block_src = first
1170 .extra
1171 .get("channel_block")
1172 .filter(|v| v.is_object())
1173 .cloned()
1174 .unwrap_or_else(|| block_from_model(first));
1175 let channel_block = inline_secrets(&channel_block_src, vault);
1176 let channel_enabled = first.extra.get("channel_enabled").and_then(Value::as_bool);
1177 if first.extra.get("account_source").and_then(Value::as_str) != Some("accounts") {
1178 let mut single = channel_block.as_object().cloned().unwrap_or_default();
1179 if first.extra.get("enabled_on").and_then(Value::as_str) == Some("channel")
1180 || !first.enabled
1181 {
1182 single.insert("enabled".into(), Value::Bool(first.enabled));
1183 }
1184 out.push((
1185 kind,
1186 ordered_object(ordered_channel_block(Value::Object(single))),
1187 ));
1188 continue;
1189 }
1190 let form = first
1191 .extra
1192 .get("accounts_form")
1193 .and_then(Value::as_str)
1194 .unwrap_or("object");
1195 let mut head = channel_block.as_object().cloned().unwrap_or_default();
1196 if let Some(e) = channel_enabled {
1197 head.insert("enabled".into(), Value::Bool(e));
1198 }
1199 let blocks: Vec<(String, Value)> = rows
1200 .iter()
1201 .map(|ch| {
1202 let block_src = ch
1203 .extra
1204 .get("account_block")
1205 .filter(|v| v.is_object())
1206 .cloned()
1207 .unwrap_or_else(|| block_from_model(ch));
1208 let block = inline_secrets(&block_src, vault);
1209 let inherited = channel_enabled.unwrap_or(true);
1210 let id = ch
1211 .extra
1212 .get("accountId")
1213 .and_then(Value::as_str)
1214 .unwrap_or("")
1215 .to_string();
1216 if ch.extra.get("enabled_on").and_then(Value::as_str) == Some("account")
1217 || ch.enabled != inherited
1218 {
1219 let mut b = vec![("enabled".to_string(), Value::Bool(ch.enabled))];
1220 b.extend(
1221 block
1222 .as_object()
1223 .map(|m| {
1224 m.iter()
1225 .map(|(k, v)| (k.clone(), v.clone()))
1226 .collect::<Vec<_>>()
1227 })
1228 .unwrap_or_default(),
1229 );
1230 (id, ordered_object(b))
1231 } else {
1232 (
1233 id,
1234 ordered_object(
1235 block
1236 .as_object()
1237 .map(|m| m.iter().map(|(k, v)| (k.clone(), v.clone())).collect())
1238 .unwrap_or_default(),
1239 ),
1240 )
1241 }
1242 })
1243 .collect();
1244 let mut pairs: Vec<(String, Value)> = head.into_iter().collect();
1245 if form == "array" {
1246 pairs.push((
1247 "accounts".into(),
1248 Value::Array(
1249 blocks
1250 .into_iter()
1251 .map(|(id, b)| {
1252 let mut p = vec![("id".to_string(), Value::String(id))];
1253 if let Some(items) = is_ordered_pairs(&b) {
1254 p.extend(items);
1255 }
1256 ordered_object(p)
1257 })
1258 .collect(),
1259 ),
1260 ));
1261 } else {
1262 pairs.push(("accounts".into(), ordered_object(blocks)));
1263 }
1264 out.push((kind, ordered_object(pairs)));
1265 }
1266 out
1267}
1268
1269fn is_ordered_pairs(value: &Value) -> Option<Vec<(String, Value)>> {
1270 let arr = value.as_array()?;
1271 if arr.len() == 2 && arr[0].as_str() == Some("__ordered__") {
1272 return arr[1].as_array().map(|items| {
1273 items
1274 .iter()
1275 .map(|p| {
1276 (
1277 p["__k"].as_str().unwrap_or("").to_string(),
1278 p["__v"].clone(),
1279 )
1280 })
1281 .collect()
1282 });
1283 }
1284 None
1285}
1286
1287const BINDING_MATCH_MAPPED: &[&str] = &["channel", "guildId", "peer"];
1290
1291fn decode_routes(config: &Value, name_for: &dyn Fn(&str) -> String) -> Vec<Route> {
1292 let mut routes = Vec::new();
1293 let Some(list) = config.get("bindings").and_then(Value::as_array) else {
1294 return routes;
1295 };
1296 for (index, binding) in list.iter().enumerate() {
1297 let (Some(bmap), Some(agent_id)) = (
1298 binding.as_object(),
1299 binding.get("agentId").and_then(Value::as_str),
1300 ) else {
1301 continue;
1302 };
1303 let m = binding
1304 .get("match")
1305 .and_then(Value::as_object)
1306 .cloned()
1307 .unwrap_or_default();
1308 let text = |v: &Value| match v {
1309 Value::String(s) => Some(s.clone()),
1310 Value::Number(n) => Some(n.to_string()),
1311 _ => None,
1312 };
1313 let matches = RouteMatch {
1314 platform: m
1315 .get("channel")
1316 .and_then(Value::as_str)
1317 .unwrap_or("")
1318 .to_string(),
1319 guild_id: m.get("guildId").filter(|v| !v.is_null()).and_then(text),
1320 chat_id: m
1321 .get("peer")
1322 .and_then(|p| p.get("id"))
1323 .filter(|v| !v.is_null())
1324 .and_then(text),
1325 thread_id: None,
1326 };
1327 let mut match_residue = Map::new();
1328 for (k, v) in &m {
1329 if !BINDING_MATCH_MAPPED.contains(&k.as_str()) {
1330 match_residue.insert(k.clone(), v.clone());
1331 }
1332 }
1333 if let Some(peer) = m.get("peer").and_then(Value::as_object) {
1334 let mut pr = peer.clone();
1335 pr.remove("id");
1336 if !pr.is_empty() {
1337 match_residue.insert("peer".into(), Value::Object(pr));
1338 }
1339 }
1340 let mut residue = Residue::default();
1341 residue.keep("agent_id", Value::String(agent_id.into()));
1342 residue.keep("index", Value::from(index));
1343 for (k, v) in bmap {
1344 if k != "agentId" && k != "match" {
1345 residue.keep(k.clone(), v.clone());
1346 }
1347 }
1348 if !match_residue.is_empty() {
1349 residue.keep("match", Value::Object(match_residue));
1350 }
1351 routes.push(Route {
1352 name: None,
1353 matches,
1354 profile: name_for(agent_id),
1355 residue,
1356 });
1357 }
1358 routes
1359}
1360
1361fn encode_routes(profile: &Profile, id_for_name: &dyn Fn(&str) -> String) -> Vec<Value> {
1362 profile
1363 .routes
1364 .iter()
1365 .map(|r| {
1366 let mut residue = r.residue.0.clone();
1367 let agent_id = residue
1368 .remove("agent_id")
1369 .and_then(|v| v.as_str().map(str::to_string))
1370 .unwrap_or_else(|| id_for_name(&r.profile));
1371 residue.remove("index");
1372 let match_residue = residue
1373 .remove("match")
1374 .and_then(|v| v.as_object().cloned())
1375 .unwrap_or_default();
1376 let mut m: Vec<(String, Value)> = Vec::new();
1377 if !r.matches.platform.is_empty() {
1378 m.push(("channel".into(), Value::String(r.matches.platform.clone())));
1379 }
1380 for (k, v) in &match_residue {
1381 if k != "peer" {
1382 m.push((k.clone(), v.clone()));
1383 }
1384 }
1385 if let Some(g) = &r.matches.guild_id {
1386 m.push(("guildId".into(), Value::String(g.clone())));
1387 }
1388 if r.matches.chat_id.is_some()
1389 || match_residue.get("peer").is_some_and(Value::is_object)
1390 {
1391 let mut peer: Vec<(String, Value)> = match_residue
1392 .get("peer")
1393 .and_then(Value::as_object)
1394 .map(|p| p.iter().map(|(k, v)| (k.clone(), v.clone())).collect())
1395 .unwrap_or_default();
1396 if let Some(c) = &r.matches.chat_id {
1397 peer.push(("id".into(), Value::String(c.clone())));
1398 }
1399 m.push(("peer".into(), ordered_object(peer)));
1400 }
1401 let mut pairs: Vec<(String, Value)> = residue.into_iter().collect();
1402 pairs.push(("agentId".into(), Value::String(agent_id)));
1403 pairs.push(("match".into(), ordered_object(m)));
1404 ordered_object(pairs)
1405 })
1406 .collect()
1407}
1408
1409#[derive(Debug, Clone, Default, PartialEq)]
1413struct HooksMeta {
1414 block: Map<String, Value>,
1415 has_token: bool,
1416}
1417
1418fn decode_hooks(
1419 config: &Value,
1420 vault: &mut BTreeMap<String, String>,
1421 name_for: &dyn Fn(&str) -> String,
1422 profiles: &mut BTreeMap<String, Profile>,
1423) -> Option<HooksMeta> {
1424 let hooks = config.get("hooks").and_then(Value::as_object)?;
1425 let mut secret = None;
1426 if let Some(Value::String(token)) = hooks.get("token") {
1427 let r = credential_ref_name("hooks", "token");
1428 vault.insert(r.clone(), token.clone());
1429 secret = Some(SecretRef::Dotenv(r));
1430 }
1431 if let Some(mappings) = hooks.get("mappings").and_then(Value::as_array) {
1432 for (index, mapping) in mappings.iter().enumerate() {
1433 let Some(mm) = mapping.as_object() else {
1434 continue;
1435 };
1436 let name = mm
1437 .get("id")
1438 .and_then(Value::as_str)
1439 .filter(|s| !s.is_empty())
1440 .map(str::to_string)
1441 .unwrap_or_else(|| format!("hook-{index}"));
1442 let mut residue_v = redact_secrets(mapping, &format!("hook_{name}"), vault)
1443 .as_object()
1444 .cloned()
1445 .unwrap_or_default();
1446 residue_v.remove("deliver");
1447 residue_v.remove("to");
1448 residue_v.insert("__index".into(), Value::from(index));
1449 let owner = mm
1450 .get("agentId")
1451 .and_then(Value::as_str)
1452 .map(name_for)
1453 .filter(|n| profiles.contains_key(n))
1454 .unwrap_or_else(|| "default".into());
1455 let deliver =
1456 mm.get("deliver")
1457 .and_then(Value::as_str)
1458 .map(|platform| Target::Explicit {
1459 platform: platform.into(),
1460 chat_id: mm.get("to").and_then(|t| match t {
1461 Value::String(s) => Some(s.clone()),
1462 Value::Number(n) => Some(n.to_string()),
1463 _ => None,
1464 }),
1465 thread_id: None,
1466 });
1467 profiles.get_mut(&owner).unwrap().subscriptions.insert(
1468 name.clone(),
1469 WebhookSubscription {
1470 name,
1471 secret: secret.clone(),
1472 events: None,
1473 prompt_template: String::new(),
1474 deliver,
1475 skills: Vec::new(),
1476 description: None,
1477 created_at: None,
1478 residue: Residue(residue_v.into_iter().collect()),
1479 },
1480 );
1481 }
1482 }
1483 let mut block = hooks.clone();
1484 block.remove("token");
1485 block.remove("mappings");
1486 Some(HooksMeta {
1487 block,
1488 has_token: secret.is_some(),
1489 })
1490}
1491
1492fn encode_hooks(
1493 orchestration: &Orchestration,
1494 meta: Option<&HooksMeta>,
1495 vault: &BTreeMap<String, String>,
1496) -> Option<Value> {
1497 let mut subs: Vec<&WebhookSubscription> = orchestration
1498 .profiles
1499 .values()
1500 .flat_map(|p| p.subscriptions.values())
1501 .collect();
1502 if meta.is_none() && subs.is_empty() {
1503 return None;
1504 }
1505 subs.sort_by_key(|s| {
1506 s.residue
1507 .0
1508 .get("__index")
1509 .and_then(Value::as_i64)
1510 .unwrap_or(0)
1511 });
1512 let mut pairs: Vec<(String, Value)> = meta
1513 .map(|m| {
1514 inline_secrets(&Value::Object(m.block.clone()), vault)
1515 .as_object()
1516 .map(|o| o.iter().map(|(k, v)| (k.clone(), v.clone())).collect())
1517 .unwrap_or_default()
1518 })
1519 .unwrap_or_default();
1520 if meta.is_some_and(|m| m.has_token) {
1521 let r = subs
1522 .iter()
1523 .find_map(|s| s.secret.clone())
1524 .unwrap_or_else(|| SecretRef::Dotenv(credential_ref_name("hooks", "token")));
1525 pairs.push((
1526 "token".into(),
1527 inline_secrets(&serde_json::to_value(&r).unwrap(), vault),
1528 ));
1529 }
1530 if !subs.is_empty() {
1531 pairs.push((
1532 "mappings".into(),
1533 Value::Array(
1534 subs.iter()
1535 .map(|sub| {
1536 let mut residue = inline_secrets(
1537 &Value::Object(sub.residue.0.clone().into_iter().collect()),
1538 vault,
1539 )
1540 .as_object()
1541 .cloned()
1542 .unwrap_or_default();
1543 residue.remove("__index");
1544 let mut mp: Vec<(String, Value)> =
1545 vec![("id".into(), Value::String(sub.name.clone()))];
1546 mp.extend(residue);
1547 if let Some(Target::Explicit {
1548 platform, chat_id, ..
1549 }) = &sub.deliver
1550 {
1551 mp.push(("deliver".into(), Value::String(platform.clone())));
1552 if let Some(c) = chat_id {
1553 mp.push(("to".into(), Value::String(c.clone())));
1554 }
1555 }
1556 ordered_object(mp)
1557 })
1558 .collect(),
1559 ),
1560 ));
1561 }
1562 Some(ordered_object(pairs))
1563}
1564
1565fn json_object_of(text: Option<&Value>) -> Map<String, Value> {
1568 text.and_then(Value::as_str)
1569 .and_then(|s| serde_json::from_str::<Value>(s).ok())
1570 .and_then(|v| v.as_object().cloned())
1571 .unwrap_or_default()
1572}
1573
1574fn opt_text(v: Option<&Value>) -> Option<String> {
1575 match v {
1576 Some(Value::String(s)) if !s.is_empty() => Some(s.clone()),
1577 Some(Value::Number(n)) => Some(n.to_string()),
1578 _ => None,
1579 }
1580}
1581
1582fn delivery_target(
1583 delivery: Option<&Map<String, Value>>,
1584 row: Option<&Map<String, Value>>,
1585) -> Option<Target> {
1586 let d = |k: &str| delivery.and_then(|d| d.get(k)).filter(|v| !v.is_null());
1587 let r = |k: &str| row.and_then(|r| r.get(k)).filter(|v| !v.is_null());
1588 let mode = d("mode")
1589 .or_else(|| d("kind"))
1590 .or_else(|| d("type"))
1591 .or_else(|| r("delivery_mode"))
1592 .and_then(Value::as_str)
1593 .map(str::to_string);
1594 let channel = d("channel")
1595 .or_else(|| r("delivery_channel"))
1596 .and_then(Value::as_str)
1597 .map(str::to_string);
1598 let to = opt_text(d("to").or_else(|| r("delivery_to")));
1599 let thread = opt_text(
1600 d("threadId")
1601 .or_else(|| d("thread_id"))
1602 .or_else(|| r("delivery_thread_id")),
1603 );
1604 if mode.as_deref() == Some("none") {
1605 return Some(Target::Local);
1606 }
1607 if matches!(channel.as_deref(), Some("last") | Some("origin")) {
1608 return Some(Target::Origin);
1609 }
1610 let Some(channel) = channel else {
1611 return if mode.is_some() {
1612 Some(Target::Local)
1613 } else {
1614 None
1615 };
1616 };
1617 Some(Target::Explicit {
1618 platform: channel,
1619 chat_id: to,
1620 thread_id: thread,
1621 })
1622}
1623
1624fn failure_target(record: &Map<String, Value>, row: &Map<String, Value>) -> Option<Target> {
1625 let f = record.get("failureDelivery").and_then(Value::as_object);
1626 let channel = f
1627 .and_then(|f| f.get("channel"))
1628 .or_else(|| row.get("failure_delivery_channel"))
1629 .filter(|v| !v.is_null())
1630 .cloned();
1631 let mode = f
1632 .and_then(|f| f.get("mode"))
1633 .or_else(|| row.get("failure_delivery_mode"))
1634 .filter(|v| !v.is_null())
1635 .cloned();
1636 if channel.is_none() && mode.is_none() {
1637 return None;
1638 }
1639 let mut synth = Map::new();
1640 if let Some(m) = mode {
1641 synth.insert("mode".into(), m);
1642 }
1643 if let Some(c) = channel {
1644 synth.insert("channel".into(), c);
1645 }
1646 if let Some(t) = f
1647 .and_then(|f| f.get("to"))
1648 .or_else(|| row.get("failure_delivery_to"))
1649 .filter(|v| !v.is_null())
1650 {
1651 synth.insert("to".into(), t.clone());
1652 }
1653 delivery_target(Some(&synth), None)
1654}
1655
1656fn origin_of(row: &Map<String, Value>) -> Option<JobOrigin> {
1657 let key = row.get("owner_session_key").and_then(Value::as_str)?;
1658 let parsed = parse_openclaw_session_key(key)?;
1659 Some(JobOrigin {
1660 platform: parsed.key.platform.unwrap_or_default(),
1661 chat_type: parsed.key.kind,
1662 chat_id: parsed.key.chat_id,
1663 thread_id: parsed.key.thread_id,
1664 })
1665}
1666
1667pub fn decode_job(row: &Map<String, Value>) -> Job {
1669 let record = json_object_of(row.get("job_json"));
1670 let state_json = json_object_of(row.get("state_json"));
1671 let raw_schedule = record
1672 .get("schedule")
1673 .and_then(Value::as_object)
1674 .cloned()
1675 .unwrap_or_default();
1676 let mut schedule_residue = raw_schedule.clone();
1677 let kind = raw_schedule.get("kind").and_then(Value::as_str);
1678 let schedule =
1679 if kind == Some("every") && raw_schedule.get("everyMs").is_some_and(Value::is_number) {
1680 schedule_residue.remove("kind");
1681 schedule_residue.remove("everyMs");
1682 Schedule::Interval {
1683 minutes: raw_schedule["everyMs"].as_f64().unwrap() / 60000.0,
1684 }
1685 } else if kind == Some("cron") && raw_schedule.get("expr").is_some_and(Value::is_string) {
1686 schedule_residue.remove("kind");
1687 schedule_residue.remove("expr");
1688 schedule_residue.remove("tz");
1689 Schedule::Cron {
1690 expr: raw_schedule["expr"].as_str().unwrap().into(),
1691 tz: raw_schedule
1692 .get("tz")
1693 .and_then(Value::as_str)
1694 .map(str::to_string),
1695 }
1696 } else if kind == Some("at") && raw_schedule.get("at").is_some_and(Value::is_string) {
1697 schedule_residue.remove("kind");
1698 schedule_residue.remove("at");
1699 Schedule::Once {
1700 run_at: raw_schedule["at"].as_str().unwrap().into(),
1701 }
1702 } else {
1703 schedule_residue.insert("__unmapped".into(), Value::Bool(true));
1704 Schedule::Cron {
1705 expr: String::new(),
1706 tz: None,
1707 }
1708 };
1709 let payload = record
1710 .get("payload")
1711 .and_then(Value::as_object)
1712 .cloned()
1713 .unwrap_or_default();
1714 let prompt = ["message", "text", "command", "script"]
1715 .iter()
1716 .find_map(|k| payload.get(*k).and_then(Value::as_str))
1717 .map(str::to_string);
1718 let delivery = record
1722 .get("delivery")
1723 .and_then(Value::as_object)
1724 .cloned()
1725 .or_else(|| {
1726 let mut d = Map::new();
1727 for (column, key) in [
1728 ("delivery_mode", "mode"),
1729 ("delivery_channel", "channel"),
1730 ("delivery_to", "to"),
1731 ("delivery_thread_id", "threadId"),
1732 ("delivery_account_id", "accountId"),
1733 ] {
1734 if let Some(v) = row.get(column).filter(|v| !v.is_null()) {
1735 if v.as_str().is_some_and(str::is_empty) {
1736 continue;
1737 }
1738 d.insert(key.into(), v.clone());
1739 }
1740 }
1741 (!d.is_empty()).then_some(d)
1742 });
1743 let session_target = record
1744 .get("sessionTarget")
1745 .and_then(Value::as_str)
1746 .map(str::to_string)
1747 .or_else(|| {
1748 row.get("session_target")
1749 .and_then(Value::as_str)
1750 .map(str::to_string)
1751 });
1752 let job_id = opt_text(row.get("job_id")).unwrap_or_default();
1753 let mut residue = Residue::default();
1754 residue.keep(
1755 "__session_target",
1756 session_target.map(Value::String).unwrap_or(Value::Null),
1757 );
1758 residue.keep(
1759 "name",
1760 record
1761 .get("name")
1762 .and_then(Value::as_str)
1763 .map(|s| Value::String(s.into()))
1764 .unwrap_or_else(|| {
1765 row.get("name")
1766 .filter(|v| !v.is_null())
1767 .cloned()
1768 .unwrap_or(Value::String(job_id.clone()))
1769 }),
1770 );
1771 for (k, v) in &record {
1772 if [
1773 "id",
1774 "name",
1775 "enabled",
1776 "schedule",
1777 "payload",
1778 "delivery",
1779 "sessionTarget",
1780 "createdAtMs",
1781 "state",
1782 ]
1783 .contains(&k.as_str())
1784 {
1785 continue;
1786 }
1787 residue.keep(k.clone(), v.clone());
1788 }
1789 if !schedule_residue.is_empty() {
1790 residue.keep("__schedule", Value::Object(schedule_residue));
1791 }
1792 if !payload.is_empty() {
1793 residue.keep("__payload", Value::Object(payload.clone()));
1794 }
1795 if let Some(d) = &delivery {
1796 residue.keep("__delivery", Value::Object(d.clone()));
1797 }
1798 if let Some(f) = record.get("failureDelivery").filter(|v| v.is_object()) {
1799 residue.keep("__failure_delivery", f.clone());
1800 }
1801 if !state_json.is_empty() {
1802 residue.keep("__state", Value::Object(state_json.clone()));
1803 }
1804 let enabled = match row.get("enabled") {
1805 None | Some(Value::Null) => true,
1806 Some(Value::Bool(b)) => *b,
1807 Some(Value::Number(n)) => n.as_f64() != Some(0.0),
1808 Some(other) => !matches!(other, Value::String(s) if s.is_empty()),
1809 };
1810 Job {
1811 id: job_id,
1812 schedule,
1813 prompt,
1814 workdir: None,
1815 model: payload
1816 .get("model")
1817 .and_then(Value::as_str)
1818 .map(str::to_string)
1819 .or_else(|| {
1820 row.get("payload_model")
1821 .and_then(Value::as_str)
1822 .map(str::to_string)
1823 }),
1824 skills: Vec::new(),
1825 context_from: None,
1826 deliver: delivery_target(delivery.as_ref(), Some(row)).unwrap_or(Target::Local),
1827 failure_deliver: failure_target(&record, row),
1828 origin: origin_of(row),
1829 attach_to_session: None,
1830 machine: None,
1831 repeat: None,
1832 enabled,
1833 next_run_at: iso_seconds(
1834 ms_of(row.get("next_run_at_ms").filter(|v| !v.is_null()))
1835 .or_else(|| ms_of(state_json.get("nextRunAtMs"))),
1836 ),
1837 last_run_at: iso_seconds(
1838 ms_of(row.get("last_run_at_ms").filter(|v| !v.is_null()))
1839 .or_else(|| ms_of(state_json.get("lastRunAtMs"))),
1840 ),
1841 last_status: row
1842 .get("last_run_status")
1843 .and_then(Value::as_str)
1844 .map(str::to_string),
1845 created_at: iso_seconds(
1846 ms_of(row.get("created_at_ms").filter(|v| !v.is_null()))
1847 .or_else(|| ms_of(record.get("createdAtMs"))),
1848 ),
1849 residue,
1850 }
1851}
1852
1853fn empty_row(columns: &[&str]) -> Map<String, Value> {
1854 columns
1855 .iter()
1856 .map(|c| (c.to_string(), Value::Null))
1857 .collect()
1858}
1859
1860pub fn encode_job_row(
1862 job: &Job,
1863 raw: Option<&Map<String, Value>>,
1864 store_key: &str,
1865) -> Map<String, Value> {
1866 let residue = &job.residue.0;
1867 let mut row = raw.cloned().unwrap_or_else(|| empty_row(CRON_JOB_COLUMNS));
1868 let mut schedule = Map::new();
1869 match &job.schedule {
1870 Schedule::Interval { minutes } => {
1871 schedule.insert("kind".into(), "every".into());
1872 schedule.insert(
1873 "everyMs".into(),
1874 Value::from((minutes * 60000.0).round() as i64),
1875 );
1876 }
1877 Schedule::Cron { expr, tz } => {
1878 schedule.insert("kind".into(), "cron".into());
1879 schedule.insert("expr".into(), Value::String(expr.clone()));
1880 if let Some(tz) = tz {
1881 schedule.insert("tz".into(), Value::String(tz.clone()));
1882 }
1883 }
1884 Schedule::Once { run_at } => {
1885 schedule.insert("kind".into(), "at".into());
1886 schedule.insert("at".into(), Value::String(run_at.clone()));
1887 }
1888 }
1889 if let Some(Value::Object(extra)) = residue.get("__schedule") {
1890 for (k, v) in extra {
1891 schedule.insert(k.clone(), v.clone());
1892 }
1893 }
1894 schedule.remove("__unmapped");
1895 let mut record = Map::new();
1896 record.insert("id".into(), Value::String(job.id.clone()));
1897 record.insert(
1898 "name".into(),
1899 residue
1900 .get("name")
1901 .cloned()
1902 .unwrap_or(Value::String(job.id.clone())),
1903 );
1904 record.insert("enabled".into(), Value::Bool(job.enabled));
1905 if let Some(ms) = ms_from_iso(job.created_at.as_deref()) {
1906 record.insert("createdAtMs".into(), Value::from(ms));
1907 }
1908 record.insert("schedule".into(), Value::Object(schedule.clone()));
1909 if let Some(st) = residue.get("__session_target").filter(|v| !v.is_null()) {
1910 record.insert("sessionTarget".into(), st.clone());
1911 }
1912 match residue.get("__payload").and_then(Value::as_object) {
1913 Some(p) => {
1914 let mut p = p.clone();
1915 if let Some(prompt) = &job.prompt {
1916 for k in ["message", "text", "command", "script"] {
1917 if p.contains_key(k) {
1918 p.insert(k.into(), Value::String(prompt.clone()));
1919 break;
1920 }
1921 }
1922 }
1923 record.insert("payload".into(), Value::Object(p));
1924 }
1925 None => {
1926 if let Some(prompt) = &job.prompt {
1927 record.insert(
1928 "payload".into(),
1929 serde_json::json!({"kind": "agentTurn", "message": prompt}),
1930 );
1931 }
1932 }
1933 }
1934 if let Some(d) = residue.get("__delivery") {
1935 record.insert("delivery".into(), d.clone());
1936 }
1937 if let Some(f) = residue.get("__failure_delivery") {
1938 record.insert("failureDelivery".into(), f.clone());
1939 }
1940 record.insert(
1941 "state".into(),
1942 residue
1943 .get("__state")
1944 .cloned()
1945 .unwrap_or_else(|| Value::Object(Map::new())),
1946 );
1947 for (k, v) in residue {
1948 if !k.starts_with("__") {
1949 record.insert(k.clone(), v.clone());
1950 }
1951 }
1952 let get = |m: &Map<String, Value>, k: &str| m.get(k).filter(|v| !v.is_null()).cloned();
1953 let delivery = record
1954 .get("delivery")
1955 .and_then(Value::as_object)
1956 .cloned()
1957 .unwrap_or_default();
1958 let payload = record
1959 .get("payload")
1960 .and_then(Value::as_object)
1961 .cloned()
1962 .unwrap_or_default();
1963 row.insert(
1964 "store_key".into(),
1965 get(&row, "store_key").unwrap_or(Value::String(store_key.into())),
1966 );
1967 row.insert("job_id".into(), Value::String(job.id.clone()));
1968 row.insert("name".into(), record["name"].clone());
1969 row.insert("enabled".into(), Value::from(i64::from(job.enabled)));
1970 row.insert(
1971 "created_at_ms".into(),
1972 record
1973 .get("createdAtMs")
1974 .cloned()
1975 .or_else(|| get(&row, "created_at_ms"))
1976 .unwrap_or(Value::from(0)),
1977 );
1978 row.insert(
1979 "schedule_kind".into(),
1980 schedule
1981 .get("kind")
1982 .cloned()
1983 .unwrap_or(Value::String("cron".into())),
1984 );
1985 row.insert(
1986 "schedule_expr".into(),
1987 schedule.get("expr").cloned().unwrap_or(Value::Null),
1988 );
1989 row.insert(
1990 "schedule_tz".into(),
1991 schedule.get("tz").cloned().unwrap_or(Value::Null),
1992 );
1993 row.insert(
1994 "every_ms".into(),
1995 schedule.get("everyMs").cloned().unwrap_or(Value::Null),
1996 );
1997 row.insert(
1998 "anchor_ms".into(),
1999 schedule.get("anchorMs").cloned().unwrap_or(Value::Null),
2000 );
2001 row.insert(
2002 "at".into(),
2003 schedule.get("at").cloned().unwrap_or(Value::Null),
2004 );
2005 row.insert(
2006 "session_target".into(),
2007 record
2008 .get("sessionTarget")
2009 .cloned()
2010 .or_else(|| get(&row, "session_target"))
2011 .unwrap_or(Value::String("isolated".into())),
2012 );
2013 row.insert(
2014 "wake_mode".into(),
2015 get(&row, "wake_mode").unwrap_or(Value::String("now".into())),
2016 );
2017 row.insert(
2018 "payload_kind".into(),
2019 payload
2020 .get("kind")
2021 .cloned()
2022 .or_else(|| get(&row, "payload_kind"))
2023 .unwrap_or(Value::String("agentTurn".into())),
2024 );
2025 row.insert(
2026 "payload_message".into(),
2027 job.prompt.clone().map(Value::String).unwrap_or(Value::Null),
2028 );
2029 row.insert(
2030 "payload_model".into(),
2031 job.model.clone().map(Value::String).unwrap_or(Value::Null),
2032 );
2033 row.insert(
2034 "delivery_mode".into(),
2035 delivery
2036 .get("mode")
2037 .or_else(|| delivery.get("kind"))
2038 .cloned()
2039 .unwrap_or(Value::Null),
2040 );
2041 row.insert(
2042 "delivery_channel".into(),
2043 delivery.get("channel").cloned().unwrap_or(Value::Null),
2044 );
2045 row.insert(
2046 "delivery_to".into(),
2047 delivery.get("to").cloned().unwrap_or(Value::Null),
2048 );
2049 row.insert(
2050 "delivery_thread_id".into(),
2051 delivery.get("threadId").cloned().unwrap_or(Value::Null),
2052 );
2053 row.insert(
2054 "delivery_account_id".into(),
2055 delivery.get("accountId").cloned().unwrap_or(Value::Null),
2056 );
2057 row.insert(
2058 "next_run_at_ms".into(),
2059 ms_from_iso(job.next_run_at.as_deref())
2060 .map(Value::from)
2061 .unwrap_or(Value::Null),
2062 );
2063 row.insert(
2064 "last_run_at_ms".into(),
2065 ms_from_iso(job.last_run_at.as_deref())
2066 .map(Value::from)
2067 .unwrap_or(Value::Null),
2068 );
2069 row.insert(
2070 "last_run_status".into(),
2071 job.last_status
2072 .clone()
2073 .map(Value::String)
2074 .unwrap_or(Value::Null),
2075 );
2076 row.insert(
2077 "job_json".into(),
2078 Value::String(serde_json::to_string(&Value::Object(record.clone())).unwrap()),
2079 );
2080 row.insert(
2081 "state_json".into(),
2082 Value::String(
2083 serde_json::to_string(record.get("state").unwrap_or(&Value::Object(Map::new())))
2084 .unwrap(),
2085 ),
2086 );
2087 row.insert(
2088 "sort_order".into(),
2089 get(&row, "sort_order").unwrap_or(Value::from(0)),
2090 );
2091 row.insert(
2092 "updated_at".into(),
2093 get(&row, "updated_at")
2094 .or_else(|| get(&row, "created_at_ms"))
2095 .unwrap_or(Value::from(0)),
2096 );
2097 row
2098}
2099
2100fn run_status(word: Option<&str>) -> FireStatus {
2103 match word {
2104 Some("ok") => FireStatus::Succeeded,
2105 Some("error") => FireStatus::Failed,
2106 _ => FireStatus::Unknown,
2107 }
2108}
2109
2110fn run_status_back(status: FireStatus) -> &'static str {
2111 match status {
2112 FireStatus::Succeeded | FireStatus::Claimed | FireStatus::Running => "ok",
2113 FireStatus::Failed | FireStatus::Timeout => "error",
2114 FireStatus::Unknown => "skipped",
2115 }
2116}
2117
2118const FIRE_MAPPED: &[&str] = &[
2119 "job_id",
2120 "seq",
2121 "ts",
2122 "error",
2123 "run_id",
2124 "run_at_ms",
2125 "session_id",
2126];
2127
2128pub fn decode_fire(row: &Map<String, Value>) -> Fire {
2130 let job_id = opt_text(row.get("job_id")).unwrap_or_default();
2131 let id = opt_text(row.get("run_id")).unwrap_or_else(|| {
2132 format!(
2133 "{job_id}#{}",
2134 row.get("seq").map(|v| v.to_string()).unwrap_or_default()
2135 )
2136 });
2137 let started = iso_millis(ms_of(row.get("run_at_ms").filter(|v| !v.is_null())));
2138 let finished = iso_millis(ms_of(row.get("ts").filter(|v| !v.is_null())));
2139 let mut residue = Residue::default();
2140 residue.keep("__no_claim", Value::Bool(true));
2141 for (k, v) in row {
2142 if !FIRE_MAPPED.contains(&k.as_str()) && !v.is_null() {
2143 residue.keep(k.clone(), v.clone());
2144 }
2145 }
2146 Fire {
2147 id,
2148 job_id,
2149 session_id: opt_text(row.get("session_id")),
2150 status: run_status(row.get("status").and_then(Value::as_str)),
2151 claimed_at: started
2152 .clone()
2153 .or_else(|| finished.clone())
2154 .unwrap_or_default(),
2155 started_at: started,
2156 finished_at: finished,
2157 error: opt_text(row.get("error")),
2158 obligation_id: None,
2159 residue,
2160 }
2161}
2162
2163pub fn encode_fire_row(
2165 fire: &Fire,
2166 raw: Option<&Map<String, Value>>,
2167 store_key: &str,
2168) -> Map<String, Value> {
2169 let residue = &fire.residue.0;
2170 let mut row = raw
2171 .cloned()
2172 .unwrap_or_else(|| empty_row(CRON_RUN_LOG_COLUMNS));
2173 let get = |m: &Map<String, Value>, k: &str| m.get(k).filter(|v| !v.is_null()).cloned();
2174 row.insert(
2175 "store_key".into(),
2176 get(&row, "store_key")
2177 .or_else(|| residue.get("store_key").cloned())
2178 .unwrap_or(Value::String(store_key.into())),
2179 );
2180 row.insert("job_id".into(), Value::String(fire.job_id.clone()));
2181 row.insert(
2182 "seq".into(),
2183 get(&row, "seq")
2184 .or_else(|| residue.get("seq").cloned())
2185 .unwrap_or(Value::Null),
2186 );
2187 let ts = ms_from_iso(fire.finished_at.as_deref())
2188 .map(Value::from)
2189 .or_else(|| get(&row, "ts"))
2190 .unwrap_or(Value::from(0));
2191 row.insert("ts".into(), ts.clone());
2192 row.insert(
2193 "status".into(),
2194 residue
2195 .get("status")
2196 .cloned()
2197 .unwrap_or(Value::String(run_status_back(fire.status).into())),
2198 );
2199 row.insert(
2200 "error".into(),
2201 fire.error.clone().map(Value::String).unwrap_or(Value::Null),
2202 );
2203 row.insert(
2204 "delivery_status".into(),
2205 residue
2206 .get("delivery_status")
2207 .cloned()
2208 .unwrap_or(Value::Null),
2209 );
2210 row.insert(
2211 "delivery_error".into(),
2212 residue
2213 .get("delivery_error")
2214 .cloned()
2215 .unwrap_or(Value::Null),
2216 );
2217 row.insert(
2218 "delivered".into(),
2219 residue.get("delivered").cloned().unwrap_or(Value::Null),
2220 );
2221 row.insert(
2222 "session_id".into(),
2223 fire.session_id
2224 .clone()
2225 .map(Value::String)
2226 .unwrap_or(Value::Null),
2227 );
2228 row.insert(
2229 "session_key".into(),
2230 residue.get("session_key").cloned().unwrap_or(Value::Null),
2231 );
2232 row.insert(
2233 "run_id".into(),
2234 if fire.id.contains('#') {
2235 Value::Null
2236 } else {
2237 Value::String(fire.id.clone())
2238 },
2239 );
2240 row.insert(
2241 "run_at_ms".into(),
2242 ms_from_iso(fire.started_at.as_deref())
2243 .map(Value::from)
2244 .unwrap_or(Value::Null),
2245 );
2246 row.insert("entry_json".into(), residue.get("entry_json").cloned().unwrap_or_else(|| Value::String(serde_json::json!({"action": "finished", "ts": ts, "jobId": fire.job_id, "status": row["status"]}).to_string())));
2247 row.insert(
2248 "created_at".into(),
2249 residue.get("created_at").cloned().unwrap_or(ts),
2250 );
2251 for c in CRON_RUN_LOG_COLUMNS {
2252 if !row.contains_key(*c) {
2253 row.insert(
2254 c.to_string(),
2255 residue.get(*c).cloned().unwrap_or(Value::Null),
2256 );
2257 }
2258 }
2259 row
2260}
2261
2262fn obl_state(native: &str) -> ObligationState {
2265 match native {
2266 "pending" | "queued" | "sending" | "in_flight" | "retrying" => ObligationState::Pending,
2267 "delivered" | "sent" => ObligationState::Sent,
2268 "failed" => ObligationState::Failed,
2269 "dropped" | "dead" | "cancelled" | "canceled" => ObligationState::Dropped,
2270 _ => ObligationState::Pending,
2271 }
2272}
2273
2274const OBL_MAPPED: &[&str] = &[
2275 "id",
2276 "status",
2277 "session_key",
2278 "channel",
2279 "target",
2280 "retry_count",
2281 "last_error",
2282 "enqueued_at",
2283 "updated_at",
2284];
2285
2286pub fn decode_obligation(row: &Map<String, Value>) -> Obligation {
2288 let parsed = row
2289 .get("session_key")
2290 .and_then(Value::as_str)
2291 .and_then(parse_openclaw_session_key);
2292 let native = row
2293 .get("status")
2294 .and_then(Value::as_str)
2295 .map(|s| s.to_ascii_lowercase())
2296 .unwrap_or_default();
2297 let state = obl_state(&native);
2298 let mut content = OutboundContent {
2299 text: String::new(),
2300 attachments: None,
2301 reply_to: None,
2302 format: None,
2303 };
2304 let mut residue = Residue::default();
2305 residue.keep("status", row.get("status").cloned().unwrap_or(Value::Null));
2306 if let Some(entry) = row
2307 .get("entry_json")
2308 .and_then(Value::as_str)
2309 .and_then(|s| serde_json::from_str::<Value>(s).ok())
2310 .and_then(|v| v.as_object().cloned())
2311 {
2312 if let Some(t) = ["text", "message", "content", "body"]
2313 .iter()
2314 .find_map(|k| entry.get(*k).and_then(Value::as_str))
2315 {
2316 content.text = t.into();
2317 }
2318 }
2319 for (k, v) in row {
2320 if !OBL_MAPPED.contains(&k.as_str()) && !v.is_null() {
2321 residue.keep(k.clone(), v.clone());
2322 }
2323 }
2324 let text = |k: &str| {
2325 row.get(k).and_then(|v| match v {
2326 Value::String(s) => Some(s.clone()),
2327 Value::Number(n) => Some(n.to_string()),
2328 _ => None,
2329 })
2330 };
2331 let created_at = text("enqueued_at").unwrap_or_else(|| "0".into());
2332 let updated_at = text("updated_at").unwrap_or_else(|| created_at.clone());
2333 Obligation {
2334 id: opt_text(row.get("id")).unwrap_or_default(),
2335 target: SurfaceKey {
2336 key: None,
2337 platform: Some(
2338 text("channel")
2339 .filter(|s| !s.is_empty())
2340 .or_else(|| parsed.as_ref().and_then(|p| p.key.platform.clone()))
2341 .unwrap_or_default(),
2342 ),
2343 kind: parsed.as_ref().and_then(|p| p.key.kind.clone()),
2344 chat_id: Some(text("target").unwrap_or_default()),
2345 thread_id: parsed.as_ref().and_then(|p| p.key.thread_id.clone()),
2346 participant_id: None,
2347 },
2348 session_key: row
2349 .get("session_key")
2350 .and_then(Value::as_str)
2351 .map(str::to_string),
2352 content,
2353 state,
2354 attempts: ms_of(row.get("retry_count")).unwrap_or(0).max(0) as u64,
2355 last_error: row
2356 .get("last_error")
2357 .and_then(Value::as_str)
2358 .map(str::to_string),
2359 delivered_at: if state == ObligationState::Sent {
2360 iso_millis(
2361 ms_of(row.get("updated_at").filter(|v| !v.is_null()))
2362 .or_else(|| ms_of(row.get("enqueued_at"))),
2363 )
2364 } else {
2365 None
2366 },
2367 created_at,
2368 updated_at,
2369 posted: None,
2370 source: match parsed {
2371 Some(p) => ObligationSource::Turn { key: Some(p.key) },
2372 None => ObligationSource::Turn { key: None },
2373 },
2374 residue,
2375 }
2376}
2377
2378pub fn encode_obligation_row(
2380 o: &Obligation,
2381 raw: Option<&Map<String, Value>>,
2382) -> Map<String, Value> {
2383 let residue = &o.residue.0;
2384 let mut row = raw
2385 .cloned()
2386 .unwrap_or_else(|| empty_row(DELIVERY_QUEUE_COLUMNS));
2387 let get = |m: &Map<String, Value>, k: &str| m.get(k).filter(|v| !v.is_null()).cloned();
2388 row.insert(
2389 "queue_name".into(),
2390 get(&row, "queue_name")
2391 .or_else(|| residue.get("queue_name").cloned())
2392 .unwrap_or(Value::String("default".into())),
2393 );
2394 row.insert("id".into(), Value::String(o.id.clone()));
2395 row.insert(
2396 "status".into(),
2397 residue
2398 .get("status")
2399 .filter(|v| !v.is_null())
2400 .cloned()
2401 .unwrap_or(Value::String(o.state.hermes_word().into())),
2402 );
2403 row.insert(
2404 "session_key".into(),
2405 o.session_key
2406 .clone()
2407 .map(Value::String)
2408 .unwrap_or(Value::Null),
2409 );
2410 row.insert(
2411 "channel".into(),
2412 o.target
2413 .platform
2414 .clone()
2415 .filter(|p| !p.is_empty())
2416 .map(Value::String)
2417 .unwrap_or(Value::Null),
2418 );
2419 row.insert(
2420 "target".into(),
2421 o.target
2422 .chat_id
2423 .clone()
2424 .filter(|c| !c.is_empty())
2425 .map(Value::String)
2426 .unwrap_or(Value::Null),
2427 );
2428 row.insert(
2429 "last_error".into(),
2430 o.last_error
2431 .clone()
2432 .map(Value::String)
2433 .unwrap_or(Value::Null),
2434 );
2435 let enq = o.created_at.parse::<f64>().map(|f| f as i64).unwrap_or(0);
2436 row.insert("enqueued_at".into(), Value::from(enq));
2437 row.insert(
2438 "updated_at".into(),
2439 Value::from(o.updated_at.parse::<f64>().map(|f| f as i64).unwrap_or(enq)),
2440 );
2441 row.insert("entry_json".into(), residue.get("entry_json").cloned().unwrap_or_else(|| Value::String(serde_json::json!({"kind": residue.get("entry_kind").cloned().unwrap_or(Value::String("message".into())), "text": o.content.text}).to_string())));
2442 for c in DELIVERY_QUEUE_COLUMNS {
2443 if row.get(*c).is_none_or(Value::is_null) {
2444 row.insert(
2445 c.to_string(),
2446 residue.get(*c).cloned().unwrap_or(Value::Null),
2447 );
2448 }
2449 }
2450 row.insert("retry_count".into(), Value::from(o.attempts));
2451 row
2452}
2453
2454fn decode_conversation_binding(row: &Map<String, Value>) -> Option<Binding> {
2463 let text = |k: &str| opt_text(row.get(k));
2464 let channel = text("channel")?;
2465 let conversation_id = text("conversation_id")?;
2466 let kind = text("conversation_kind");
2467 let parent = text("parent_conversation_id");
2468 let (chat_id, thread_id) = if kind.as_deref() == Some("thread") && parent.is_some() {
2469 (parent.clone(), Some(conversation_id.clone()))
2470 } else {
2471 (Some(conversation_id.clone()), None)
2472 };
2473 let status = text("status");
2474 let started_at = iso_millis(ms_of(row.get("bound_at").filter(|v| !v.is_null())));
2475 let updated_at = iso_millis(ms_of(row.get("updated_at").filter(|v| !v.is_null())));
2476 let mut residue = Residue::default();
2477 residue.keep("__observed", Value::Bool(true));
2478 for k in [
2479 "binding_id",
2480 "binding_key",
2481 "account_id",
2482 "target_kind",
2483 "target_session_key",
2484 "status",
2485 "expires_at",
2486 ] {
2487 if let Some(v) = row.get(k).filter(|v| !v.is_null()) {
2488 residue.keep(k.to_string(), v.clone());
2489 }
2490 }
2491 for k in ["metadata_json", "record_json"] {
2492 if let Some(t) = text(k) {
2493 let parsed = serde_json::from_str::<Value>(&t).unwrap_or(Value::String(t));
2494 residue.keep(k.trim_end_matches("_json").to_string(), parsed);
2495 }
2496 }
2497 if let Some(agent) = text("target_agent_id") {
2498 residue.keep("agent_id", Value::String(agent));
2499 }
2500 Some(Binding {
2501 key: SurfaceKey {
2502 key: text("target_session_key"),
2503 platform: Some(channel),
2504 kind,
2505 chat_id,
2506 thread_id,
2507 participant_id: None,
2508 },
2509 profile: None,
2510 worker: Worker {
2511 harness: HarnessId::new(HarnessId::OPENCLAW),
2512 session_id: text("target_session_id"),
2513 locator: None,
2514 },
2515 trigger: Trigger::Channel,
2516 recurrence: None,
2517 handoff: None,
2518 started_at: started_at.clone(),
2519 last_activity_at: updated_at.clone().or(started_at),
2520 ended_at: if status.as_deref().is_some_and(|s| s != "active") {
2521 updated_at
2522 } else {
2523 None
2524 },
2525 end_reason: None,
2526 residue,
2527 })
2528}
2529
2530fn bindings_for_agent(agent_dir: &Path, agent_id: &str) -> Result<Vec<Binding>> {
2531 let mut out = Vec::new();
2532 let sessions = agent_dir.join("sessions");
2533 if !sessions.is_dir() {
2534 return Ok(out);
2535 }
2536 let mut entries: Vec<_> = fs::read_dir(&sessions)?.flatten().collect();
2537 entries.sort_by_key(|e| e.file_name());
2538 for entry in entries {
2539 let name = entry.file_name().to_string_lossy().into_owned();
2540 if !name.ends_with(".jsonl") || name.ends_with(".trajectory.jsonl") {
2541 continue;
2542 }
2543 let path = entry.path();
2544 let Ok(text) = fs::read_to_string(&path) else {
2545 continue;
2546 };
2547 let lines: Vec<&str> = text.lines().filter(|l| !l.trim().is_empty()).collect();
2548 let Some(first) = lines.first() else { continue };
2549 let Ok(header) = serde_json::from_str::<Value>(first) else {
2550 continue;
2551 };
2552 if header.get("type").and_then(Value::as_str) != Some("session") {
2553 continue;
2554 }
2555 let Some(key) = header
2556 .get("sessionKey")
2557 .or_else(|| header.get("__openclaw").and_then(|o| o.get("sessionKey")))
2558 .and_then(Value::as_str)
2559 else {
2560 continue;
2561 };
2562 let Some(parsed) = parse_openclaw_session_key(key) else {
2563 continue;
2564 };
2565 let header_ts = header
2566 .get("timestamp")
2567 .and_then(Value::as_str)
2568 .map(str::to_string);
2569 let mut last = header_ts.clone();
2570 for line in lines.iter().rev() {
2571 if let Ok(rec) = serde_json::from_str::<Value>(line) {
2572 if let Some(ts) = rec.get("timestamp").and_then(Value::as_str) {
2573 last = Some(ts.into());
2574 break;
2575 }
2576 }
2577 }
2578 let mut residue = Residue(parsed.residue.into_iter().collect());
2579 residue.keep(
2580 "agent_id",
2581 Value::String(parsed.agent.clone().unwrap_or_else(|| agent_id.into())),
2582 );
2583 out.push(Binding {
2584 trigger: if parsed.recurrence.is_some() {
2585 Trigger::Cron
2586 } else {
2587 Trigger::Channel
2588 },
2589 key: parsed.key,
2590 profile: None,
2591 worker: Worker {
2592 harness: HarnessId::new(HarnessId::OPENCLAW),
2593 session_id: Some(
2594 header
2595 .get("id")
2596 .and_then(Value::as_str)
2597 .map(str::to_string)
2598 .unwrap_or_else(|| name.trim_end_matches(".jsonl").into()),
2599 ),
2600 locator: Some(path.display().to_string()),
2601 },
2602 recurrence: parsed.recurrence,
2603 handoff: None,
2604 started_at: header_ts.clone().or_else(|| last.clone()),
2605 last_activity_at: last.or(header_ts),
2606 ended_at: None,
2607 end_reason: None,
2608 residue,
2609 });
2610 }
2611 Ok(out)
2612}
2613
2614fn list_unmodeled(state_dir: &Path) -> Result<Vec<String>> {
2615 let mut out = Vec::new();
2616 fn walk(base: &Path, dir: &Path, out: &mut Vec<String>) -> Result<()> {
2617 let mut entries: Vec<_> = fs::read_dir(dir)?.flatten().collect();
2618 entries.sort_by_key(|e| e.file_name());
2619 for entry in entries {
2620 let name = entry.file_name().to_string_lossy().into_owned();
2621 if name == "node_modules" || name == ".git" {
2622 continue;
2623 }
2624 let p = entry.path();
2625 let rel = p
2626 .strip_prefix(base)
2627 .unwrap_or(&p)
2628 .to_string_lossy()
2629 .replace('\\', "/");
2630 let Ok(st) = fs::symlink_metadata(&p) else {
2631 continue;
2632 };
2633 if st.is_dir() {
2634 walk(base, &p, out)?;
2635 continue;
2636 }
2637 if rel == OPENCLAW_CONFIG
2638 || rel == OPENCLAW_STATE_DB
2639 || rel.starts_with(&format!("{OPENCLAW_STATE_DB}-"))
2640 || rel == "cron/jobs.json"
2641 {
2642 continue;
2643 }
2644 out.push(rel);
2645 }
2646 Ok(())
2647 }
2648 walk(state_dir, state_dir, &mut out)?;
2649 Ok(out)
2650}
2651
2652fn config_record(orchestration: &Orchestration) -> Value {
2655 let root = &orchestration.profiles["default"];
2656 serde_json::json!({
2657 "channels": root.channels, "routes": root.routes, "residue": root.residue.config,
2658 "profiles": orchestration.profiles.iter().map(|(n, p)| (n.clone(), serde_json::json!({"agent": p.residue.config.get("openclaw_agent"), "subscriptions": p.subscriptions}))).collect::<BTreeMap<_, _>>(),
2659 })
2660}
2661
2662fn store_record(profile: &Profile) -> Value {
2665 let jobs: BTreeMap<&String, &Job> = profile
2666 .jobs
2667 .iter()
2668 .filter(|(_, job)| !is_legacy_file_job(job))
2669 .collect();
2670 serde_json::json!({ "jobs": jobs, "fires": profile.fires, "obligations": profile.obligations })
2671}
2672
2673pub fn from_openclaw(state_dir: &Path) -> Result<OpenclawLoaded> {
2675 if !state_dir.is_dir() {
2676 return Err(load_error(
2677 &state_dir.display().to_string(),
2678 "",
2679 "not a directory",
2680 ));
2681 }
2682 let mut vault = BTreeMap::new();
2683 let config_path = state_dir.join(OPENCLAW_CONFIG);
2684 let config_text = if config_path.is_file() {
2685 Some(fs::read_to_string(&config_path)?)
2686 } else {
2687 None
2688 };
2689 let config = match &config_text {
2690 Some(t) => parse_json5(&config_path.display().to_string(), t)?,
2691 None => Value::Object(Map::new()),
2692 };
2693 if !config.is_object() {
2694 return Err(load_error(
2695 &config_path.display().to_string(),
2696 "",
2697 "expected a JSON5 object",
2698 ));
2699 }
2700 let entries = agent_entries(&config);
2701 let default_id = default_agent_id(&config);
2702 let declared: Vec<(String, Value)> = if entries.is_empty() {
2703 vec![(default_id.clone(), Value::Object(Map::new()))]
2704 } else {
2705 entries
2706 };
2707 let name_for = |agent_id: &str| -> String {
2708 if agent_id == default_id {
2709 "default".into()
2710 } else {
2711 agent_id.into()
2712 }
2713 };
2714
2715 let mut profiles: BTreeMap<String, Profile> = BTreeMap::new();
2716 let mut ios: BTreeMap<String, OpenclawProfileIo> = BTreeMap::new();
2717 for (id, entry) in &declared {
2718 let name = name_for(id);
2719 let dir = state_dir.join("agents").join(id);
2720 let mut profile = empty_profile(&name, &dir);
2721 profile.residue.config.insert(
2722 "openclaw_agent".into(),
2723 redact_secrets(entry, &format!("agent_{id}"), &mut vault),
2724 );
2725 if let Ok(text) = fs::read_to_string(dir.join("AGENTS.md")) {
2726 profile.persona = Some(persona_ref(&text));
2727 }
2728 for b in bindings_for_agent(&dir, id)? {
2729 profile.bindings.insert(surface_key_string(&b.key), b);
2730 }
2731 profiles.insert(name.clone(), profile);
2732 ios.insert(
2733 name,
2734 OpenclawProfileIo {
2735 agent_id: id.clone(),
2736 source_dir: dir,
2737 ..Default::default()
2738 },
2739 );
2740 }
2741 let channels = decode_channels(&config, &mut vault);
2743 let routes = decode_routes(&config, &name_for);
2744 let mut rest = Map::new();
2745 for (k, v) in config.as_object().unwrap() {
2746 if !["channels", "bindings", "agents", "hooks"].contains(&k.as_str()) {
2747 rest.insert(k.clone(), v.clone());
2748 }
2749 }
2750 let mut agents_rest = Map::new();
2751 if let Some(a) = config.get("agents").and_then(Value::as_object) {
2752 for (k, v) in a {
2753 if k != "list" && k != "entries" {
2754 agents_rest.insert(k.clone(), v.clone());
2755 }
2756 }
2757 }
2758 let hooks = decode_hooks(&config, &mut vault, &name_for, &mut profiles);
2759 {
2760 let root = profiles.get_mut("default").unwrap();
2761 root.channels = channels;
2762 root.routes = routes;
2763 root.residue.config.insert("openclaw".into(), serde_json::json!({
2764 "default_agent": default_id,
2765 "agents_form": agents_form(&config),
2766 "agents_rest": redact_secrets(&Value::Object(agents_rest), "agents", &mut vault),
2767 "hooks": hooks.as_ref().map(|h| serde_json::json!({"block": h.block, "has_token": h.has_token})),
2768 "rest": redact_secrets(&Value::Object(rest), "config", &mut vault),
2769 }));
2770 }
2771 let mut root_io = OpenclawRootIo {
2772 state_dir: state_dir.to_path_buf(),
2773 config_raw: config_text.clone().unwrap_or_default(),
2774 config_present: config_text.is_some(),
2775 default_agent: default_id.clone(),
2776 ..Default::default()
2777 };
2778
2779 let db_path = state_dir.join(OPENCLAW_STATE_DB);
2781 let mut store_key = state_dir.join("cron/jobs.json").display().to_string();
2782 let mut job_owner: BTreeMap<String, String> = BTreeMap::new();
2783 if table_exists(&db_path, "cron_jobs") {
2784 for row in read_rows(
2785 &db_path,
2786 "select * from cron_jobs order by sort_order, created_at_ms, job_id",
2787 &[],
2788 )?
2789 .unwrap_or_default()
2790 {
2791 let job_id = opt_text(row.get("job_id")).unwrap_or_default();
2792 if let Some(k) = row
2793 .get("store_key")
2794 .and_then(Value::as_str)
2795 .filter(|s| !s.is_empty())
2796 {
2797 store_key = k.into();
2798 }
2799 let record = json_object_of(row.get("job_json"));
2800 let agent_id = record
2801 .get("agentId")
2802 .or_else(|| row.get("agent_id"))
2803 .or_else(|| row.get("owner_agent_id"))
2804 .and_then(Value::as_str)
2805 .unwrap_or(&default_id)
2806 .to_string();
2807 let owner = if profiles.contains_key(&name_for(&agent_id)) {
2808 name_for(&agent_id)
2809 } else {
2810 "default".into()
2811 };
2812 let job = decode_job(&row);
2813 root_io.cron_jobs.insert(job_id.clone(), row);
2814 profiles
2815 .get_mut(&owner)
2816 .unwrap()
2817 .jobs
2818 .insert(job.id.clone(), job);
2819 job_owner.insert(job_id, owner);
2820 }
2821 }
2822 let legacy_path = state_dir.join("cron/jobs.json");
2826 if let Ok(text) = fs::read_to_string(&legacy_path) {
2827 let (records, object_form) = legacy_job_records(&text);
2828 root_io.legacy_jobs_raw = Some(text);
2829 root_io.legacy_jobs_object_form = object_form;
2830 for record in records {
2831 let Some(id) = legacy_job_id(&record) else {
2832 continue;
2833 };
2834 if profiles.values().any(|p| p.jobs.contains_key(&id)) {
2835 continue;
2836 }
2837 let agent_id = record
2838 .get("agentId")
2839 .or_else(|| record.get("agent_id"))
2840 .and_then(Value::as_str)
2841 .unwrap_or(&default_id)
2842 .to_string();
2843 let owner = if profiles.contains_key(&name_for(&agent_id)) {
2844 name_for(&agent_id)
2845 } else {
2846 "default".into()
2847 };
2848 let job = decode_legacy_job(&record);
2849 root_io.legacy_jobs.insert(id.clone(), record);
2850 profiles
2851 .get_mut(&owner)
2852 .unwrap()
2853 .jobs
2854 .insert(id.clone(), job);
2855 job_owner.insert(id, owner);
2856 }
2857 }
2858 if table_exists(&db_path, "cron_run_logs") {
2859 for row in read_rows(
2860 &db_path,
2861 "select * from cron_run_logs order by ts, seq",
2862 &[],
2863 )?
2864 .unwrap_or_default()
2865 {
2866 let fire = decode_fire(&row);
2867 root_io.cron_run_logs.insert(fire.id.clone(), row);
2868 let owner = job_owner
2869 .get(&fire.job_id)
2870 .cloned()
2871 .unwrap_or_else(|| "default".into());
2872 profiles.get_mut(&owner).unwrap().fires.push(fire);
2873 }
2874 }
2875 if table_exists(&db_path, "delivery_queue_entries") {
2876 for row in read_rows(
2877 &db_path,
2878 "select * from delivery_queue_entries order by enqueued_at, id",
2879 &[],
2880 )?
2881 .unwrap_or_default()
2882 {
2883 let o = decode_obligation(&row);
2884 root_io.delivery_queue_entries.insert(o.id.clone(), row);
2885 let owner = o
2886 .session_key
2887 .as_deref()
2888 .and_then(parse_openclaw_session_key)
2889 .and_then(|p| p.agent)
2890 .map(|a| name_for(&a))
2891 .filter(|n| profiles.contains_key(n))
2892 .unwrap_or_else(|| "default".into());
2893 profiles.get_mut(&owner).unwrap().obligations.push(o);
2894 }
2895 }
2896 if table_exists(&db_path, "current_conversation_bindings") {
2902 for row in read_rows(
2903 &db_path,
2904 "select * from current_conversation_bindings order by bound_at, binding_key",
2905 &[],
2906 )?
2907 .unwrap_or_default()
2908 {
2909 let Some(b) = decode_conversation_binding(&row) else {
2910 continue;
2911 };
2912 let agent_id =
2913 opt_text(row.get("target_agent_id")).unwrap_or_else(|| default_id.clone());
2914 let owner = if profiles.contains_key(&name_for(&agent_id)) {
2915 name_for(&agent_id)
2916 } else {
2917 "default".into()
2918 };
2919 profiles
2920 .get_mut(&owner)
2921 .unwrap()
2922 .bindings
2923 .insert(surface_key_string(&b.key), b);
2924 }
2925 }
2926 if table_exists(&db_path, "schema_meta") {
2927 root_io.schema_meta =
2928 read_rows(&db_path, "select * from schema_meta order by meta_key", &[])?
2929 .unwrap_or_default();
2930 }
2931 root_io.store_key = store_key;
2932 root_io.db_present = db_path.exists();
2933
2934 let mut fire_by_key: BTreeMap<String, (String, String, Option<String>)> = BTreeMap::new();
2937 for (name, profile) in &profiles {
2938 for fire in &profile.fires {
2939 let Some(key) = fire.residue.0.get("session_key").and_then(Value::as_str) else {
2940 continue;
2941 };
2942 let finished = fire.finished_at.clone();
2943 match fire_by_key.get(key) {
2944 Some((_, _, held))
2945 if held.clone().unwrap_or_default() > finished.clone().unwrap_or_default() => {}
2946 _ => {
2947 fire_by_key.insert(key.into(), (name.clone(), fire.id.clone(), finished));
2948 }
2949 }
2950 }
2951 }
2952 let mut links: Vec<(String, String, String)> = Vec::new(); for profile in profiles.values_mut() {
2954 for o in &mut profile.obligations {
2955 let Some((owner, fire_id, _)) =
2956 o.session_key.as_deref().and_then(|k| fire_by_key.get(k))
2957 else {
2958 continue;
2959 };
2960 o.source = ObligationSource::Fire {
2961 fire_id: fire_id.clone(),
2962 };
2963 links.push((owner.clone(), fire_id.clone(), o.id.clone()));
2964 }
2965 }
2966 for (owner, fire_id, obligation_id) in links {
2967 if let Some(fire) = profiles
2968 .get_mut(&owner)
2969 .and_then(|p| p.fires.iter_mut().find(|f| f.id == fire_id))
2970 {
2971 fire.obligation_id = Some(obligation_id);
2972 }
2973 }
2974
2975 let mut orchestration = Orchestration {
2976 root: state_dir.to_path_buf(),
2977 profiles,
2978 workflow: crate::orchestration::hermes_instance(),
2979 };
2980 root_io.config_snapshot = canonical_json(&config_record(&orchestration));
2981 for (name, profile) in orchestration.profiles.iter_mut() {
2982 let io = ios.get_mut(name).unwrap();
2983 io.store_snapshot = canonical_json(&store_record(profile));
2984 io.bindings_snapshot = canonical_json(&serde_json::to_value(&profile.bindings).unwrap());
2985 if name != "default" {
2986 profile.residue.files = Vec::new();
2987 }
2988 }
2989 orchestration
2990 .profiles
2991 .get_mut("default")
2992 .unwrap()
2993 .residue
2994 .files = list_unmodeled(state_dir)?;
2995 Ok(OpenclawLoaded {
2996 orchestration,
2997 vault,
2998 root: root_io,
2999 profiles: ios,
3000 })
3001}
3002
3003#[derive(Debug, Clone, Default)]
3007pub struct OpenclawReport {
3008 pub written: Vec<ArtifactFidelity>,
3010 pub refused: Vec<super::hermes::Refusal>,
3012 pub notes: Vec<String>,
3014 pub rows_byte: usize,
3016 pub rows_emitted: usize,
3018}
3019
3020fn write_atomic(path: &Path, text: &str) -> Result<()> {
3021 if let Some(parent) = path.parent() {
3022 fs::create_dir_all(parent)?;
3023 }
3024 let tmp = path.with_file_name(format!(
3025 "{}.tmp-{}",
3026 path.file_name().unwrap().to_string_lossy(),
3027 std::process::id()
3028 ));
3029 fs::write(&tmp, text)?;
3030 fs::rename(&tmp, path)?;
3031 Ok(())
3032}
3033
3034fn encode_config(loaded: &OpenclawLoaded) -> Value {
3035 let orchestration = &loaded.orchestration;
3036 let root = &orchestration.profiles["default"];
3037 let own = root
3038 .residue
3039 .config
3040 .get("openclaw")
3041 .and_then(Value::as_object)
3042 .cloned()
3043 .unwrap_or_default();
3044 let default_id = own
3045 .get("default_agent")
3046 .and_then(Value::as_str)
3047 .map(str::to_string)
3048 .unwrap_or_else(|| loaded.root.default_agent.clone());
3049 let id_for_name = |name: &str| -> String {
3050 if name == "default" {
3051 default_id.clone()
3052 } else {
3053 loaded
3054 .profiles
3055 .get(name)
3056 .map(|io| io.agent_id.clone())
3057 .unwrap_or_else(|| name.into())
3058 }
3059 };
3060 let mut pairs: Vec<(String, Value)> = inline_secrets(
3061 own.get("rest").unwrap_or(&Value::Object(Map::new())),
3062 &loaded.vault,
3063 )
3064 .as_object()
3065 .map(|m| m.iter().map(|(k, v)| (k.clone(), v.clone())).collect())
3066 .unwrap_or_default();
3067 pairs.push((
3068 "channels".into(),
3069 ordered_object(encode_channels(root, &loaded.vault)),
3070 ));
3071 let agent_blocks: Vec<(String, Vec<(String, Value)>)> = orchestration
3072 .profiles
3073 .iter()
3074 .map(|(name, p)| {
3075 let id = id_for_name(name);
3076 let entry = inline_secrets(
3077 p.residue
3078 .config
3079 .get("openclaw_agent")
3080 .unwrap_or(&Value::Object(Map::new())),
3081 &loaded.vault,
3082 );
3083 let mut block: Vec<(String, Value)> = vec![("id".into(), Value::String(id.clone()))];
3084 if let Some(m) = entry.as_object() {
3085 for (k, v) in m {
3086 if k != "id" {
3087 block.push((k.clone(), v.clone()));
3088 }
3089 }
3090 }
3091 (id, block)
3092 })
3093 .collect();
3094 let mut agents: Vec<(String, Value)> = inline_secrets(
3095 own.get("agents_rest").unwrap_or(&Value::Object(Map::new())),
3096 &loaded.vault,
3097 )
3098 .as_object()
3099 .map(|m| m.iter().map(|(k, v)| (k.clone(), v.clone())).collect())
3100 .unwrap_or_default();
3101 if own.get("agents_form").and_then(Value::as_str) == Some("entries") {
3102 agents.push((
3103 "entries".into(),
3104 ordered_object(
3105 agent_blocks
3106 .into_iter()
3107 .map(|(id, block)| {
3108 (
3109 id,
3110 ordered_object(block.into_iter().filter(|(k, _)| k != "id").collect()),
3111 )
3112 })
3113 .collect(),
3114 ),
3115 ));
3116 } else {
3117 agents.push((
3118 "list".into(),
3119 Value::Array(
3120 agent_blocks
3121 .into_iter()
3122 .map(|(_, block)| ordered_object(block))
3123 .collect(),
3124 ),
3125 ));
3126 }
3127 pairs.push(("agents".into(), ordered_object(agents)));
3128 pairs.push((
3129 "bindings".into(),
3130 Value::Array(encode_routes(root, &id_for_name)),
3131 ));
3132 let hooks_meta = own
3133 .get("hooks")
3134 .and_then(Value::as_object)
3135 .map(|h| HooksMeta {
3136 block: h
3137 .get("block")
3138 .and_then(Value::as_object)
3139 .cloned()
3140 .unwrap_or_default(),
3141 has_token: h.get("has_token").and_then(Value::as_bool).unwrap_or(false),
3142 });
3143 if let Some(hooks) = encode_hooks(orchestration, hooks_meta.as_ref(), &loaded.vault) {
3144 pairs.push(("hooks".into(), hooks));
3145 }
3146 ordered_object(pairs)
3147}
3148
3149fn plain(value: &Value) -> Value {
3151 if let Some(pairs) = is_ordered_pairs(value) {
3152 return Value::Object(pairs.into_iter().map(|(k, v)| (k, plain(&v))).collect());
3153 }
3154 match value {
3155 Value::Array(items) => Value::Array(items.iter().map(plain).collect()),
3156 Value::Object(m) => Value::Object(m.iter().map(|(k, v)| (k.clone(), plain(v))).collect()),
3157 other => other.clone(),
3158 }
3159}
3160
3161fn resequence(rows: &mut [Map<String, Value>]) {
3162 let key = |r: &Map<String, Value>| {
3163 format!(
3164 "{}\u{0}{}",
3165 r.get("store_key")
3166 .map(|v| v.to_string())
3167 .unwrap_or_default(),
3168 r.get("job_id").map(|v| v.to_string()).unwrap_or_default()
3169 )
3170 };
3171 let mut taken: BTreeMap<String, Vec<i64>> = BTreeMap::new();
3172 for row in rows.iter_mut() {
3173 let Some(seq) = row.get("seq").and_then(Value::as_i64) else {
3174 continue;
3175 };
3176 let seen = taken.entry(key(row)).or_default();
3177 if seen.contains(&seq) {
3178 row.insert("seq".into(), Value::Null);
3179 continue;
3180 }
3181 seen.push(seq);
3182 }
3183 for row in rows.iter_mut() {
3184 if row.get("seq").is_some_and(|v| !v.is_null()) {
3185 continue;
3186 }
3187 let seen = taken.entry(key(row)).or_default();
3188 let mut next = 1;
3189 while seen.contains(&next) {
3190 next += 1;
3191 }
3192 row.insert("seq".into(), Value::from(next));
3193 seen.push(next);
3194 }
3195}
3196
3197fn insert_for(table: &str, columns: &[&str]) -> String {
3198 format!(
3199 "insert into {table} ({}) values ({})",
3200 columns.join(", "),
3201 columns.iter().map(|_| "?").collect::<Vec<_>>().join(",")
3202 )
3203}
3204
3205fn params_of(row: &Map<String, Value>, columns: &[&str]) -> Vec<Param> {
3206 columns
3207 .iter()
3208 .map(|c| Param::from(row.get(*c).unwrap_or(&Value::Null)))
3209 .collect()
3210}
3211
3212fn legacy_job_records(text: &str) -> (Vec<Value>, bool) {
3216 match serde_json::from_str::<Value>(text) {
3217 Ok(Value::Array(items)) => (items, false),
3218 Ok(Value::Object(map)) => (
3219 map.get("jobs")
3220 .and_then(Value::as_array)
3221 .cloned()
3222 .unwrap_or_default(),
3223 true,
3224 ),
3225 _ => (Vec::new(), false),
3226 }
3227}
3228
3229fn legacy_job_id(record: &Value) -> Option<String> {
3230 ["id", "job_id", "jobId"]
3231 .iter()
3232 .find_map(|k| record.get(*k).and_then(Value::as_str))
3233 .filter(|s| !s.is_empty())
3234 .map(str::to_string)
3235}
3236
3237fn is_legacy_file_job(job: &Job) -> bool {
3238 job.residue.0.get("__legacy_file") == Some(&Value::Bool(true))
3239}
3240
3241fn decode_legacy_job(record: &Value) -> Job {
3247 let mut modern = record.as_object().cloned().unwrap_or_default();
3248 let text = |v: Option<&Value>| match v {
3249 Some(Value::String(s)) if !s.is_empty() => Some(s.clone()),
3250 Some(Value::Number(n)) => Some(n.to_string()),
3251 _ => None,
3252 };
3253 let ms = |v: &Value| match v {
3254 Value::Number(n) => n.as_i64(),
3255 Value::String(s) => iso_epoch(s).map(|secs| (secs * 1000.0) as i64),
3256 _ => None,
3257 };
3258 let schedule = match modern.get("schedule") {
3259 Some(Value::Object(o)) if o.get("kind").is_some() => Some(Value::Object(o.clone())),
3260 Some(Value::String(expr)) => Some(serde_json::json!({ "kind": "cron", "expr": expr })),
3261 _ => None,
3262 }
3263 .or_else(|| {
3264 text(modern.get("cron")).map(|expr| serde_json::json!({ "kind": "cron", "expr": expr }))
3265 })
3266 .or_else(|| {
3267 modern
3268 .get("everyMinutes")
3269 .and_then(Value::as_f64)
3270 .map(|m| serde_json::json!({ "kind": "every", "everyMs": (m * 60_000.0) as i64 }))
3271 })
3272 .or_else(|| {
3273 modern
3274 .get("everyMs")
3275 .and_then(Value::as_f64)
3276 .map(|ms| serde_json::json!({ "kind": "every", "everyMs": ms as i64 }))
3277 })
3278 .or_else(|| {
3279 text(modern.get("runAt").or_else(|| modern.get("run_at")))
3280 .map(|at| serde_json::json!({ "kind": "at", "at": at }))
3281 });
3282 for k in ["cron", "everyMinutes", "everyMs", "runAt", "run_at"] {
3283 modern.remove(k);
3284 }
3285 if let Some(schedule) = schedule {
3286 modern.insert("schedule".into(), schedule);
3287 }
3288 match modern.get("delivery").cloned() {
3289 Some(Value::String(mode)) => {
3290 modern.insert("delivery".into(), serde_json::json!({ "mode": mode }));
3291 }
3292 Some(Value::Object(mut d)) => {
3293 if !d.contains_key("mode") {
3294 if let Some(kind) = d.remove("kind").or_else(|| d.remove("type")) {
3295 d.insert("mode".into(), kind);
3296 }
3297 }
3298 modern.insert("delivery".into(), Value::Object(d));
3299 }
3300 _ => {}
3301 }
3302 let mut row = Map::new();
3303 row.insert(
3304 "job_id".into(),
3305 Value::String(legacy_job_id(record).unwrap_or_default()),
3306 );
3307 if let Some(name) = modern.get("name").cloned() {
3308 row.insert("name".into(), name);
3309 }
3310 if let Some(enabled) = modern.get("enabled").cloned() {
3311 row.insert("enabled".into(), enabled);
3312 }
3313 for (legacy, column) in [
3314 ("nextRunAt", "next_run_at_ms"),
3315 ("next_run_at", "next_run_at_ms"),
3316 ("lastRunAt", "last_run_at_ms"),
3317 ("last_run_at", "last_run_at_ms"),
3318 ] {
3319 if let Some(v) = modern.remove(legacy) {
3320 if let Some(at) = ms(&v) {
3321 row.insert(column.into(), Value::from(at));
3322 }
3323 }
3324 }
3325 if let Some(status) = modern
3326 .remove("lastStatus")
3327 .or_else(|| modern.remove("last_status"))
3328 {
3329 row.insert("last_run_status".into(), status);
3330 }
3331 if let Some(v) = modern
3332 .remove("createdAt")
3333 .or_else(|| modern.remove("created_at"))
3334 {
3335 if let Some(at) = ms(&v) {
3336 modern.insert("createdAtMs".into(), Value::from(at));
3337 }
3338 }
3339 row.insert(
3341 "job_json".into(),
3342 Value::String(serde_json::to_string(&Value::Object(modern)).unwrap()),
3343 );
3344 let mut job = decode_job(&row);
3345 job.residue.keep("__legacy_file", Value::Bool(true));
3346 job
3347}
3348
3349fn write_legacy_jobs(
3353 loaded: &OpenclawLoaded,
3354 dest: &Path,
3355 report: &mut OpenclawReport,
3356) -> Result<()> {
3357 let jobs: Vec<&Job> = loaded
3358 .orchestration
3359 .profiles
3360 .values()
3361 .flat_map(|p| p.jobs.values())
3362 .filter(|j| is_legacy_file_job(j))
3363 .collect();
3364 if jobs.is_empty() && loaded.root.legacy_jobs_raw.is_none() {
3365 return Ok(());
3366 }
3367 let target = dest.join("cron/jobs.json");
3368 if let Some(parent) = target.parent() {
3369 fs::create_dir_all(parent)?;
3370 }
3371 let unchanged = loaded.root.legacy_jobs_raw.is_some()
3372 && jobs.len() == loaded.root.legacy_jobs.len()
3373 && jobs.iter().all(|job| {
3374 loaded.root.legacy_jobs.get(&job.id).is_some_and(|record| {
3375 canonical_json(&serde_json::to_value(decode_legacy_job(record)).unwrap())
3376 == canonical_json(&serde_json::to_value(job).unwrap())
3377 })
3378 });
3379 if unchanged {
3380 fs::write(
3381 &target,
3382 loaded.root.legacy_jobs_raw.as_deref().unwrap_or(""),
3383 )?;
3384 report
3385 .written
3386 .push(ArtifactFidelity::byte("cron/jobs.json"));
3387 return Ok(());
3388 }
3389 if jobs.is_empty() {
3390 return Ok(());
3391 }
3392 let store_key = target.display().to_string();
3393 let records: Vec<Value> = jobs
3394 .iter()
3395 .map(|job| {
3396 let raw = loaded
3397 .root
3398 .legacy_jobs
3399 .get(&job.id)
3400 .and_then(Value::as_object);
3401 match encode_job_row(job, raw, &store_key).remove("job_json") {
3402 Some(Value::String(text)) => {
3403 serde_json::from_str(&text).unwrap_or(Value::String(text))
3404 }
3405 Some(other) => other,
3406 None => Value::Null,
3407 }
3408 })
3409 .collect();
3410 let body = if loaded.root.legacy_jobs_object_form {
3411 serde_json::json!({ "jobs": records })
3412 } else {
3413 Value::Array(records)
3414 };
3415 fs::write(
3416 &target,
3417 format!("{}\n", serde_json::to_string_pretty(&body).unwrap()),
3418 )?;
3419 report.written.push(ArtifactFidelity::semantic(
3420 "cron/jobs.json",
3421 vec!["re-emitted: a file job changed".into()],
3422 ));
3423 Ok(())
3424}
3425
3426fn write_store(loaded: &OpenclawLoaded, dest: &Path, report: &mut OpenclawReport) -> Result<()> {
3427 let store_key = if loaded.root.store_key.is_empty() {
3428 dest.join("cron/jobs.json").display().to_string()
3429 } else {
3430 loaded.root.store_key.clone()
3431 };
3432 let mut job_rows = Vec::new();
3433 let mut fire_rows = Vec::new();
3434 let mut obligation_rows = Vec::new();
3435 let (mut byte_rows, mut emitted) = (0usize, 0usize);
3436 for profile in loaded.orchestration.profiles.values() {
3437 for job in profile.jobs.values() {
3438 if is_legacy_file_job(job) {
3439 continue;
3440 }
3441 let original = loaded.root.cron_jobs.get(&job.id);
3442 let unchanged = original.is_some_and(|o| {
3443 canonical_json(&serde_json::to_value(decode_job(o)).unwrap())
3444 == canonical_json(&serde_json::to_value(job).unwrap())
3445 });
3446 if unchanged {
3447 job_rows.push(original.unwrap().clone());
3448 byte_rows += 1;
3449 } else {
3450 job_rows.push(encode_job_row(job, original, &store_key));
3451 emitted += 1;
3452 }
3453 }
3454 for fire in &profile.fires {
3455 let original = loaded.root.cron_run_logs.get(&fire.id);
3456 let unchanged = original.is_some_and(|o| {
3457 let mut d = decode_fire(o);
3458 d.obligation_id = fire.obligation_id.clone();
3459 canonical_json(&serde_json::to_value(d).unwrap())
3460 == canonical_json(&serde_json::to_value(fire).unwrap())
3461 });
3462 if unchanged {
3463 fire_rows.push(original.unwrap().clone());
3464 byte_rows += 1;
3465 } else {
3466 fire_rows.push(encode_fire_row(fire, original, &store_key));
3467 emitted += 1;
3468 }
3469 }
3470 resequence(&mut fire_rows);
3471 for o in &profile.obligations {
3472 let original = loaded.root.delivery_queue_entries.get(&o.id);
3473 let unchanged = original.is_some_and(|r| {
3474 canonical_json(&serde_json::to_value(decode_obligation(r)).unwrap())
3475 == canonical_json(&serde_json::to_value(o).unwrap())
3476 });
3477 if unchanged {
3478 obligation_rows.push(original.unwrap().clone());
3479 byte_rows += 1;
3480 } else {
3481 obligation_rows.push(encode_obligation_row(o, original));
3482 emitted += 1;
3483 }
3484 }
3485 }
3486 let target = dest.join(OPENCLAW_STATE_DB);
3487 fs::create_dir_all(dest.join("state"))?;
3488 let tmp = target.with_file_name(format!("openclaw.sqlite.tmp-{}", std::process::id()));
3489 let _ = fs::remove_file(&tmp);
3490 let meta: Vec<Map<String, Value>> = if loaded.root.schema_meta.is_empty() {
3491 vec![serde_json::from_value(serde_json::json!({"meta_key": "global", "role": "global", "schema_version": 1, "agent_id": null, "app_version": null, "created_at": 0, "updated_at": 0})).unwrap()]
3492 } else {
3493 loaded.root.schema_meta.clone()
3494 };
3495 write_table(
3496 &tmp,
3497 OPENCLAW_DDL,
3498 &insert_for("schema_meta", SCHEMA_META_COLUMNS),
3499 &meta
3500 .iter()
3501 .map(|r| params_of(r, SCHEMA_META_COLUMNS))
3502 .collect::<Vec<_>>(),
3503 )?;
3504 write_table(
3505 &tmp,
3506 "",
3507 &insert_for("cron_jobs", CRON_JOB_COLUMNS),
3508 &job_rows
3509 .iter()
3510 .map(|r| params_of(r, CRON_JOB_COLUMNS))
3511 .collect::<Vec<_>>(),
3512 )?;
3513 write_table(
3514 &tmp,
3515 "",
3516 &insert_for("cron_run_logs", CRON_RUN_LOG_COLUMNS),
3517 &fire_rows
3518 .iter()
3519 .map(|r| params_of(r, CRON_RUN_LOG_COLUMNS))
3520 .collect::<Vec<_>>(),
3521 )?;
3522 write_table(
3523 &tmp,
3524 "",
3525 &insert_for("delivery_queue_entries", DELIVERY_QUEUE_COLUMNS),
3526 &obligation_rows
3527 .iter()
3528 .map(|r| params_of(r, DELIVERY_QUEUE_COLUMNS))
3529 .collect::<Vec<_>>(),
3530 )?;
3531 fs::rename(&tmp, &target)?;
3532 report.written.push(ArtifactFidelity {
3533 path: OPENCLAW_STATE_DB.into(),
3534 fidelity: if emitted == 0 {
3535 Fidelity::ByteLossless
3536 } else {
3537 Fidelity::Semantic
3538 },
3539 loss: Vec::new(),
3540 });
3541 report.rows_byte = byte_rows;
3542 report.rows_emitted = emitted;
3543 Ok(())
3544}
3545
3546fn copy_unmodeled(
3547 files: &[String],
3548 src: &Path,
3549 into: &Path,
3550 report: &mut OpenclawReport,
3551 prefix: &str,
3552) -> Result<()> {
3553 for rel in files {
3554 let from = src.join(rel);
3555 if !from.exists() {
3556 continue;
3557 }
3558 let to = into.join(rel);
3559 if let Some(parent) = to.parent() {
3560 fs::create_dir_all(parent)?;
3561 }
3562 fs::copy(&from, &to)?;
3563 report
3564 .written
3565 .push(ArtifactFidelity::byte(format!("{prefix}{rel}")));
3566 }
3567 Ok(())
3568}
3569
3570pub fn to_openclaw(loaded: &OpenclawLoaded, dest: &Path) -> Result<OpenclawReport> {
3572 let mut report = OpenclawReport::default();
3573 let orchestration = &loaded.orchestration;
3574 if !orchestration.profiles.contains_key("default") {
3575 return Err(load_error(
3576 &dest.display().to_string(),
3577 "",
3578 "no `default` profile: an OpenClaw install always has a default agent",
3579 ));
3580 }
3581 fs::create_dir_all(dest)?;
3582
3583 let cfg_unchanged =
3585 loaded.root.config_snapshot == canonical_json(&config_record(orchestration));
3586 if cfg_unchanged && loaded.root.config_present {
3587 write_atomic(&dest.join(OPENCLAW_CONFIG), &loaded.root.config_raw)?;
3588 report.written.push(ArtifactFidelity::byte(OPENCLAW_CONFIG));
3589 } else {
3590 let encoded = encode_config(loaded);
3591 let missing = unresolved_refs(&plain(&encoded), "");
3592 if !missing.is_empty() {
3593 report.refused.push(super::hermes::Refusal { file: OPENCLAW_CONFIG.into(), reason: format!("the vault has no value for {}; openclaw would read the reference itself as the credential", missing.iter().map(|(p, r)| format!("{r} ({p})")).collect::<Vec<_>>().join(", ")) });
3594 } else {
3595 write_atomic(
3596 &dest.join(OPENCLAW_CONFIG),
3597 &format!("{}\n", pretty_ordered(&encoded, 0)),
3598 )?;
3599 report.written.push(ArtifactFidelity::semantic(OPENCLAW_CONFIG, vec!["re-emitted as JSON: openclaw reads it with a JSON5 parser, so it loads, but the source's comments, trailing commas and key order are gone".into()]));
3600 report.notes.push(format!("{OPENCLAW_CONFIG}: re-emitted as JSON; comments, trailing commas and key order are gone"));
3601 }
3602 }
3603
3604 let store_unchanged = orchestration.profiles.iter().all(|(n, p)| {
3606 loaded
3607 .profiles
3608 .get(n)
3609 .map(|io| io.store_snapshot == canonical_json(&store_record(p)))
3610 .unwrap_or(false)
3611 });
3612 let src_db = loaded.root.state_dir.join(OPENCLAW_STATE_DB);
3613 let has_rows = orchestration
3614 .profiles
3615 .values()
3616 .any(|p| !p.jobs.is_empty() || !p.fires.is_empty() || !p.obligations.is_empty());
3617 if store_unchanged && loaded.root.db_present && src_db.exists() {
3618 fs::create_dir_all(dest.join("state"))?;
3619 fs::copy(&src_db, dest.join(OPENCLAW_STATE_DB))?;
3620 report
3621 .written
3622 .push(ArtifactFidelity::byte(OPENCLAW_STATE_DB));
3623 } else if has_rows || loaded.root.db_present {
3624 write_store(loaded, dest, &mut report)?;
3625 }
3626 write_legacy_jobs(loaded, dest, &mut report)?;
3627
3628 for (name, profile) in &orchestration.profiles {
3630 if profile.bindings.is_empty() {
3631 continue;
3632 }
3633 let io = loaded.profiles.get(name);
3634 let snapshot = io
3635 .map(|io| io.bindings_snapshot.clone())
3636 .filter(|s| !s.is_empty());
3637 if snapshot.as_deref()
3638 == Some(canonical_json(&serde_json::to_value(&profile.bindings).unwrap()).as_str())
3639 {
3640 continue;
3641 }
3642 let agent = io
3643 .map(|io| io.agent_id.clone())
3644 .unwrap_or_else(|| name.clone());
3645 if snapshot.is_some() {
3646 report.refused.push(super::hermes::Refusal { file: format!("agents/{agent}/sessions/*.jsonl"), reason: "bindings changed since import; writing the change into an OpenClaw transcript is UNI-18 (the first write into another harness's live session store), not this codec's".into() });
3649 continue;
3650 }
3651 let sessions_dir = dest.join("agents").join(&agent).join("sessions");
3657 for (slot, b) in &profile.bindings {
3658 let id = b.worker.session_id.clone().unwrap_or_else(|| slot.clone());
3659 let file = sessions_dir.join(format!("{id}.jsonl"));
3660 let rel = format!("agents/{agent}/sessions/{id}.jsonl");
3661 if file.exists() {
3662 report.refused.push(super::hermes::Refusal { file: rel, reason: "the destination already holds this transcript; writing into a live OpenClaw session store is UNI-18, not this codec's".into() });
3663 continue;
3664 }
3665 fs::create_dir_all(&sessions_dir)?;
3666 let header = serde_json::json!({
3667 "type": "session",
3668 "version": 3,
3669 "id": id,
3670 "timestamp": b.started_at.clone().or_else(|| b.last_activity_at.clone()).unwrap_or_default(),
3671 "sessionKey": crate::ontology::render_openclaw_session_key(&agent, b),
3672 });
3673 fs::write(&file, format!("{header}\n"))?;
3674 report.written.push(ArtifactFidelity::semantic(rel, vec!["a fresh transcript header carrying the session key; the turns live in the worker's store and are not carried".into()]));
3675 }
3676 }
3677
3678 copy_unmodeled(
3680 &orchestration.profiles["default"].residue.files,
3681 &loaded.root.state_dir,
3682 dest,
3683 &mut report,
3684 "",
3685 )?;
3686 for (name, profile) in &orchestration.profiles {
3687 if name == "default" {
3688 continue;
3689 }
3690 let Some(io) = loaded.profiles.get(name) else {
3691 continue;
3692 };
3693 copy_unmodeled(
3694 &profile.residue.files,
3695 &io.source_dir,
3696 &dest.join("agents").join(&io.agent_id),
3697 &mut report,
3698 &format!("agents/{}/", io.agent_id),
3699 )?;
3700 }
3701 Ok(report)
3702}