Skip to main content

supercode_interchange/orchestration/codec/
openclaw.rs

1//! The OpenClaw codec (`openclaw.mjs`): an OpenClaw state directory ⇄ the
2//! orchestration. `agents.*` become profiles — the default agent IS the `default`
3//! profile, every other agent keeps its id — each rooted at its own
4//! `agents/<id>/`. Channels, bindings and hooks are install-wide and land on
5//! `default`; jobs, fires and obligations land on the profile their agent id
6//! or session key names. `openclaw.json` is JSON5 (comments and trailing
7//! commas survive the read; a re-emit is JSON, which OpenClaw's parser
8//! loads); `state/openclaw.sqlite` is the pinned v2026.7.1-2 schema, its
9//! rows written back column for column when their record is unchanged.
10
11use 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
32/// The config file.
33pub const OPENCLAW_CONFIG: &str = "openclaw.json";
34/// The shared state store.
35pub const OPENCLAW_STATE_DB: &str = "state/openclaw.sqlite";
36/// Keys an inline account id may sit under.
37pub const OPENCLAW_ACCOUNT_KEYS: &[&str] = &[
38    "accountId",
39    "account_id",
40    "account",
41    "teamId",
42    "appId",
43    "userId",
44];
45
46/// The pinned v2026.7.1-2 schema, transcribed column for column.
47pub 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"#;
207/// `cron_jobs` columns in order.
208pub 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];
285/// `cron_run_logs` columns in order.
286pub 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];
310/// `delivery_queue_entries` columns in order.
311pub 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];
330/// `schema_meta` columns in order.
331pub 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
341// ---------------------------------------------------------------- JSON5
342
343/// Reduce JSON5 to JSON: `//` and `/* */` comments outside strings, trailing
344/// commas before `}` / `]` (`profiles.rs::strip_json5`, which the Node codec
345/// transcribed; this is the same rule in the crate that owns the codec), and
346/// JSON5's key and string forms OpenClaw's own configs use: identifier keys
347/// (`agents: {…}`) are quoted and single-quoted strings become double-quoted.
348pub 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            // A single-quoted string: its `\'` is a quote, its `"` must be escaped.
356            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            // An identifier: a key when it follows `{` or `,` and a `:` follows it.
380            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
488/// Parse a JSON5 file's text.
489pub 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
494// ---------------------------------------------------------------- secrets
495
496/// Credential-SHAPED names OpenClaw uses for things that are not secrets.
497pub 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
512/// OpenClaw's own credential rule: the lower-cased key ENDS WITH a marker; only the name decides.
513pub 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
523/// The `.env` name for a credential: `OPENCLAW_<scope>_<key>`.
524pub 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
537/// Deep copy with every string under a credential-shaped key replaced by a `{dotenv}` ref, the value handed to the vault.
538pub 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
581/// Inverse of [`redact_secrets`]: `${NAME}` references become their values.
582pub 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
600/// Every `${NAME}` reference a write left unresolved, with its path.
601fn 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
628// ---------------------------------------------------------------- time
629
630fn 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
656/// Unix ms → RFC 3339 at SECOND precision (`jobs.rs::iso_from_ms`).
657pub 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
664/// Unix ms → RFC 3339 with millis (`sidecar::ms_to_rfc3339`).
665pub 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
672/// RFC 3339 (`YYYY-MM-DDTHH:MM:SS[.fff][Z|±HH:MM]`) → Unix ms; `None` for anything else.
673pub 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    // days from civil (Howard Hinnant)
712    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
733// ---------------------------------------------------------------- agents
734
735/// The declared agents as `[id, entry]` in declaration order.
736pub 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
762/// `list` | `entries` | none.
763pub 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
774/// The default agent: the entry flagged `default: true`, else the first declared, else `main`.
775pub 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// ---------------------------------------------------------------- session keys
786
787/// What an OpenClaw session key mapped to, in the orchestration's binding form.
788#[derive(Debug, Clone, PartialEq)]
789pub struct ParsedKey {
790    /// The agent the key names, if any.
791    pub agent: Option<String>,
792    /// The surface (a `main` key is the DM COLLAPSE: `{main, dm, main}`).
793    pub key: SurfaceKey,
794    /// What the mapping had to translate (`session_key`, `dm_collapse`, `chat_kind`, `thread_word`).
795    pub residue: Map<String, Value>,
796    /// The job a `cron:` key names.
797    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
813/// An OpenClaw gateway session key → the orchestration's surface key plus what the
814/// mapping had to translate (`openclaw.mjs::parseOpenclawSessionKey`). This is
815/// the WORLD's form; the spine's discovery nouns keep the key's own shape
816/// through `Binding::from_openclaw_key`.
817pub 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// ---------------------------------------------------------------- io
896
897/// Root bookkeeping for an OpenClaw state dir.
898#[derive(Debug, Clone, Default)]
899pub struct OpenclawRootIo {
900    /// The state dir.
901    pub state_dir: PathBuf,
902    /// `openclaw.json`'s bytes as read (empty when absent).
903    pub config_raw: String,
904    /// Whether `openclaw.json` existed.
905    pub config_present: bool,
906    /// Canonical JSON of the config record as loaded.
907    pub config_snapshot: String,
908    /// Store rows as read, by table and id.
909    pub cron_jobs: BTreeMap<String, Map<String, Value>>,
910    /// Run-log rows by fire id.
911    pub cron_run_logs: BTreeMap<String, Map<String, Value>>,
912    /// Queue rows by id.
913    pub delivery_queue_entries: BTreeMap<String, Map<String, Value>>,
914    /// `schema_meta` rows.
915    pub schema_meta: Vec<Map<String, Value>>,
916    /// The `store_key` the jobs carry.
917    pub store_key: String,
918    /// Whether the store existed.
919    pub db_present: bool,
920    /// The default agent id.
921    pub default_agent: String,
922    /// The legacy `cron/jobs.json` beside the store, byte for byte, when
923    /// one exists (an install the pin has not migrated yet).
924    pub legacy_jobs_raw: Option<String>,
925    /// Its records by id, as read.
926    pub legacy_jobs: BTreeMap<String, Value>,
927    /// Whether the file was `{"jobs": [...]}` rather than a bare array.
928    pub legacy_jobs_object_form: bool,
929}
930
931/// Per-profile bookkeeping.
932#[derive(Debug, Clone, Default)]
933pub struct OpenclawProfileIo {
934    /// The agent id.
935    pub agent_id: String,
936    /// `agents/<id>/`.
937    pub source_dir: PathBuf,
938    /// Canonical JSON of the store record (jobs, fires, obligations) as loaded.
939    pub store_snapshot: String,
940    /// Canonical JSON of the bindings as loaded.
941    pub bindings_snapshot: String,
942}
943
944/// A loaded OpenClaw state directory.
945#[derive(Debug, Clone)]
946pub struct OpenclawLoaded {
947    /// The orchestration.
948    pub orchestration: Orchestration,
949    /// Secret values by `.env` name; never in the orchestration.
950    pub vault: BTreeMap<String, String>,
951    /// Root bookkeeping.
952    pub root: OpenclawRootIo,
953    /// Per-profile bookkeeping by profile name.
954    pub profiles: BTreeMap<String, OpenclawProfileIo>,
955}
956
957impl OpenclawLoaded {
958    /// An orchestration that did not come from an OpenClaw store (our folder on its way
959    /// out): nothing to reuse, every artifact re-emitted, every binding refused.
960    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
1002// ---------------------------------------------------------------- channels
1003
1004/// `channels.<name>.accounts` as `[accountId, entry]`, sorted by id.
1005pub 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
1344// ---------------------------------------------------------------- routes
1345
1346const 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// ---------------------------------------------------------------- hooks
1467
1468/// What the `hooks` block carried beside its mappings.
1469#[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
1622// ---------------------------------------------------------------- jobs
1623
1624fn 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
1724/// `cron_jobs` row → Job (`openclaw.mjs::decodeJob`).
1725pub 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    // the store's writer keeps delivery in `job_json` AND in columns; a row
1776    // that has it only in the columns still has to re-emit as a delivery, so
1777    // the columns' object is what the residue keeps when `job_json` has none
1778    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
1914/// A job as a `cron_jobs` row, over the original row when one exists.
1915pub 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
2154// ---------------------------------------------------------------- fires
2155
2156fn 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
2182/// `cron_run_logs` row → Fire.
2183pub 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
2217/// A fire as a `cron_run_logs` row, over the original when one exists.
2218pub 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
2316// ---------------------------------------------------------------- obligations
2317
2318fn 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
2340/// `delivery_queue_entries` row → Obligation.
2341pub 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
2432/// An obligation as a `delivery_queue_entries` row, over the original when one exists.
2433pub 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
2508// ---------------------------------------------------------------- bindings
2509
2510/// One `current_conversation_bindings` row as a [`Binding`]: the conversation
2511/// ref is the surface (a thread's parent is the chat, the thread is the
2512/// thread), the target session is the worker, `bound_at` is when it started,
2513/// and a binding that is not `active` ended when the row last changed. The
2514/// row's own facts (binding id and key, account, target kind, status, expiry,
2515/// metadata, the full record) ride as residue.
2516fn 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
2710// ---------------------------------------------------------------- from
2711
2712fn 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
2720/// The rows the sqlite store holds for a profile; a legacy file's jobs are
2721/// the file's rows, not the store's.
2722fn 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
2731/// Compile an OpenClaw state directory into the orchestration.
2732pub 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    // `agents.defaults` is each agent's worker unless it sets its own: its model (a string, or an object whose
2774    // `primary` is the model) and its workspace land on every agent's worker; what the IR has no field for stays.
2775    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        // `agents.*` → a profile plus its agent (C8): the id names both, `name` is the label, `model` and
2795        // `workspace` are the profile's worker; anything else stays on the agent, verbatim.
2796        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    // install-wide config: channels, bindings, and everything else verbatim
2869    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                // what the agents' workers took (model's primary, workspace) is not residue
2882                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    // the shared state DB
2927    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    // the legacy file beside the store: jobs the pin's migration has not
2970    // moved yet, read by the vocabulary that file was written in; a store
2971    // row with the same id wins
2972    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    // the store's own binding record (the pin's `current_conversation_bindings`,
3044    // ORC-12 finding 7): which agent and session a conversation is bound to,
3045    // observed rather than inferred from transcript headers. A row replaces
3046    // the inferred binding on the same surface; a store without the table
3047    // keeps the inferred ones. Read, never written back (UNI-22).
3048    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    // a fire's delivery OUTCOME rides on its own run-log row; the queue entry on
3082    // the same conversation key is the obligation that outcome is about
3083    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(); // (profile, fire id, obligation id)
3100    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// ---------------------------------------------------------------- to
3156
3157/// What an OpenClaw decompile did.
3158#[derive(Debug, Clone, Default)]
3159pub struct OpenclawReport {
3160    /// Every artifact written, with its tier.
3161    pub written: Vec<ArtifactFidelity>,
3162    /// Every write refused, with the gate named.
3163    pub refused: Vec<super::hermes::Refusal>,
3164    /// What a semantic write gave up.
3165    pub notes: Vec<String>,
3166    /// Store rows written back column for column vs re-encoded.
3167    pub rows_byte: usize,
3168    /// Store rows re-encoded.
3169    pub rows_emitted: usize,
3170    /// What the export folded or dropped, listed before anything was written and kept beside it.
3171    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            // the entry from its agent (C8): label, the profile's worker, then what the IR has no field for
3231            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                // an object-form model's other keys (fallbacks) join the worker's model as its `primary`
3254                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
3323/// Flatten an ordered value back into plain JSON (for ref scans).
3324fn 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
3386// ---------------------------------------------------------------- legacy cron/jobs.json
3387
3388/// The records of a legacy `cron/jobs.json` (a bare array, or `{"jobs": [...]}`).
3389fn 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
3415/// A legacy file record as a `cron_jobs` row, in the vocabulary the pin's
3416/// migration reads: a bare cron string / `cron` / `everyMinutes` / `everyMs`
3417/// / `runAt` schedule becomes the typed `schedule` object, `delivery.kind`
3418/// is the mode, and the ISO state stamps become the row's millisecond
3419/// columns. The record is then one [`decode_job`], marked `__legacy_file`.
3420fn 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    // the column form: the store keeps the record as JSON text
3514    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
3523/// The legacy `cron/jobs.json` on the way out: the source bytes when every
3524/// file job is as it was read, else the file re-emitted from the jobs
3525/// (`job_json` records, in the file's own form).
3526fn 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
3744/// Write an OpenClaw state directory from the orchestration.
3745pub 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    // openclaw.json
3761    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    // state/openclaw.sqlite
3782    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    // bindings: read, never written back (UNI-22)
3806    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            // the source transcripts hold what the orchestration does not model; the
3824            // change has to be written INTO them, which is UNI-18's
3825            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        // our own bindings: an OpenClaw conversation's surface lives in the
3829        // session KEY inside the transcript header, so each binding becomes a
3830        // FRESH transcript whose header carries that key (the discovery door
3831        // reads `sessionKey`); the turns live in the worker's store and are
3832        // not carried. An existing transcript is never overwritten.
3833        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    // everything under the state dir the IR does not model
3856    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
3881/// An OpenClaw model setting's model: the string itself, or an object's `primary`.
3882fn 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
3889/// An object-form model setting without its `primary` (its fallbacks, say), or None when nothing else is left.
3890fn 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}