Skip to main content

supercode_interchange/ontology/
binding.rs

1//! The one conversation record (`docs/ONTOLOGY.md` §2.2): a conversation
2//! reached on a surface, routed to a profile, run by a worker, with a
3//! lifecycle. A Hermes `sessions` row, an OpenClaw session key and an
4//! orchestrator `bindings` row each decode to it ONCE, here; the discovery
5//! row's nouns are its projection ([`Binding::nouns`]).
6
7use schemars::JsonSchema;
8use serde::{Deserialize, Serialize};
9
10use super::residue::Residue;
11use super::surface::{CrossSurface, Recurrence, SurfaceKey, Trigger};
12use super::HarnessId;
13use crate::session::OrchestrationNouns;
14
15/// Why a binding ended (`docs/ORCHESTRATOR-IR.md` §2.5).
16#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
17#[serde(rename_all = "snake_case")]
18pub enum EndReason {
19    /// Idle expiry.
20    Idle,
21    /// The daily reset boundary.
22    Daily,
23    /// `/reset` (or the operator verb).
24    Reset,
25    /// `/new` (or the operator verb).
26    New,
27    /// Moved to another surface.
28    Handoff,
29    /// The conversation was transferred to another agent (G2).
30    Transfer,
31    /// The agent left the conversation.
32    Left,
33    /// The worker failed.
34    Error,
35}
36
37impl EndReason {
38    /// Parse the wire word; anything else is not an end reason.
39    pub fn parse(word: &str) -> Option<Self> {
40        Some(match word {
41            "idle" => Self::Idle,
42            "daily" => Self::Daily,
43            "reset" => Self::Reset,
44            "new" => Self::New,
45            "handoff" => Self::Handoff,
46            "transfer" => Self::Transfer,
47            "left" => Self::Left,
48            "error" => Self::Error,
49            _ => return None,
50        })
51    }
52
53    /// The wire word.
54    pub fn as_str(self) -> &'static str {
55        match self {
56            Self::Idle => "idle",
57            Self::Daily => "daily",
58            Self::Reset => "reset",
59            Self::New => "new",
60            Self::Handoff => "handoff",
61            Self::Transfer => "transfer",
62            Self::Left => "left",
63            Self::Error => "error",
64        }
65    }
66}
67
68/// The worker session a binding points at.
69#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
70pub struct Worker {
71    /// Which harness runs the conversation.
72    pub harness: HarnessId,
73    /// The harness's own session id, when one exists yet.
74    #[serde(default, skip_serializing_if = "Option::is_none")]
75    pub session_id: Option<String>,
76    /// Where that session's transcript can be read (a path or store address).
77    #[serde(default, skip_serializing_if = "Option::is_none")]
78    pub locator: Option<String>,
79}
80
81/// A conversation moved (or moving) to another surface — Hermes `handoff_*`.
82#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
83pub struct Handoff {
84    /// The destination platform.
85    #[serde(default, skip_serializing_if = "Option::is_none")]
86    pub to: Option<String>,
87    /// `pending` | `done` | `failed` (the source's own word, verbatim).
88    pub state: String,
89    #[serde(default, skip_serializing_if = "Option::is_none")]
90    /// The failure, when `state` is `failed`.
91    pub error: Option<String>,
92}
93
94impl Default for Worker {
95    fn default() -> Self {
96        Self {
97            harness: HarnessId::new(""),
98            session_id: None,
99            locator: None,
100        }
101    }
102}
103
104/// The conversation record.
105#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
106pub struct Binding {
107    /// The surface; degenerate (all `None`) on a terminal.
108    pub key: SurfaceKey,
109    /// The profile that owns it (`None` = the home's default).
110    #[serde(default)]
111    pub profile: Option<String>,
112    /// The conversation it holds a session in (G1/G5); `None` where the source records none.
113    #[serde(default, skip_serializing_if = "Option::is_none")]
114    pub conversation: Option<String>,
115    /// The agent whose session it is; with `conversation`, its key. `None` where the source records none.
116    #[serde(default, skip_serializing_if = "Option::is_none")]
117    pub agent: Option<String>,
118    /// The worker session.
119    pub worker: Worker,
120    /// Why the conversation exists. The orchestrator's own store keeps no
121    /// trigger column (it is derived from the key and recurrence), so the
122    /// wire may omit it.
123    #[serde(default)]
124    pub trigger: Trigger,
125    #[serde(default)]
126    /// The job this conversation is a fire of.
127    pub recurrence: Option<Recurrence>,
128    #[serde(default)]
129    /// Moved (or moving) to another surface.
130    pub handoff: Option<Handoff>,
131    /// RFC3339 instants where the source has them.
132    #[serde(default)]
133    pub started_at: Option<String>,
134    #[serde(default)]
135    /// Last activity, RFC3339.
136    pub last_activity_at: Option<String>,
137    #[serde(default)]
138    /// End, RFC3339, when ended.
139    pub ended_at: Option<String>,
140    #[serde(default)]
141    /// Why it ended.
142    pub end_reason: Option<EndReason>,
143    /// Source fields the record does not model, verbatim.
144    #[serde(default)]
145    pub residue: Residue,
146}
147
148impl Default for Binding {
149    fn default() -> Self {
150        Self {
151            key: SurfaceKey::default(),
152            profile: None,
153            worker: Worker::default(),
154            trigger: Trigger::Unknown,
155            recurrence: None,
156            conversation: None,
157            agent: None,
158            handoff: None,
159            started_at: None,
160            last_activity_at: None,
161            ended_at: None,
162            end_reason: None,
163            residue: Residue::default(),
164        }
165    }
166}
167
168impl Binding {
169    /// The surface as a discovery row shows it: `None` on a terminal, where the
170    /// key carries no key string, platform or chat id.
171    pub fn surface(&self) -> Option<SurfaceKey> {
172        let k = &self.key;
173        if k.key.is_some() || k.platform.is_some() || k.chat_id.is_some() {
174            Some(k.clone())
175        } else {
176            None
177        }
178    }
179
180    /// The discovery row's nouns. `workspace` is left for the row's own cwd
181    /// rule (`SessionMeta::workspace`), the one derivation that is not this
182    /// record's to make.
183    pub fn nouns(&self) -> OrchestrationNouns {
184        OrchestrationNouns {
185            trigger: Some(self.trigger),
186            surface: self.surface(),
187            profile: self.profile.clone(),
188            recurrence: self.recurrence.clone(),
189            cross_surface: self.handoff.as_ref().map(|h| CrossSurface {
190                state: h.state.clone(),
191                platform: h.to.clone(),
192                error: h.error.clone(),
193            }),
194            workspace: None,
195        }
196    }
197}
198
199// ------------------------------------------------------------- Hermes rows
200
201/// One Hermes `sessions` row, the columns a binding is made of. Empty strings
202/// are read as absent by the decoder.
203#[derive(Debug, Clone, Default)]
204pub struct HermesSessionRow {
205    /// `sessions.id`.
206    pub id: String,
207    /// `sessions.source`: `cli` | `tui` | `acp` | `api_server` | `cron` | `webhook` | a platform.
208    pub source: Option<String>,
209    /// `delegate` for a subagent child (Hermes lineage), else anything.
210    pub lineage_kind: Option<String>,
211    /// `sessions.session_key`, the gateway conversation key.
212    pub session_key: Option<String>,
213    /// `sessions.chat_id`.
214    pub chat_id: Option<String>,
215    /// `sessions.chat_type`.
216    pub chat_type: Option<String>,
217    /// `sessions.thread_id`.
218    pub thread_id: Option<String>,
219    /// `sessions.user_id` (the participant).
220    pub user_id: Option<String>,
221    /// `sessions.profile_name`, the routed profile.
222    pub profile_name: Option<String>,
223    /// `sessions.handoff_state`.
224    pub handoff_state: Option<String>,
225    /// `sessions.handoff_platform`.
226    pub handoff_platform: Option<String>,
227    /// `sessions.handoff_error`.
228    pub handoff_error: Option<String>,
229    /// Epoch seconds, as the store keeps them.
230    pub started_at: Option<f64>,
231    /// Epoch seconds when the session ended, if it has.
232    pub ended_at: Option<f64>,
233    /// `sessions.end_reason`, the store's own word.
234    pub end_reason: Option<String>,
235}
236
237/// Hermes `sessions.source` → trigger. Cron fires are tagged `cron`; the CLI,
238/// TUI and ACP adapter are human surfaces; `api_server` is the HTTP API;
239/// `webhook` is inbound; every other value is a messaging platform.
240pub fn hermes_trigger_for_source(source: &str) -> Trigger {
241    match source {
242        "" => Trigger::Unknown,
243        "cron" => Trigger::Cron,
244        "webhook" => Trigger::Webhook,
245        "cli" | "tui" | "acp" | "console" => Trigger::Human,
246        "api_server" | "api" => Trigger::Api,
247        "kanban" => Trigger::Task,
248        _ => Trigger::Channel,
249    }
250}
251
252/// Hermes cron fire session ids are minted as `cron_<job_id>_<YYYYMMDD_HHMMSS>`
253/// (`cron/scheduler.py`); recover the job id.
254pub fn hermes_cron_job_id(session_id: &str) -> Option<String> {
255    let rest = session_id.strip_prefix("cron_")?;
256    let (job, stamp) = rest.rsplit_once('_')?;
257    let (job, date) = job.rsplit_once('_')?;
258    let ok = date.len() == 8
259        && stamp.len() == 6
260        && date.chars().all(|c| c.is_ascii_digit())
261        && stamp.chars().all(|c| c.is_ascii_digit());
262    if ok && !job.is_empty() {
263        Some(job.to_string())
264    } else {
265        None
266    }
267}
268
269/// Parse a Hermes gateway session key
270/// (`agent:<profile|main>:<platform>:<chat_type>[:<chat_id>][:<thread_id>][:<participant>]`).
271/// Returns the surface and the profile namespace (`None` for `main`).
272pub fn parse_hermes_session_key(key: &str) -> Option<(SurfaceKey, Option<String>)> {
273    let parts: Vec<&str> = key.split(':').collect();
274    if parts.len() < 4 || parts[0] != "agent" {
275        return None;
276    }
277    let profile = match parts[1] {
278        "" | "main" | "default" => None,
279        p => Some(p.to_string()),
280    };
281    let surface = SurfaceKey {
282        key: Some(key.to_string()),
283        platform: Some(parts[2].to_string()),
284        kind: Some(parts[3].to_string()),
285        chat_id: parts.get(4).map(|s| s.to_string()),
286        thread_id: parts.get(5).map(|s| s.to_string()),
287        participant_id: parts.get(6).map(|s| s.to_string()),
288    };
289    Some((surface, profile))
290}
291
292/// Render a surface as Hermes's `build_session_key` form, which
293/// [`parse_hermes_session_key`] reads back unchanged.
294pub fn render_hermes_session_key(profile: &str, key: &SurfaceKey) -> String {
295    let mut parts = vec![
296        "agent".to_string(),
297        if profile.is_empty() {
298            "main".to_string()
299        } else {
300            profile.to_string()
301        },
302        key.platform.clone().unwrap_or_default(),
303        key.kind.clone().unwrap_or_default(),
304    ];
305    parts.extend(
306        [
307            key.chat_id.clone(),
308            key.thread_id.clone(),
309            key.participant_id.clone(),
310        ]
311        .into_iter()
312        .flatten(),
313    );
314    parts.join(":")
315}
316
317fn epoch_to_rfc3339(seconds: f64) -> String {
318    let millis = (seconds * 1000.0).round() as i64;
319    let secs = millis.div_euclid(1000);
320    let sub = millis.rem_euclid(1000) as u32;
321    // civil-from-days (Howard Hinnant), enough for a timestamp string
322    let days = secs.div_euclid(86_400);
323    let sod = secs.rem_euclid(86_400);
324    let z = days + 719_468;
325    let era = z.div_euclid(146_097);
326    let doe = z - era * 146_097;
327    let yoe = (doe - doe / 1460 + doe / 36_524 - doe / 146_096) / 365;
328    let y = yoe + era * 400;
329    let doy = doe - (365 * yoe + yoe / 4 - yoe / 100);
330    let mp = (5 * doy + 2) / 153;
331    let d = doy - (153 * mp + 2) / 5 + 1;
332    let m = if mp < 10 { mp + 3 } else { mp - 9 };
333    let y = if m <= 2 { y + 1 } else { y };
334    format!(
335        "{y:04}-{m:02}-{d:02}T{:02}:{:02}:{:02}.{sub:03}Z",
336        sod / 3600,
337        (sod % 3600) / 60,
338        sod % 60
339    )
340}
341
342impl Binding {
343    /// Decode a Hermes `sessions` row. Always yields a record: a terminal
344    /// session is a binding with a degenerate key. The columns win over the
345    /// parsed key when both are present (ORC-8 finding: `api_server`
346    /// conversations carry the surface only in `session_key`).
347    pub fn from_hermes_row(row: &HermesSessionRow, locator: Option<&str>) -> Self {
348        let nonempty = |v: &Option<String>| v.clone().filter(|s| !s.is_empty());
349        let source = nonempty(&row.source).unwrap_or_default();
350        let mut trigger = hermes_trigger_for_source(&source);
351        if row.lineage_kind.as_deref() == Some("delegate") {
352            trigger = Trigger::Parent;
353        }
354        let mut recurrence = None;
355        if let Some(job_id) = hermes_cron_job_id(&row.id) {
356            recurrence = Some(Recurrence {
357                job_id,
358                kind: "cron".into(),
359            });
360            trigger = Trigger::Cron;
361        }
362        let mut profile = None;
363        let mut key = nonempty(&row.session_key)
364            .and_then(|k| parse_hermes_session_key(&k))
365            .map(|(surface, key_profile)| {
366                profile = key_profile;
367                surface
368            })
369            .unwrap_or_default();
370        if key.key.is_none() {
371            key.key = nonempty(&row.session_key);
372        }
373        if let Some(v) = nonempty(&row.chat_id) {
374            key.chat_id = Some(v);
375        }
376        if let Some(v) = nonempty(&row.chat_type) {
377            key.kind = Some(v);
378        }
379        if let Some(v) = nonempty(&row.thread_id) {
380            key.thread_id = Some(v);
381        }
382        if let Some(v) = nonempty(&row.user_id) {
383            key.participant_id = Some(v);
384        }
385        if key.platform.is_none() && trigger == Trigger::Channel {
386            key.platform = Some(source.clone());
387        }
388        if let Some(p) = nonempty(&row.profile_name) {
389            profile = Some(p);
390        }
391        let handoff = nonempty(&row.handoff_state).map(|state| Handoff {
392            to: nonempty(&row.handoff_platform),
393            state,
394            error: nonempty(&row.handoff_error),
395        });
396        let mut residue = Residue::default();
397        let end_reason = match nonempty(&row.end_reason) {
398            Some(word) => match EndReason::parse(&word) {
399                Some(r) => Some(r),
400                None => {
401                    residue.keep("end_reason", serde_json::Value::String(word));
402                    None
403                }
404            },
405            None => None,
406        };
407        Self {
408            key,
409            profile,
410            worker: Worker {
411                harness: HarnessId::new(HarnessId::HERMES),
412                session_id: Some(row.id.clone()),
413                locator: locator.map(str::to_string),
414            },
415            trigger,
416            recurrence,
417            conversation: None,
418            agent: None,
419            handoff,
420            started_at: row.started_at.map(epoch_to_rfc3339),
421            last_activity_at: row.ended_at.or(row.started_at).map(epoch_to_rfc3339),
422            ended_at: row.ended_at.map(epoch_to_rfc3339),
423            end_reason,
424            residue,
425        }
426    }
427}
428
429/// The Hermes `sessions.source` word for a binding (inverse of
430/// [`hermes_trigger_for_source`]): a channel conversation's source is its
431/// platform; the rest are the words Hermes's own doors mint.
432pub fn hermes_source_for_binding(binding: &Binding) -> String {
433    match binding.trigger {
434        Trigger::Cron => "cron".into(),
435        Trigger::Webhook => "webhook".into(),
436        Trigger::Api => "api_server".into(),
437        Trigger::Task => "kanban".into(),
438        Trigger::Human => "cli".into(),
439        Trigger::Parent => "delegate".into(),
440        Trigger::Heartbeat => "heartbeat".into(),
441        Trigger::Channel | Trigger::Unknown => {
442            binding.key.platform.clone().unwrap_or_else(|| "cli".into())
443        }
444    }
445}
446
447// ----------------------------------------------------------- OpenClaw keys
448
449/// Render an OpenClaw gateway session key for a binding under an agent
450/// (inverse of [`parse_openclaw_session_key`]): a cron fire is `cron:<jobId>`,
451/// a conversation is the agent shape.
452pub fn render_openclaw_session_key(agent: &str, binding: &Binding) -> String {
453    let key = &binding.key;
454    if let Some(k) = &key.key {
455        return k.clone();
456    }
457    if let Some(r) = &binding.recurrence {
458        return format!("cron:{}", r.job_id);
459    }
460    if key.kind.as_deref() == Some("main") || key.platform.is_none() {
461        return format!("agent:{agent}:main");
462    }
463    let mut out = format!(
464        "agent:{agent}:{}:{}:{}",
465        key.platform.clone().unwrap_or_default(),
466        key.kind.clone().unwrap_or_else(|| "dm".into()),
467        key.chat_id.clone().unwrap_or_default()
468    );
469    if let Some(t) = &key.thread_id {
470        out.push_str(&format!(":thread:{t}"));
471    }
472    out
473}
474
475/// Parse an OpenClaw gateway session key. Shapes (`docs/channels/channel-routing.md`,
476/// `docs/automation/cron-jobs.md`, `docs/cli/acp.md` upstream):
477/// `agent:<id>:main`, `agent:<id>:<channel>:<group|channel>:<cid>[:thread|topic:<tid>]`,
478/// `cron:<jobId>`, `hook:<name>:<id>`, `acp-bridge:<uuid>`.
479pub fn parse_openclaw_session_key(
480    key: &str,
481) -> Option<(Option<String>, SurfaceKey, Trigger, Option<Recurrence>)> {
482    let parts: Vec<&str> = key.split(':').collect();
483    match parts.first().copied() {
484        Some("agent") if parts.len() >= 3 => {
485            let agent = Some(parts[1].to_string());
486            if parts[2] == "main" {
487                let surface = SurfaceKey {
488                    key: Some(key.to_string()),
489                    kind: Some("main".to_string()),
490                    ..SurfaceKey::default()
491                };
492                return Some((agent, surface, Trigger::Unknown, None));
493            }
494            if parts.len() < 5 {
495                return None;
496            }
497            let thread_id = match (parts.get(5), parts.get(6)) {
498                (Some(&"thread"), Some(t)) | (Some(&"topic"), Some(t)) => Some(t.to_string()),
499                _ => None,
500            };
501            let surface = SurfaceKey {
502                key: Some(key.to_string()),
503                platform: Some(parts[2].to_string()),
504                kind: Some(parts[3].to_string()),
505                chat_id: Some(parts[4].to_string()),
506                thread_id,
507                participant_id: None,
508            };
509            Some((agent, surface, Trigger::Channel, None))
510        }
511        Some("cron") if parts.len() >= 2 => Some((
512            None,
513            SurfaceKey {
514                key: Some(key.to_string()),
515                ..SurfaceKey::default()
516            },
517            Trigger::Cron,
518            Some(Recurrence {
519                job_id: parts[1..].join(":"),
520                kind: "cron".into(),
521            }),
522        )),
523        Some("hook") if parts.len() >= 2 => Some((
524            None,
525            SurfaceKey {
526                key: Some(key.to_string()),
527                ..SurfaceKey::default()
528            },
529            Trigger::Webhook,
530            None,
531        )),
532        Some("acp-bridge") => Some((
533            None,
534            SurfaceKey {
535                key: Some(key.to_string()),
536                platform: Some("acp".into()),
537                ..SurfaceKey::default()
538            },
539            Trigger::Api,
540            None,
541        )),
542        _ => None,
543    }
544}
545
546impl Binding {
547    /// Decode an OpenClaw session key. `None` when the key has no known shape
548    /// (nothing is claimed, as before). `agent_from_path` is the profile the
549    /// file's own `agents/<id>/` directory names, used when the key names none.
550    pub fn from_openclaw_key(
551        key: &str,
552        agent_from_path: Option<&str>,
553        session_id: Option<&str>,
554        locator: Option<&str>,
555    ) -> Option<Self> {
556        let (agent, surface, trigger, recurrence) = parse_openclaw_session_key(key)?;
557        Some(Self {
558            key: surface,
559            profile: agent.or_else(|| agent_from_path.map(str::to_string)),
560            worker: Worker {
561                harness: HarnessId::new(HarnessId::OPENCLAW),
562                session_id: session_id.map(str::to_string),
563                locator: locator.map(str::to_string),
564            },
565            trigger,
566            recurrence,
567            ..Self::default()
568        })
569    }
570}
571
572// ------------------------------------------------------ orchestrator rows
573
574/// One row of an orchestrator profile's `bindings` table, as read.
575#[derive(Debug, Clone, Default)]
576pub struct OrchestratorBindingRow {
577    /// Surface platform.
578    pub platform: String,
579    /// Surface chat type (`dm` | `group` | `channel` | `thread`).
580    pub chat_type: String,
581    /// Surface chat id.
582    pub chat_id: Option<String>,
583    /// Surface thread id.
584    pub thread_id: Option<String>,
585    /// Surface participant id.
586    pub participant_id: Option<String>,
587    /// The worker harness.
588    pub worker_harness: String,
589    /// The worker session id; `None` until the worker has reported one.
590    pub worker_session_id: Option<String>,
591    /// Where the worker transcript can be read.
592    pub worker_locator: Option<String>,
593    /// RFC3339 start.
594    pub started_at: Option<String>,
595    /// RFC3339 last activity.
596    pub last_activity_at: Option<String>,
597    /// RFC3339 end, if ended.
598    pub ended_at: Option<String>,
599    /// Why it ended, the store's own word.
600    pub end_reason: Option<String>,
601    /// Handoff destination platform.
602    pub handoff_to: Option<String>,
603    /// Handoff state.
604    pub handoff_state: Option<String>,
605    /// Handoff error.
606    pub handoff_error: Option<String>,
607    /// The job a fire binding belongs to.
608    pub recurrence_job_id: Option<String>,
609}
610
611impl Binding {
612    /// Decode an orchestrator `bindings` row under `profile`. The key string is
613    /// the orchestrator's own rendering (`docs/ORCHESTRATOR-IR.md` §2.3), which
614    /// is Hermes's form, so [`parse_hermes_session_key`] reads it back unchanged.
615    pub fn from_orchestrator_row(profile: &str, row: &OrchestratorBindingRow) -> Self {
616        let mut key = SurfaceKey {
617            key: None,
618            platform: Some(row.platform.clone()),
619            kind: Some(row.chat_type.clone()),
620            chat_id: row.chat_id.clone(),
621            thread_id: row.thread_id.clone(),
622            participant_id: row.participant_id.clone(),
623        };
624        key.key = Some(render_hermes_session_key(profile, &key));
625        let trigger = if row.recurrence_job_id.is_some() {
626            Trigger::Cron
627        } else if row.platform == "webhook" {
628            Trigger::Webhook
629        } else {
630            Trigger::Channel
631        };
632        let mut residue = Residue::default();
633        let end_reason = match row.end_reason.as_deref() {
634            Some(word) => match EndReason::parse(word) {
635                Some(r) => Some(r),
636                None => {
637                    residue.keep("end_reason", serde_json::Value::String(word.to_string()));
638                    None
639                }
640            },
641            None => None,
642        };
643        Self {
644            key,
645            profile: Some(profile.to_string()),
646            worker: Worker {
647                harness: HarnessId::new(&row.worker_harness),
648                session_id: row.worker_session_id.clone().filter(|s| !s.is_empty()),
649                locator: row.worker_locator.clone(),
650            },
651            trigger,
652            recurrence: row.recurrence_job_id.clone().map(|job_id| Recurrence {
653                job_id,
654                kind: "cron".into(),
655            }),
656            conversation: None,
657            agent: None,
658            handoff: row.handoff_state.clone().map(|state| Handoff {
659                to: row.handoff_to.clone(),
660                state,
661                error: row.handoff_error.clone(),
662            }),
663            started_at: row.started_at.clone(),
664            last_activity_at: row.last_activity_at.clone(),
665            ended_at: row.ended_at.clone(),
666            end_reason,
667            residue,
668        }
669    }
670}