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