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