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