1use std::collections::BTreeMap;
6
7use serde_json::{Map, Value};
8
9use crate::ontology::{HarnessId, Residue, SecretRef, SurfaceKey};
10use crate::orchestration::{
11 ChannelConfig, EnvValue, Fire, FireStatus, Job, JobOrigin, Obligation, ObligationSource,
12 ObligationState, OutboundContent, PermissionDefault, PermissionPolicy, PermissionUnattended,
13 Posted, Repeat, Route, RouteMatch, Schedule, Target, WebhookSubscription, WorkerSpec,
14};
15use crate::Result;
16
17#[derive(Debug, Clone, PartialEq, Eq)]
19pub struct LoadError {
20 pub file: String,
22 pub key: Option<String>,
24 pub message: String,
26}
27
28impl std::fmt::Display for LoadError {
29 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
30 match &self.key {
31 Some(key) => write!(f, "{} [{key}]: {}", self.file, self.message),
32 None => write!(f, "{}: {}", self.file, self.message),
33 }
34 }
35}
36
37impl From<LoadError> for crate::Error {
38 fn from(error: LoadError) -> Self {
39 crate::Error::Other(error.to_string())
40 }
41}
42
43pub(crate) fn load_error(file: &str, key: &str, message: impl Into<String>) -> crate::Error {
44 LoadError {
45 file: file.to_string(),
46 key: if key.is_empty() {
47 None
48 } else {
49 Some(key.to_string())
50 },
51 message: message.into(),
52 }
53 .into()
54}
55
56fn expect_keys(file: &str, key: &str, obj: &Value, allowed: &[&str]) -> Result<()> {
57 let Some(map) = obj.as_object() else {
58 return Err(load_error(file, key, "expected an object"));
59 };
60 for k in map.keys() {
61 if !allowed.contains(&k.as_str()) {
62 let path = if key.is_empty() {
63 k.clone()
64 } else {
65 format!("{key}.{k}")
66 };
67 return Err(load_error(file, &path, "unknown key"));
68 }
69 }
70 Ok(())
71}
72
73fn req_str(file: &str, key: &str, v: Option<&Value>) -> Result<String> {
74 match v {
75 Some(Value::String(s)) if !s.is_empty() => Ok(s.clone()),
76 _ => Err(load_error(file, key, "expected a non-empty string")),
77 }
78}
79
80fn opt_str(file: &str, key: &str, v: Option<&Value>) -> Result<Option<String>> {
81 match v {
82 None | Some(Value::Null) => Ok(None),
83 Some(Value::String(s)) => Ok(Some(s.clone())),
84 _ => Err(load_error(file, key, "expected a string")),
85 }
86}
87
88fn opt_bool(
89 file: &str,
90 key: &str,
91 v: Option<&Value>,
92 default: Option<bool>,
93) -> Result<Option<bool>> {
94 match v {
95 None | Some(Value::Null) => Ok(default),
96 Some(Value::Bool(b)) => Ok(Some(*b)),
97 _ => Err(load_error(file, key, "expected a boolean")),
98 }
99}
100
101fn residue_of(obj: &Map<String, Value>, mapped: &[&str]) -> Residue {
102 let mut out = Residue::default();
103 for (k, v) in obj {
104 if !mapped.contains(&k.as_str()) {
105 out.keep(k.clone(), v.clone());
106 }
107 }
108 out
109}
110
111fn row_residue_of(row: &Map<String, Value>, mapped: &[&str]) -> Residue {
113 let mut out = Residue::default();
114 for (k, v) in row {
115 if !mapped.contains(&k.as_str()) && !v.is_null() {
116 out.keep(k.clone(), v.clone());
117 }
118 }
119 out
120}
121
122fn text_of(v: &Value) -> Option<String> {
124 match v {
125 Value::String(s) => Some(s.clone()),
126 Value::Number(n) => Some(match n.as_f64() {
127 Some(f) if n.is_f64() && f.fract() == 0.0 && f.abs() < 1e21 => format!("{}", f as i64),
128 _ => n.to_string(),
129 }),
130 _ => None,
131 }
132}
133
134fn number_of(v: &Value) -> Option<f64> {
135 match v {
136 Value::Number(n) => n.as_f64(),
137 Value::String(s) => s.parse().ok(),
138 _ => None,
139 }
140}
141
142pub const CHAT_TYPES: &[&str] = &["dm", "group", "channel", "thread"];
146
147pub fn surface_key_string(key: &SurfaceKey) -> String {
149 format!(
150 "{}|{}|{}|{}|{}",
151 key.platform.as_deref().unwrap_or(""),
152 key.kind.as_deref().unwrap_or(""),
153 key.chat_id.as_deref().unwrap_or(""),
154 key.thread_id.as_deref().unwrap_or(""),
155 key.participant_id.as_deref().unwrap_or("")
156 )
157}
158
159pub fn decode_surface_key(file: &str, key: &str, raw: &Value) -> Result<SurfaceKey> {
161 let map = raw
162 .as_object()
163 .ok_or_else(|| load_error(file, key, "expected an object"))?;
164 for k in map.keys() {
165 if ![
166 "platform",
167 "chat_type",
168 "kind",
169 "chat_id",
170 "thread_id",
171 "participant_id",
172 "key",
173 ]
174 .contains(&k.as_str())
175 {
176 return Err(load_error(file, &format!("{key}.{k}"), "unknown key"));
177 }
178 }
179 let chat_type = map
180 .get("chat_type")
181 .or_else(|| map.get("kind"))
182 .and_then(Value::as_str)
183 .map(str::to_string);
184 if !chat_type
185 .as_deref()
186 .is_some_and(|c| CHAT_TYPES.contains(&c))
187 {
188 return Err(load_error(
189 file,
190 key,
191 "expected platform, chat_type, chat_id",
192 ));
193 }
194 Ok(SurfaceKey {
195 key: None,
196 platform: Some(req_str(
197 file,
198 &format!("{key}.platform"),
199 map.get("platform"),
200 )?),
201 kind: chat_type,
202 chat_id: map.get("chat_id").and_then(text_of),
203 thread_id: map
204 .get("thread_id")
205 .and_then(text_of)
206 .filter(|s| !s.is_empty()),
207 participant_id: map
208 .get("participant_id")
209 .and_then(text_of)
210 .filter(|s| !s.is_empty()),
211 })
212}
213
214pub fn encode_surface_key(key: &SurfaceKey) -> Value {
216 let mut out = Map::new();
217 out.insert(
218 "platform".into(),
219 Value::String(key.platform.clone().unwrap_or_default()),
220 );
221 out.insert(
222 "kind".into(),
223 Value::String(key.kind.clone().unwrap_or_default()),
224 );
225 if let Some(c) = &key.chat_id {
226 out.insert("chat_id".into(), Value::String(c.clone()));
227 }
228 if let Some(t) = &key.thread_id {
229 out.insert("thread_id".into(), Value::String(t.clone()));
230 }
231 if let Some(p) = &key.participant_id {
232 out.insert("participant_id".into(), Value::String(p.clone()));
233 }
234 Value::Object(out)
235}
236
237pub fn parse_target(
241 file: &str,
242 key: &str,
243 word: Option<&Value>,
244 extra: Option<&Value>,
245) -> Result<Option<Target>> {
246 let word = match word {
247 None | Some(Value::Null) => return Ok(None),
248 Some(Value::String(s)) if !s.is_empty() => s.as_str(),
249 _ => return Err(load_error(file, key, "expected a delivery word")),
250 };
251 Ok(Some(match word {
252 "origin" => Target::Origin,
253 "local" => Target::Local,
254 "home" => Target::Home,
255 _ => {
256 let parts: Vec<&str> = word.split(':').collect();
257 let extra_get = |name: &str| extra.and_then(|e| e.get(name)).and_then(text_of);
258 let chat_id = parts
259 .get(1)
260 .map(|s| s.to_string())
261 .or_else(|| extra_get("chat_id"))
262 .filter(|s| !s.is_empty());
263 let thread_id = parts
264 .get(2)
265 .map(|s| s.to_string())
266 .or_else(|| extra_get("thread_id"))
267 .filter(|s| !s.is_empty());
268 Target::Explicit {
269 platform: parts[0].to_string(),
270 chat_id,
271 thread_id,
272 }
273 }
274 }))
275}
276
277pub fn decode_env_value(file: &str, key: &str, v: &Value) -> Result<EnvValue> {
281 match v {
282 Value::String(s) => Ok(match placeholder_ref(s) {
283 Some(name) => EnvValue::Secret(SecretRef::Dotenv(name.to_string())),
284 None => EnvValue::Literal(s.clone()),
285 }),
286 _ => Err(load_error(
287 file,
288 key,
289 "expected a string (a secret is written `${NAME}`)",
290 )),
291 }
292}
293
294pub fn decode_worker(file: &str, raw: Option<&Value>) -> Result<Option<WorkerSpec>> {
296 let Some(raw) = raw.filter(|v| !v.is_null()) else {
297 return Ok(None);
298 };
299 expect_keys(
300 file,
301 "worker",
302 raw,
303 &[
304 "harness",
305 "model",
306 "reasoning_effort",
307 "preset",
308 "cwd",
309 "env",
310 "permission",
311 "home",
312 "surface",
313 "capacity",
314 ],
315 )?;
316 let cwd = match raw.get("cwd") {
317 None => ".".to_string(),
318 v => req_str(file, "worker.cwd", v)?,
319 };
320 if cwd.starts_with('/') || cwd.split('/').any(|p| p == "..") {
321 return Err(load_error(
322 file,
323 "worker.cwd",
324 "must be a relative path inside the profile",
325 ));
326 }
327 let mut env = BTreeMap::new();
328 if let Some(e) = raw.get("env") {
329 let map = e
330 .as_object()
331 .ok_or_else(|| load_error(file, "worker.env", "expected a map"))?;
332 for (k, v) in map {
333 env.insert(
334 k.clone(),
335 decode_env_value(file, &format!("worker.env.{k}"), v)?,
336 );
337 }
338 }
339 let mut permission = PermissionPolicy::default();
340 if let Some(p) = raw.get("permission") {
341 expect_keys(
342 file,
343 "worker.permission",
344 p,
345 &["timeout_seconds", "default", "unattended"],
346 )?;
347 if let Some(t) = p.get("timeout_seconds") {
348 permission.timeout_seconds = t.as_u64().filter(|v| *v >= 1).ok_or_else(|| {
349 load_error(
350 file,
351 "worker.permission.timeout_seconds",
352 "expected an integer >= 1",
353 )
354 })? as u32;
355 }
356 if let Some(d) = p.get("default") {
357 permission.default = match d.as_str() {
358 Some("deny") => PermissionDefault::Deny,
359 Some("allow") => PermissionDefault::Allow,
360 _ => {
361 return Err(load_error(
362 file,
363 "worker.permission.default",
364 "expected deny|allow",
365 ))
366 }
367 };
368 }
369 if let Some(u) = p.get("unattended") {
370 permission.unattended = match u.as_str() {
371 Some("deny") => PermissionUnattended::Deny,
372 Some("approve") => PermissionUnattended::Approve,
373 _ => {
374 return Err(load_error(
375 file,
376 "worker.permission.unattended",
377 "expected deny|approve",
378 ))
379 }
380 };
381 }
382 }
383 let home = match opt_str(file, "worker.home", raw.get("home"))?.as_deref() {
384 None | Some("profile") => crate::orchestration::WorkerHome::Profile,
385 Some("user") => crate::orchestration::WorkerHome::User,
386 Some(_) => {
387 return Err(load_error(
388 file,
389 "worker.home",
390 "expected `profile` or `user`",
391 ))
392 }
393 };
394 let surface = match opt_str(file, "worker.surface", raw.get("surface"))?.as_deref() {
395 None | Some("headless") => crate::orchestration::WorkerSurface::Headless,
396 Some("pane") => crate::orchestration::WorkerSurface::Pane,
397 Some(_) => {
398 return Err(load_error(
399 file,
400 "worker.surface",
401 "expected `headless` or `pane`",
402 ))
403 }
404 };
405 let capacity = match raw.get("capacity") {
406 None => None,
407 Some(v) => Some(
408 v.as_u64()
409 .filter(|n| *n >= 1 && *n <= u64::from(u32::MAX))
410 .ok_or_else(|| load_error(file, "worker.capacity", "expected an integer >= 1"))?
411 as u32,
412 ),
413 };
414 Ok(Some(WorkerSpec {
415 harness: HarnessId::new(req_str(file, "worker.harness", raw.get("harness"))?),
416 model: opt_str(file, "worker.model", raw.get("model"))?,
417 reasoning_effort: opt_str(file, "worker.reasoning_effort", raw.get("reasoning_effort"))?,
418 preset: opt_str(file, "worker.preset", raw.get("preset"))?,
419 cwd,
420 env,
421 permission,
422 home,
423 surface,
424 capacity,
425 }))
426}
427
428const JOB_MAPPED: &[&str] = &[
431 "id",
432 "schedule",
433 "prompt",
434 "workdir",
435 "model",
436 "skills",
437 "context_from",
438 "deliver",
439 "failure_deliver",
440 "origin",
441 "attach_to_session",
442 "session_target",
443 "machine",
444 "repeat",
445 "enabled",
446 "next_run_at",
447 "last_run_at",
448 "last_status",
449 "created_at",
450];
451
452pub const HERMES_JOB_ORDER: &[&str] = &[
454 "id",
455 "schedule",
456 "prompt",
457 "skills",
458 "script",
459 "no_agent",
460 "model",
461 "provider",
462 "workdir",
463 "enabled_toolsets",
464 "context_from",
465 "deliver",
466 "failure_deliver",
467 "attach_to_session",
468 "origin",
469 "repeat",
470 "enabled",
471 "next_run_at",
472 "last_run_at",
473 "last_status",
474 "created_at",
475 "fire_claim",
476];
477
478pub fn decode_context_from(
480 file: &str,
481 key: &str,
482 v: Option<&Value>,
483) -> Result<Option<Vec<String>>> {
484 let items: Vec<String> = match v {
485 None | Some(Value::Null) => return Ok(None),
486 Some(Value::String(s)) => vec![s.clone()],
487 Some(Value::Array(a)) => a
488 .iter()
489 .map(|x| match x {
490 Value::String(s) => s.clone(),
491 other => other.to_string(),
492 })
493 .collect(),
494 _ => {
495 return Err(load_error(
496 file,
497 key,
498 "expected a job id, a list of job ids, or \"self\"",
499 ))
500 }
501 };
502 let refs: Vec<String> = items
503 .into_iter()
504 .map(|s| s.trim().to_string())
505 .filter(|s| !s.is_empty())
506 .collect();
507 Ok(if refs.is_empty() { None } else { Some(refs) })
508}
509
510pub fn decode_repeat(file: &str, key: &str, raw: Option<&Value>) -> Result<Option<Repeat>> {
512 Ok(match raw {
513 None | Some(Value::Null) => None,
514 Some(Value::Bool(true)) => Some(Repeat {
515 times: None,
516 completed: 0,
517 }),
518 Some(Value::Bool(false)) => Some(Repeat {
519 times: Some(1),
520 completed: 0,
521 }),
522 Some(Value::Number(n)) => Some(Repeat {
523 times: n.as_f64().filter(|f| *f > 0.0).map(|f| f.floor() as u32),
524 completed: 0,
525 }),
526 Some(Value::Object(map)) => {
527 let times = match map.get("times") {
528 None | Some(Value::Null) => None,
529 Some(Value::Number(n)) if n.as_u64().is_some_and(|v| v >= 1) => {
530 Some(n.as_u64().unwrap() as u32)
531 }
532 _ => {
533 return Err(load_error(
534 file,
535 &format!("{key}.times"),
536 "expected null or an integer >= 1",
537 ))
538 }
539 };
540 let completed = match map.get("completed") {
541 None => 0,
542 Some(Value::Number(n)) if n.as_u64().is_some() => n.as_u64().unwrap() as u32,
543 _ => {
544 return Err(load_error(
545 file,
546 &format!("{key}.completed"),
547 "expected an integer >= 0",
548 ))
549 }
550 };
551 Some(Repeat { times, completed })
552 }
553 _ => {
554 return Err(load_error(
555 file,
556 key,
557 "expected null, {times, completed}, a number or a boolean",
558 ))
559 }
560 })
561}
562
563pub fn decode_job(file: &str, raw: &Value) -> Result<Job> {
565 let map = raw
566 .as_object()
567 .ok_or_else(|| load_error(file, "", "expected a job object"))?;
568 let id = req_str(file, "id", map.get("id"))?;
569 let k = |s: &str| format!("{id}.{s}");
570 let sched = map
571 .get("schedule")
572 .and_then(Value::as_object)
573 .ok_or_else(|| load_error(file, &k("schedule"), "expected an object"))?;
574 let schedule = match sched.get("kind").and_then(Value::as_str) {
575 Some("once") => Schedule::Once {
576 run_at: req_str(file, &k("schedule.run_at"), sched.get("run_at"))?,
577 },
578 Some("interval") => {
579 let minutes = sched
580 .get("minutes")
581 .and_then(Value::as_f64)
582 .filter(|m| *m > 0.0)
583 .ok_or_else(|| {
584 load_error(file, &k("schedule.minutes"), "expected a positive number")
585 })?;
586 Schedule::Interval { minutes }
587 }
588 Some("cron") => Schedule::Cron {
589 expr: req_str(file, &k("schedule.expr"), sched.get("expr"))?,
590 tz: opt_str(file, &k("schedule.tz"), sched.get("tz"))?,
591 },
592 _ => {
593 return Err(load_error(
594 file,
595 &k("schedule.kind"),
596 "expected once|interval|cron",
597 ))
598 }
599 };
600 let schedule_residue = residue_of(sched, &["kind", "run_at", "minutes", "expr", "tz"]);
601 let mut residue = residue_of(map, JOB_MAPPED);
602 if !schedule_residue.is_empty() {
603 residue.keep(
604 "__schedule",
605 Value::Object(schedule_residue.0.into_iter().collect()),
606 );
607 }
608 let origin = match map.get("origin") {
609 Some(Value::Object(o)) => {
610 let r = residue_of(o, &["platform", "chat_id", "thread_id"]);
611 if !r.is_empty() {
612 residue.keep("__origin", Value::Object(r.0.into_iter().collect()));
613 }
614 Some(JobOrigin {
615 platform: req_str(file, &k("origin.platform"), o.get("platform"))?,
616 chat_type: None,
617 chat_id: opt_str(file, &k("origin.chat_id"), o.get("chat_id"))?,
618 thread_id: opt_str(file, &k("origin.thread_id"), o.get("thread_id"))?,
619 })
620 }
621 _ => None,
622 };
623 Ok(Job {
624 schedule,
625 prompt: opt_str(file, &k("prompt"), map.get("prompt"))?,
626 workdir: opt_str(file, &k("workdir"), map.get("workdir"))?,
627 model: opt_str(file, &k("model"), map.get("model"))?,
628 skills: map
629 .get("skills")
630 .and_then(Value::as_array)
631 .map(|a| {
632 a.iter()
633 .map(|v| {
634 v.as_str()
635 .map(str::to_string)
636 .unwrap_or_else(|| v.to_string())
637 })
638 .collect()
639 })
640 .unwrap_or_default(),
641 context_from: decode_context_from(file, &k("context_from"), map.get("context_from"))?,
642 deliver: parse_target(file, &k("deliver"), map.get("deliver"), None)?
643 .unwrap_or(Target::Local),
644 failure_deliver: parse_target(
645 file,
646 &k("failure_deliver"),
647 map.get("failure_deliver"),
648 None,
649 )?,
650 origin,
651 attach_to_session: opt_bool(
652 file,
653 &k("attach_to_session"),
654 map.get("attach_to_session"),
655 None,
656 )?,
657 session_target: match opt_str(file, &k("session_target"), map.get("session_target"))? {
658 Some(t)
659 if !["main", "current", "isolated"].contains(&t.as_str())
660 && !t.starts_with("session:") =>
661 {
662 return Err(load_error(
663 file,
664 &k("session_target"),
665 "expected main|current|isolated|session:<id>",
666 ))
667 }
668 other => other,
669 },
670 machine: opt_str(file, &k("machine"), map.get("machine"))?,
671 repeat: decode_repeat(file, &k("repeat"), map.get("repeat"))?,
672 enabled: opt_bool(file, &k("enabled"), map.get("enabled"), Some(true))?.unwrap_or(true),
673 next_run_at: opt_str(file, &k("next_run_at"), map.get("next_run_at"))?,
674 last_run_at: opt_str(file, &k("last_run_at"), map.get("last_run_at"))?,
675 last_status: opt_str(file, &k("last_status"), map.get("last_status"))?,
676 created_at: opt_str(file, &k("created_at"), map.get("created_at"))?,
677 residue,
678 id,
679 })
680}
681
682pub fn encode_job(job: &Job) -> Vec<(String, Value)> {
685 let mut sched = Map::new();
686 sched.insert("kind".into(), Value::String(job.schedule.kind().into()));
687 match &job.schedule {
688 Schedule::Once { run_at } => {
689 sched.insert("run_at".into(), Value::String(run_at.clone()));
690 }
691 Schedule::Interval { minutes } => {
692 sched.insert("minutes".into(), serde_json::json!(*minutes));
693 }
694 Schedule::Cron { expr, tz } => {
695 sched.insert("expr".into(), Value::String(expr.clone()));
696 if let Some(tz) = tz {
697 sched.insert("tz".into(), Value::String(tz.clone()));
698 }
699 }
700 }
701 if let Some(Value::Object(extra)) = job.residue.0.get("__schedule") {
702 for (k, v) in extra {
703 sched.insert(k.clone(), v.clone());
704 }
705 }
706 let origin = job
707 .origin
708 .as_ref()
709 .map(|o| {
710 let mut m = Map::new();
711 m.insert("platform".into(), Value::String(o.platform.clone()));
712 m.insert(
713 "chat_id".into(),
714 o.chat_id.clone().map(Value::String).unwrap_or(Value::Null),
715 );
716 m.insert(
717 "thread_id".into(),
718 o.thread_id
719 .clone()
720 .map(Value::String)
721 .unwrap_or(Value::Null),
722 );
723 if let Some(Value::Object(extra)) = job.residue.0.get("__origin") {
724 for (k, v) in extra {
725 m.insert(k.clone(), v.clone());
726 }
727 }
728 Value::Object(m)
729 })
730 .unwrap_or(Value::Null);
731 let opt = |s: &Option<String>| s.clone().map(Value::String).unwrap_or(Value::Null);
732 let mut mapped: Vec<(String, Value)> = vec![
733 ("id".into(), Value::String(job.id.clone())),
734 ("schedule".into(), Value::Object(sched)),
735 ("prompt".into(), opt(&job.prompt)),
736 (
737 "skills".into(),
738 Value::Array(
739 job.skills
740 .iter()
741 .map(|s| Value::String(s.clone()))
742 .collect(),
743 ),
744 ),
745 ("model".into(), opt(&job.model)),
746 ("workdir".into(), opt(&job.workdir)),
747 (
748 "context_from".into(),
749 job.context_from
750 .as_ref()
751 .map(|l| Value::Array(l.iter().map(|s| Value::String(s.clone())).collect()))
752 .unwrap_or(Value::Null),
753 ),
754 ("deliver".into(), Value::String(job.deliver.render())),
755 (
756 "failure_deliver".into(),
757 job.failure_deliver
758 .as_ref()
759 .map(|t| Value::String(t.render()))
760 .unwrap_or(Value::Null),
761 ),
762 (
763 "attach_to_session".into(),
764 job.attach_to_session
765 .map(Value::Bool)
766 .unwrap_or(Value::Null),
767 ),
768 ("origin".into(), origin),
769 (
770 "repeat".into(),
771 job.repeat
772 .as_ref()
773 .map(|r| serde_json::json!({"times": r.times, "completed": r.completed}))
774 .unwrap_or(Value::Null),
775 ),
776 ("enabled".into(), Value::Bool(job.enabled)),
777 ("next_run_at".into(), opt(&job.next_run_at)),
778 ("last_run_at".into(), opt(&job.last_run_at)),
779 ("last_status".into(), opt(&job.last_status)),
780 ("created_at".into(), opt(&job.created_at)),
781 ];
782 if let Some(target) = &job.session_target {
784 mapped.push(("session_target".into(), Value::String(target.clone())));
785 }
786 if let Some(machine) = &job.machine {
788 mapped.push(("machine".into(), Value::String(machine.clone())));
789 }
790 let mut residue: Vec<(String, Value)> = job
791 .residue
792 .0
793 .iter()
794 .filter(|(k, _)| k.as_str() != "__schedule" && k.as_str() != "__origin")
800 .map(|(k, v)| (k.clone(), v.clone()))
801 .collect();
802 let mut out: Vec<(String, Value)> = Vec::new();
803 for key in HERMES_JOB_ORDER {
804 if let Some(pos) = mapped.iter().position(|(k, _)| k == key) {
805 out.push(mapped.remove(pos));
806 } else if let Some(pos) = residue.iter().position(|(k, _)| k == key) {
807 out.push(residue.remove(pos));
808 }
809 }
810 out.extend(mapped);
811 out.extend(residue);
812 out
813}
814
815pub fn decode_fire_row(file: &str, row: &Map<String, Value>) -> Result<Fire> {
817 let id = req_str(file, "id", row.get("id"))?;
818 let status_word = row.get("status").and_then(Value::as_str).unwrap_or("");
819 let status = FireStatus::from_hermes_word(status_word).ok_or_else(|| {
820 load_error(
821 file,
822 &format!("{id}.status"),
823 format!(
824 "unknown fire status {}",
825 serde_json::to_string(status_word).unwrap()
826 ),
827 )
828 })?;
829 Ok(Fire {
830 job_id: req_str(file, &format!("{id}.job_id"), row.get("job_id"))?,
831 session_id: row
832 .get("session_id")
833 .and_then(text_of)
834 .filter(|s| !s.is_empty()),
835 status,
836 claimed_at: req_str(file, &format!("{id}.claimed_at"), row.get("claimed_at"))?,
837 started_at: opt_str(file, &format!("{id}.started_at"), row.get("started_at"))?,
838 finished_at: opt_str(file, &format!("{id}.finished_at"), row.get("finished_at"))?,
839 error: opt_str(file, &format!("{id}.error"), row.get("error"))?,
840 obligation_id: row
841 .get("obligation_id")
842 .and_then(text_of)
843 .filter(|s| !s.is_empty()),
844 residue: {
845 let mut residue = row_residue_of(
846 row,
847 &[
848 "id",
849 "job_id",
850 "status",
851 "claimed_at",
852 "started_at",
853 "finished_at",
854 "error",
855 "residue_json",
856 "session_id",
857 "obligation_id",
858 ],
859 );
860 if let Some(extra) = row
863 .get("residue_json")
864 .and_then(text_of)
865 .and_then(|t| serde_json::from_str::<Map<String, Value>>(&t).ok())
866 {
867 for (k, v) in extra {
868 residue.keep(k, v);
869 }
870 }
871 residue
872 },
873 id,
874 })
875}
876
877pub fn encode_fire_row(fire: &Fire) -> Vec<Value> {
879 let r = &fire.residue.0;
880 let opt = |s: &Option<String>| s.clone().map(Value::String).unwrap_or(Value::Null);
881 vec![
882 Value::String(fire.id.clone()),
883 Value::String(fire.job_id.clone()),
884 r.get("source")
885 .cloned()
886 .unwrap_or(Value::String("scheduler".into())),
887 r.get("process_id")
888 .cloned()
889 .unwrap_or(Value::String(String::new())),
890 r.get("pid").cloned().unwrap_or(Value::from(0)),
891 r.get("process_started_at").cloned().unwrap_or(Value::Null),
892 Value::String(fire.status.hermes_word().into()),
893 Value::String(fire.claimed_at.clone()),
894 opt(&fire.started_at),
895 opt(&fire.finished_at),
896 opt(&fire.error),
897 ]
898}
899
900pub const EXECUTION_COLUMNS: &[&str] = &[
902 "id",
903 "job_id",
904 "source",
905 "process_id",
906 "pid",
907 "process_started_at",
908 "status",
909 "claimed_at",
910 "started_at",
911 "finished_at",
912 "error",
913];
914
915pub fn decode_obligation_row(file: &str, row: &Map<String, Value>) -> Result<Obligation> {
917 let id = req_str(file, "obligation_id", row.get("obligation_id"))?;
918 let state_word = row.get("state").and_then(Value::as_str).unwrap_or("");
919 let state = ObligationState::from_hermes_word(state_word).ok_or_else(|| {
920 load_error(
921 file,
922 &format!("{id}.state"),
923 format!(
924 "unknown obligation state {}",
925 serde_json::to_string(state_word).unwrap()
926 ),
927 )
928 })?;
929 let session_key = row
930 .get("session_key")
931 .and_then(Value::as_str)
932 .filter(|s| !s.is_empty())
933 .map(str::to_string);
934 let parsed = session_key
935 .as_deref()
936 .and_then(crate::ontology::parse_hermes_session_key)
937 .map(|(k, _)| k);
938 let target = SurfaceKey {
939 key: None,
940 platform: Some(req_str(
941 file,
942 &format!("{id}.platform"),
943 row.get("platform"),
944 )?),
945 kind: parsed.as_ref().and_then(|k| k.kind.clone()),
946 chat_id: Some(row.get("chat_id").and_then(text_of).unwrap_or_default()),
947 thread_id: row
948 .get("thread_id")
949 .and_then(text_of)
950 .filter(|s| !s.is_empty()),
951 participant_id: None,
952 };
953 let text = |k: &str| row.get(k).and_then(text_of);
954 let created_at = text("created_at").unwrap_or_else(|| "undefined".into());
955 let updated_at = text("updated_at").unwrap_or_else(|| created_at.clone());
956 Ok(Obligation {
957 target,
958 session_key,
959 content: OutboundContent {
960 text: row.get("content").and_then(text_of).unwrap_or_default(),
961 attachments: None,
962 reply_to: None,
963 format: None,
964 },
965 state,
966 attempts: row.get("attempts").and_then(number_of).unwrap_or(0.0) as u64,
967 last_error: row.get("last_error").and_then(text_of),
968 delivered_at: if state == ObligationState::Sent {
969 Some(updated_at.clone())
970 } else {
971 None
972 },
973 created_at,
974 updated_at,
975 posted: row
979 .get("posted_message_id")
980 .and_then(text_of)
981 .filter(|s| !s.is_empty())
982 .map(|message_id| Posted { message_id }),
983 source: row
984 .get("source_json")
985 .and_then(text_of)
986 .and_then(|text| serde_json::from_str::<ObligationSource>(&text).ok())
987 .unwrap_or(ObligationSource::Turn { key: None }),
988 residue: row_residue_of(
989 row,
990 &[
991 "obligation_id",
992 "session_key",
993 "platform",
994 "chat_id",
995 "thread_id",
996 "content",
997 "state",
998 "attempts",
999 "created_at",
1000 "updated_at",
1001 "last_error",
1002 "posted_message_id",
1003 "source_json",
1004 ],
1005 ),
1006 id,
1007 })
1008}
1009
1010pub const FOLDER_OBLIGATION_EXTRA_COLUMNS: &[&str] = &["posted_message_id", "source_json"];
1012
1013pub fn encode_obligation_folder_extras(o: &Obligation) -> Vec<Value> {
1015 vec![
1016 o.posted
1017 .as_ref()
1018 .map(|p| Value::String(p.message_id.clone()))
1019 .unwrap_or(Value::Null),
1020 Value::String(serde_json::to_string(&o.source).unwrap()),
1021 ]
1022}
1023
1024pub fn encode_obligation_row(o: &Obligation) -> Vec<Value> {
1026 let r = &o.residue.0;
1027 let num = |s: &str| {
1028 s.parse::<f64>()
1029 .map(|f| serde_json::json!(f))
1030 .unwrap_or(Value::from(0))
1031 };
1032 vec![
1033 Value::String(o.id.clone()),
1034 Value::String(o.session_key.clone().unwrap_or_default()),
1035 Value::String(o.target.platform.clone().unwrap_or_default()),
1036 Value::String(o.target.chat_id.clone().unwrap_or_default()),
1037 o.target
1038 .thread_id
1039 .clone()
1040 .map(Value::String)
1041 .unwrap_or(Value::Null),
1042 Value::String(o.content.text.clone()),
1043 Value::String(o.state.hermes_word().into()),
1044 Value::from(o.attempts),
1045 num(&o.created_at),
1046 num(o
1047 .delivered_at
1048 .as_deref()
1049 .map(|_| o.updated_at.as_str())
1050 .unwrap_or(&o.updated_at)),
1051 r.get("owner_pid").cloned().unwrap_or(Value::Null),
1052 r.get("owner_started_at").cloned().unwrap_or(Value::Null),
1053 o.last_error
1054 .clone()
1055 .map(Value::String)
1056 .unwrap_or(Value::Null),
1057 r.get("adapter_profile").cloned().unwrap_or(Value::Null),
1058 ]
1059}
1060
1061pub const OBLIGATION_COLUMNS: &[&str] = &[
1063 "obligation_id",
1064 "session_key",
1065 "platform",
1066 "chat_id",
1067 "thread_id",
1068 "content",
1069 "state",
1070 "attempts",
1071 "created_at",
1072 "updated_at",
1073 "owner_pid",
1074 "owner_started_at",
1075 "last_error",
1076 "adapter_profile",
1077];
1078
1079const ROUTE_MATCH: &[&str] = &["platform", "guild_id", "chat_id", "thread_id"];
1080
1081pub fn decode_route(file: &str, index: usize, raw: &Value) -> Result<Route> {
1083 let map = raw.as_object().ok_or_else(|| {
1084 load_error(
1085 file,
1086 &format!("profile_routes[{index}]"),
1087 "expected an object",
1088 )
1089 })?;
1090 let text = |k: &str| map.get(k).filter(|v| !v.is_null()).and_then(text_of);
1091 Ok(Route {
1092 name: opt_str(
1093 file,
1094 &format!("profile_routes[{index}].name"),
1095 map.get("name"),
1096 )?,
1097 matches: RouteMatch {
1098 platform: req_str(
1099 file,
1100 &format!("profile_routes[{index}].platform"),
1101 map.get("platform"),
1102 )?,
1103 guild_id: text("guild_id"),
1104 chat_id: text("chat_id"),
1105 thread_id: text("thread_id"),
1106 },
1107 agent: match map.get("agent").filter(|v| !v.is_null()) {
1109 Some(agent) => req_str(file, &format!("profile_routes[{index}].agent"), Some(agent))?,
1110 None => req_str(
1111 file,
1112 &format!("profile_routes[{index}].profile"),
1113 map.get("profile"),
1114 )?,
1115 },
1116 residue: residue_of(
1117 map,
1118 &[
1119 "name",
1120 "agent",
1121 "profile",
1122 "platform",
1123 "guild_id",
1124 "chat_id",
1125 "thread_id",
1126 ],
1127 ),
1128 })
1129}
1130
1131pub fn encode_route(
1134 r: &Route,
1135 agents: &std::collections::BTreeMap<String, crate::orchestration::AgentDecl>,
1136 hermes: bool,
1137) -> Value {
1138 let mut out = Map::new();
1139 if let Some(n) = &r.name {
1140 out.insert("name".into(), Value::String(n.clone()));
1141 }
1142 let m = &r.matches;
1143 for (k, v) in [
1144 ("platform", Some(&m.platform)),
1145 ("guild_id", m.guild_id.as_ref()),
1146 ("chat_id", m.chat_id.as_ref()),
1147 ("thread_id", m.thread_id.as_ref()),
1148 ] {
1149 if let Some(v) = v {
1150 out.insert(k.into(), Value::String(v.clone()));
1151 }
1152 }
1153 let _ = ROUTE_MATCH;
1154 let profile = agents
1155 .get(&r.agent)
1156 .map(|decl| decl.profile.clone())
1157 .unwrap_or_else(|| r.agent.clone());
1158 if !hermes && profile != r.agent {
1159 out.insert("agent".into(), Value::String(r.agent.clone()));
1160 }
1161 out.insert("profile".into(), Value::String(profile));
1162 for (k, v) in &r.residue.0 {
1163 out.insert(k.clone(), v.clone());
1164 }
1165 Value::Object(out)
1166}
1167
1168pub fn is_credential_key(key: &str) -> bool {
1170 let k = key.to_ascii_lowercase();
1171 [
1172 "token",
1173 "secret",
1174 "key",
1175 "password",
1176 "passwd",
1177 "api_key",
1178 "app_secret",
1179 "signing",
1180 "webhook_url",
1181 "private",
1182 ]
1183 .iter()
1184 .any(|w| k.contains(w))
1185}
1186
1187pub fn placeholder_ref(text: &str) -> Option<&str> {
1194 let name = text.strip_prefix("${")?.strip_suffix('}')?;
1195 let mut chars = name.chars();
1196 let first = chars.next()?;
1197 ((first.is_ascii_alphabetic() || first == '_')
1198 && chars.all(|c| c.is_ascii_alphanumeric() || c == '_'))
1199 .then_some(name)
1200}
1201
1202pub fn render_ref(r: &SecretRef) -> Value {
1205 Value::String(format!("${{{}}}", r.name()))
1206}
1207
1208pub fn credential_ref_name(platform: &str, key: &str) -> String {
1210 format!("{platform}_{key}")
1211 .chars()
1212 .map(|c| {
1213 if c.is_ascii_alphanumeric() {
1214 c.to_ascii_uppercase()
1215 } else {
1216 '_'
1217 }
1218 })
1219 .collect()
1220}
1221
1222pub fn decode_channel(
1224 file: &str,
1225 platform: &str,
1226 raw: &Value,
1227 vault: &mut BTreeMap<String, String>,
1228) -> Result<ChannelConfig> {
1229 let map = raw
1230 .as_object()
1231 .ok_or_else(|| load_error(file, &format!("platforms.{platform}"), "expected an object"))?;
1232 let mut ch = ChannelConfig {
1233 platform: platform.to_string(),
1234 enabled: opt_bool(
1235 file,
1236 &format!("platforms.{platform}.enabled"),
1237 map.get("enabled"),
1238 Some(true),
1239 )?
1240 .unwrap_or(true),
1241 credentials: BTreeMap::new(),
1242 extra: BTreeMap::new(),
1243 };
1244 fn walk(
1245 map: &Map<String, Value>,
1246 prefix: &str,
1247 platform: &str,
1248 ch: &mut ChannelConfig,
1249 vault: &mut BTreeMap<String, String>,
1250 ) {
1251 for (k, v) in map {
1252 if k == "enabled" && prefix.is_empty() {
1253 continue;
1254 }
1255 let name = format!("{prefix}{k}");
1256 if let Value::String(s) = v {
1257 if let Some(r) = placeholder_ref(s) {
1258 ch.credentials
1259 .insert(name, SecretRef::Dotenv(r.to_string()));
1260 continue;
1261 }
1262 }
1263 if let Some(obj) = v.as_object() {
1264 if k == "extra" && prefix.is_empty() {
1265 walk(obj, "extra.", platform, ch, vault);
1266 continue;
1267 }
1268 }
1269 if let Value::String(s) = v {
1270 if is_credential_key(k) {
1271 let r = credential_ref_name(platform, &name);
1272 vault.insert(r.clone(), s.clone());
1273 ch.credentials.insert(name, SecretRef::Dotenv(r));
1274 continue;
1275 }
1276 }
1277 ch.extra.insert(name, v.clone());
1278 }
1279 }
1280 walk(map, "", platform, &mut ch, vault);
1281 Ok(ch)
1282}
1283
1284pub fn encode_channel(ch: &ChannelConfig, vault: Option<&BTreeMap<String, String>>) -> Value {
1286 let mut out = Map::new();
1287 out.insert("enabled".into(), Value::Bool(ch.enabled));
1288 let mut extra = Map::new();
1289 for (k, v) in &ch.extra {
1290 match k.strip_prefix("extra.") {
1291 Some(inner) => {
1292 extra.insert(inner.to_string(), v.clone());
1293 }
1294 None => {
1295 out.insert(k.clone(), v.clone());
1296 }
1297 }
1298 }
1299 for (k, r) in &ch.credentials {
1300 let rendered = match vault.and_then(|v| v.get(r.name())) {
1301 Some(value) => Value::String(value.clone()),
1302 None => render_ref(r),
1303 };
1304 match k.strip_prefix("extra.") {
1305 Some(inner) => {
1306 extra.insert(inner.to_string(), rendered);
1307 }
1308 None => {
1309 out.insert(k.clone(), rendered);
1310 }
1311 }
1312 }
1313 if !extra.is_empty() {
1314 out.insert("extra".into(), Value::Object(extra));
1315 }
1316 Value::Object(out)
1317}
1318
1319const SUB_MAPPED: &[&str] = &[
1320 "events",
1321 "prompt",
1322 "skills",
1323 "deliver",
1324 "deliver_extra",
1325 "secret",
1326 "description",
1327 "created_at",
1328];
1329
1330pub fn decode_subscription(
1332 file: &str,
1333 name: &str,
1334 raw: &Value,
1335 vault: &mut BTreeMap<String, String>,
1336) -> Result<WebhookSubscription> {
1337 let map = raw
1338 .as_object()
1339 .ok_or_else(|| load_error(file, name, "expected an object"))?;
1340 let mut sub = WebhookSubscription {
1341 name: name.to_string(),
1342 secret: None,
1343 events: map.get("events").and_then(Value::as_array).map(|a| {
1344 a.iter()
1345 .map(|v| {
1346 v.as_str()
1347 .map(str::to_string)
1348 .unwrap_or_else(|| v.to_string())
1349 })
1350 .collect()
1351 }),
1352 prompt_template: opt_str(file, &format!("{name}.prompt"), map.get("prompt"))?
1353 .unwrap_or_default(),
1354 deliver: parse_target(
1355 file,
1356 &format!("{name}.deliver"),
1357 map.get("deliver"),
1358 map.get("deliver_extra"),
1359 )?,
1360 skills: map
1361 .get("skills")
1362 .and_then(Value::as_array)
1363 .map(|a| {
1364 a.iter()
1365 .map(|v| {
1366 v.as_str()
1367 .map(str::to_string)
1368 .unwrap_or_else(|| v.to_string())
1369 })
1370 .collect()
1371 })
1372 .unwrap_or_default(),
1373 description: opt_str(file, &format!("{name}.description"), map.get("description"))?,
1374 created_at: opt_str(file, &format!("{name}.created_at"), map.get("created_at"))?,
1375 residue: residue_of(map, SUB_MAPPED),
1376 };
1377 match map.get("secret") {
1378 Some(Value::String(secret)) => {
1381 let r = match placeholder_ref(secret) {
1382 Some(name) => name.to_string(),
1383 None => {
1384 let r = credential_ref_name("webhook", &format!("{name}_secret"));
1385 vault.insert(r.clone(), secret.clone());
1386 r
1387 }
1388 };
1389 sub.secret = Some(SecretRef::Dotenv(r));
1390 }
1391 Some(Value::Null) | None => {}
1392 Some(_) => {
1393 return Err(load_error(
1394 file,
1395 &format!("{name}.secret"),
1396 "expected a string (a secret is written `${NAME}`)",
1397 ))
1398 }
1399 }
1400 Ok(sub)
1401}
1402
1403pub fn encode_subscription(
1405 sub: &WebhookSubscription,
1406 vault: Option<&BTreeMap<String, String>>,
1407) -> Value {
1408 let mut out = Map::new();
1409 if let Some(d) = &sub.description {
1410 out.insert("description".into(), Value::String(d.clone()));
1411 }
1412 if let Some(e) = &sub.events {
1413 out.insert(
1414 "events".into(),
1415 Value::Array(e.iter().map(|s| Value::String(s.clone())).collect()),
1416 );
1417 }
1418 out.insert("prompt".into(), Value::String(sub.prompt_template.clone()));
1419 out.insert(
1420 "skills".into(),
1421 Value::Array(
1422 sub.skills
1423 .iter()
1424 .map(|s| Value::String(s.clone()))
1425 .collect(),
1426 ),
1427 );
1428 if let Some(t) = &sub.deliver {
1429 match t {
1430 Target::Explicit {
1431 platform,
1432 chat_id,
1433 thread_id,
1434 } => {
1435 out.insert("deliver".into(), Value::String(platform.clone()));
1436 let mut extra = Map::new();
1437 if let Some(c) = chat_id {
1438 extra.insert("chat_id".into(), Value::String(c.clone()));
1439 }
1440 if let Some(th) = thread_id {
1441 extra.insert("thread_id".into(), Value::String(th.clone()));
1442 }
1443 if !extra.is_empty() {
1444 out.insert("deliver_extra".into(), Value::Object(extra));
1445 }
1446 }
1447 other => {
1448 out.insert("deliver".into(), Value::String(other.render()));
1449 }
1450 }
1451 }
1452 if let Some(s) = &sub.secret {
1453 let value = vault
1454 .and_then(|v| v.get(s.name()).cloned())
1455 .unwrap_or_else(|| format!("${{{}}}", s.name()));
1456 out.insert("secret".into(), Value::String(value));
1457 }
1458 if let Some(c) = &sub.created_at {
1459 out.insert("created_at".into(), Value::String(c.clone()));
1460 }
1461 for (k, v) in &sub.residue.0 {
1462 out.insert(k.clone(), v.clone());
1463 }
1464 Value::Object(out)
1465}