agentplane 0.30.0

Durable, replayable agent runtime — the journal is the plan of record
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
1001
1002
1003
1004
1005
1006
1007
1008
1009
1010
1011
1012
1013
1014
1015
1016
1017
1018
1019
1020
1021
1022
1023
1024
1025
1026
1027
1028
1029
1030
1031
1032
1033
1034
1035
1036
1037
1038
1039
1040
1041
1042
1043
1044
1045
1046
1047
1048
1049
1050
1051
1052
1053
1054
1055
1056
1057
1058
1059
1060
1061
1062
1063
1064
1065
1066
1067
1068
1069
1070
1071
1072
1073
1074
1075
1076
1077
//! Memory an agent keeps between runs, and the rules that keep it survivable.
//!
//! # Writable memory is delayed code
//!
//! A vector store bolted onto an agent looks like a cache and behaves like a
//! program. Whatever is written today is read back tomorrow *into a context
//! window*, where a model treats it as established fact — so a single poisoned
//! write becomes a standing instruction that fires on every later session. The
//! literature calls this the attack that waits, and its distinguishing property
//! is that nothing at read time looks wrong.
//!
//! Three rules follow, and none of them is about the storage engine.
//!
//! # Trust comes from provenance, never from content
//!
//! A recalled item is labelled by **where it came from**, not by what it says.
//! Content-inferred trust is adversarially gameable by construction: text that
//! asserts its own reliability is the cheapest thing an attacker can write.
//!
//! So [`MemoryItem`] carries its sources, and [`StepCtx::recall`] hands back a
//! [`Tainted`] value whose label is the join of them. An item written from a
//! model's output is untrusted forever, however many times it is re-read, and
//! reaching a mutating sink with it takes the same journaled release as any
//! other untrusted value.
//!
//! [`StepCtx::recall`]: crate::runtime::StepCtx::recall
//! [`Tainted`]: crate::core::Tainted
//!
//! # Retrieval is an effect, not a lookup
//!
//! Memory is mutable state outside the journal, so reading it from inside the
//! deterministic zone would make replay depend on what the store happens to hold
//! *now*. A run replayed after a later write would retrieve different items,
//! reach different conclusions, and produce a history that disagrees with itself
//! — the exact failure the effect protocol exists to prevent.
//!
//! So a recall is journaled: the query, the filters, and the **selection** — item
//! ids with their exact versions and content digests. Replay reads that record
//! and re-materialises those versions rather than re-running the search, so the
//! ranking is not re-computed and cannot drift with the corpus.
//!
//! # Content is versioned and supersedable, never edited in place
//!
//! A memory that can be rewritten in place cannot be audited, and cannot be
//! repaired: there is no way to ask what the agent believed last Tuesday, and no
//! way to undo one bad write without guessing what it replaced. Writes append a
//! new version and mark the old version's lifecycle metadata superseded. Its
//! content and security metadata remain unchanged, forgetting is selective, and
//! lineage survives correction so a later erasure can still traverse it.

use std::fmt::Debug;

use async_trait::async_trait;
use serde::{Deserialize, Serialize};
use serde_json::Value;

use crate::core::{Digest, Label, Sensitivity, SourceId, StoreError, Timestamp, Trust};

/// One remembered thing.
///
/// The fields beyond `content` are not bookkeeping. Each answers a question that
/// an agent acting on a memory has to be able to ask, and that a store holding
/// only text cannot answer.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct MemoryItem {
    /// Stable across versions: this is *the same memory*, revised. Subject and
    /// purpose cannot change under it, and it is not reusable after forgetting;
    /// use a new id for another scope or unrelated content.
    pub id: String,
    /// What this is about — an account, a customer, a matter. The primary
    /// retrieval axis, and the unit selective forgetting works on.
    pub subject: String,
    /// Why it was kept. Retrieval filters on it, so a memory written for
    /// support triage is not silently read into a payments decision.
    pub purpose: String,
    /// The remembered content.
    pub content: Value,
    /// Where it came from.
    ///
    /// The label is derived from this and never from `content` — see the module
    /// docs on why content-inferred trust is gameable.
    pub provenance: Vec<SourceId>,
    /// How far this may travel.
    pub sensitivity: Sensitivity,
    /// Whether the content may be believed without a release.
    ///
    /// Set by the writer from the source, not inferred. Anything derived from a
    /// model, a peer, or an inbound message is untrusted.
    pub trust: Trust,
    /// The run that wrote it, so a bad write is attributable to a history.
    pub written_by: String,
    /// Monotonic per `id`. A write appends; nothing is edited in place.
    pub version: u64,
    #[serde(with = "time::serde::rfc3339")]
    pub created_at: Timestamp,
    /// When this version stops being eligible for fresh recall.
    ///
    /// Exact-version reads remain available for replay until an explicit
    /// lifecycle sweep erases the memory. The cutoff is evaluated against the
    /// journaled `Recall::as_of`, never an ambient store clock.
    ///
    /// This is a **hard ceiling** and it is immutable: sliding access
    /// retention may shorten a memory's life below it, and nothing — no touch,
    /// however recent — extends a life past it. The effective expiry is the
    /// *earlier* of this and the access window.
    #[serde(
        default,
        with = "time::serde::rfc3339::option",
        skip_serializing_if = "Option::is_none"
    )]
    pub expires_at: Option<Timestamp>,
    /// Sliding retention window, in seconds of disuse.
    ///
    /// The window opens **at the write** — the write is itself an access — so
    /// an untouched memory expires `created_at + window` rather than living
    /// forever waiting for its first touch, and each explicit journaled touch
    /// slides it forward. It only ever operates *inside* `expires_at`: sliding
    /// retention answers "collect this once nobody reads it", the fixed expiry
    /// answers "this is gone on Thursday whatever happens", and where both are
    /// set the earlier instant wins.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub access_retention_seconds: Option<u64>,
    /// Set when a later version replaced this one.
    #[serde(
        default,
        with = "time::serde::rfc3339::option",
        skip_serializing_if = "Option::is_none"
    )]
    pub superseded_at: Option<Timestamp>,
    /// The memories this one was derived from, at the versions actually read.
    ///
    /// Empty for an ordinary write. Populated by compaction, and it is what
    /// makes a summary repairable: without it, forgetting a poisoned memory
    /// leaves every summary that absorbed it in place, and the attack survives
    /// its own remedy. Stores validate every source commitment and require the
    /// derived memory to remain in the same subject, so subject erasure reaches
    /// the whole derivation graph.
    #[serde(default, skip_serializing_if = "Vec::is_empty")]
    pub derived_from: Vec<Selected>,
}

/// Where a runtime memory write lands.
///
/// Security metadata is intentionally absent. [`StepCtx::remember`] derives
/// trust, provenance and sensitivity from the `Tainted<Value>` being stored, so
/// an untrusted model result cannot be promoted by constructing metadata that
/// says otherwise.
///
/// [`StepCtx::remember`]: crate::runtime::StepCtx::remember
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct MemoryWrite {
    pub id: String,
    pub subject: String,
    pub purpose: String,
    pub expires_at: Option<Timestamp>,
    pub access_retention_seconds: Option<u64>,
}

impl MemoryWrite {
    #[must_use]
    pub fn new(
        id: impl Into<String>,
        subject: impl Into<String>,
        purpose: impl Into<String>,
    ) -> Self {
        Self {
            id: id.into(),
            subject: subject.into(),
            purpose: purpose.into(),
            expires_at: None,
            access_retention_seconds: None,
        }
    }

    /// Expire this memory from fresh recall at a deterministic instant.
    ///
    /// A hard ceiling: a sliding window set beside it may collect the memory
    /// earlier, and no touch carries it past this.
    #[must_use]
    pub const fn expires_at(mut self, at: Timestamp) -> Self {
        self.expires_at = Some(at);
        self
    }

    /// Keep the memory for this long after each explicitly journaled recall.
    ///
    /// The window opens at the write itself, so an untouched memory is
    /// collected `seconds` after it was stored rather than kept forever
    /// waiting for a first touch — and it never extends past an `expires_at`
    /// set beside it.
    #[must_use]
    pub const fn retain_after_access(mut self, seconds: u64) -> Self {
        self.access_retention_seconds = Some(seconds);
        self
    }
}

impl MemoryItem {
    /// The label a recalled value carries.
    ///
    /// Derived from the declared trust, sensitivity and sources — never from
    /// reading the content. That is the whole defence: an item that says
    /// "verified by the security team" is a string, and a string cannot promote
    /// itself.
    #[must_use]
    pub fn label(&self) -> Label {
        let mut label = if self.trust == Trust::Trusted {
            Label::trusted()
        } else {
            Label::untrusted(SourceId::new(format!("memory:{}", self.id)))
        };
        for source in &self.provenance {
            label.provenance.insert(source.clone());
        }
        label.sensitivity = self.sensitivity;
        label
    }

    /// What the journal records instead of the content.
    ///
    /// Personal data belongs in an erasable store, not in a hash chain that
    /// cannot be redacted — so a recall journals this and re-materialises the
    /// content on replay.
    #[must_use]
    pub fn digest(&self) -> Digest {
        Digest::of(&crate::core::canon::value_bytes(&self.content))
    }

    /// Commitment used when this version is selected for a run.
    ///
    /// Content alone is insufficient: changing an item's trust, provenance,
    /// sensitivity, scope, or attribution changes what a replay may do with the
    /// same bytes. `id` and `version` travel beside this digest in [`Selected`];
    /// `superseded_at` is deliberately excluded because it is lifecycle state
    /// that may change after a run selected the version.
    #[must_use]
    pub fn selection_digest(&self) -> Digest {
        Digest::of(&crate::core::canon::value_bytes(&serde_json::json!({
            "subject": self.subject,
            "purpose": self.purpose,
            "content": self.digest().to_hex(),
            "provenance": self.provenance,
            "sensitivity": self.sensitivity,
            "trust": self.trust,
            "written_by": self.written_by,
            "created_at": crate::core::format_timestamp(self.created_at),
            "expires_at": self.expires_at.map(crate::core::format_timestamp),
            "access_retention_seconds": self.access_retention_seconds,
            "derived_from": self.derived_from,
        })))
    }
}

/// What to recall.
///
/// Deliberately not a free-text similarity query alone. `subject` and `purpose`
/// are the axes an operator can reason about and a policy can be written
/// against.
///
/// # The selection rule is part of the contract, not a backend choice
///
/// Recall truncates, so the order decides what an agent sees. Every backend
/// must select the same way: **most trusted first, then newest, then id** —
/// and when `purpose` is `None`, that ordering is **global across purposes**,
/// not per-purpose. Two backends ranking differently is an agent that recalls
/// different facts depending on which store a deployment happened to wire,
/// which no test above the store can see.
///
/// Trust leads recency because truncating by recency alone is an eviction an
/// attacker steers: anything able to write an untrusted memory writes `limit`
/// of them and the trusted ones silently lose their place. What this rule does
/// **not** cover: relevance. Within one trust rank the order is recency, so a
/// subject with more memories than `limit` can still push an older, more
/// pertinent memory out of the window — semantic retrieval, not recall, is the
/// tool for that.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct Recall {
    pub subject: String,
    /// Restrict to memories kept for this purpose, when given.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub purpose: Option<String>,
    /// At most this many, newest first.
    pub limit: usize,
    /// Deterministic lifecycle cutoff. Set by [`StepCtx::recall`] from its
    /// journaled clock; direct store callers may choose one explicitly.
    ///
    /// [`StepCtx::recall`]: crate::runtime::StepCtx::recall
    #[serde(
        default,
        with = "time::serde::rfc3339::option",
        skip_serializing_if = "Option::is_none"
    )]
    pub as_of: Option<Timestamp>,
    /// Refresh sliding retention for selected memories as a second effect.
    #[serde(default)]
    pub refresh_access: bool,
}

impl Recall {
    #[must_use]
    pub fn about(subject: impl Into<String>) -> Self {
        Self {
            subject: subject.into(),
            purpose: None,
            limit: 10,
            as_of: None,
            refresh_access: false,
        }
    }

    #[must_use]
    pub fn for_purpose(mut self, purpose: impl Into<String>) -> Self {
        self.purpose = Some(purpose.into());
        self
    }

    #[must_use]
    pub const fn limit(mut self, n: usize) -> Self {
        self.limit = n;
        self
    }

    #[must_use]
    pub const fn at(mut self, at: Timestamp) -> Self {
        self.as_of = Some(at);
        self
    }

    #[must_use]
    pub const fn refresh_access(mut self) -> Self {
        self.refresh_access = true;
        self
    }
}

/// What a caller asks a semantic index for.
///
/// Deliberately no vector, no model name and no snapshot: those are facts about
/// the wiring, and the runtime reads them from the seams. See [`IndexIdentity`]
/// for what a caller stating them wrongly would cost.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SemanticSearch {
    pub subject: String,
    /// Restrict to memories kept for this purpose, when given.
    pub purpose: Option<String>,
    pub limit: usize,
    /// Highest sensitivity the embedder and the retriever may be shown.
    ///
    /// One ceiling for both: the text goes to one and its coordinates to the
    /// other, so a deployment refusing the first has no reason to allow the
    /// second. `Public` by default.
    pub max_sensitivity: Sensitivity,
}

impl SemanticSearch {
    #[must_use]
    pub fn about(subject: impl Into<String>) -> Self {
        Self {
            subject: subject.into(),
            purpose: None,
            limit: 10,
            max_sensitivity: Sensitivity::Public,
        }
    }

    #[must_use]
    pub fn for_purpose(mut self, purpose: impl Into<String>) -> Self {
        self.purpose = Some(purpose.into());
        self
    }

    #[must_use]
    pub const fn limit(mut self, n: usize) -> Self {
        self.limit = n;
        self
    }

    #[must_use]
    pub const fn max_sensitivity(mut self, sensitivity: Sensitivity) -> Self {
        self.max_sensitivity = sensitivity;
        self
    }
}

/// Which vector space an index lives in, and which snapshot it holds.
///
/// Cosine similarity is defined between any two vectors of equal width, so a
/// query embedded by one revision and searched against an index built by
/// another does not fail — it ranks unrelated memories confidently, with no
/// exception to catch and nothing in the output that looks wrong. An index
/// therefore states the space it accepts, and the pairing is refused at wiring
/// rather than per query.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct IndexIdentity {
    /// Immutable snapshot of the corpus this index holds.
    ///
    /// In the retrieval effect's key, so re-indexing is replay divergence
    /// rather than a silently different ranking.
    pub snapshot: String,
    /// The exact [`Embedder::revision`] a **query** vector must come from.
    ///
    /// Not the revision the documents were embedded with: asymmetric embedders
    /// embed a query and a document deliberately differently, so an index built
    /// from `…/search_document` names `…/search_query` here.
    pub query_revision: String,
}

/// A vector and the space it lives in.
///
/// The revision is read from [`Embedder::revision`], never supplied beside the
/// floats: it is the only thing that makes them comparable to anything, and a
/// stated space is a space that can be stated wrongly ([`IndexIdentity`]).
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct Embedding {
    pub vector: Vec<f32>,
    pub revision: String,
}

/// A semantic query as the journal records it.
///
/// Assembled by the runtime from a [`SemanticSearch`], the wired
/// [`Embedder`] and the wired [`SemanticRetriever`]'s [`IndexIdentity`]. Every
/// field is in the retrieval effect's key, so a re-indexed corpus or a changed
/// embedding revision is divergence a replay reports rather than a different
/// answer nothing explains.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct SemanticQuery {
    pub subject: String,
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub purpose: Option<String>,
    /// The query as written — **a search input, not only provenance.**
    ///
    /// It is carried beside the vector rather than derived from it because a
    /// retriever may legitimately want both: dense similarity for meaning and
    /// the literal terms for exact matches that embeddings famously lose —
    /// identifiers, error codes, product names. A retriever that fuses the two
    /// is doing *hybrid* retrieval, and everything it needs is already in this
    /// struct.
    ///
    /// Where the fusion is *declared* is [`SemanticRetriever::profile`], which
    /// is in the effect key: changing a weighting or switching fusion off is
    /// replay divergence rather than a silently different ranking. That is the
    /// supported axis, and it is per-retriever on purpose — a per-call knob
    /// would let one run rank two ways with nothing on the record saying which.
    ///
    /// The shipped [`InMemorySemanticRetriever`] ignores this field and ranks on
    /// the vector alone. It is a reference implementation, not a statement about
    /// what the seam permits.
    pub text: String,
    /// Exact query vector, obtained through
    /// [`StepCtx::embed`](crate::runtime::StepCtx::embed).
    ///
    /// It is in the retrieval effect's key, so a replay must reproduce it
    /// exactly — which is why producing it is itself a journaled effect rather
    /// than a call a skill makes on its own. An embedding API is a network
    /// observation: two calls with the same text are not obliged to return the
    /// same floats, and a model revision guarantees they will not. Computing one
    /// inside the deterministic zone therefore quarantines a healthy run at the
    /// next replay, for a reason nothing on the record explains.
    pub embedding: Vec<f32>,
    /// The space and snapshot this query was resolved against.
    pub index: IndexIdentity,
    pub limit: usize,
    /// Highest sensitivity this retriever may receive.
    pub max_sensitivity: Sensitivity,
    /// The lifecycle cutoff this selection was screened at.
    ///
    /// A semantic index is derived and therefore stale by construction: it
    /// keeps naming versions after they are superseded, expire, or are
    /// erased. Live dispatch screens every hit against the authoritative
    /// store at this instant — the run's journaled clock, so the cutoff is on
    /// the record beside the selection it shaped rather than being an ambient
    /// wall clock two machines disagree about.
    #[serde(with = "time::serde::rfc3339")]
    pub as_of: Timestamp,
}

/// One ranked semantic selection as journaled.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct SemanticHit {
    pub selected: Selected,
    pub score: f32,
}

/// Derived semantic index, never durable memory truth.
///
/// # A static reference corpus is in scope, as memory items
///
/// Worth stating because the answer is narrower than "a vector-database seam"
/// suggests, and an evaluator reached for a bespoke retrieval tool before
/// finding this. A hit is a [`Selected`] — an id, a version and a content
/// digest — so what comes back is always a **memory item**, and the index is
/// derived from `MemoryStore`, which stays the durable truth.
///
/// Operator-ingested regulatory or reference material is therefore a first-class
/// use, written once as memory items in a scope nothing else writes to. It is
/// not *an agent's* learned memory, and it does not have to be: `memory.formation`
/// governs what a model may add, retention is optional, and a corpus with no
/// expiry simply never expires. What this seam cannot do is return documents
/// that are not memory items — an external corpus has to be ingested rather than
/// federated, so that a replay can re-materialise the exact versions a run read.
#[async_trait]
pub trait SemanticRetriever: Send + Sync + Debug {
    /// Stable non-secret implementation configuration for effect identity.
    fn profile(&self) -> Value;

    /// The space this index lives in and the snapshot it holds.
    ///
    /// Read at wiring, where a mismatch with the plane's embedder is refused,
    /// and again per query to stamp the journaled record — so an index that
    /// moves under a running plane shows up as a changed effect key.
    fn index(&self) -> IndexIdentity;

    async fn search(&self, query: &SemanticQuery) -> Result<Vec<SemanticHit>, StoreError>;
}

/// Turns text into a vector.
///
/// # Why this is a seam and not a helper
///
/// [`SemanticQuery::embedding`] enters the retrieval effect's key, so a strict
/// replay has to arrive at the same floats. An embedding service is a network
/// call: batching, hardware and a silent model revision all move the last bits,
/// and nothing about *the same text* obliges the same answer. A skill that
/// computed its own vector would therefore be making a nondeterministic
/// observation inside the deterministic zone, which is the one thing the effect
/// protocol exists to forbid — and the symptom would be a quarantine on replay
/// with nothing on the record to explain it.
///
/// So embedding crosses the effect protocol like every other observation:
/// [`StepCtx::embed`](crate::runtime::StepCtx::embed) journals the vector, and a
/// replay reads it back rather than asking again. That also makes the cost
/// visible — an embedding call is metered like any other effect — and puts the
/// model revision on the record beside the vector it produced.
///
/// [`OpenAiEmbedder`](crate::model::embeddings::OpenAiEmbedder) is the shipped
/// one, behind `providers`. Shipping it is not in tension with the crate's
/// refusal to ship a content *classifier*: that refusal is specific to
/// classifiers, and an embedder is a model driver by another name — the crate
/// ships five of those and a policy-evaluator adapter.
///
/// The consequence settled it. Without one, `StepCtx::embed` and therefore the
/// whole semantic-retrieval tier could not be used without writing Rust — the
/// same language barrier that kept the A2A server, the MCP host and push out of
/// reach of a declarative agent.
#[async_trait]
pub trait Embedder: Send + Sync + Debug {
    /// Stable, non-secret identity of the model and revision producing vectors.
    ///
    /// It goes in the effect key beside the text, so a revision change is
    /// replay divergence rather than a silently different vector, and it is
    /// what [`IndexIdentity::query_revision`] is matched against at wiring.
    /// Anything that changes what the floats *mean* belongs here — the shipped
    /// drivers fold in width and input type. Secrets never do.
    fn revision(&self) -> String;

    /// Embed one text.
    ///
    /// # Errors
    ///
    /// [`StoreError::Backend`] when the service cannot be reached or answers
    /// with something that is not a vector.
    async fn embed(&self, text: &str) -> Result<Vec<f32>, StoreError>;
}

/// One immutable vector record for [`InMemorySemanticRetriever`].
#[derive(Debug, Clone)]
pub struct SemanticVector {
    pub subject: String,
    pub purpose: String,
    pub selected: Selected,
    pub embedding: Vec<f32>,
}

/// Deterministic exact cosine retriever for tests and small corpora.
#[derive(Debug, Clone)]
pub struct InMemorySemanticRetriever {
    index: IndexIdentity,
    vectors: Vec<SemanticVector>,
}

impl InMemorySemanticRetriever {
    #[must_use]
    pub fn new(index: IndexIdentity, vectors: Vec<SemanticVector>) -> Self {
        Self { index, vectors }
    }
}

#[async_trait]
impl SemanticRetriever for InMemorySemanticRetriever {
    fn profile(&self) -> Value {
        serde_json::json!({ "driver": "in-memory-exact-cosine/v1" })
    }

    fn index(&self) -> IndexIdentity {
        self.index.clone()
    }

    async fn search(&self, query: &SemanticQuery) -> Result<Vec<SemanticHit>, StoreError> {
        validate_vector(&query.embedding)?;
        let mut hits = Vec::new();
        for candidate in &self.vectors {
            if candidate.subject != query.subject
                || query
                    .purpose
                    .as_ref()
                    .is_some_and(|purpose| purpose != &candidate.purpose)
            {
                continue;
            }
            validate_vector(&candidate.embedding)?;
            if candidate.embedding.len() != query.embedding.len() {
                return Err(StoreError::Backend(format!(
                    "semantic vector dimension {} does not match query dimension {}",
                    candidate.embedding.len(),
                    query.embedding.len()
                )));
            }
            let dot: f32 = candidate
                .embedding
                .iter()
                .zip(&query.embedding)
                .map(|(a, b)| *a * *b)
                .sum();
            let left = candidate
                .embedding
                .iter()
                .map(|value| value.powi(2))
                .sum::<f32>()
                .sqrt();
            let right = query
                .embedding
                .iter()
                .map(|value| value.powi(2))
                .sum::<f32>()
                .sqrt();
            let score = if left == 0.0 || right == 0.0 {
                0.0
            } else {
                dot / (left * right)
            };
            hits.push(SemanticHit {
                selected: candidate.selected.clone(),
                score,
            });
        }
        hits.sort_by(|a, b| {
            b.score
                .total_cmp(&a.score)
                .then_with(|| a.selected.id.cmp(&b.selected.id))
                .then_with(|| a.selected.version.cmp(&b.selected.version))
        });
        hits.truncate(query.limit);
        Ok(hits)
    }
}

fn validate_vector(vector: &[f32]) -> Result<(), StoreError> {
    if vector.is_empty() || vector.iter().any(|value| !value.is_finite()) {
        return Err(StoreError::Backend(
            "semantic vectors must be non-empty and finite".to_owned(),
        ));
    }
    Ok(())
}

/// Where a summary should land.
///
/// The parts of a memory a caller genuinely decides. Everything else about a
/// summary — its provenance, trust, sensitivity, and what it was made from — is
/// derived from the inputs, because those are not matters of opinion.
#[derive(Debug, Clone, PartialEq)]
pub struct Compaction {
    pub id: String,
    pub subject: String,
    pub purpose: String,
    pub at: Timestamp,
    /// What to tell the model to do with the memories.
    pub instruction: String,
    /// The highest sensitivity the summarising model may be shown.
    ///
    /// `Public` by default, which refuses to summarise anything above it. That
    /// is deliberate: compaction *sends the memories to a model*, so it is an
    /// egress decision and not a storage one. Without a ceiling here, summarising
    /// would be the way to move confidential content past a limit that stops
    /// every other path — and it would look like housekeeping.
    pub max_sensitivity: Sensitivity,
}

/// Governed model-assisted formation of durable memories.
#[derive(Debug, Clone, PartialEq)]
pub struct Formation {
    pub subject: String,
    pub purpose: String,
    pub instruction: String,
    pub max_items: usize,
    pub expires_at: Option<Timestamp>,
    pub access_retention_seconds: Option<u64>,
    pub max_sensitivity: Sensitivity,
}

/// One selected memory, as the journal records it.
///
/// Ids and versions rather than content: this is what makes a replay reproduce
/// the *selection* without re-running the search, and it keeps the content out
/// of a chain that cannot be redacted.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct Selected {
    pub id: String,
    pub version: u64,
    /// What the immutable content and security metadata were when selected.
    ///
    /// Checked on replay: an item whose content, label inputs, scope, lineage,
    /// or attribution changed under a version that is supposed to be immutable
    /// is a store that cannot reproduce history. Lifecycle-only
    /// `superseded_at` is excluded so a later legitimate revision does not make
    /// an earlier run unreplayable.
    pub digest: Digest,
}

/// Why a memory operation could not be completed.
#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
pub enum MemoryError {
    /// A version named by the journal is no longer in the store.
    ///
    /// Distinguished from a plain miss because the two mean opposite things: a
    /// miss is an empty result, and this is a **history that can no longer be
    /// reproduced** — usually because the memory was deliberately forgotten.
    #[error(
        "memory '{id}' version {version} was recalled by this run and is no longer \
         stored — it was forgotten, so this history cannot be replayed as it happened"
    )]
    Forgotten { id: String, version: u64 },

    /// The stored content no longer hashes to what was recorded.
    #[error(
        "memory '{id}' version {version} has different content or security metadata \
            than when it was recalled — a version is supposed to be immutable, so this \
            store cannot reproduce its own history"
    )]
    Rewritten { id: String, version: u64 },

    #[error("the memory store could not be reached: {0}")]
    Unavailable(String),
}

impl From<StoreError> for MemoryError {
    fn from(e: StoreError) -> Self {
        Self::Unavailable(e.to_string())
    }
}

/// Where memories live.
#[async_trait]
pub trait MemoryStore: Send + Sync + Debug {
    /// Whose rows this handle can reach.
    ///
    /// Defaults to [`TenantId::DEFAULT`](crate::core::TenantId::DEFAULT), the
    /// tenant a store serves until told otherwise. Override it with the tenant
    /// the handle is actually scoped to.
    ///
    /// This exists so a mismatch with the plane's tenant is a **startup
    /// refusal**. When a key ring is wired, `build()` seals this state under
    /// the plane's tenant while the store writes rows under its own; the two
    /// disagreeing is not a leak — the scopes simply differ — but it seals the
    /// state under a scope erasure will never destroy. That is an erasure that
    /// reports success and misses, which is the one failure a deletion
    /// guarantee cannot have.
    fn tenant(&self) -> &str {
        crate::core::TenantId::DEFAULT
    }

    /// Whether this store's erasure lifecycle lock spans instances.
    ///
    /// `None` — the default, and the honest answer for most stores — means
    /// there is no lifecycle lock because there is no cryptographic erasure to
    /// serialise. `Some(false)` means there is one and it is process-local, so
    /// a plane sharing its durable state with another instance would have a
    /// window between an erasure's hold check and its key destruction in which
    /// the other instance can write. The runtime refuses that pairing at
    /// `build`, which is why this is a question a store can be asked rather
    /// than a fact an operator has to remember.
    fn erasure_is_distributed(&self) -> Option<bool> {
        None
    }

    /// Append a new version of a memory.
    ///
    /// Appends rather than replaces: the previous version is marked superseded
    /// and kept, so lineage survives and one bad write can be undone without
    /// guessing what it replaced.
    ///
    /// # Errors
    ///
    /// If the store cannot be reached.
    async fn remember(&self, item: &MemoryItem) -> Result<u64, StoreError>;

    /// The current versions matching a query — most trusted first, then
    /// newest, then id, globally across purposes (see [`Recall`] for why the
    /// ordering is contractual).
    ///
    /// Superseded and forgotten versions are not returned.
    ///
    /// # Errors
    ///
    /// If the store cannot be reached.
    async fn recall(&self, query: &Recall) -> Result<Vec<MemoryItem>, StoreError>;

    /// Every id currently belonging to a subject — the erasure path's
    /// enumeration.
    ///
    /// A separate operation rather than a `recall` with a huge limit, because
    /// erasure must not ride a bounded, content-returning query: a backend is
    /// free to cap or refuse extreme recall limits (the `PostgreSQL` store
    /// refuses limits beyond `BIGINT`), and a cryptographic erasure that
    /// enumerated through one silently missed whatever the cap cut off — while
    /// reporting success. This returns ids only, unbounded, in stable order.
    ///
    /// What it does **not** cover: forgotten and swept ids. They have no
    /// current version, no content, and their tombstones are not subject-keyed
    /// — an erasure that needs them already erased them.
    ///
    /// # Errors
    ///
    /// If the store cannot be reached.
    async fn subject_ids(&self, subject: &str) -> Result<Vec<String>, StoreError>;

    /// One exact version, superseded or not.
    ///
    /// This is what replay uses: it names the version the run actually read, not
    /// whatever is current now.
    ///
    /// # Errors
    ///
    /// If the store cannot be reached.
    async fn version(&self, id: &str, version: u64) -> Result<Option<MemoryItem>, StoreError>;

    /// The one version of `id` a fresh recall at `as_of` would be allowed to
    /// see — or `None`, which covers a memory that was never written, was
    /// forgotten or swept, or whose effective expiry (`min(expires_at, access
    /// window)`) has passed the cutoff. `None` for `as_of` skips the expiry
    /// check, exactly as it does on [`recall`](Self::recall).
    ///
    /// The by-id twin of [`recall`](Self::recall)'s lifecycle rule, and a
    /// question [`version`](Self::version) deliberately cannot answer:
    /// `version` serves replay, which must keep reading superseded and
    /// expired state, so nothing built on it can tell *still current* from
    /// *still stored*. The semantic tier's lifecycle screen is the consumer —
    /// see [`SemanticRecall`](crate::runtime::effects::SemanticRecall) for
    /// why a stale hit leaves the selection rather than being served or
    /// failing the query.
    ///
    /// # Errors
    ///
    /// If the store cannot be reached.
    async fn current(
        &self,
        id: &str,
        as_of: Option<Timestamp>,
    ) -> Result<Option<MemoryItem>, StoreError>;

    /// Forget a memory and every version of it.
    ///
    /// Selective by construction: one id, not a purge. A repair that can only
    /// drop everything is one nobody performs. The id remains reserved after
    /// erasure: recycling it could make an old journal selection or derivation
    /// edge refer to unrelated new content.
    ///
    /// # Errors
    ///
    /// If the store cannot be reached.
    async fn forget(&self, id: &str) -> Result<(), StoreError>;

    /// Forget everything about a subject.
    ///
    /// The unit an erasure request names — a person, an account, a matter.
    ///
    /// # Errors
    ///
    /// If the store cannot be reached.
    async fn forget_subject(&self, subject: &str) -> Result<usize, StoreError>;

    /// Memories derived from this one, directly.
    ///
    /// The repair path. A summary absorbs what it summarised, so a poisoned
    /// memory does not stop being a problem when it is forgotten — it stops
    /// being *visible* while its content continues to arrive in every summary
    /// that read it.
    ///
    /// # Errors
    ///
    /// If the store cannot be reached.
    async fn derivatives(&self, id: &str) -> Result<Vec<MemoryItem>, StoreError>;

    /// Forget a memory and everything derived from it, transitively.
    ///
    /// The form an **erasure** request needs. [`forget`](Self::forget) is the
    /// form a *correction* needs: a stale memory whose summaries are still
    /// legitimate should not take them with it. The two are separate calls
    /// because they answer different questions, and defaulting either way would
    /// be wrong half the time — silently.
    ///
    /// This is a required store operation rather than a default assembled from
    /// [`derivatives`](Self::derivatives) and [`forget`](Self::forget). Those
    /// calls are individually atomic but leave a gap in which another writer
    /// can add a derivative that the erasure never sees. Implementations must
    /// serialize derivative creation with the complete traversal and deletion.
    ///
    /// The traversal passes **through tombstones**: a memory forgotten
    /// individually keeps its derivation edges in both directions, so a later
    /// cascade from further upstream still reaches everything transitively
    /// derived — A → B → C with B corrected away must not shelter C from A's
    /// erasure.
    ///
    /// It is also **version-granular**, because supersession does not
    /// un-absorb anything. A derivative whose *current* version read a doomed
    /// node is erased as a whole id. A derivative that has since been honestly
    /// re-derived from clean sources keeps its current version — but the
    /// **superseded versions that named the doomed source are erased**, along
    /// with anything transitively derived from exactly those versions.
    /// Otherwise a rolling summary's v1, which absorbed the poisoned memory
    /// and remained readable through [`version`](Self::version), would outlive
    /// the erasure that claimed to reach everything derived.
    ///
    /// Returns the number of ids whose state this call actually removed — an
    /// id that lost only superseded versions counts once; a tombstone passed
    /// through is routed, not counted.
    ///
    /// # Errors
    ///
    /// If the store cannot be reached.
    async fn forget_cascading(&self, id: &str) -> Result<usize, StoreError>;

    /// Place or release a legal hold. A held id cannot be forgotten, swept, or
    /// removed as part of subject/cascading erasure.
    async fn set_legal_hold(&self, id: &str, held: bool) -> Result<(), StoreError>;

    /// Whether an id is currently protected by legal hold.
    async fn legal_hold(&self, id: &str) -> Result<bool, StoreError>;

    /// Atomically erase current memories whose effective expiry has passed,
    /// unless held. Returns the number of memory ids erased.
    ///
    /// The effective expiry is `min(expires_at, access window)` — the hard
    /// ceiling wins, and a lapsed sliding window collects an item its ceiling
    /// would still have kept. The cutoff is inclusive (`<= at`), and the sweep
    /// keeps derivation edges in both directions exactly as `forget` does, so
    /// a later cascading erasure still routes through the tombstone.
    async fn sweep_expired(&self, at: Timestamp) -> Result<usize, StoreError>;

    /// Refresh sliding access retention for current ids at a journaled instant.
    ///
    /// Slides each touched id's window to `at + window`. It cannot extend a
    /// life past an `expires_at` — that ceiling is immutable, and both recall
    /// and the sweep take the earlier of the two.
    ///
    /// A touch on an id whose window has **lapsed but not yet been swept**
    /// deliberately *resurrects* it: expiry becomes fact at the sweep, not at
    /// the instant the window closes, and stores take `max(existing, at +
    /// window)` so a late touch simply opens a new window. The alternative —
    /// refusing the touch — would make the answer depend on how recently the
    /// sweeper ran, which is an ambient race no journal records. What this
    /// choice does not cover: the `expires_at` ceiling, which no touch moves,
    /// and a swept id, which is a tombstone no touch revives.
    async fn touch(&self, ids: &[String], at: Timestamp) -> Result<(), StoreError>;
}

#[cfg(test)]
mod write_tests {
    use super::*;

    /// The two lifecycle builders set their own field and no other.
    ///
    /// Fixed expiry and sliding retention answer opposite questions — "this
    /// stops being recallable on Thursday whatever happens" versus "this stays
    /// as long as it keeps being read" — and a memory carrying both by accident
    /// is one whose disposal date nobody can state. The store paths were tested
    /// through YAML and direct field assignment; the builders a caller reaches
    /// for had no test, so a swapped assignment would have been invisible.
    #[test]
    fn each_lifecycle_builder_sets_only_what_it_names() {
        let plain = MemoryWrite::new("m-1", "account-1", "support");
        assert_eq!(plain.expires_at, None, "neither is set by default");
        assert_eq!(plain.access_retention_seconds, None);

        let sliding = MemoryWrite::new("m-1", "account-1", "support").retain_after_access(600);
        assert_eq!(sliding.access_retention_seconds, Some(600));
        assert_eq!(
            sliding.expires_at, None,
            "a sliding window is not also a fixed expiry"
        );

        let at = Timestamp::UNIX_EPOCH;
        let fixed = MemoryWrite::new("m-1", "account-1", "support").expires_at(at);
        assert_eq!(fixed.expires_at, Some(at));
        assert_eq!(
            fixed.access_retention_seconds, None,
            "a fixed expiry is not also a sliding window"
        );

        // Identity is untouched by either: the id, subject and purpose are what
        // erasure and retrieval key on.
        for built in [&sliding, &fixed] {
            assert_eq!(built.id, plain.id);
            assert_eq!(built.subject, plain.subject);
            assert_eq!(built.purpose, plain.purpose);
        }
    }
}

#[cfg(test)]
mod semantic_tests {
    use super::*;

    fn selected(id: &str) -> Selected {
        Selected {
            id: id.to_owned(),
            version: 1,
            digest: Digest::of(id.as_bytes()),
        }
    }

    #[tokio::test]
    async fn exact_cosine_retrieval_is_scoped_and_ranked() {
        let retriever = InMemorySemanticRetriever::new(
            IndexIdentity {
                snapshot: "snapshot-7".to_owned(),
                query_revision: "embed-v3@2026-07-01".to_owned(),
            },
            vec![
                SemanticVector {
                    subject: "account-1".to_owned(),
                    purpose: "support".to_owned(),
                    selected: selected("near"),
                    embedding: vec![1.0, 0.0],
                },
                SemanticVector {
                    subject: "account-1".to_owned(),
                    purpose: "support".to_owned(),
                    selected: selected("far"),
                    embedding: vec![0.0, 1.0],
                },
                SemanticVector {
                    subject: "account-2".to_owned(),
                    purpose: "support".to_owned(),
                    selected: selected("wrong-subject"),
                    embedding: vec![1.0, 0.0],
                },
            ],
        );
        let query = SemanticQuery {
            subject: "account-1".to_owned(),
            purpose: Some("support".to_owned()),
            text: "query".to_owned(),
            embedding: vec![1.0, 0.0],
            index: retriever.index(),
            limit: 2,
            max_sensitivity: Sensitivity::Internal,
            as_of: Timestamp::UNIX_EPOCH,
        };
        let hits = retriever.search(&query).await.expect("semantic search");
        assert_eq!(
            hits.iter()
                .map(|hit| hit.selected.id.as_str())
                .collect::<Vec<_>>(),
            vec!["near", "far"]
        );
        assert!(hits[0].score > hits[1].score);

        // A vector of another width is a different space wearing the same
        // struct, and comparing the two would rank on whichever prefix happened
        // to line up.
        let mut mismatched = query;
        mismatched.embedding = vec![1.0, 0.0, 0.0];
        assert!(retriever.search(&mismatched).await.is_err());
    }
}