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 "preset",
307 "cwd",
308 "env",
309 "permission",
310 "home",
311 "surface",
312 "capacity",
313 ],
314 )?;
315 let cwd = match raw.get("cwd") {
316 None => ".".to_string(),
317 v => req_str(file, "worker.cwd", v)?,
318 };
319 if cwd.starts_with('/') || cwd.split('/').any(|p| p == "..") {
320 return Err(load_error(
321 file,
322 "worker.cwd",
323 "must be a relative path inside the profile",
324 ));
325 }
326 let mut env = BTreeMap::new();
327 if let Some(e) = raw.get("env") {
328 let map = e
329 .as_object()
330 .ok_or_else(|| load_error(file, "worker.env", "expected a map"))?;
331 for (k, v) in map {
332 env.insert(
333 k.clone(),
334 decode_env_value(file, &format!("worker.env.{k}"), v)?,
335 );
336 }
337 }
338 let mut permission = PermissionPolicy::default();
339 if let Some(p) = raw.get("permission") {
340 expect_keys(
341 file,
342 "worker.permission",
343 p,
344 &["timeout_seconds", "default", "unattended"],
345 )?;
346 if let Some(t) = p.get("timeout_seconds") {
347 permission.timeout_seconds = t.as_u64().filter(|v| *v >= 1).ok_or_else(|| {
348 load_error(
349 file,
350 "worker.permission.timeout_seconds",
351 "expected an integer >= 1",
352 )
353 })? as u32;
354 }
355 if let Some(d) = p.get("default") {
356 permission.default = match d.as_str() {
357 Some("deny") => PermissionDefault::Deny,
358 Some("allow") => PermissionDefault::Allow,
359 _ => {
360 return Err(load_error(
361 file,
362 "worker.permission.default",
363 "expected deny|allow",
364 ))
365 }
366 };
367 }
368 if let Some(u) = p.get("unattended") {
369 permission.unattended = match u.as_str() {
370 Some("deny") => PermissionUnattended::Deny,
371 Some("approve") => PermissionUnattended::Approve,
372 _ => {
373 return Err(load_error(
374 file,
375 "worker.permission.unattended",
376 "expected deny|approve",
377 ))
378 }
379 };
380 }
381 }
382 let home = match opt_str(file, "worker.home", raw.get("home"))?.as_deref() {
383 None | Some("profile") => crate::orchestration::WorkerHome::Profile,
384 Some("user") => crate::orchestration::WorkerHome::User,
385 Some(_) => {
386 return Err(load_error(
387 file,
388 "worker.home",
389 "expected `profile` or `user`",
390 ))
391 }
392 };
393 let surface = match opt_str(file, "worker.surface", raw.get("surface"))?.as_deref() {
394 None | Some("headless") => crate::orchestration::WorkerSurface::Headless,
395 Some("pane") => crate::orchestration::WorkerSurface::Pane,
396 Some(_) => {
397 return Err(load_error(
398 file,
399 "worker.surface",
400 "expected `headless` or `pane`",
401 ))
402 }
403 };
404 let capacity = match raw.get("capacity") {
405 None => None,
406 Some(v) => Some(
407 v.as_u64()
408 .filter(|n| *n >= 1 && *n <= u64::from(u32::MAX))
409 .ok_or_else(|| load_error(file, "worker.capacity", "expected an integer >= 1"))?
410 as u32,
411 ),
412 };
413 Ok(Some(WorkerSpec {
414 harness: HarnessId::new(req_str(file, "worker.harness", raw.get("harness"))?),
415 model: opt_str(file, "worker.model", raw.get("model"))?,
416 preset: opt_str(file, "worker.preset", raw.get("preset"))?,
417 cwd,
418 env,
419 permission,
420 home,
421 surface,
422 capacity,
423 }))
424}
425
426const JOB_MAPPED: &[&str] = &[
429 "id",
430 "schedule",
431 "prompt",
432 "workdir",
433 "model",
434 "skills",
435 "context_from",
436 "deliver",
437 "failure_deliver",
438 "origin",
439 "attach_to_session",
440 "repeat",
441 "enabled",
442 "next_run_at",
443 "last_run_at",
444 "last_status",
445 "created_at",
446];
447
448pub const HERMES_JOB_ORDER: &[&str] = &[
450 "id",
451 "schedule",
452 "prompt",
453 "skills",
454 "script",
455 "no_agent",
456 "model",
457 "provider",
458 "workdir",
459 "enabled_toolsets",
460 "context_from",
461 "deliver",
462 "failure_deliver",
463 "attach_to_session",
464 "origin",
465 "repeat",
466 "enabled",
467 "next_run_at",
468 "last_run_at",
469 "last_status",
470 "created_at",
471 "fire_claim",
472];
473
474pub fn decode_context_from(
476 file: &str,
477 key: &str,
478 v: Option<&Value>,
479) -> Result<Option<Vec<String>>> {
480 let items: Vec<String> = match v {
481 None | Some(Value::Null) => return Ok(None),
482 Some(Value::String(s)) => vec![s.clone()],
483 Some(Value::Array(a)) => a
484 .iter()
485 .map(|x| match x {
486 Value::String(s) => s.clone(),
487 other => other.to_string(),
488 })
489 .collect(),
490 _ => {
491 return Err(load_error(
492 file,
493 key,
494 "expected a job id, a list of job ids, or \"self\"",
495 ))
496 }
497 };
498 let refs: Vec<String> = items
499 .into_iter()
500 .map(|s| s.trim().to_string())
501 .filter(|s| !s.is_empty())
502 .collect();
503 Ok(if refs.is_empty() { None } else { Some(refs) })
504}
505
506pub fn decode_repeat(file: &str, key: &str, raw: Option<&Value>) -> Result<Option<Repeat>> {
508 Ok(match raw {
509 None | Some(Value::Null) => None,
510 Some(Value::Bool(true)) => Some(Repeat {
511 times: None,
512 completed: 0,
513 }),
514 Some(Value::Bool(false)) => Some(Repeat {
515 times: Some(1),
516 completed: 0,
517 }),
518 Some(Value::Number(n)) => Some(Repeat {
519 times: n.as_f64().filter(|f| *f > 0.0).map(|f| f.floor() as u32),
520 completed: 0,
521 }),
522 Some(Value::Object(map)) => {
523 let times = match map.get("times") {
524 None | Some(Value::Null) => None,
525 Some(Value::Number(n)) if n.as_u64().is_some_and(|v| v >= 1) => {
526 Some(n.as_u64().unwrap() as u32)
527 }
528 _ => {
529 return Err(load_error(
530 file,
531 &format!("{key}.times"),
532 "expected null or an integer >= 1",
533 ))
534 }
535 };
536 let completed = match map.get("completed") {
537 None => 0,
538 Some(Value::Number(n)) if n.as_u64().is_some() => n.as_u64().unwrap() as u32,
539 _ => {
540 return Err(load_error(
541 file,
542 &format!("{key}.completed"),
543 "expected an integer >= 0",
544 ))
545 }
546 };
547 Some(Repeat { times, completed })
548 }
549 _ => {
550 return Err(load_error(
551 file,
552 key,
553 "expected null, {times, completed}, a number or a boolean",
554 ))
555 }
556 })
557}
558
559pub fn decode_job(file: &str, raw: &Value) -> Result<Job> {
561 let map = raw
562 .as_object()
563 .ok_or_else(|| load_error(file, "", "expected a job object"))?;
564 let id = req_str(file, "id", map.get("id"))?;
565 let k = |s: &str| format!("{id}.{s}");
566 let sched = map
567 .get("schedule")
568 .and_then(Value::as_object)
569 .ok_or_else(|| load_error(file, &k("schedule"), "expected an object"))?;
570 let schedule = match sched.get("kind").and_then(Value::as_str) {
571 Some("once") => Schedule::Once {
572 run_at: req_str(file, &k("schedule.run_at"), sched.get("run_at"))?,
573 },
574 Some("interval") => {
575 let minutes = sched
576 .get("minutes")
577 .and_then(Value::as_f64)
578 .filter(|m| *m > 0.0)
579 .ok_or_else(|| {
580 load_error(file, &k("schedule.minutes"), "expected a positive number")
581 })?;
582 Schedule::Interval { minutes }
583 }
584 Some("cron") => Schedule::Cron {
585 expr: req_str(file, &k("schedule.expr"), sched.get("expr"))?,
586 tz: opt_str(file, &k("schedule.tz"), sched.get("tz"))?,
587 },
588 _ => {
589 return Err(load_error(
590 file,
591 &k("schedule.kind"),
592 "expected once|interval|cron",
593 ))
594 }
595 };
596 let schedule_residue = residue_of(sched, &["kind", "run_at", "minutes", "expr", "tz"]);
597 let mut residue = residue_of(map, JOB_MAPPED);
598 if !schedule_residue.is_empty() {
599 residue.keep(
600 "__schedule",
601 Value::Object(schedule_residue.0.into_iter().collect()),
602 );
603 }
604 let origin = match map.get("origin") {
605 Some(Value::Object(o)) => {
606 let r = residue_of(o, &["platform", "chat_id", "thread_id"]);
607 if !r.is_empty() {
608 residue.keep("__origin", Value::Object(r.0.into_iter().collect()));
609 }
610 Some(JobOrigin {
611 platform: req_str(file, &k("origin.platform"), o.get("platform"))?,
612 chat_type: None,
613 chat_id: opt_str(file, &k("origin.chat_id"), o.get("chat_id"))?,
614 thread_id: opt_str(file, &k("origin.thread_id"), o.get("thread_id"))?,
615 })
616 }
617 _ => None,
618 };
619 Ok(Job {
620 schedule,
621 prompt: opt_str(file, &k("prompt"), map.get("prompt"))?,
622 workdir: opt_str(file, &k("workdir"), map.get("workdir"))?,
623 model: opt_str(file, &k("model"), map.get("model"))?,
624 skills: map
625 .get("skills")
626 .and_then(Value::as_array)
627 .map(|a| {
628 a.iter()
629 .map(|v| {
630 v.as_str()
631 .map(str::to_string)
632 .unwrap_or_else(|| v.to_string())
633 })
634 .collect()
635 })
636 .unwrap_or_default(),
637 context_from: decode_context_from(file, &k("context_from"), map.get("context_from"))?,
638 deliver: parse_target(file, &k("deliver"), map.get("deliver"), None)?
639 .unwrap_or(Target::Local),
640 failure_deliver: parse_target(
641 file,
642 &k("failure_deliver"),
643 map.get("failure_deliver"),
644 None,
645 )?,
646 origin,
647 attach_to_session: opt_bool(
648 file,
649 &k("attach_to_session"),
650 map.get("attach_to_session"),
651 None,
652 )?,
653 repeat: decode_repeat(file, &k("repeat"), map.get("repeat"))?,
654 enabled: opt_bool(file, &k("enabled"), map.get("enabled"), Some(true))?.unwrap_or(true),
655 next_run_at: opt_str(file, &k("next_run_at"), map.get("next_run_at"))?,
656 last_run_at: opt_str(file, &k("last_run_at"), map.get("last_run_at"))?,
657 last_status: opt_str(file, &k("last_status"), map.get("last_status"))?,
658 created_at: opt_str(file, &k("created_at"), map.get("created_at"))?,
659 residue,
660 id,
661 })
662}
663
664pub fn encode_job(job: &Job) -> Vec<(String, Value)> {
667 let mut sched = Map::new();
668 sched.insert("kind".into(), Value::String(job.schedule.kind().into()));
669 match &job.schedule {
670 Schedule::Once { run_at } => {
671 sched.insert("run_at".into(), Value::String(run_at.clone()));
672 }
673 Schedule::Interval { minutes } => {
674 sched.insert("minutes".into(), serde_json::json!(*minutes));
675 }
676 Schedule::Cron { expr, tz } => {
677 sched.insert("expr".into(), Value::String(expr.clone()));
678 if let Some(tz) = tz {
679 sched.insert("tz".into(), Value::String(tz.clone()));
680 }
681 }
682 }
683 if let Some(Value::Object(extra)) = job.residue.0.get("__schedule") {
684 for (k, v) in extra {
685 sched.insert(k.clone(), v.clone());
686 }
687 }
688 let origin = job
689 .origin
690 .as_ref()
691 .map(|o| {
692 let mut m = Map::new();
693 m.insert("platform".into(), Value::String(o.platform.clone()));
694 m.insert(
695 "chat_id".into(),
696 o.chat_id.clone().map(Value::String).unwrap_or(Value::Null),
697 );
698 m.insert(
699 "thread_id".into(),
700 o.thread_id
701 .clone()
702 .map(Value::String)
703 .unwrap_or(Value::Null),
704 );
705 if let Some(Value::Object(extra)) = job.residue.0.get("__origin") {
706 for (k, v) in extra {
707 m.insert(k.clone(), v.clone());
708 }
709 }
710 Value::Object(m)
711 })
712 .unwrap_or(Value::Null);
713 let opt = |s: &Option<String>| s.clone().map(Value::String).unwrap_or(Value::Null);
714 let mut mapped: Vec<(String, Value)> = vec![
715 ("id".into(), Value::String(job.id.clone())),
716 ("schedule".into(), Value::Object(sched)),
717 ("prompt".into(), opt(&job.prompt)),
718 (
719 "skills".into(),
720 Value::Array(
721 job.skills
722 .iter()
723 .map(|s| Value::String(s.clone()))
724 .collect(),
725 ),
726 ),
727 ("model".into(), opt(&job.model)),
728 ("workdir".into(), opt(&job.workdir)),
729 (
730 "context_from".into(),
731 job.context_from
732 .as_ref()
733 .map(|l| Value::Array(l.iter().map(|s| Value::String(s.clone())).collect()))
734 .unwrap_or(Value::Null),
735 ),
736 ("deliver".into(), Value::String(job.deliver.render())),
737 (
738 "failure_deliver".into(),
739 job.failure_deliver
740 .as_ref()
741 .map(|t| Value::String(t.render()))
742 .unwrap_or(Value::Null),
743 ),
744 (
745 "attach_to_session".into(),
746 job.attach_to_session
747 .map(Value::Bool)
748 .unwrap_or(Value::Null),
749 ),
750 ("origin".into(), origin),
751 (
752 "repeat".into(),
753 job.repeat
754 .as_ref()
755 .map(|r| serde_json::json!({"times": r.times, "completed": r.completed}))
756 .unwrap_or(Value::Null),
757 ),
758 ("enabled".into(), Value::Bool(job.enabled)),
759 ("next_run_at".into(), opt(&job.next_run_at)),
760 ("last_run_at".into(), opt(&job.last_run_at)),
761 ("last_status".into(), opt(&job.last_status)),
762 ("created_at".into(), opt(&job.created_at)),
763 ];
764 let mut residue: Vec<(String, Value)> = job
765 .residue
766 .0
767 .iter()
768 .filter(|(k, _)| k.as_str() != "__schedule" && k.as_str() != "__origin")
774 .map(|(k, v)| (k.clone(), v.clone()))
775 .collect();
776 let mut out: Vec<(String, Value)> = Vec::new();
777 for key in HERMES_JOB_ORDER {
778 if let Some(pos) = mapped.iter().position(|(k, _)| k == key) {
779 out.push(mapped.remove(pos));
780 } else if let Some(pos) = residue.iter().position(|(k, _)| k == key) {
781 out.push(residue.remove(pos));
782 }
783 }
784 out.extend(mapped);
785 out.extend(residue);
786 out
787}
788
789pub fn decode_fire_row(file: &str, row: &Map<String, Value>) -> Result<Fire> {
791 let id = req_str(file, "id", row.get("id"))?;
792 let status_word = row.get("status").and_then(Value::as_str).unwrap_or("");
793 let status = FireStatus::from_hermes_word(status_word).ok_or_else(|| {
794 load_error(
795 file,
796 &format!("{id}.status"),
797 format!(
798 "unknown fire status {}",
799 serde_json::to_string(status_word).unwrap()
800 ),
801 )
802 })?;
803 Ok(Fire {
804 job_id: req_str(file, &format!("{id}.job_id"), row.get("job_id"))?,
805 session_id: row
806 .get("session_id")
807 .and_then(text_of)
808 .filter(|s| !s.is_empty()),
809 status,
810 claimed_at: req_str(file, &format!("{id}.claimed_at"), row.get("claimed_at"))?,
811 started_at: opt_str(file, &format!("{id}.started_at"), row.get("started_at"))?,
812 finished_at: opt_str(file, &format!("{id}.finished_at"), row.get("finished_at"))?,
813 error: opt_str(file, &format!("{id}.error"), row.get("error"))?,
814 obligation_id: row
815 .get("obligation_id")
816 .and_then(text_of)
817 .filter(|s| !s.is_empty()),
818 residue: {
819 let mut residue = row_residue_of(
820 row,
821 &[
822 "id",
823 "job_id",
824 "status",
825 "claimed_at",
826 "started_at",
827 "finished_at",
828 "error",
829 "residue_json",
830 "session_id",
831 "obligation_id",
832 ],
833 );
834 if let Some(extra) = row
837 .get("residue_json")
838 .and_then(text_of)
839 .and_then(|t| serde_json::from_str::<Map<String, Value>>(&t).ok())
840 {
841 for (k, v) in extra {
842 residue.keep(k, v);
843 }
844 }
845 residue
846 },
847 id,
848 })
849}
850
851pub fn encode_fire_row(fire: &Fire) -> Vec<Value> {
853 let r = &fire.residue.0;
854 let opt = |s: &Option<String>| s.clone().map(Value::String).unwrap_or(Value::Null);
855 vec![
856 Value::String(fire.id.clone()),
857 Value::String(fire.job_id.clone()),
858 r.get("source")
859 .cloned()
860 .unwrap_or(Value::String("scheduler".into())),
861 r.get("process_id")
862 .cloned()
863 .unwrap_or(Value::String(String::new())),
864 r.get("pid").cloned().unwrap_or(Value::from(0)),
865 r.get("process_started_at").cloned().unwrap_or(Value::Null),
866 Value::String(fire.status.hermes_word().into()),
867 Value::String(fire.claimed_at.clone()),
868 opt(&fire.started_at),
869 opt(&fire.finished_at),
870 opt(&fire.error),
871 ]
872}
873
874pub const EXECUTION_COLUMNS: &[&str] = &[
876 "id",
877 "job_id",
878 "source",
879 "process_id",
880 "pid",
881 "process_started_at",
882 "status",
883 "claimed_at",
884 "started_at",
885 "finished_at",
886 "error",
887];
888
889pub fn decode_obligation_row(file: &str, row: &Map<String, Value>) -> Result<Obligation> {
891 let id = req_str(file, "obligation_id", row.get("obligation_id"))?;
892 let state_word = row.get("state").and_then(Value::as_str).unwrap_or("");
893 let state = ObligationState::from_hermes_word(state_word).ok_or_else(|| {
894 load_error(
895 file,
896 &format!("{id}.state"),
897 format!(
898 "unknown obligation state {}",
899 serde_json::to_string(state_word).unwrap()
900 ),
901 )
902 })?;
903 let session_key = row
904 .get("session_key")
905 .and_then(Value::as_str)
906 .filter(|s| !s.is_empty())
907 .map(str::to_string);
908 let parsed = session_key
909 .as_deref()
910 .and_then(crate::ontology::parse_hermes_session_key)
911 .map(|(k, _)| k);
912 let target = SurfaceKey {
913 key: None,
914 platform: Some(req_str(
915 file,
916 &format!("{id}.platform"),
917 row.get("platform"),
918 )?),
919 kind: parsed.as_ref().and_then(|k| k.kind.clone()),
920 chat_id: Some(row.get("chat_id").and_then(text_of).unwrap_or_default()),
921 thread_id: row
922 .get("thread_id")
923 .and_then(text_of)
924 .filter(|s| !s.is_empty()),
925 participant_id: None,
926 };
927 let text = |k: &str| row.get(k).and_then(text_of);
928 let created_at = text("created_at").unwrap_or_else(|| "undefined".into());
929 let updated_at = text("updated_at").unwrap_or_else(|| created_at.clone());
930 Ok(Obligation {
931 target,
932 session_key,
933 content: OutboundContent {
934 text: row.get("content").and_then(text_of).unwrap_or_default(),
935 attachments: None,
936 reply_to: None,
937 format: None,
938 },
939 state,
940 attempts: row.get("attempts").and_then(number_of).unwrap_or(0.0) as u64,
941 last_error: row.get("last_error").and_then(text_of),
942 delivered_at: if state == ObligationState::Sent {
943 Some(updated_at.clone())
944 } else {
945 None
946 },
947 created_at,
948 updated_at,
949 posted: row
953 .get("posted_message_id")
954 .and_then(text_of)
955 .filter(|s| !s.is_empty())
956 .map(|message_id| Posted { message_id }),
957 source: row
958 .get("source_json")
959 .and_then(text_of)
960 .and_then(|text| serde_json::from_str::<ObligationSource>(&text).ok())
961 .unwrap_or(ObligationSource::Turn { key: None }),
962 residue: row_residue_of(
963 row,
964 &[
965 "obligation_id",
966 "session_key",
967 "platform",
968 "chat_id",
969 "thread_id",
970 "content",
971 "state",
972 "attempts",
973 "created_at",
974 "updated_at",
975 "last_error",
976 "posted_message_id",
977 "source_json",
978 ],
979 ),
980 id,
981 })
982}
983
984pub const FOLDER_OBLIGATION_EXTRA_COLUMNS: &[&str] = &["posted_message_id", "source_json"];
986
987pub fn encode_obligation_folder_extras(o: &Obligation) -> Vec<Value> {
989 vec![
990 o.posted
991 .as_ref()
992 .map(|p| Value::String(p.message_id.clone()))
993 .unwrap_or(Value::Null),
994 Value::String(serde_json::to_string(&o.source).unwrap()),
995 ]
996}
997
998pub fn encode_obligation_row(o: &Obligation) -> Vec<Value> {
1000 let r = &o.residue.0;
1001 let num = |s: &str| {
1002 s.parse::<f64>()
1003 .map(|f| serde_json::json!(f))
1004 .unwrap_or(Value::from(0))
1005 };
1006 vec![
1007 Value::String(o.id.clone()),
1008 Value::String(o.session_key.clone().unwrap_or_default()),
1009 Value::String(o.target.platform.clone().unwrap_or_default()),
1010 Value::String(o.target.chat_id.clone().unwrap_or_default()),
1011 o.target
1012 .thread_id
1013 .clone()
1014 .map(Value::String)
1015 .unwrap_or(Value::Null),
1016 Value::String(o.content.text.clone()),
1017 Value::String(o.state.hermes_word().into()),
1018 Value::from(o.attempts),
1019 num(&o.created_at),
1020 num(o
1021 .delivered_at
1022 .as_deref()
1023 .map(|_| o.updated_at.as_str())
1024 .unwrap_or(&o.updated_at)),
1025 r.get("owner_pid").cloned().unwrap_or(Value::Null),
1026 r.get("owner_started_at").cloned().unwrap_or(Value::Null),
1027 o.last_error
1028 .clone()
1029 .map(Value::String)
1030 .unwrap_or(Value::Null),
1031 r.get("adapter_profile").cloned().unwrap_or(Value::Null),
1032 ]
1033}
1034
1035pub const OBLIGATION_COLUMNS: &[&str] = &[
1037 "obligation_id",
1038 "session_key",
1039 "platform",
1040 "chat_id",
1041 "thread_id",
1042 "content",
1043 "state",
1044 "attempts",
1045 "created_at",
1046 "updated_at",
1047 "owner_pid",
1048 "owner_started_at",
1049 "last_error",
1050 "adapter_profile",
1051];
1052
1053const ROUTE_MATCH: &[&str] = &["platform", "guild_id", "chat_id", "thread_id"];
1054
1055pub fn decode_route(file: &str, index: usize, raw: &Value) -> Result<Route> {
1057 let map = raw.as_object().ok_or_else(|| {
1058 load_error(
1059 file,
1060 &format!("profile_routes[{index}]"),
1061 "expected an object",
1062 )
1063 })?;
1064 let text = |k: &str| map.get(k).filter(|v| !v.is_null()).and_then(text_of);
1065 Ok(Route {
1066 name: opt_str(
1067 file,
1068 &format!("profile_routes[{index}].name"),
1069 map.get("name"),
1070 )?,
1071 matches: RouteMatch {
1072 platform: req_str(
1073 file,
1074 &format!("profile_routes[{index}].platform"),
1075 map.get("platform"),
1076 )?,
1077 guild_id: text("guild_id"),
1078 chat_id: text("chat_id"),
1079 thread_id: text("thread_id"),
1080 },
1081 profile: req_str(
1082 file,
1083 &format!("profile_routes[{index}].profile"),
1084 map.get("profile"),
1085 )?,
1086 residue: residue_of(
1087 map,
1088 &[
1089 "name",
1090 "profile",
1091 "platform",
1092 "guild_id",
1093 "chat_id",
1094 "thread_id",
1095 ],
1096 ),
1097 })
1098}
1099
1100pub fn encode_route(r: &Route) -> Value {
1102 let mut out = Map::new();
1103 if let Some(n) = &r.name {
1104 out.insert("name".into(), Value::String(n.clone()));
1105 }
1106 let m = &r.matches;
1107 for (k, v) in [
1108 ("platform", Some(&m.platform)),
1109 ("guild_id", m.guild_id.as_ref()),
1110 ("chat_id", m.chat_id.as_ref()),
1111 ("thread_id", m.thread_id.as_ref()),
1112 ] {
1113 if let Some(v) = v {
1114 out.insert(k.into(), Value::String(v.clone()));
1115 }
1116 }
1117 let _ = ROUTE_MATCH;
1118 out.insert("profile".into(), Value::String(r.profile.clone()));
1119 for (k, v) in &r.residue.0 {
1120 out.insert(k.clone(), v.clone());
1121 }
1122 Value::Object(out)
1123}
1124
1125pub fn is_credential_key(key: &str) -> bool {
1127 let k = key.to_ascii_lowercase();
1128 [
1129 "token",
1130 "secret",
1131 "key",
1132 "password",
1133 "passwd",
1134 "api_key",
1135 "app_secret",
1136 "signing",
1137 "webhook_url",
1138 "private",
1139 ]
1140 .iter()
1141 .any(|w| k.contains(w))
1142}
1143
1144pub fn placeholder_ref(text: &str) -> Option<&str> {
1151 let name = text.strip_prefix("${")?.strip_suffix('}')?;
1152 let mut chars = name.chars();
1153 let first = chars.next()?;
1154 ((first.is_ascii_alphabetic() || first == '_')
1155 && chars.all(|c| c.is_ascii_alphanumeric() || c == '_'))
1156 .then_some(name)
1157}
1158
1159pub fn render_ref(r: &SecretRef) -> Value {
1162 Value::String(format!("${{{}}}", r.name()))
1163}
1164
1165pub fn credential_ref_name(platform: &str, key: &str) -> String {
1167 format!("{platform}_{key}")
1168 .chars()
1169 .map(|c| {
1170 if c.is_ascii_alphanumeric() {
1171 c.to_ascii_uppercase()
1172 } else {
1173 '_'
1174 }
1175 })
1176 .collect()
1177}
1178
1179pub fn decode_channel(
1181 file: &str,
1182 platform: &str,
1183 raw: &Value,
1184 vault: &mut BTreeMap<String, String>,
1185) -> Result<ChannelConfig> {
1186 let map = raw
1187 .as_object()
1188 .ok_or_else(|| load_error(file, &format!("platforms.{platform}"), "expected an object"))?;
1189 let mut ch = ChannelConfig {
1190 platform: platform.to_string(),
1191 enabled: opt_bool(
1192 file,
1193 &format!("platforms.{platform}.enabled"),
1194 map.get("enabled"),
1195 Some(true),
1196 )?
1197 .unwrap_or(true),
1198 credentials: BTreeMap::new(),
1199 extra: BTreeMap::new(),
1200 };
1201 fn walk(
1202 map: &Map<String, Value>,
1203 prefix: &str,
1204 platform: &str,
1205 ch: &mut ChannelConfig,
1206 vault: &mut BTreeMap<String, String>,
1207 ) {
1208 for (k, v) in map {
1209 if k == "enabled" && prefix.is_empty() {
1210 continue;
1211 }
1212 let name = format!("{prefix}{k}");
1213 if let Value::String(s) = v {
1214 if let Some(r) = placeholder_ref(s) {
1215 ch.credentials
1216 .insert(name, SecretRef::Dotenv(r.to_string()));
1217 continue;
1218 }
1219 }
1220 if let Some(obj) = v.as_object() {
1221 if k == "extra" && prefix.is_empty() {
1222 walk(obj, "extra.", platform, ch, vault);
1223 continue;
1224 }
1225 }
1226 if let Value::String(s) = v {
1227 if is_credential_key(k) {
1228 let r = credential_ref_name(platform, &name);
1229 vault.insert(r.clone(), s.clone());
1230 ch.credentials.insert(name, SecretRef::Dotenv(r));
1231 continue;
1232 }
1233 }
1234 ch.extra.insert(name, v.clone());
1235 }
1236 }
1237 walk(map, "", platform, &mut ch, vault);
1238 Ok(ch)
1239}
1240
1241pub fn encode_channel(ch: &ChannelConfig, vault: Option<&BTreeMap<String, String>>) -> Value {
1243 let mut out = Map::new();
1244 out.insert("enabled".into(), Value::Bool(ch.enabled));
1245 let mut extra = Map::new();
1246 for (k, v) in &ch.extra {
1247 match k.strip_prefix("extra.") {
1248 Some(inner) => {
1249 extra.insert(inner.to_string(), v.clone());
1250 }
1251 None => {
1252 out.insert(k.clone(), v.clone());
1253 }
1254 }
1255 }
1256 for (k, r) in &ch.credentials {
1257 let rendered = match vault.and_then(|v| v.get(r.name())) {
1258 Some(value) => Value::String(value.clone()),
1259 None => render_ref(r),
1260 };
1261 match k.strip_prefix("extra.") {
1262 Some(inner) => {
1263 extra.insert(inner.to_string(), rendered);
1264 }
1265 None => {
1266 out.insert(k.clone(), rendered);
1267 }
1268 }
1269 }
1270 if !extra.is_empty() {
1271 out.insert("extra".into(), Value::Object(extra));
1272 }
1273 Value::Object(out)
1274}
1275
1276const SUB_MAPPED: &[&str] = &[
1277 "events",
1278 "prompt",
1279 "skills",
1280 "deliver",
1281 "deliver_extra",
1282 "secret",
1283 "description",
1284 "created_at",
1285];
1286
1287pub fn decode_subscription(
1289 file: &str,
1290 name: &str,
1291 raw: &Value,
1292 vault: &mut BTreeMap<String, String>,
1293) -> Result<WebhookSubscription> {
1294 let map = raw
1295 .as_object()
1296 .ok_or_else(|| load_error(file, name, "expected an object"))?;
1297 let mut sub = WebhookSubscription {
1298 name: name.to_string(),
1299 secret: None,
1300 events: map.get("events").and_then(Value::as_array).map(|a| {
1301 a.iter()
1302 .map(|v| {
1303 v.as_str()
1304 .map(str::to_string)
1305 .unwrap_or_else(|| v.to_string())
1306 })
1307 .collect()
1308 }),
1309 prompt_template: opt_str(file, &format!("{name}.prompt"), map.get("prompt"))?
1310 .unwrap_or_default(),
1311 deliver: parse_target(
1312 file,
1313 &format!("{name}.deliver"),
1314 map.get("deliver"),
1315 map.get("deliver_extra"),
1316 )?,
1317 skills: map
1318 .get("skills")
1319 .and_then(Value::as_array)
1320 .map(|a| {
1321 a.iter()
1322 .map(|v| {
1323 v.as_str()
1324 .map(str::to_string)
1325 .unwrap_or_else(|| v.to_string())
1326 })
1327 .collect()
1328 })
1329 .unwrap_or_default(),
1330 description: opt_str(file, &format!("{name}.description"), map.get("description"))?,
1331 created_at: opt_str(file, &format!("{name}.created_at"), map.get("created_at"))?,
1332 residue: residue_of(map, SUB_MAPPED),
1333 };
1334 match map.get("secret") {
1335 Some(Value::String(secret)) => {
1338 let r = match placeholder_ref(secret) {
1339 Some(name) => name.to_string(),
1340 None => {
1341 let r = credential_ref_name("webhook", &format!("{name}_secret"));
1342 vault.insert(r.clone(), secret.clone());
1343 r
1344 }
1345 };
1346 sub.secret = Some(SecretRef::Dotenv(r));
1347 }
1348 Some(Value::Null) | None => {}
1349 Some(_) => {
1350 return Err(load_error(
1351 file,
1352 &format!("{name}.secret"),
1353 "expected a string (a secret is written `${NAME}`)",
1354 ))
1355 }
1356 }
1357 Ok(sub)
1358}
1359
1360pub fn encode_subscription(
1362 sub: &WebhookSubscription,
1363 vault: Option<&BTreeMap<String, String>>,
1364) -> Value {
1365 let mut out = Map::new();
1366 if let Some(d) = &sub.description {
1367 out.insert("description".into(), Value::String(d.clone()));
1368 }
1369 if let Some(e) = &sub.events {
1370 out.insert(
1371 "events".into(),
1372 Value::Array(e.iter().map(|s| Value::String(s.clone())).collect()),
1373 );
1374 }
1375 out.insert("prompt".into(), Value::String(sub.prompt_template.clone()));
1376 out.insert(
1377 "skills".into(),
1378 Value::Array(
1379 sub.skills
1380 .iter()
1381 .map(|s| Value::String(s.clone()))
1382 .collect(),
1383 ),
1384 );
1385 if let Some(t) = &sub.deliver {
1386 match t {
1387 Target::Explicit {
1388 platform,
1389 chat_id,
1390 thread_id,
1391 } => {
1392 out.insert("deliver".into(), Value::String(platform.clone()));
1393 let mut extra = Map::new();
1394 if let Some(c) = chat_id {
1395 extra.insert("chat_id".into(), Value::String(c.clone()));
1396 }
1397 if let Some(th) = thread_id {
1398 extra.insert("thread_id".into(), Value::String(th.clone()));
1399 }
1400 if !extra.is_empty() {
1401 out.insert("deliver_extra".into(), Value::Object(extra));
1402 }
1403 }
1404 other => {
1405 out.insert("deliver".into(), Value::String(other.render()));
1406 }
1407 }
1408 }
1409 if let Some(s) = &sub.secret {
1410 let value = vault
1411 .and_then(|v| v.get(s.name()).cloned())
1412 .unwrap_or_else(|| format!("${{{}}}", s.name()));
1413 out.insert("secret".into(), Value::String(value));
1414 }
1415 if let Some(c) = &sub.created_at {
1416 out.insert("created_at".into(), Value::String(c.clone()));
1417 }
1418 for (k, v) in &sub.residue.0 {
1419 out.insert(k.clone(), v.clone());
1420 }
1421 Value::Object(out)
1422}