Skip to main content

a3s_code_core/evaluation/
evidence.rs

1//! Bounded evidence reads over the existing Code run and artifact stores.
2
3use super::identity::{
4    digest_bytes, digest_json, validate_digest, ExecutionFrameV1, ExecutionTargetV1,
5};
6use super::journal::{ExecutionFactInputV1, ExecutionFactJournal, ExecutionFactV1};
7use crate::agent::AgentEvent;
8use crate::event_protocol::{run_event_envelope_v1, EventEnvelopeV1};
9use crate::run::{InMemoryRunStore, RunSnapshot, RunStatus};
10use crate::tools::ArtifactStore;
11use async_trait::async_trait;
12use serde::{Deserialize, Serialize};
13use serde_json::Value;
14use std::collections::BTreeSet;
15use std::sync::Arc;
16use thiserror::Error;
17
18pub const EVIDENCE_SNAPSHOT_SCHEMA_V1: &str = "a3s.code.evidence-snapshot.v1";
19pub const EVIDENCE_MAX_EVENTS: usize = 4096;
20pub const EVIDENCE_MAX_EVENT_BYTES: usize = 16 * 1024 * 1024;
21pub const EVIDENCE_MAX_ARTIFACTS: usize = 256;
22pub const EVIDENCE_MAX_ARTIFACT_BYTES: usize = 16 * 1024 * 1024;
23pub const EVIDENCE_MAX_PROMPT_BYTES: usize = 1024 * 1024;
24pub const EVIDENCE_MAX_RESULT_BYTES: usize = 1024 * 1024;
25const MAX_EVENT_PAYLOAD_BYTES: u64 = 4 * 1024 * 1024;
26const MAX_REFERENCED_ARTIFACTS: usize = 256;
27const MAX_ARTIFACT_URI_BYTES: usize = 1024;
28const MAX_TOOL_NAME_BYTES: usize = 256;
29const MAX_EVENT_TYPE_BYTES: usize = 256;
30const MAX_EVENT_METADATA_BYTES: usize = 16 * 1024;
31
32#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, Default)]
33#[serde(rename_all = "snake_case")]
34pub enum EvidenceContentModeV1 {
35    /// Return event payloads as digest/size markers only.
36    #[default]
37    DigestOnly,
38    /// Return event payloads only while each encoded event remains within the
39    /// request's byte budget; oversized payloads become digest markers.
40    BoundedPayload,
41}
42
43#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
44#[serde(deny_unknown_fields)]
45pub struct EvidenceLimitsV1 {
46    pub max_events: usize,
47    pub max_event_bytes: usize,
48    pub max_artifacts: usize,
49    pub max_artifact_bytes: usize,
50    pub max_prompt_bytes: usize,
51    pub max_result_bytes: usize,
52}
53
54impl Default for EvidenceLimitsV1 {
55    fn default() -> Self {
56        Self {
57            max_events: 128,
58            max_event_bytes: 256 * 1024,
59            max_artifacts: 32,
60            max_artifact_bytes: 1024 * 1024,
61            max_prompt_bytes: 16 * 1024,
62            max_result_bytes: 64 * 1024,
63        }
64    }
65}
66
67impl EvidenceLimitsV1 {
68    pub fn validate(&self) -> Result<(), EvidenceError> {
69        if self.max_events == 0
70            || self.max_event_bytes == 0
71            || self.max_artifact_bytes == 0
72            || self.max_prompt_bytes == 0
73            || self.max_result_bytes == 0
74            || self.max_events > EVIDENCE_MAX_EVENTS
75            || self.max_event_bytes > EVIDENCE_MAX_EVENT_BYTES
76            || self.max_artifacts > EVIDENCE_MAX_ARTIFACTS
77            || self.max_artifact_bytes > EVIDENCE_MAX_ARTIFACT_BYTES
78            || self.max_prompt_bytes > EVIDENCE_MAX_PROMPT_BYTES
79            || self.max_result_bytes > EVIDENCE_MAX_RESULT_BYTES
80        {
81            return Err(EvidenceError::InvalidLimit);
82        }
83        Ok(())
84    }
85}
86
87#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
88#[serde(deny_unknown_fields)]
89pub struct EvidenceReadRequestV1 {
90    pub target: ExecutionTargetV1,
91    #[serde(default, skip_serializing_if = "Option::is_none")]
92    pub after_sequence: Option<u64>,
93    #[serde(default)]
94    pub limits: EvidenceLimitsV1,
95    #[serde(default)]
96    pub content_mode: EvidenceContentModeV1,
97    #[serde(default)]
98    pub include_prompt: bool,
99    /// Include bounded terminal result/error text in addition to their
100    /// digests. The default remains digest-only for cross-tenant safety.
101    #[serde(default)]
102    pub include_terminal_text: bool,
103    #[serde(default)]
104    pub include_artifact_content: bool,
105}
106
107impl EvidenceReadRequestV1 {
108    pub fn new(target: ExecutionTargetV1) -> Self {
109        Self {
110            target,
111            after_sequence: None,
112            limits: EvidenceLimitsV1::default(),
113            content_mode: EvidenceContentModeV1::DigestOnly,
114            include_prompt: false,
115            include_terminal_text: false,
116            include_artifact_content: false,
117        }
118    }
119
120    pub fn validate(&self) -> Result<(), EvidenceError> {
121        self.target
122            .validate()
123            .map_err(|_| EvidenceError::InvalidTarget)?;
124        self.limits.validate()
125    }
126}
127
128#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
129#[serde(deny_unknown_fields)]
130pub struct EvidenceRunStateV1 {
131    pub schema: String,
132    pub target: ExecutionTargetV1,
133    pub status: RunStatus,
134    pub created_at_ms: u64,
135    pub updated_at_ms: u64,
136    pub event_count: u64,
137    pub prompt_bytes: u64,
138    pub prompt_digest: String,
139    #[serde(default, skip_serializing_if = "Option::is_none")]
140    pub prompt: Option<String>,
141    #[serde(default, skip_serializing_if = "Option::is_none")]
142    pub result_bytes: Option<u64>,
143    #[serde(default, skip_serializing_if = "Option::is_none")]
144    pub result: Option<String>,
145    #[serde(default, skip_serializing_if = "Option::is_none")]
146    pub error_bytes: Option<u64>,
147    #[serde(default, skip_serializing_if = "Option::is_none")]
148    pub error: Option<String>,
149    #[serde(default, skip_serializing_if = "Option::is_none")]
150    pub result_digest: Option<String>,
151    #[serde(default, skip_serializing_if = "Option::is_none")]
152    pub error_digest: Option<String>,
153    #[serde(default, skip_serializing_if = "Option::is_none")]
154    pub workspace_change_set_digest: Option<String>,
155    #[serde(default, skip_serializing_if = "Option::is_none")]
156    pub capability_binding_digest: Option<String>,
157    #[serde(default, skip_serializing_if = "Option::is_none")]
158    pub cognitive_binding_digest: Option<String>,
159}
160
161impl EvidenceRunStateV1 {
162    fn from_snapshot(
163        snapshot: &RunSnapshot,
164        target: &ExecutionTargetV1,
165        limits: EvidenceLimitsV1,
166        include_prompt: bool,
167        include_terminal_text: bool,
168    ) -> Result<Self, EvidenceError> {
169        let prompt = if include_prompt {
170            bounded_string(&snapshot.prompt, limits.max_prompt_bytes)
171        } else {
172            None
173        };
174        let prompt_bytes =
175            u64::try_from(snapshot.prompt.len()).map_err(|_| EvidenceError::NumericOverflow)?;
176        let prompt_digest = digest_bytes("a3s.code.evidence.prompt.v1", snapshot.prompt.as_bytes());
177        let result_bytes = snapshot
178            .result_text
179            .as_ref()
180            .map(|value| u64::try_from(value.len()))
181            .transpose()
182            .map_err(|_| EvidenceError::NumericOverflow)?;
183        let result_digest = snapshot
184            .result_text
185            .as_deref()
186            .map(|value| digest_bytes("a3s.code.evidence.result.v1", value.as_bytes()));
187        let error_bytes = snapshot
188            .error
189            .as_ref()
190            .map(|value| u64::try_from(value.len()))
191            .transpose()
192            .map_err(|_| EvidenceError::NumericOverflow)?;
193        let error_digest = snapshot
194            .error
195            .as_deref()
196            .map(|value| digest_bytes("a3s.code.evidence.error.v1", value.as_bytes()));
197        let result = if include_terminal_text {
198            snapshot
199                .result_text
200                .as_deref()
201                .and_then(|value| bounded_string(value, limits.max_result_bytes))
202        } else {
203            None
204        };
205        let error = if include_terminal_text {
206            snapshot
207                .error
208                .as_deref()
209                .and_then(|value| bounded_string(value, limits.max_result_bytes))
210        } else {
211            None
212        };
213        let workspace_change_set_digest = snapshot
214            .workspace_change_set
215            .as_ref()
216            .map(|value| digest_json("a3s.code.evidence.workspace-change-set.v1", value))
217            .transpose()
218            .map_err(EvidenceError::Serialization)?;
219        let capability_binding_digest = snapshot
220            .capability_binding
221            .as_ref()
222            .map(|value| digest_json("a3s.code.evidence.capability-binding.v1", value))
223            .transpose()
224            .map_err(EvidenceError::Serialization)?;
225        let cognitive_binding_digest = snapshot
226            .cognitive_package_binding
227            .as_ref()
228            .map(|value| digest_json("a3s.code.evidence.cognitive-binding.v1", value))
229            .transpose()
230            .map_err(EvidenceError::Serialization)?;
231        Ok(Self {
232            schema: "a3s.code.evidence-run-state.v1".to_string(),
233            target: target.clone(),
234            status: snapshot.status,
235            created_at_ms: snapshot.created_at_ms,
236            updated_at_ms: snapshot.updated_at_ms,
237            event_count: u64::try_from(snapshot.event_count)
238                .map_err(|_| EvidenceError::NumericOverflow)?,
239            prompt_bytes,
240            prompt_digest,
241            prompt,
242            result_bytes,
243            result,
244            error_bytes,
245            error,
246            result_digest,
247            error_digest,
248            workspace_change_set_digest,
249            capability_binding_digest,
250            cognitive_binding_digest,
251        })
252    }
253}
254
255#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
256#[serde(deny_unknown_fields)]
257pub struct EvidenceEventV1 {
258    pub sequence: u64,
259    pub occurred_at_ms: u64,
260    pub event: EventEnvelopeV1,
261    pub payload_digest: String,
262    pub payload_bytes: u64,
263}
264
265#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
266#[serde(deny_unknown_fields)]
267pub struct EvidenceArtifactV1 {
268    pub artifact_uri: String,
269    pub tool_name: String,
270    pub content_digest: String,
271    pub content_bytes: u64,
272    #[serde(default, skip_serializing_if = "Option::is_none")]
273    pub content: Option<String>,
274}
275
276#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
277#[serde(deny_unknown_fields)]
278pub struct EvidenceSnapshotV1 {
279    pub schema: String,
280    pub target: ExecutionTargetV1,
281    pub observed_at_ms: u64,
282    #[serde(default, skip_serializing_if = "Option::is_none")]
283    pub after_sequence: Option<u64>,
284    pub first_available_sequence: Option<u64>,
285    pub latest_sequence_exclusive: u64,
286    pub state: EvidenceRunStateV1,
287    pub facts: Vec<ExecutionFactV1>,
288    pub events: Vec<EvidenceEventV1>,
289    pub artifacts: Vec<EvidenceArtifactV1>,
290    pub complete: bool,
291    pub retention_gap: bool,
292    pub snapshot_digest: String,
293}
294
295impl EvidenceSnapshotV1 {
296    pub fn validate(&self) -> Result<(), EvidenceError> {
297        if self.schema != EVIDENCE_SNAPSHOT_SCHEMA_V1 {
298            return Err(EvidenceError::UnsupportedSchema);
299        }
300        self.target
301            .validate()
302            .map_err(|_| EvidenceError::InvalidTarget)?;
303        if self.state.target != self.target {
304            return Err(EvidenceError::TargetMismatch);
305        }
306        if self
307            .first_available_sequence
308            .is_some_and(|first| first >= self.latest_sequence_exclusive)
309        {
310            return Err(EvidenceError::InvalidField("first_available_sequence"));
311        }
312        if self.state.schema != "a3s.code.evidence-run-state.v1" {
313            return Err(EvidenceError::UnsupportedSchema);
314        }
315        if self.state.event_count != self.latest_sequence_exclusive {
316            return Err(EvidenceError::InvalidField("state.event_count"));
317        }
318        if self.facts.len() > EVIDENCE_MAX_EVENTS
319            || self.events.len() > EVIDENCE_MAX_EVENTS
320            || self.artifacts.len() > EVIDENCE_MAX_ARTIFACTS
321        {
322            return Err(EvidenceError::InvalidLimit);
323        }
324        if self.state.created_at_ms > self.state.updated_at_ms
325            || self.observed_at_ms < self.state.updated_at_ms
326        {
327            return Err(EvidenceError::InvalidField("state.timestamps"));
328        }
329        let requested_start = self
330            .after_sequence
331            .map(|sequence| sequence.saturating_add(1))
332            .unwrap_or(0);
333        let expected_gap = requested_start < self.latest_sequence_exclusive
334            && self
335                .first_available_sequence
336                .is_none_or(|first| requested_start < first);
337        if self.retention_gap != expected_gap {
338            return Err(EvidenceError::InvalidField("retention_gap"));
339        }
340        validate_digest(&self.state.prompt_digest)
341            .map_err(|_| EvidenceError::InvalidDigest("state.prompt_digest"))?;
342        if let Some(prompt) = self.state.prompt.as_deref() {
343            let prompt_bytes =
344                u64::try_from(prompt.len()).map_err(|_| EvidenceError::NumericOverflow)?;
345            if self.state.prompt_bytes != prompt_bytes {
346                return Err(EvidenceError::InvalidField("state.prompt_bytes"));
347            }
348            if self.state.prompt_digest
349                != digest_bytes("a3s.code.evidence.prompt.v1", prompt.as_bytes())
350            {
351                return Err(EvidenceError::DigestMismatch("state.prompt_digest"));
352            }
353        }
354        if self.state.result.is_some() && self.state.result_digest.is_none() {
355            return Err(EvidenceError::InvalidField("state.result_digest"));
356        }
357        if self.state.error.is_some() && self.state.error_digest.is_none() {
358            return Err(EvidenceError::InvalidField("state.error_digest"));
359        }
360        if self.state.result_bytes.is_some() != self.state.result_digest.is_some()
361            || self.state.error_bytes.is_some() != self.state.error_digest.is_some()
362        {
363            return Err(EvidenceError::InvalidField("state.text_bytes"));
364        }
365        for (field, digest) in [
366            ("state.result_digest", self.state.result_digest.as_deref()),
367            ("state.error_digest", self.state.error_digest.as_deref()),
368        ] {
369            if let Some(digest) = digest {
370                validate_digest(digest).map_err(|_| EvidenceError::InvalidDigest(field))?;
371            }
372        }
373        if let (Some(result), Some(digest)) = (
374            self.state.result.as_deref(),
375            self.state.result_digest.as_deref(),
376        ) {
377            if self.state.result_bytes != u64::try_from(result.len()).ok() {
378                return Err(EvidenceError::InvalidField("state.result_bytes"));
379            }
380            if digest != digest_bytes("a3s.code.evidence.result.v1", result.as_bytes()) {
381                return Err(EvidenceError::DigestMismatch("state.result_digest"));
382            }
383        }
384        if let (Some(error), Some(digest)) = (
385            self.state.error.as_deref(),
386            self.state.error_digest.as_deref(),
387        ) {
388            if self.state.error_bytes != u64::try_from(error.len()).ok() {
389                return Err(EvidenceError::InvalidField("state.error_bytes"));
390            }
391            if digest != digest_bytes("a3s.code.evidence.error.v1", error.as_bytes()) {
392                return Err(EvidenceError::DigestMismatch("state.error_digest"));
393            }
394        }
395        let mut previous_fact: Option<(u64, u64)> = None;
396        for fact in &self.facts {
397            fact.validate().map_err(EvidenceError::Journal)?;
398            if fact.frame.target != self.target {
399                return Err(EvidenceError::TargetMismatch);
400            }
401            if fact.observed_at_ms > self.observed_at_ms {
402                return Err(EvidenceError::InvalidField("facts.observed_at_ms"));
403            }
404            if self
405                .after_sequence
406                .is_some_and(|cursor| fact.sequence <= cursor)
407                || fact.sequence >= self.latest_sequence_exclusive
408                || previous_fact.is_some_and(|(sequence, timestamp)| {
409                    fact.sequence != sequence.saturating_add(1) || fact.observed_at_ms < timestamp
410                })
411            {
412                return Err(EvidenceError::InvalidField("facts.sequence"));
413            }
414            previous_fact = Some((fact.sequence, fact.observed_at_ms));
415        }
416        let mut previous_event: Option<(u64, u64)> = None;
417        for (index, event) in self.events.iter().enumerate() {
418            validate_digest(&event.payload_digest)
419                .map_err(|_| EvidenceError::InvalidDigest("payload_digest"))?;
420            let payload = serde_json::to_vec(&event.event).map_err(EvidenceError::Serialization)?;
421            if event.payload_bytes == 0 || event.payload_bytes > MAX_EVENT_PAYLOAD_BYTES {
422                return Err(EvidenceError::InvalidField("payload_bytes"));
423            }
424            if event.event.version != 1
425                || event.event.event_type.is_empty()
426                || event.event.event_type.len() > MAX_EVENT_TYPE_BYTES
427                || event.event.event_type.contains('\0')
428                || event.event.event_type.lines().count() != 1
429                || payload.is_empty()
430            {
431                return Err(EvidenceError::InvalidField("event"));
432            }
433            validate_event_metadata(
434                &event.event,
435                &self.target,
436                event.sequence,
437                event.occurred_at_ms,
438            )?;
439            if event.occurred_at_ms > self.observed_at_ms {
440                return Err(EvidenceError::InvalidField("events.occurred_at_ms"));
441            }
442            if self
443                .after_sequence
444                .is_some_and(|cursor| event.sequence <= cursor)
445                || event.sequence >= self.latest_sequence_exclusive
446                || previous_event.is_some_and(|(sequence, timestamp)| {
447                    event.sequence != sequence.saturating_add(1) || event.occurred_at_ms < timestamp
448                })
449            {
450                return Err(EvidenceError::InvalidField("events.sequence"));
451            }
452            let redacted = is_redacted_payload(
453                &event.event.payload,
454                &event.payload_digest,
455                event.payload_bytes,
456            );
457            if redacted {
458                let marker = event
459                    .event
460                    .payload
461                    .as_object()
462                    .ok_or(EvidenceError::InvalidField("event.payload"))?;
463                let marker_digest = marker
464                    .get("digest")
465                    .and_then(Value::as_str)
466                    .ok_or(EvidenceError::InvalidField("event.payload"))?;
467                let marker_bytes = marker
468                    .get("bytes")
469                    .and_then(Value::as_u64)
470                    .ok_or(EvidenceError::InvalidField("event.payload"))?;
471                if marker_digest != event.payload_digest || marker_bytes != event.payload_bytes {
472                    return Err(EvidenceError::DigestMismatch("payload_digest"));
473                }
474            } else if event.payload_bytes != u64::try_from(payload.len()).unwrap_or(u64::MAX)
475                || event.payload_digest
476                    != digest_bytes("a3s.code.evidence.event-payload.v1", &payload)
477            {
478                return Err(EvidenceError::DigestMismatch("payload_digest"));
479            }
480            if let Some(fact) = self.facts.get(index) {
481                // A fact and an event at the same cursor must describe the
482                // same source observation. A snapshot may be marked
483                // incomplete while two independently captured stores are
484                // converging, but it must never claim that mismatched data is
485                // a complete evidence window.
486                let pair_matches = fact.sequence == event.sequence
487                    && fact.event_type == event.event.event_type
488                    && fact.observed_at_ms == event.occurred_at_ms
489                    && fact_payload_matches_event(fact, event);
490                if !pair_matches && self.complete {
491                    return Err(EvidenceError::InvalidField("facts.events"));
492                }
493            }
494            previous_event = Some((event.sequence, event.occurred_at_ms));
495        }
496        if self.complete {
497            let event_sequences = self
498                .events
499                .iter()
500                .map(|event| event.sequence)
501                .collect::<Vec<_>>();
502            let fact_sequences = self
503                .facts
504                .iter()
505                .map(|fact| fact.sequence)
506                .collect::<Vec<_>>();
507            if event_sequences != fact_sequences {
508                return Err(EvidenceError::InvalidField("facts.events"));
509            }
510        }
511        for artifact in &self.artifacts {
512            if artifact.artifact_uri.is_empty()
513                || artifact.artifact_uri.len() > MAX_ARTIFACT_URI_BYTES
514                || artifact.artifact_uri.contains('\0')
515                || artifact.artifact_uri.lines().count() != 1
516                || artifact.tool_name.is_empty()
517                || artifact.tool_name.len() > MAX_TOOL_NAME_BYTES
518                || artifact.tool_name.contains('\0')
519                || artifact.tool_name.lines().count() != 1
520            {
521                return Err(EvidenceError::InvalidField("artifact"));
522            }
523            validate_digest(&artifact.content_digest)
524                .map_err(|_| EvidenceError::InvalidDigest("content_digest"))?;
525            if artifact.content.as_ref().is_some_and(|content| {
526                u64::try_from(content.len()).ok() != Some(artifact.content_bytes)
527            }) {
528                return Err(EvidenceError::InvalidField("artifact.content"));
529            }
530            if let Some(content) = artifact.content.as_deref() {
531                if artifact.content_digest
532                    != digest_bytes("a3s.code.evidence.artifact-content.v1", content.as_bytes())
533                {
534                    return Err(EvidenceError::DigestMismatch("artifact.content_digest"));
535                }
536            }
537        }
538        if self
539            .artifacts
540            .windows(2)
541            .any(|window| window[0].artifact_uri >= window[1].artifact_uri)
542        {
543            return Err(EvidenceError::InvalidField("artifacts"));
544        }
545        if self.retention_gap && self.complete {
546            return Err(EvidenceError::InvalidField("complete"));
547        }
548        if let Some(first) = self.events.first() {
549            if (self.retention_gap && self.first_available_sequence != Some(first.sequence))
550                || (!self.retention_gap
551                    && requested_start < self.latest_sequence_exclusive
552                    && first.sequence != requested_start)
553                || self
554                    .first_available_sequence
555                    .is_some_and(|available| first.sequence < available)
556            {
557                return Err(EvidenceError::InvalidField("events"));
558            }
559        } else if self.retention_gap && self.first_available_sequence.is_some() {
560            return Err(EvidenceError::InvalidField("events"));
561        }
562        let last_sequence = self.events.last().map(|event| event.sequence);
563        if requested_start < self.latest_sequence_exclusive {
564            let page_complete = last_sequence.is_some_and(|sequence| {
565                sequence.saturating_add(1) == self.latest_sequence_exclusive
566            });
567            if self.complete && !page_complete {
568                return Err(EvidenceError::InvalidField("complete"));
569            }
570            if self.events.is_empty() && self.complete {
571                return Err(EvidenceError::InvalidField("complete"));
572            }
573        }
574        if self.snapshot_digest != self.expected_digest()? {
575            return Err(EvidenceError::DigestMismatch("snapshot_digest"));
576        }
577        Ok(())
578    }
579
580    fn expected_digest(&self) -> Result<String, EvidenceError> {
581        #[derive(Serialize)]
582        struct Identity<'a> {
583            schema: &'a str,
584            target: &'a ExecutionTargetV1,
585            observed_at_ms: u64,
586            after_sequence: Option<u64>,
587            first_available_sequence: Option<u64>,
588            latest_sequence_exclusive: u64,
589            state: &'a EvidenceRunStateV1,
590            facts: &'a [ExecutionFactV1],
591            events: &'a [EvidenceEventV1],
592            artifacts: &'a [EvidenceArtifactV1],
593            complete: bool,
594            retention_gap: bool,
595        }
596        digest_json(
597            "a3s.code.evidence-snapshot.identity.v1",
598            &Identity {
599                schema: &self.schema,
600                target: &self.target,
601                observed_at_ms: self.observed_at_ms,
602                after_sequence: self.after_sequence,
603                first_available_sequence: self.first_available_sequence,
604                latest_sequence_exclusive: self.latest_sequence_exclusive,
605                state: &self.state,
606                facts: &self.facts,
607                events: &self.events,
608                artifacts: &self.artifacts,
609                complete: self.complete,
610                retention_gap: self.retention_gap,
611            },
612        )
613        .map_err(EvidenceError::Serialization)
614    }
615}
616
617#[derive(Debug, Error)]
618pub enum EvidenceError {
619    #[error("evidence target is invalid or unknown")]
620    InvalidTarget,
621    #[error("evidence target does not match its source")]
622    TargetMismatch,
623    #[error("evidence schema is unsupported")]
624    UnsupportedSchema,
625    #[error("evidence field `{0}` is invalid")]
626    InvalidField(&'static str),
627    #[error("evidence digest for `{0}` is invalid")]
628    InvalidDigest(&'static str),
629    #[error("evidence digest for `{0}` does not match")]
630    DigestMismatch(&'static str),
631    #[error("evidence limit is invalid")]
632    InvalidLimit,
633    #[error("evidence numeric value does not fit the wire type")]
634    NumericOverflow,
635    #[error("run was not found")]
636    RunNotFound,
637    #[error("evidence serialization failed: {0}")]
638    Serialization(#[from] serde_json::Error),
639    #[error("execution fact error: {0}")]
640    Journal(#[from] super::journal::JournalError),
641}
642
643#[async_trait]
644pub trait EvidenceReader: Send + Sync {
645    async fn read(
646        &self,
647        request: EvidenceReadRequestV1,
648    ) -> Result<EvidenceSnapshotV1, EvidenceError>;
649}
650
651/// Reader over the native in-memory run journal and optional artifact/fact
652/// stores.  It captures state and the event window from one RunStore
653/// observation generation.
654#[derive(Clone)]
655pub struct RunEvidenceReader {
656    runs: Arc<InMemoryRunStore>,
657    facts: Option<Arc<dyn ExecutionFactJournal>>,
658    artifacts: Option<ArtifactStore>,
659}
660
661impl std::fmt::Debug for RunEvidenceReader {
662    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
663        formatter
664            .debug_struct("RunEvidenceReader")
665            .field("facts_bound", &self.facts.is_some())
666            .field("artifacts_bound", &self.artifacts.is_some())
667            .finish()
668    }
669}
670
671impl RunEvidenceReader {
672    pub fn new(runs: Arc<InMemoryRunStore>) -> Self {
673        Self {
674            runs,
675            facts: None,
676            artifacts: None,
677        }
678    }
679
680    pub fn with_facts(mut self, facts: Arc<dyn ExecutionFactJournal>) -> Self {
681        self.facts = Some(facts);
682        self
683    }
684
685    pub fn with_artifacts(mut self, artifacts: ArtifactStore) -> Self {
686        self.artifacts = Some(artifacts);
687        self
688    }
689
690    pub async fn read(
691        &self,
692        request: EvidenceReadRequestV1,
693    ) -> Result<EvidenceSnapshotV1, EvidenceError> {
694        self.read_inner(request).await
695    }
696
697    async fn read_inner(
698        &self,
699        request: EvidenceReadRequestV1,
700    ) -> Result<EvidenceSnapshotV1, EvidenceError> {
701        request.validate()?;
702        let after_sequence = request
703            .after_sequence
704            .map(|sequence| usize::try_from(sequence).map_err(|_| EvidenceError::NumericOverflow))
705            .transpose()?;
706        let observation = self
707            .runs
708            .event_observation(
709                &request.target.run_id,
710                after_sequence,
711                request.limits.max_events,
712            )
713            .await
714            .ok_or(EvidenceError::RunNotFound)?;
715        if observation.snapshot.session_id != request.target.session_id {
716            return Err(EvidenceError::TargetMismatch);
717        }
718
719        let state_content_truncated = state_content_truncated(
720            &observation.snapshot,
721            request.limits,
722            request.include_prompt,
723            request.include_terminal_text,
724        );
725        let state = EvidenceRunStateV1::from_snapshot(
726            &observation.snapshot,
727            &request.target,
728            request.limits,
729            request.include_prompt,
730            request.include_terminal_text,
731        )?;
732        let frame = ExecutionFrameV1::root(request.target.clone());
733        let observation_first_available = observation
734            .page
735            .first_available_sequence
736            .map(|sequence| u64::try_from(sequence).map_err(|_| EvidenceError::NumericOverflow))
737            .transpose()?;
738        let observation_latest_sequence_exclusive =
739            u64::try_from(observation.page.latest_sequence_exclusive)
740                .map_err(|_| EvidenceError::NumericOverflow)?;
741        let mut events = Vec::new();
742        let mut event_bytes = 0usize;
743        // A bounded page with more source events is not a complete evidence
744        // window, even when no retention gap exists. Callers can request the
745        // next page with the last returned sequence.
746        let mut complete = !observation.page.retention_gap
747            && !observation.page.has_more
748            && !state_content_truncated;
749        let mut referenced_artifacts = BTreeSet::new();
750        let mut references_truncated = false;
751
752        for record in &observation.page.events {
753            let sequence =
754                u64::try_from(record.sequence).map_err(|_| EvidenceError::NumericOverflow)?;
755            let envelope =
756                run_event_envelope_v1(record, &request.target.run_id, &request.target.session_id)
757                    .map_err(|_| EvidenceError::InvalidField("event"))?;
758            let encoded = serde_json::to_vec(&envelope).map_err(EvidenceError::Serialization)?;
759            let payload_digest = digest_bytes("a3s.code.evidence.event-payload.v1", &encoded);
760            let payload_bytes =
761                u64::try_from(encoded.len()).map_err(|_| EvidenceError::NumericOverflow)?;
762            let mut projected = envelope.clone();
763            let over_budget = encoded.len() > request.limits.max_event_bytes
764                || event_bytes.saturating_add(encoded.len()) > request.limits.max_event_bytes;
765            if request.content_mode == EvidenceContentModeV1::DigestOnly || over_budget {
766                projected.payload = serde_json::json!({
767                    "digest": payload_digest,
768                    "bytes": payload_bytes,
769                    "content": "redacted"
770                });
771                // Digest-only is an intentional, complete representation of
772                // the event window.  Only an omitted payload caused by a
773                // bounded-payload budget makes the source incomplete.
774                if request.content_mode == EvidenceContentModeV1::BoundedPayload {
775                    complete &= !over_budget;
776                }
777            } else {
778                event_bytes = event_bytes.saturating_add(encoded.len());
779            }
780            references_truncated |=
781                collect_refs_from_value(&envelope.payload, &mut referenced_artifacts);
782            events.push(EvidenceEventV1 {
783                sequence,
784                occurred_at_ms: record.timestamp_ms,
785                event: projected,
786                payload_digest,
787                payload_bytes,
788            });
789        }
790
791        let facts = if let Some(journal) = &self.facts {
792            let page = journal
793                .page(
794                    &request.target,
795                    request.after_sequence,
796                    request.limits.max_events,
797                )
798                .ok_or(EvidenceError::RunNotFound)?;
799            complete &= !page.retention_gap && !page.has_more;
800            if page.latest_sequence_exclusive != observation_latest_sequence_exclusive {
801                complete = false;
802            }
803            page.facts
804        } else {
805            observation
806                .page
807                .events
808                .iter()
809                .map(|record| {
810                    let input = ExecutionFactInputV1::from_run_event(frame.clone(), record)?;
811                    ExecutionFactV1::from_input(input).map_err(EvidenceError::Journal)
812                })
813                .collect::<Result<Vec<_>, EvidenceError>>()?
814        };
815        let event_sequences = events
816            .iter()
817            .map(|event| event.sequence)
818            .collect::<Vec<_>>();
819        let fact_sequences = facts.iter().map(|fact| fact.sequence).collect::<Vec<_>>();
820        if event_sequences != fact_sequences {
821            // The run store and an independently durable fact journal do not
822            // share a transaction. Never claim a complete evidence window
823            // when their generations disagree.
824            complete = false;
825        }
826        if events.len() != facts.len()
827            || events.iter().zip(&facts).any(|(event, fact)| {
828                event.sequence != fact.sequence
829                    || event.event.event_type != fact.event_type
830                    || event.occurred_at_ms != fact.observed_at_ms
831                    || !fact_payload_matches_event(fact, event)
832            })
833        {
834            // The fact journal and the run store are separate observations.
835            // Preserve the bounded data for diagnostics, but make the
836            // generation unusable for a gate until the host reconciles it.
837            complete = false;
838        }
839        if references_truncated {
840            complete = false;
841        }
842
843        let artifacts = self.read_artifacts(
844            &referenced_artifacts,
845            request.limits,
846            request.include_artifact_content,
847            &mut complete,
848        )?;
849        let observed_at_ms = observation.snapshot.updated_at_ms;
850        let mut snapshot = EvidenceSnapshotV1 {
851            schema: EVIDENCE_SNAPSHOT_SCHEMA_V1.to_string(),
852            target: request.target,
853            observed_at_ms,
854            after_sequence: request.after_sequence,
855            first_available_sequence: observation_first_available,
856            latest_sequence_exclusive: observation_latest_sequence_exclusive,
857            state,
858            facts,
859            events,
860            artifacts,
861            complete,
862            retention_gap: observation.page.retention_gap,
863            snapshot_digest: String::new(),
864        };
865        snapshot.snapshot_digest = snapshot.expected_digest()?;
866        snapshot.validate()?;
867        Ok(snapshot)
868    }
869
870    fn read_artifacts(
871        &self,
872        references: &BTreeSet<String>,
873        limits: EvidenceLimitsV1,
874        include_content: bool,
875        complete: &mut bool,
876    ) -> Result<Vec<EvidenceArtifactV1>, EvidenceError> {
877        let Some(store) = &self.artifacts else {
878            if !references.is_empty() {
879                *complete = false;
880            }
881            return Ok(Vec::new());
882        };
883        let mut artifacts = Vec::new();
884        let mut total_bytes = 0usize;
885        // ArtifactStore preserves insertion order, while the evidence wire
886        // contract is canonical URI order. Sort before applying byte/count
887        // limits so the selected projection is deterministic as well as
888        // validatable.
889        let mut stored_artifacts = store.artifacts();
890        stored_artifacts.sort_by(|left, right| left.artifact_uri.cmp(&right.artifact_uri));
891        for artifact in stored_artifacts {
892            if !references.contains(&artifact.artifact_uri) {
893                continue;
894            }
895            if artifacts.len() >= limits.max_artifacts {
896                *complete = false;
897                break;
898            }
899            let content_bytes = artifact.content.len();
900            let digest = digest_bytes(
901                "a3s.code.evidence.artifact-content.v1",
902                artifact.content.as_bytes(),
903            );
904            let content = if include_content
905                && content_bytes <= limits.max_artifact_bytes
906                && total_bytes.saturating_add(content_bytes) <= limits.max_artifact_bytes
907            {
908                total_bytes = total_bytes.saturating_add(content_bytes);
909                Some(artifact.content.clone())
910            } else {
911                if include_content {
912                    *complete = false;
913                }
914                None
915            };
916            artifacts.push(EvidenceArtifactV1 {
917                artifact_uri: artifact.artifact_uri,
918                tool_name: artifact.tool_name,
919                content_digest: digest,
920                content_bytes: u64::try_from(content_bytes)
921                    .map_err(|_| EvidenceError::NumericOverflow)?,
922                content,
923            });
924        }
925        if artifacts.len() < references.len() {
926            *complete = false;
927        }
928        Ok(artifacts)
929    }
930}
931
932#[async_trait]
933impl EvidenceReader for RunEvidenceReader {
934    async fn read(
935        &self,
936        request: EvidenceReadRequestV1,
937    ) -> Result<EvidenceSnapshotV1, EvidenceError> {
938        self.read(request).await
939    }
940}
941
942fn bounded_string(value: &str, limit: usize) -> Option<String> {
943    (value.len() <= limit).then(|| value.to_string())
944}
945
946fn state_content_truncated(
947    snapshot: &RunSnapshot,
948    limits: EvidenceLimitsV1,
949    include_prompt: bool,
950    include_terminal_text: bool,
951) -> bool {
952    (include_prompt && snapshot.prompt.len() > limits.max_prompt_bytes)
953        || (include_terminal_text
954            && (snapshot
955                .result_text
956                .as_ref()
957                .is_some_and(|value| value.len() > limits.max_result_bytes)
958                || snapshot
959                    .error
960                    .as_ref()
961                    .is_some_and(|value| value.len() > limits.max_result_bytes)))
962}
963
964fn is_redacted_payload(value: &Value, digest: &str, bytes: u64) -> bool {
965    let Some(object) = value.as_object() else {
966        return false;
967    };
968    object.len() == 3
969        && object.get("content") == Some(&Value::String("redacted".to_string()))
970        && object.get("digest").and_then(Value::as_str) == Some(digest)
971        && object.get("bytes").and_then(Value::as_u64) == Some(bytes)
972}
973
974fn validate_event_metadata(
975    event: &EventEnvelopeV1,
976    target: &ExecutionTargetV1,
977    sequence: u64,
978    occurred_at_ms: u64,
979) -> Result<(), EvidenceError> {
980    let metadata = event
981        .metadata
982        .as_ref()
983        .ok_or(EvidenceError::InvalidField("event.metadata"))?;
984    let encoded = serde_json::to_vec(metadata).map_err(EvidenceError::Serialization)?;
985    if encoded.len() > MAX_EVENT_METADATA_BYTES {
986        return Err(EvidenceError::InvalidLimit);
987    }
988    let object = metadata
989        .as_object()
990        .ok_or(EvidenceError::InvalidField("event.metadata"))?;
991    let exact = object.get("run_id").and_then(Value::as_str) == Some(target.run_id.as_str())
992        && object.get("session_id").and_then(Value::as_str) == Some(target.session_id.as_str())
993        && object.get("sequence").and_then(Value::as_u64) == Some(sequence)
994        && object.get("timestamp_ms").and_then(Value::as_u64) == Some(occurred_at_ms);
995    if !exact {
996        return Err(EvidenceError::TargetMismatch);
997    }
998    Ok(())
999}
1000
1001/// Compare an unredacted event projection with the digest-only fact that was
1002/// captured from the same runtime event. The fact and evidence layers use
1003/// different wire domains, so rebuild the original AgentEvent JSON shape by
1004/// restoring its top-level `type` field. Redacted payload markers deliberately
1005/// skip this check: their digest is the retained source commitment, not the
1006/// marker's digest.
1007fn fact_payload_matches_event(fact: &ExecutionFactV1, event: &EvidenceEventV1) -> bool {
1008    if is_redacted_payload(
1009        &event.event.payload,
1010        &event.payload_digest,
1011        event.payload_bytes,
1012    ) {
1013        return true;
1014    }
1015    let Some(mut object) = event.event.payload.as_object().cloned() else {
1016        return false;
1017    };
1018    object.insert(
1019        "type".to_string(),
1020        Value::String(event.event.event_type.clone()),
1021    );
1022    // Deserialize back through the runtime enum before serializing. This
1023    // restores the exact variant field order used when the fact digest was
1024    // captured; serializing the intermediate JSON map directly can reorder
1025    // keys and produce a different byte digest for the same event.
1026    let Ok(runtime_event) = serde_json::from_value::<AgentEvent>(Value::Object(object)) else {
1027        return false;
1028    };
1029    let Ok(encoded) = serde_json::to_vec(&runtime_event) else {
1030        return false;
1031    };
1032    u64::try_from(encoded.len()).ok() == Some(fact.payload_bytes)
1033        && digest_bytes("a3s.code.execution-fact.payload.v1", &encoded) == fact.payload_digest
1034}
1035
1036fn collect_refs_from_value(value: &Value, refs: &mut BTreeSet<String>) -> bool {
1037    let mut truncated = false;
1038    collect_refs_from_value_inner(value, refs, &mut truncated);
1039    truncated
1040}
1041
1042fn collect_refs_from_value_inner(value: &Value, refs: &mut BTreeSet<String>, truncated: &mut bool) {
1043    match value {
1044        Value::Object(object) => {
1045            for key in ["artifact_uri", "content_ref", "content_uri"] {
1046                if let Some(uri) = object.get(key).and_then(Value::as_str) {
1047                    if uri.is_empty() || uri.len() > MAX_ARTIFACT_URI_BYTES {
1048                        *truncated = true;
1049                    } else if refs.len() < MAX_REFERENCED_ARTIFACTS {
1050                        refs.insert(uri.to_string());
1051                    } else if !refs.contains(uri) {
1052                        *truncated = true;
1053                    }
1054                }
1055            }
1056            for child in object.values() {
1057                collect_refs_from_value_inner(child, refs, truncated);
1058            }
1059        }
1060        Value::Array(items) => {
1061            for child in items {
1062                collect_refs_from_value_inner(child, refs, truncated);
1063            }
1064        }
1065        _ => {}
1066    }
1067}
1068
1069#[cfg(test)]
1070mod tests;