1use 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#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
17#[serde(rename_all = "snake_case")]
18pub enum EndReason {
19 Idle,
21 Daily,
23 Reset,
25 New,
27 Handoff,
29 Transfer,
31 Left,
33 Error,
35}
36
37impl EndReason {
38 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 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#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
70pub struct Worker {
71 pub harness: HarnessId,
73 #[serde(default, skip_serializing_if = "Option::is_none")]
75 pub session_id: Option<String>,
76 #[serde(default, skip_serializing_if = "Option::is_none")]
78 pub locator: Option<String>,
79}
80
81#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
83pub struct Handoff {
84 #[serde(default, skip_serializing_if = "Option::is_none")]
86 pub to: Option<String>,
87 pub state: String,
89 #[serde(default, skip_serializing_if = "Option::is_none")]
90 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#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
106pub struct Binding {
107 pub key: SurfaceKey,
109 #[serde(default)]
111 pub profile: Option<String>,
112 #[serde(default, skip_serializing_if = "Option::is_none")]
114 pub conversation: Option<String>,
115 #[serde(default, skip_serializing_if = "Option::is_none")]
117 pub agent: Option<String>,
118 pub worker: Worker,
120 #[serde(default)]
124 pub trigger: Trigger,
125 #[serde(default)]
126 pub recurrence: Option<Recurrence>,
128 #[serde(default)]
129 pub handoff: Option<Handoff>,
131 #[serde(default)]
133 pub started_at: Option<String>,
134 #[serde(default)]
135 pub last_activity_at: Option<String>,
137 #[serde(default)]
138 pub ended_at: Option<String>,
140 #[serde(default)]
141 pub end_reason: Option<EndReason>,
143 #[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 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 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#[derive(Debug, Clone, Default)]
204pub struct HermesSessionRow {
205 pub id: String,
207 pub source: Option<String>,
209 pub lineage_kind: Option<String>,
211 pub session_key: Option<String>,
213 pub chat_id: Option<String>,
215 pub chat_type: Option<String>,
217 pub thread_id: Option<String>,
219 pub user_id: Option<String>,
221 pub profile_name: Option<String>,
223 pub handoff_state: Option<String>,
225 pub handoff_platform: Option<String>,
227 pub handoff_error: Option<String>,
229 pub started_at: Option<f64>,
231 pub ended_at: Option<f64>,
233 pub end_reason: Option<String>,
235}
236
237pub 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
252pub 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
269pub 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
292pub 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 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 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
429pub 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
447pub 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
475pub 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 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#[derive(Debug, Clone, Default)]
576pub struct OrchestratorBindingRow {
577 pub platform: String,
579 pub chat_type: String,
581 pub chat_id: Option<String>,
583 pub thread_id: Option<String>,
585 pub participant_id: Option<String>,
587 pub worker_harness: String,
589 pub worker_session_id: Option<String>,
591 pub worker_locator: Option<String>,
593 pub started_at: Option<String>,
595 pub last_activity_at: Option<String>,
597 pub ended_at: Option<String>,
599 pub end_reason: Option<String>,
601 pub handoff_to: Option<String>,
603 pub handoff_state: Option<String>,
605 pub handoff_error: Option<String>,
607 pub recurrence_job_id: Option<String>,
609}
610
611impl Binding {
612 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}