Skip to main content

areev_loop/
substrate.rs

1//! The `OmsSubstrate` trait — the engine's only contact with a store. It is
2//! defined in terms of the OMS Level-2 protocol (CAL text ↔ JSON rows, grain
3//! get/put/supersede) plus curated typed reads the built-in analyzers use.
4//! Areev is the first substrate; the in-repo `ReferenceSubstrate` lets engine
5//! CI run with zero Areev, and doubles as the third-party conformance kit.
6
7use crate::error::Result;
8use crate::model::GrainRecord;
9use serde_json::{Map, Value};
10
11/// Optional substrate capabilities, declared once and matched against each
12/// analyzer manifest's `requires` list. A missing capability degrades an
13/// analyzer to an activation-ladder entry, never a silent no-op (§8).
14#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
15pub struct Capabilities {
16    /// Multiple concurrent heads per entity are tracked and queryable
17    /// (fork surfacing needs this).
18    pub forks: bool,
19    /// A telemetry sidecar records recall/access history.
20    pub telemetry: bool,
21    /// An embedder is installed (upgrades T0 analyzers to T1).
22    pub embeddings: bool,
23}
24
25/// Read filters for curated grain reads.
26#[derive(Debug, Clone, Copy)]
27pub struct ReadOpts {
28    /// When true (default), only live (non-superseded) grains are returned.
29    pub live_only: bool,
30    /// When set, only grains created at or after this epoch-ms are returned
31    /// (the incremental watermark scan, §8).
32    pub since_ms: Option<i64>,
33}
34
35impl Default for ReadOpts {
36    fn default() -> Self {
37        ReadOpts {
38            live_only: true,
39            since_ms: None,
40        }
41    }
42}
43
44/// A grain to be written by an apply. `derived_from` and other provenance go
45/// in `fields`; the substrate computes the content address.
46#[derive(Debug, Clone, PartialEq)]
47pub struct GrainSpec {
48    pub grain_type: String,
49    pub namespace: String,
50    pub fields: Map<String, Value>,
51}
52
53impl GrainSpec {
54    pub fn new(grain_type: impl Into<String>, namespace: impl Into<String>) -> Self {
55        GrainSpec {
56            grain_type: grain_type.into(),
57            namespace: namespace.into(),
58            fields: Map::new(),
59        }
60    }
61
62    pub fn with_field(mut self, key: impl Into<String>, value: impl Into<Value>) -> Self {
63        self.fields.insert(key.into(), value.into());
64        self
65    }
66}
67
68/// One entity holding more than one live head (fork surfacing input).
69#[derive(Debug, Clone, PartialEq)]
70pub struct HeadGroup {
71    /// Entity identity, e.g. `"caller/john"` (namespace-qualified subject).
72    pub entity: String,
73    /// The competing head hashes.
74    pub heads: Vec<String>,
75}
76
77/// A snapshot of the recall-telemetry sidecar's rollups (§8). Telemetry-fed
78/// analyzers (`cold_grains`, `coverage_gap`, `budget_pressure`) read this via
79/// [`crate::analyzer::AnalyzeCtx::telemetry`]. It is an owned snapshot, not a
80/// live handle — analyzers stay read-only and can't reach the sidecar. A
81/// substrate without a sidecar returns `None`, so the analyzer degrades to an
82/// activation-ladder entry rather than firing on absent evidence.
83#[derive(Debug, Clone, Default, PartialEq)]
84pub struct TelemetryView {
85    /// Per-grain recall rollups. A grain **absent** here has never been
86    /// recalled (the `cold_grains` signal).
87    pub access: Vec<GrainAccess>,
88    /// Per-question rollups over free-text recalls (the `coverage_gap` signal).
89    pub queries: Vec<QueryUsage>,
90    /// Assembly-budget pressure rollup (the `budget_pressure` signal).
91    pub budget: BudgetUsage,
92}
93
94/// How often one grain has been surfaced by recall.
95#[derive(Debug, Clone, PartialEq)]
96pub struct GrainAccess {
97    pub hash: String,
98    pub recall_count: i64,
99    pub last_ms: i64,
100}
101
102/// How a recurring recall question has fared.
103#[derive(Debug, Clone, PartialEq)]
104pub struct QueryUsage {
105    /// A short human-readable sample of the query intent.
106    pub sample: String,
107    pub run_count: i64,
108    /// How many of those runs returned nothing — the coverage-gap signal.
109    pub empty_count: i64,
110    pub sum_results: i64,
111    pub last_ms: i64,
112}
113
114/// Assembly-budget pressure rollup.
115#[derive(Debug, Clone, Default, PartialEq)]
116pub struct BudgetUsage {
117    pub sample_count: i64,
118    pub overflow_count: i64,
119}
120
121/// The read-only slice of the substrate. Analyzers receive this (via
122/// `AnalyzeCtx`) and nothing else — the trust floor's "analyzers execute
123/// read-only" is enforced by the type system: a `&dyn SubstrateRead` cannot
124/// reach any mutating method. It is object-safe (no generics) so
125/// `builtin_analyzers()` can hand out `Box<dyn Analyzer>`.
126pub trait SubstrateRead {
127    /// Declared optional capabilities.
128    fn capabilities(&self) -> Capabilities;
129
130    /// Curated read: all grains of one OMS type, optionally namespace-scoped
131    /// and watermark/liveness filtered.
132    fn grains_of_type(
133        &self,
134        grain_type: &str,
135        namespace: Option<&str>,
136        opts: ReadOpts,
137    ) -> Result<Vec<GrainRecord>>;
138
139    /// Fetch one grain by content address.
140    fn grain(&self, hash: &str) -> Result<Option<GrainRecord>>;
141
142    /// Entities with more than one live head. Requires the `forks` capability;
143    /// the default impl reports it missing so non-fork substrates degrade
144    /// cleanly rather than pretend.
145    fn heads(&self, _namespace: Option<&str>) -> Result<Vec<HeadGroup>> {
146        Err(crate::error::Error::CapabilityMissing("forks".into()))
147    }
148
149    /// A snapshot of the recall-telemetry rollups (§8). Requires the
150    /// `telemetry` capability; the default returns `None` so substrates without
151    /// a sidecar degrade cleanly. `namespace` scopes the snapshot when set.
152    fn telemetry(&self, _namespace: Option<&str>) -> Result<Option<TelemetryView>> {
153        Ok(None)
154    }
155}
156
157/// The full store protocol the engine binds to: reads (via the supertrait)
158/// plus governed writes, CAL, and state persistence. All methods are fallible;
159/// a substrate fault surfaces as [`crate::error::Error::Substrate`].
160pub trait OmsSubstrate: SubstrateRead {
161    /// Append a new grain; returns its content address.
162    fn put_grain(&mut self, spec: &GrainSpec) -> Result<String>;
163
164    /// Supersede `target_hash` with a new grain carrying `justification`;
165    /// returns the new grain's address. Atomic and distinct from put
166    /// (OMS §28.4).
167    fn supersede(
168        &mut self,
169        target_hash: &str,
170        spec: &GrainSpec,
171        justification: &str,
172    ) -> Result<String>;
173
174    /// Index-layer retraction (`verification_status = retracted`) — the
175    /// inverse of an applied ADD, used by rollback. Not destructive (the grain
176    /// stays content-addressed; only the index marks it retracted). The
177    /// default reports it unsupported so substrates opt in.
178    fn retract(&mut self, hash: &str, reason: &str) -> Result<()> {
179        Err(crate::error::Error::Substrate(format!(
180            "retract not supported by this substrate ({hash}: {reason})"
181        )))
182    }
183
184    /// Store an opaque blob (candidate tool CODE, evalset payloads) in the
185    /// substrate's CAS, returning its address. §7.4's blob seam —
186    /// CAPABILITY-GATED: the default refuses, so a loop can only carry code
187    /// on substrates that explicitly opt in. Code enters the substrate only
188    /// through this seam or an authored add — never from a git mirror.
189    fn put_blob(&mut self, bytes: &[u8]) -> Result<String> {
190        let _ = bytes;
191        Err(crate::error::Error::Substrate(
192            "put_blob not supported by this substrate (code-carrying loops \
193             need an opted-in blob seam)"
194            .into(),
195        ))
196    }
197
198    /// Fetch a blob by the address `put_blob` returned. Same capability gate.
199    fn get_blob(&mut self, address: &str) -> Result<Vec<u8>> {
200        Err(crate::error::Error::Substrate(format!(
201            "get_blob not supported by this substrate ({address})"
202        )))
203    }
204
205    /// Execute CAL text, returning result rows as JSON. Used to regenerate
206    /// evidence sets (`evidence_query`) and to apply `proposal_cal`. A
207    /// substrate MAY reject CAL it cannot run with [`Error::CalUnsupported`].
208    ///
209    /// [`Error::CalUnsupported`]: crate::error::Error::CalUnsupported
210    fn execute_cal(&mut self, cal: &str) -> Result<Vec<Value>>;
211
212    /// Validate a CAL batch without executing it (statement classification,
213    /// destructive-op detection). Delegated to the substrate — the engine
214    /// contains a CAL *writer*, never a parser.
215    fn validate_cal(&self, cal: &str) -> Result<()>;
216
217    /// Load the persisted loop state blob (config + watermarks/cooldowns).
218    /// Returns `Value::Null` when nothing has been stored yet.
219    fn load_state(&self) -> Result<Value>;
220
221    /// Persist the loop state blob (a file-truth, so it travels with the
222    /// file on sync).
223    fn store_state(&mut self, state: &Value) -> Result<()>;
224}