1use 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 #[default]
37 DigestOnly,
38 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 #[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 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#[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 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 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 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 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 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
1001fn 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 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;