1pub mod ulid;
14
15use crate::obs::log::Logger;
16use crate::store::{Envelope, KeySeq, PutOutcome, SharedStore, StoreError};
17use serde::{Deserialize, Serialize};
18use serde_json::{Value, json};
19use std::collections::BTreeSet;
20use std::collections::{BTreeMap, HashMap};
21use std::sync::atomic::{AtomicBool, Ordering};
22use std::sync::{Mutex, OnceLock};
23use std::time::{Duration, Instant};
24
25pub(crate) use crate::store::now_ms;
26
27#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord)]
30pub enum Kind {
31 Manifest,
32 Inbox,
33 Context,
34 Run,
35 Subagent,
36 Task,
37 Memory,
38 Artifact,
39 Timer,
40 Event,
44 Audit,
45 Cred,
49}
50
51impl Kind {
52 pub fn as_str(self) -> &'static str {
53 match self {
54 Kind::Manifest => "manifest",
55 Kind::Inbox => "inbox",
56 Kind::Context => "context",
57 Kind::Run => "run",
58 Kind::Subagent => "subagent",
59 Kind::Task => "task",
60 Kind::Memory => "memory",
61 Kind::Artifact => "artifact",
62 Kind::Timer => "timer",
63 Kind::Event => "event",
64 Kind::Audit => "audit",
65 Kind::Cred => "cred",
66 }
67 }
68 pub fn parse(s: &str) -> Option<Kind> {
69 Some(match s {
70 "manifest" => Kind::Manifest,
71 "inbox" => Kind::Inbox,
72 "event" => Kind::Event,
73 "context" => Kind::Context,
74 "run" => Kind::Run,
75 "subagent" => Kind::Subagent,
76 "task" => Kind::Task,
77 "memory" => Kind::Memory,
78 "artifact" => Kind::Artifact,
79 "timer" => Kind::Timer,
80 "audit" => Kind::Audit,
81 "cred" => Kind::Cred,
82 _ => return None,
83 })
84 }
85 pub fn indexed(self) -> bool {
89 !matches!(
90 self,
91 Kind::Manifest | Kind::Memory | Kind::Audit | Kind::Cred | Kind::Event
92 )
93 }
94}
95
96#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
98pub struct EntityRef {
99 pub kind: String,
100 pub id: String,
101 pub seq: u64,
102}
103
104#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, Default)]
106pub struct StreamMeta {
107 pub seq: u64,
109 pub first: u64,
111}
112
113#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, Default)]
116pub struct Manifest {
117 #[serde(default)]
118 pub generation: u64,
119 #[serde(default)]
120 pub created: u64,
121 #[serde(default)]
122 pub updated: u64,
123 #[serde(default)]
124 pub entities: Vec<EntityRef>,
125 #[serde(default)]
127 pub starts: BTreeMap<String, Value>,
128 #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
132 pub streams: BTreeMap<String, StreamMeta>,
133 #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
138 pub breakers: BTreeMap<String, Value>,
139 #[serde(default)]
142 pub budget: Value,
143 #[serde(default)]
144 pub lifecycle: Value,
145 #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
151 pub config_digest: BTreeMap<String, String>,
152 #[serde(default, skip_serializing_if = "Vec::is_empty")]
157 pub retired: Vec<EntityRef>,
158}
159
160impl Manifest {
161 fn upsert(&mut self, kind: &str, id: &str, seq: u64) {
162 match self
163 .entities
164 .iter_mut()
165 .find(|e| e.kind == kind && e.id == id)
166 {
167 Some(e) => e.seq = seq,
168 None => self.entities.push(EntityRef {
169 kind: kind.to_string(),
170 id: id.to_string(),
171 seq,
172 }),
173 }
174 }
175 fn remove(&mut self, kind: &str, id: &str) {
176 self.entities.retain(|e| !(e.kind == kind && e.id == id));
177 }
178}
179
180#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
183pub struct InboxEvent {
184 pub id: String,
185 pub kind: String,
186 pub ts: u64,
187 #[serde(default, skip_serializing_if = "Option::is_none")]
188 pub principal: Option<String>,
189 pub payload: Value,
190 #[serde(default)]
191 pub status: InboxStatus,
192}
193
194#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, Default)]
195#[serde(rename_all = "lowercase")]
196pub enum InboxStatus {
197 #[default]
198 Pending,
199 Done,
200}
201
202impl InboxEvent {
203 pub fn new(kind: &str, principal: Option<String>, payload: Value) -> InboxEvent {
204 InboxEvent {
205 id: ulid::new(),
206 kind: kind.to_string(),
207 ts: now_ms(),
208 principal,
209 payload,
210 status: InboxStatus::Pending,
211 }
212 }
213}
214
215#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
219pub struct TimerRecord {
220 pub id: String,
221 pub deadline_ms: u64,
222 pub owner: Value,
223 #[serde(default)]
224 pub payload: Value,
225}
226
227#[derive(Debug, Clone)]
230pub struct Policy {
231 pub debounce: Duration,
232 pub on_error: crate::config::v2::StoreOnError,
233 pub retries: u32,
234}
235
236impl Default for Policy {
237 fn default() -> Self {
238 Policy {
239 debounce: Duration::from_millis(250),
240 on_error: crate::config::v2::StoreOnError::Halt,
241 retries: 3,
242 }
243 }
244}
245
246impl Policy {
247 pub fn from_settings(s: &crate::config::v2::Store) -> Policy {
248 Policy {
249 debounce: Duration::from_millis(s.checkpoint.debounce_ms.unwrap_or(250)),
250 on_error: s.on_error,
251 retries: 3,
252 }
253 }
254}
255
256static FRESH: AtomicBool = AtomicBool::new(false);
271static CONFIG_DIGEST: OnceLock<BTreeMap<String, String>> = OnceLock::new();
272
273pub fn request_fresh() {
276 FRESH.store(true, Ordering::Relaxed);
277}
278
279pub fn fresh_requested() -> bool {
281 FRESH.load(Ordering::Relaxed)
282}
283
284pub fn record_config_digest(settings: &crate::config::v2::Settings) {
288 let _ = CONFIG_DIGEST.set(config_digest(settings));
289}
290
291fn recorded_config_digest() -> BTreeMap<String, String> {
292 CONFIG_DIGEST.get().cloned().unwrap_or_default()
293}
294
295pub fn config_digest(settings: &crate::config::v2::Settings) -> BTreeMap<String, String> {
311 let mut out = BTreeMap::new();
312 let digest = |v: &Value| crate::sha::sha256_hex(v.to_string().as_bytes());
315 out.insert(
316 "workflows".to_string(),
317 digest(&Value::Array(settings.workflows.clone())),
318 );
319 out.insert("store".to_string(), digest(&store_shape(&settings.store)));
320 out.insert(
321 "limits".to_string(),
322 digest(&limits_shape(&settings.limits)),
323 );
324 out
325}
326
327fn store_shape(s: &crate::config::v2::Store) -> Value {
332 json!({
333 "kind": format!("{:?}", s.kind),
334 "prefix": s.prefix(),
335 "on_error": format!("{:?}", s.on_error),
336 "audit": s.audit,
337 "checkpoint_debounce_ms": s.checkpoint.debounce_ms,
338 "durability": format!("{:?}", s.durability),
339 "timeout_ms": s.timeout.map(|d| d.0.as_millis() as u64),
340 "mcp_server": s.mcp.as_ref().map(|m| m.server.clone()),
344 "http": s.http.is_some(),
345 })
346}
347
348fn limits_shape(s: &crate::config::v2::Limits) -> Value {
352 json!({
353 "max_runs": s.max_runs,
354 "run_steps": s.run.steps(),
355 "run_tokens": s.run.tokens(),
356 "run_deadline_ms": s.run.deadline().as_millis() as u64,
357 "subagents": format!("{:?}", s.subagents),
358 "inline_max_bytes": s.inline_max_bytes,
359 "step_timeout_ms": s.step_timeout.map(|d| d.0.as_millis() as u64),
360 })
361}
362
363fn changed_sections(
370 recorded: &BTreeMap<String, String>,
371 current: &BTreeMap<String, String>,
372) -> Vec<String> {
373 if recorded.is_empty() || current.is_empty() {
374 return Vec::new();
375 }
376 let mut out: Vec<String> = current
377 .iter()
378 .filter(|(k, v)| recorded.get(*k) != Some(*v))
379 .map(|(k, _)| k.clone())
380 .collect();
381 out.extend(
382 recorded
383 .keys()
384 .filter(|k| !current.contains_key(*k))
385 .cloned(),
386 );
387 out.sort();
388 out.dedup();
389 out
390}
391
392const RECONCILED: [Kind; 7] = [
395 Kind::Inbox,
396 Kind::Context,
397 Kind::Run,
398 Kind::Subagent,
399 Kind::Task,
400 Kind::Timer,
401 Kind::Artifact,
402];
403
404#[derive(Debug, Default)]
406pub struct Restored {
407 pub manifest: Option<Manifest>,
409 pub entities: BTreeMap<String, Vec<Envelope>>,
411 pub lost: Vec<EntityRef>,
413 pub unindexed: Vec<EntityRef>,
416}
417
418impl Restored {
419 pub fn inbox_pending(&self) -> Vec<InboxEvent> {
420 let mut out: Vec<InboxEvent> = self
421 .entities
422 .get("inbox")
423 .map(|v| {
424 v.iter()
425 .filter_map(|e| serde_json::from_value::<InboxEvent>(e.state.clone()).ok())
426 .filter(|e| e.status == InboxStatus::Pending)
427 .collect()
428 })
429 .unwrap_or_default();
430 out.sort_by(|a, b| a.ts.cmp(&b.ts).then(a.id.cmp(&b.id)));
431 out
432 }
433 pub fn timers(&self) -> Vec<TimerRecord> {
434 self.entities
435 .get("timer")
436 .map(|v| {
437 v.iter()
438 .filter_map(|e| serde_json::from_value(e.state.clone()).ok())
439 .collect()
440 })
441 .unwrap_or_default()
442 }
443 pub fn of(&self, kind: Kind) -> &[Envelope] {
444 self.entities
445 .get(kind.as_str())
446 .map(Vec::as_slice)
447 .unwrap_or(&[])
448 }
449 pub fn count(&self) -> usize {
450 self.entities.values().map(Vec::len).sum()
451 }
452}
453
454pub struct Durable {
456 store: SharedStore,
457 prefix: String,
458 instance: String,
459 policy: Policy,
460 seqs: Mutex<HashMap<String, u64>>,
463 manifest: Mutex<Manifest>,
464 manifest_dirty: AtomicBool,
465 last_flush: Mutex<Instant>,
466 degraded: AtomicBool,
467 log: Option<Logger>,
468 fresh: bool,
470 config_digest: BTreeMap<String, String>,
473}
474
475impl Durable {
476 pub fn new(
477 store: SharedStore,
478 prefix: &str,
479 instance: &str,
480 policy: Policy,
481 log: Option<Logger>,
482 ) -> Durable {
483 Durable {
484 store,
485 prefix: prefix.to_string(),
486 instance: instance.to_string(),
487 policy,
488 seqs: Mutex::new(HashMap::new()),
489 manifest: Mutex::new(Manifest::default()),
490 manifest_dirty: AtomicBool::new(false),
491 last_flush: Mutex::new(Instant::now()),
492 degraded: AtomicBool::new(false),
493 log,
494 fresh: fresh_requested(),
495 config_digest: recorded_config_digest(),
496 }
497 }
498
499 pub fn with_fresh(mut self, fresh: bool) -> Durable {
503 self.fresh = fresh;
504 self
505 }
506
507 pub fn with_config_digest(mut self, digest: BTreeMap<String, String>) -> Durable {
509 self.config_digest = digest;
510 self
511 }
512
513 pub fn store_kind(&self) -> &'static str {
514 self.store.kind()
515 }
516 pub fn instance(&self) -> &str {
517 &self.instance
518 }
519 pub fn prefix(&self) -> &str {
520 &self.prefix
521 }
522 pub fn key(&self, kind: Kind, id: &str) -> String {
523 crate::store::key(&self.prefix, &self.instance, kind.as_str(), id)
524 }
525 pub fn is_degraded(&self) -> bool {
527 self.degraded.load(Ordering::Relaxed)
528 }
529 pub fn policy(&self) -> &Policy {
530 &self.policy
531 }
532
533 pub fn put(
540 &self,
541 kind: Kind,
542 id: &str,
543 state: Value,
544 hash: Option<String>,
545 ) -> Result<u64, StoreError> {
546 let key = self.key(kind, id);
547 let mut adopted = false;
548 let started = std::time::Instant::now();
549 loop {
550 let (seq, warmed) = {
551 let seqs = self.seqs.lock().unwrap_or_else(|e| e.into_inner());
552 match seqs.get(&key).copied() {
553 Some(s) => (s + 1, true),
554 None => (1, false),
555 }
556 };
557 let env = Envelope::new(
558 kind.as_str(),
559 id,
560 seq,
561 &self.instance,
562 hash.clone(),
563 state.clone(),
564 );
565 kill_point("state.before_put");
566 let outcome = crate::store::with_retry(
567 || self.store.put(&key, seq, &env.to_value()),
568 self.policy.retries,
569 );
570 match outcome {
571 Ok(PutOutcome::Ok) => {
572 self.seqs
573 .lock()
574 .unwrap_or_else(|e| e.into_inner())
575 .insert(key.clone(), seq);
576 if kind.indexed() {
577 let mut m = self.manifest.lock().unwrap_or_else(|e| e.into_inner());
578 m.upsert(kind.as_str(), id, seq);
579 m.updated = now_ms();
580 self.manifest_dirty.store(true, Ordering::Relaxed);
581 }
582 self.degraded.store(false, Ordering::Relaxed);
583 kill_point("state.after_put");
584 crate::obs::metrics::record_store_op(
585 "ok",
586 started.elapsed().as_millis() as u64,
587 );
588 return Ok(seq);
589 }
590 Ok(PutOutcome::Conflict { latest_seq }) => {
591 if !warmed && !adopted {
592 if let Some(l) = latest_seq {
596 self.seqs
597 .lock()
598 .unwrap_or_else(|e| e.into_inner())
599 .insert(key.clone(), l);
600 adopted = true;
601 self.log_event("store.seq_adopted", json!({"key": key, "latest": l}));
602 continue;
603 }
604 }
605 self.log_event(
606 "store.conflict",
607 json!({"key": key, "seq": seq, "latest": latest_seq}),
608 );
609 crate::obs::metrics::record_store_op(
610 "conflict",
611 started.elapsed().as_millis() as u64,
612 );
613 return Err(StoreError::Conflict(format!(
614 "key {key}: another writer owns it (our seq {seq}, latest {latest_seq:?})"
615 )));
616 }
617 Err(e) => {
618 self.log_event("store.put.fail", json!({"key": key, "err": e.to_string()}));
619 crate::obs::metrics::record_store_op(
620 "error",
621 started.elapsed().as_millis() as u64,
622 );
623 if self.policy.on_error == crate::config::v2::StoreOnError::Degrade {
624 self.degraded.store(true, Ordering::Relaxed);
625 self.seqs
628 .lock()
629 .unwrap_or_else(|e| e.into_inner())
630 .insert(key.clone(), seq);
631 return Ok(seq);
632 }
633 return Err(e);
634 }
635 }
636 }
637 }
638
639 pub fn get(&self, kind: Kind, id: &str) -> Result<Option<Envelope>, StoreError> {
641 let key = self.key(kind, id);
642 let v = crate::store::with_retry(|| self.store.get(&key, None), self.policy.retries)?;
643 match v {
644 None => Ok(None),
645 Some(v) => {
646 let env = Envelope::from_value(v)?;
647 self.seqs
648 .lock()
649 .unwrap_or_else(|e| e.into_inner())
650 .insert(key, env.seq);
651 Ok(if env.is_tombstone() { None } else { Some(env) })
652 }
653 }
654 }
655
656 pub fn delete(&self, kind: Kind, id: &str) -> Result<(), StoreError> {
659 let key = self.key(kind, id);
660 match crate::store::with_retry(|| self.store.delete(&key), self.policy.retries) {
661 Ok(()) => {
662 self.seqs
663 .lock()
664 .unwrap_or_else(|e| e.into_inner())
665 .remove(&key);
666 }
667 Err(StoreError::Unsupported(_)) => {
668 self.put(kind, id, Value::Null, None)?;
669 }
670 Err(e) => return Err(e),
671 }
672 if kind.indexed() {
673 let mut m = self.manifest.lock().unwrap_or_else(|e| e.into_inner());
674 m.remove(kind.as_str(), id);
675 m.updated = now_ms();
676 self.manifest_dirty.store(true, Ordering::Relaxed);
677 }
678 Ok(())
679 }
680
681 pub fn list(&self, kind: Kind) -> Result<Vec<KeySeq>, StoreError> {
683 let prefix = format!("{}/{}/{}/", self.prefix, self.instance, kind.as_str());
684 self.store.list(&prefix)
685 }
686
687 pub fn inbox_put(&self, ev: &InboxEvent) -> Result<u64, StoreError> {
691 let seq = self.put(
692 Kind::Inbox,
693 &ev.id,
694 serde_json::to_value(ev).unwrap_or(Value::Null),
695 None,
696 )?;
697 kill_point("inbox.after_put");
698 Ok(seq)
699 }
700
701 pub fn inbox_done(&self, id: &str) -> Result<(), StoreError> {
703 self.delete(Kind::Inbox, id)
704 }
705
706 pub fn timer_arm(&self, t: &TimerRecord) -> Result<u64, StoreError> {
707 self.put(
708 Kind::Timer,
709 &t.id,
710 serde_json::to_value(t).unwrap_or(Value::Null),
711 None,
712 )
713 }
714
715 pub fn timer_disarm(&self, id: &str) -> Result<(), StoreError> {
716 self.delete(Kind::Timer, id)
717 }
718
719 pub fn manifest(&self) -> Manifest {
722 self.manifest
723 .lock()
724 .unwrap_or_else(|e| e.into_inner())
725 .clone()
726 }
727
728 pub fn manifest_update(&self, f: impl FnOnce(&mut Manifest)) {
731 let mut m = self.manifest.lock().unwrap_or_else(|e| e.into_inner());
732 f(&mut m);
733 m.updated = now_ms();
734 self.manifest_dirty.store(true, Ordering::Relaxed);
735 }
736
737 pub fn flush(&self, force: bool) -> Result<bool, StoreError> {
739 if !self.manifest_dirty.load(Ordering::Relaxed) {
740 return Ok(false);
741 }
742 {
743 let last = self.last_flush.lock().unwrap_or_else(|e| e.into_inner());
744 if !force && last.elapsed() < self.policy.debounce {
745 return Ok(false);
746 }
747 }
748 let snapshot = self.manifest();
749 self.put(
750 Kind::Manifest,
751 "agent",
752 serde_json::to_value(&snapshot).unwrap_or(Value::Null),
753 None,
754 )?;
755 self.manifest_dirty.store(false, Ordering::Relaxed);
756 *self.last_flush.lock().unwrap_or_else(|e| e.into_inner()) = Instant::now();
757 Ok(true)
758 }
759
760 pub fn restore(&self) -> Result<Restored, StoreError> {
770 let mut out = Restored::default();
771 let (manifest, fresh) = match self.get(Kind::Manifest, "agent")? {
772 None => (
773 Manifest {
774 generation: 0,
775 created: now_ms(),
776 updated: now_ms(),
777 ..Manifest::default()
778 },
779 true,
780 ),
781 Some(env) => (
782 serde_json::from_value::<Manifest>(env.state.clone())
783 .map_err(|e| StoreError::Corrupt(format!("manifest does not parse: {e}")))?,
784 false,
785 ),
786 };
787 if self.fresh {
792 return self.restore_fresh(manifest, fresh);
793 }
794 let moved = changed_sections(&manifest.config_digest, &self.config_digest);
801 if !moved.is_empty() {
802 self.log_event(
803 "store.config_changed",
804 json!({
805 "sections": moved,
806 "msg": "state was written under a different configuration — resuming anyway; --fresh to start a new generation",
807 }),
808 );
809 }
810 for r in &manifest.entities {
812 let Some(kind) = Kind::parse(&r.kind) else {
813 out.lost.push(r.clone());
814 continue;
815 };
816 match self.get(kind, &r.id)? {
817 Some(env) => out.entities.entry(r.kind.clone()).or_default().push(env),
818 None => out.lost.push(r.clone()),
819 }
820 }
821 let mut seen: Option<BTreeSet<(String, String)>> = None;
826 for kind in RECONCILED {
827 match self.list(kind) {
828 Ok(keys) => {
829 let seen = seen.get_or_insert_with(BTreeSet::new);
830 for ks in keys {
831 let Some((_, id)) =
832 crate::store::parse_key(&self.prefix, &self.instance, &ks.key)
833 else {
834 continue;
835 };
836 seen.insert((kind.as_str().to_string(), id.to_string()));
837 let indexed = manifest
838 .entities
839 .iter()
840 .any(|e| e.kind == kind.as_str() && e.id == id);
841 if indexed {
842 continue;
843 }
844 if manifest
849 .retired
850 .iter()
851 .any(|r| r.kind == kind.as_str() && r.id == id)
852 {
853 continue;
854 }
855 if let Some(env) = self.get(kind, id)? {
856 out.unindexed.push(EntityRef {
857 kind: kind.as_str().to_string(),
858 id: id.to_string(),
859 seq: env.seq,
860 });
861 out.entities
862 .entry(kind.as_str().to_string())
863 .or_default()
864 .push(env);
865 }
866 }
867 }
868 Err(StoreError::Unsupported(_)) => {}
869 Err(e) => return Err(e),
870 }
871 }
872 let mut m = manifest.clone();
876 m.entities.retain(|e| !out.lost.iter().any(|l| l == e));
877 for u in &out.unindexed {
878 m.upsert(&u.kind, &u.id, u.seq);
879 }
880 m.generation += 1;
881 m.updated = now_ms();
882 if !self.config_digest.is_empty() {
887 m.config_digest = self.config_digest.clone();
888 }
889 if let Some(seen) = &seen {
893 m.retired
894 .retain(|r| seen.contains(&(r.kind.clone(), r.id.clone())));
895 }
896 *self.manifest.lock().unwrap_or_else(|e| e.into_inner()) = m.clone();
897 self.manifest_dirty.store(true, Ordering::Relaxed);
898 self.flush(true)?;
899 if fresh && out.count() == 0 {
900 self.log_event("restore.fresh", json!({"generation": m.generation}));
901 return Ok(out);
902 }
903 self.log_event(
904 "restore.done",
905 json!({
906 "generation": m.generation,
907 "fresh_manifest": fresh,
908 "entities": out.count(),
909 "lost": out.lost.len(),
910 "unindexed": out.unindexed.len(),
911 "inbox_pending": out.inbox_pending().len(),
912 }),
913 );
914 out.manifest = Some(m);
915 Ok(out)
916 }
917
918 fn restore_fresh(
934 &self,
935 prior: Manifest,
936 no_prior_manifest: bool,
937 ) -> Result<Restored, StoreError> {
938 let mut retired: Vec<EntityRef> = Vec::new();
941 let mut push = |kind: &str, id: &str, seq: u64| {
942 if !retired.iter().any(|r| r.kind == kind && r.id == id) {
943 retired.push(EntityRef {
944 kind: kind.to_string(),
945 id: id.to_string(),
946 seq,
947 });
948 }
949 };
950 for kind in RECONCILED {
951 match self.list(kind) {
952 Ok(keys) => {
953 for ks in keys {
954 if let Some((_, id)) =
955 crate::store::parse_key(&self.prefix, &self.instance, &ks.key)
956 {
957 push(kind.as_str(), id, ks.seq.unwrap_or(0));
958 }
959 }
960 }
961 Err(StoreError::Unsupported(_)) => {}
962 Err(e) => return Err(e),
963 }
964 }
965 for e in &prior.entities {
966 push(&e.kind, &e.id, e.seq);
967 }
968 for e in &prior.retired {
969 push(&e.kind, &e.id, e.seq);
970 }
971 if !no_prior_manifest {
974 self.put(
975 Kind::Manifest,
976 &format!("agent.gen{}", prior.generation),
977 serde_json::to_value(&prior).unwrap_or(Value::Null),
978 None,
979 )?;
980 }
981 let m = Manifest {
984 generation: prior.generation + 1,
985 created: if prior.created == 0 {
986 now_ms()
987 } else {
988 prior.created
989 },
990 updated: now_ms(),
991 entities: Vec::new(),
992 starts: BTreeMap::new(),
993 streams: BTreeMap::new(),
994 breakers: BTreeMap::new(),
995 budget: Value::Null,
996 lifecycle: Value::Null,
997 config_digest: self.config_digest.clone(),
998 retired,
999 };
1000 *self.manifest.lock().unwrap_or_else(|e| e.into_inner()) = m.clone();
1001 self.manifest_dirty.store(true, Ordering::Relaxed);
1002 self.flush(true)?;
1003 self.log_event(
1004 "restore.fresh",
1005 json!({
1006 "generation": m.generation,
1007 "superseded": prior.generation,
1008 "retired": m.retired.len(),
1009 "msg": "--fresh: this generation starts empty; the previous one's records were kept, not deleted",
1010 }),
1011 );
1012 Ok(Restored::default())
1013 }
1014
1015 fn log_event(&self, event: &str, fields: Value) {
1016 if let Some(l) = &self.log {
1017 match event {
1018 e if e.ends_with(".fail")
1019 || e == "store.conflict"
1020 || e == "store.config_changed" =>
1023 {
1024 l.warn(event, fields)
1025 }
1026 _ => l.info(event, fields),
1027 }
1028 }
1029 }
1030}
1031
1032pub fn kill_point(name: &str) {
1038 #[cfg(any(feature = "internal-mocks", debug_assertions))]
1039 {
1040 if std::env::var("AGENTD_TEST_KILL_AT").as_deref() == Ok(name) {
1041 #[cfg(unix)]
1042 unsafe {
1043 libc::raise(libc::SIGKILL);
1044 }
1045 std::process::abort();
1046 }
1047 }
1048 #[cfg(not(any(feature = "internal-mocks", debug_assertions)))]
1049 {
1050 let _ = name;
1051 }
1052}
1053
1054#[cfg(test)]
1055mod tests {
1056 use super::*;
1057 use crate::store::Store;
1058 use crate::store::memory::MemoryStore;
1059 use std::sync::Arc;
1060
1061 fn durable(store: Arc<MemoryStore>) -> Durable {
1062 Durable::new(
1063 store,
1064 "agentd",
1065 "inst",
1066 Policy {
1067 debounce: Duration::from_millis(0),
1068 ..Policy::default()
1069 },
1070 None,
1071 )
1072 }
1073
1074 #[test]
1075 fn put_allocates_seqs_indexes_and_flushes_manifest() {
1076 let mem = Arc::new(MemoryStore::new());
1077 let d = durable(mem.clone());
1078 assert!(d.restore().unwrap().manifest.is_none(), "fresh");
1079 assert_eq!(
1080 d.put(
1081 Kind::Run,
1082 "r1",
1083 json!({"status": "running"}),
1084 Some("h".into())
1085 )
1086 .unwrap(),
1087 1
1088 );
1089 assert_eq!(
1090 d.put(Kind::Run, "r1", json!({"status": "done"}), Some("h".into()))
1091 .unwrap(),
1092 2
1093 );
1094 assert_eq!(
1095 d.put(Kind::Context, "root", json!({"v": 1}), None).unwrap(),
1096 1
1097 );
1098 let env = d.get(Kind::Run, "r1").unwrap().unwrap();
1099 assert_eq!(env.seq, 2);
1100 assert_eq!(env.state["status"], json!("done"));
1101 assert_eq!(env.hash.as_deref(), Some("h"));
1102 assert!(d.flush(true).unwrap());
1104 let m = d.manifest();
1105 assert_eq!(m.entities.len(), 2);
1106 assert!(
1107 m.entities
1108 .iter()
1109 .any(|e| e.kind == "run" && e.id == "r1" && e.seq == 2)
1110 );
1111 assert!(!d.flush(true).unwrap(), "clean after a flush");
1112 d.delete(Kind::Context, "root").unwrap();
1114 assert!(d.get(Kind::Context, "root").unwrap().is_none());
1115 assert_eq!(d.manifest().entities.len(), 1);
1116 }
1117
1118 #[test]
1119 fn conflicts_are_fatal_on_owned_keys_but_adopted_on_first_touch() {
1120 let mem = Arc::new(MemoryStore::new());
1121 let stale = Envelope::new("run", "old", 5, "inst", None, json!({"x": 1}));
1123 mem.put("agentd/inst/run/old", 5, &stale.to_value())
1124 .unwrap();
1125 let d = durable(mem.clone());
1126 assert_eq!(d.put(Kind::Run, "old", json!({"x": 2}), None).unwrap(), 6);
1128 let other = Envelope::new("run", "old", 7, "other", None, json!({"x": 3}));
1130 mem.put("agentd/inst/run/old", 7, &other.to_value())
1131 .unwrap();
1132 assert!(matches!(
1133 d.put(Kind::Run, "old", json!({"x": 4}), None),
1134 Err(StoreError::Conflict(_))
1135 ));
1136 }
1137
1138 #[test]
1139 fn inbox_write_ahead_timers_and_restore() {
1140 let mem = Arc::new(MemoryStore::new());
1141 {
1142 let d = durable(mem.clone());
1143 d.restore().unwrap();
1144 let e1 = InboxEvent::new(
1145 "a2a_message",
1146 Some("user:andrii".into()),
1147 json!({"text": "hi"}),
1148 );
1149 let e2 = InboxEvent::new("start_fired", None, json!({"workflow": "w"}));
1150 d.inbox_put(&e1).unwrap();
1151 d.inbox_put(&e2).unwrap();
1152 d.inbox_done(&e1.id).unwrap();
1153 d.timer_arm(&TimerRecord {
1154 id: "t1".into(),
1155 deadline_ms: 42,
1156 owner: json!({"run": "r"}),
1157 payload: Value::Null,
1158 })
1159 .unwrap();
1160 d.put(
1161 Kind::Run,
1162 "r",
1163 json!({"status": "running"}),
1164 Some("hash".into()),
1165 )
1166 .unwrap();
1167 d.put(Kind::Task, "task-1", json!({"state": "working"}), None)
1168 .unwrap();
1169 d.manifest_update(|m| {
1170 m.starts.insert("w.s".into(), json!({"last_fired": 1}));
1171 });
1172 d.put(Kind::Subagent, "gone", json!({}), None).unwrap();
1174 d.flush(true).unwrap();
1175 mem.delete("agentd/inst/subagent/gone").unwrap();
1176 d.put(Kind::Run, "r2", json!({"status": "running"}), None)
1179 .unwrap();
1180 }
1181 let d2 = durable(mem.clone());
1183 let r = d2.restore().unwrap();
1184 let m = r.manifest.as_ref().unwrap();
1185 assert_eq!(m.generation, 2, "generation bumped");
1186 assert_eq!(m.starts["w.s"]["last_fired"], json!(1));
1187 let pending = r.inbox_pending();
1188 assert_eq!(pending.len(), 1, "the done event does not replay");
1189 assert_eq!(pending[0].kind, "start_fired");
1190 assert_eq!(r.timers().len(), 1);
1191 assert_eq!(r.timers()[0].deadline_ms, 42);
1192 assert_eq!(
1193 r.of(Kind::Run).len(),
1194 2,
1195 "indexed + unindexed runs restored"
1196 );
1197 assert!(r.unindexed.iter().any(|u| u.id == "r2"));
1198 assert!(
1199 r.lost
1200 .iter()
1201 .any(|l| l.kind == "subagent" && l.id == "gone")
1202 );
1203 assert_eq!(r.of(Kind::Task).len(), 1);
1204 assert_eq!(
1206 d2.put(
1207 Kind::Run,
1208 "r",
1209 json!({"status": "done"}),
1210 Some("hash".into())
1211 )
1212 .unwrap(),
1213 2
1214 );
1215 assert!(!d2.manifest().entities.iter().any(|e| e.id == "gone"));
1217 }
1218
1219 #[test]
1220 fn degrade_policy_keeps_going_and_flags_it() {
1221 let mem = Arc::new(MemoryStore::new());
1222 let d = Durable::new(
1223 mem.clone(),
1224 "agentd",
1225 "inst",
1226 Policy {
1227 debounce: Duration::from_millis(0),
1228 on_error: crate::config::v2::StoreOnError::Degrade,
1229 retries: 1,
1230 },
1231 None,
1232 );
1233 mem.fail_next(1);
1234 assert_eq!(
1235 d.put(Kind::Run, "r", json!({}), None).unwrap(),
1236 1,
1237 "degraded write reports the intended seq"
1238 );
1239 assert!(d.is_degraded());
1240 assert_eq!(
1241 d.put(Kind::Run, "r", json!({}), None).unwrap(),
1242 2,
1243 "seq not reused"
1244 );
1245 assert!(!d.is_degraded(), "a successful write clears the flag");
1246 let d2 = durable(mem.clone());
1248 mem.fail_next(5);
1249 assert!(matches!(
1250 d2.put(Kind::Run, "x", json!({}), None),
1251 Err(StoreError::Io(_))
1252 ));
1253 }
1254
1255 #[test]
1260 fn fresh_opens_a_new_generation_without_resuming_and_deletes_nothing() {
1261 let mem = Arc::new(MemoryStore::new());
1262
1263 let d = durable(mem.clone());
1265 assert!(d.restore().unwrap().manifest.is_none());
1266 d.put(Kind::Run, "r1", json!({"status": "running"}), None)
1267 .unwrap();
1268 let ev = InboxEvent::new("a2a.message", None, json!({"n": 1}));
1269 d.inbox_put(&ev).unwrap();
1270 assert_eq!(d.manifest().generation, 1);
1271
1272 let f = durable(mem.clone()).with_fresh(true);
1274 let r = f.restore().unwrap();
1275 assert!(
1276 r.manifest.is_none(),
1277 "a new generation reports no prior life"
1278 );
1279 assert_eq!(r.count(), 0);
1280 assert!(r.inbox_pending().is_empty(), "the inbox does not replay");
1281 let m = f.manifest();
1282 assert_eq!(m.generation, 2, "the counter is inherited, not reset");
1283 assert!(m.entities.is_empty());
1284
1285 assert!(
1287 f.get(Kind::Run, "r1").unwrap().is_some(),
1288 "--fresh keeps the abandoned records"
1289 );
1290 let kept = f
1291 .get(Kind::Manifest, "agent.gen1")
1292 .unwrap()
1293 .expect("the outgoing manifest is preserved");
1294 let kept: Manifest = serde_json::from_value(kept.state).unwrap();
1295 assert_eq!(kept.generation, 1);
1296 assert!(m.retired.iter().any(|e| e.kind == "run" && e.id == "r1"));
1297 assert!(m.retired.iter().any(|e| e.kind == "inbox" && e.id == ev.id));
1298
1299 let d3 = durable(mem.clone());
1302 let r3 = d3.restore().unwrap();
1303 assert_eq!(r3.count(), 0, "retired records stay retired");
1304 assert!(r3.unindexed.is_empty());
1305 assert_eq!(d3.manifest().generation, 3);
1306 }
1307
1308 #[test]
1312 fn a_moved_config_digest_reports_but_never_gates_the_resume() {
1313 let before: BTreeMap<String, String> = [
1314 ("workflows".to_string(), "aaa".to_string()),
1315 ("store".to_string(), "sss".to_string()),
1316 ]
1317 .into_iter()
1318 .collect();
1319 let mut after = before.clone();
1320 after.insert("workflows".to_string(), "bbb".to_string());
1321 assert_eq!(changed_sections(&before, &after), vec!["workflows"]);
1322 assert!(changed_sections(&BTreeMap::new(), &after).is_empty());
1325 assert!(changed_sections(&before, &BTreeMap::new()).is_empty());
1326
1327 let mem = Arc::new(MemoryStore::new());
1328 let d = durable(mem.clone()).with_config_digest(before.clone());
1329 d.restore().unwrap();
1330 d.put(Kind::Run, "r1", json!({"status": "running"}), None)
1331 .unwrap();
1332 assert_eq!(d.manifest().config_digest, before);
1333
1334 let d2 = durable(mem.clone()).with_config_digest(after.clone());
1335 let r = d2.restore().unwrap();
1336 assert_eq!(r.of(Kind::Run).len(), 1, "state is still resumed");
1337 assert_eq!(d2.manifest().generation, 2);
1338 assert_eq!(
1339 d2.manifest().config_digest,
1340 after,
1341 "the next life compares against what this one ran under"
1342 );
1343 }
1344
1345 #[test]
1348 fn the_digest_covers_workflows_store_and_limits_only() {
1349 let doc = json!({
1350 "config_version": "1",
1351 "agent": {"name": "a", "instruction": "one"},
1352 "workflows": [{"name": "w", "version": 3, "steps": {"s": {"kind": "once"}}}],
1353 "limits": {"run": {"steps": 10}},
1354 });
1355 let settings = |patch: &dyn Fn(&mut Value)| {
1356 let mut d = doc.clone();
1357 patch(&mut d);
1358 serde_json::from_value::<crate::config::v2::Settings>(d).expect("settings")
1359 };
1360 let base = config_digest(&settings(&|_| {}));
1361 assert_eq!(
1362 base.keys().collect::<Vec<_>>(),
1363 ["limits", "store", "workflows"]
1364 );
1365
1366 let elsewhere = config_digest(&settings(&|d| d["agent"]["instruction"] = json!("two")));
1367 assert_eq!(base, elsewhere, "an instruction edit is not a state change");
1368
1369 let wf = config_digest(&settings(&|d| {
1370 d["workflows"][0]["steps"]["t"] = json!({"kind": "noop"})
1371 }));
1372 assert_eq!(changed_sections(&base, &wf), vec!["workflows"]);
1373
1374 let lim = config_digest(&settings(&|d| d["limits"]["run"]["steps"] = json!(11)));
1375 assert_eq!(changed_sections(&base, &lim), vec!["limits"]);
1376 }
1377}