use std::collections::BTreeMap;
use serde::{Deserialize, Serialize};
use serde_json::Value;
use crate::core::{
AttestError, Attestation, CaseId, Compensation, DeadlineState, Digest, Disposition,
EffectDescriptor, EffectKey, Epoch, Label, Phase, PolicyBundleIdentity, Principal, Recovery,
RunId, Seq, Signer, Spend, StepId, StoreError, Timestamp, Verifier, canon,
};
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct AgentIdentity {
pub name: String,
pub version: String,
pub digest: Digest,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub publisher: Option<crate::core::KeyId>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(tag = "kind", rename_all = "PascalCase", deny_unknown_fields)]
#[non_exhaustive]
pub enum RecordKind {
RunAdmitted {
capability: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
governed_by: Option<Box<AgentIdentity>>,
input: Value,
input_label: crate::core::Label,
#[serde(default, skip_serializing_if = "Option::is_none")]
policy_bundle: Option<Box<PolicyBundleIdentity>>,
canon: u16,
#[serde(default, skip_serializing_if = "Option::is_none")]
idempotency_key: Option<String>,
},
QuotaPassStarted {
#[serde(default, skip_serializing_if = "Option::is_none")]
period: Option<String>,
release_slot: bool,
},
PlanFrozen {
steps: Vec<String>,
plan: Value,
},
StepStarted {
skill: String,
},
StepFinished {
outcome: String,
},
CaseBound {
case_kind: String,
opened: bool,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
correlation: Vec<crate::core::CorrelationKey>,
},
DeadlineRegistered {
name: String,
#[serde(with = "time::serde::rfc3339")]
resolved_at: Timestamp,
calendar_digest: Digest,
},
DeadlineTransition {
name: String,
from: DeadlineState,
to: DeadlineState,
},
RunSuspended {
reason: crate::core::SuspendReason,
},
Note {
text: String,
},
EffectStarted {
descriptor: EffectDescriptor,
recovery: Recovery,
mutates: bool,
attempt: u32,
backoff_ms: u64,
#[serde(default, skip_serializing_if = "Option::is_none")]
outbound_label: Option<Label>,
#[serde(default, skip_serializing_if = "Option::is_none")]
outbound_bytes: Option<u64>,
},
EffectDone {
output: Value,
#[serde(default, skip_serializing_if = "Option::is_none")]
source: Option<String>,
#[serde(default, skip_serializing_if = "Spend::is_free_ref")]
spend: Spend,
declared: crate::core::DeclaredOutput,
},
EffectFailed {
error: String,
#[serde(default, skip_serializing_if = "Spend::is_free_ref")]
spend: Spend,
disposition: Disposition,
#[serde(default, skip_serializing_if = "std::ops::Not::not")]
permanent: bool,
},
BudgetRefused {
limit: String,
used: String,
},
BudgetReadmitted {
limit: String,
},
AuthorityWithheld {
subject: String,
reason: String,
by: crate::core::Operator,
},
AuthorityRestored {
subject: String,
},
IdentityBound {
chain: Vec<Principal>,
},
PolicyDenied {
reason: String,
action: String,
resource: String,
},
GroupOpened {
group: String,
resources: Vec<String>,
},
GroupSettled {
group: String,
outcome: crate::core::GroupOutcome,
#[serde(default, skip_serializing_if = "Option::is_none")]
detail: Option<String>,
},
QuarantineDecided {
decider: crate::core::Operator,
reason: String,
decision: crate::core::QuarantineDecision,
},
StepCompensated {
compensation: Compensation,
outcome: String,
},
EffectReconciled {
disposition: Disposition,
#[serde(default, skip_serializing_if = "Option::is_none")]
output: Option<Value>,
#[serde(default, skip_serializing_if = "Spend::is_free_ref")]
spend: Spend,
#[serde(default, skip_serializing_if = "Option::is_none")]
detail: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
declared: Option<crate::core::DeclaredOutput>,
#[serde(default, skip_serializing_if = "Option::is_none")]
asserted_by: Option<crate::core::Operator>,
},
Released {
releaser: String,
release: crate::core::Release,
label: Label,
field_labels: BTreeMap<String, Label>,
value: Digest,
},
RunCancelled {
actor: crate::core::Operator,
reason: String,
},
RunConcluded {
outcome: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
reason: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
exhaustion: Option<crate::core::BudgetExceeded>,
#[serde(default, skip_serializing_if = "Spend::is_free_ref")]
live_spend: Spend,
chain_head: Digest,
},
BreakGlass {
actor: crate::core::Operator,
roles: Vec<String>,
reason: String,
},
Swept {
subject: String,
action: crate::core::SweptAction,
#[serde(default, skip_serializing_if = "Option::is_none")]
detail: Option<String>,
},
}
impl RecordKind {
#[must_use]
pub fn kind_str(&self) -> &'static str {
match self {
Self::RunAdmitted { .. } => "RunAdmitted",
Self::QuotaPassStarted { .. } => "QuotaPassStarted",
Self::PlanFrozen { .. } => "PlanFrozen",
Self::StepStarted { .. } => "StepStarted",
Self::StepFinished { .. } => "StepFinished",
Self::Note { .. } => "Note",
Self::RunSuspended { .. } => "RunSuspended",
Self::CaseBound { .. } => "CaseBound",
Self::DeadlineRegistered { .. } => "DeadlineRegistered",
Self::DeadlineTransition { .. } => "DeadlineTransition",
Self::EffectStarted { .. } => "EffectStarted",
Self::EffectDone { .. } => "EffectDone",
Self::EffectFailed { .. } => "EffectFailed",
Self::EffectReconciled { .. } => "EffectReconciled",
Self::StepCompensated { .. } => "StepCompensated",
Self::QuarantineDecided { .. } => "QuarantineDecided",
Self::GroupOpened { .. } => "GroupOpened",
Self::GroupSettled { .. } => "GroupSettled",
Self::BudgetRefused { .. } => "BudgetRefused",
Self::BudgetReadmitted { .. } => "BudgetReadmitted",
Self::AuthorityWithheld { .. } => "AuthorityWithheld",
Self::AuthorityRestored { .. } => "AuthorityRestored",
Self::IdentityBound { .. } => "IdentityBound",
Self::PolicyDenied { .. } => "PolicyDenied",
Self::Released { .. } => "Released",
Self::RunCancelled { .. } => "RunCancelled",
Self::RunConcluded { .. } => "RunConcluded",
Self::BreakGlass { .. } => "BreakGlass",
Self::Swept { .. } => "Swept",
}
}
#[must_use]
pub fn version(&self) -> u16 {
1
}
}
fn lift(
upcaster: &dyn super::Upcaster,
raw: &[u8],
kind: &str,
version: u16,
) -> Result<RecordBody, StoreError> {
let written: Value = serde_json::from_slice(raw)?;
let lifted = upcaster.upcast(kind, version, written)?;
serde_json::from_value(lifted).map_err(StoreError::Encoding)
}
fn version_claimed(raw: &[u8]) -> Option<(String, u16)> {
let value: Value = serde_json::from_slice(raw).ok()?;
let kind = value.get("kind")?.as_str()?.to_owned();
let version = u16::try_from(value.get("v")?.as_u64()?).ok()?;
Some((kind, version))
}
pub(crate) fn unreadable(raw: &[u8], parse: serde_json::Error) -> StoreError {
match version_claimed(raw) {
Some((kind, version)) => StoreError::UnreadableRecordShape {
kind,
version,
detail: parse.to_string(),
},
None => StoreError::Encoding(parse),
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct RecordBody {
pub seq: Seq,
pub run: RunId,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub case: Option<CaseId>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub step: Option<StepId>,
#[serde(default, skip_serializing_if = "Phase::is_forward_ref")]
pub phase: Phase,
pub epoch: Epoch,
pub v: u16,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub effect_key: Option<EffectKey>,
#[serde(flatten)]
pub kind: RecordKind,
}
#[derive(Debug, Clone, PartialEq)]
pub struct Record {
pub body: RecordBody,
pub prev_hash: Digest,
pub hash: Digest,
pub attestation: Option<Attestation>,
raw: Vec<u8>,
}
fn record_signing_input(hash: Digest) -> Digest {
crate::core::signing_hash(crate::core::DOMAIN_RECORD, &hash)
}
impl Record {
pub fn seal(body: RecordBody, prev_hash: Digest) -> Result<Self, StoreError> {
Self::seal_signed(body, prev_hash, None)
}
pub const MAX_RECORD_BYTES: usize = 1 << 20;
pub fn seal_signed(
body: RecordBody,
prev_hash: Digest,
signer: Option<&dyn Signer>,
) -> Result<Self, StoreError> {
let raw = canon::to_bytes(&body)?;
if raw.len() > Self::MAX_RECORD_BYTES {
return Err(StoreError::RecordTooLarge {
bytes: raw.len(),
limit: Self::MAX_RECORD_BYTES,
});
}
let hash = Digest::chain(prev_hash, &raw);
Ok(Self {
body,
prev_hash,
hash,
attestation: signer.map(|s| s.attest(&record_signing_input(hash))),
raw,
})
}
pub fn from_stored(raw: Vec<u8>, prev_hash: Digest, hash: Digest) -> Result<Self, StoreError> {
Self::from_stored_attested(raw, prev_hash, hash, None)
}
pub fn from_stored_attested(
raw: Vec<u8>,
prev_hash: Digest,
hash: Digest,
attestation: Option<Attestation>,
) -> Result<Self, StoreError> {
Self::from_stored_with(&super::Identity, raw, prev_hash, hash, attestation)
}
pub fn from_stored_with(
upcaster: &dyn super::Upcaster,
raw: Vec<u8>,
prev_hash: Digest,
hash: Digest,
attestation: Option<Attestation>,
) -> Result<Self, StoreError> {
let recomputed = Digest::chain(prev_hash, &raw);
if recomputed != hash {
let seq = serde_json::from_slice::<serde_json::Map<String, Value>>(&raw)
.ok()
.and_then(|m| m.get("seq").and_then(Value::as_u64))
.unwrap_or(0);
return Err(StoreError::Corrupt {
seq,
detail: format!(
"hash mismatch: stored {hash:?}, recomputed {recomputed:?} — \
record was altered after it was written"
),
});
}
let body = match serde_json::from_slice::<RecordBody>(&raw) {
Ok(body) if body.v == upcaster.current_version(body.kind.kind_str()) => body,
Ok(body) => lift(upcaster, &raw, body.kind.kind_str(), body.v)?,
Err(parse) => match version_claimed(&raw) {
Some((kind, v)) if v != upcaster.current_version(&kind) => {
lift(upcaster, &raw, &kind, v)?
}
_ => return Err(unreadable(&raw, parse)),
},
};
Ok(Self {
body,
prev_hash,
hash,
attestation,
raw,
})
}
#[must_use]
pub fn raw(&self) -> &[u8] {
&self.raw
}
#[must_use]
pub fn seq(&self) -> Seq {
self.body.seq
}
#[must_use]
pub fn kind(&self) -> &RecordKind {
&self.body.kind
}
#[must_use]
pub fn effect_key(&self) -> Option<EffectKey> {
self.body.effect_key
}
#[cfg(feature = "keyring")]
#[must_use]
pub(crate) fn with_opened_kind(mut self, kind: RecordKind) -> Self {
self.body.kind = kind;
self
}
pub fn verify_chain(records: &[Self], from: Digest) -> Result<Digest, StoreError> {
let mut prev = from;
let start = records.first().map_or(1, Self::seq);
for (expect_seq, r) in (start..).zip(records.iter()) {
if r.seq() != expect_seq {
return Err(StoreError::Corrupt {
seq: r.seq(),
detail: format!("sequence gap: expected {expect_seq}, found {}", r.seq()),
});
}
if r.prev_hash != prev {
return Err(StoreError::Corrupt {
seq: r.seq(),
detail: "broken link: prev_hash does not match predecessor".into(),
});
}
let recomputed = Digest::chain(prev, &r.raw);
if recomputed != r.hash {
return Err(StoreError::Corrupt {
seq: r.seq(),
detail: "hash does not cover the stored bytes".into(),
});
}
prev = r.hash;
}
Ok(prev)
}
pub fn verify_attested(
records: &[Self],
from: Digest,
verifier: &dyn Verifier,
require_signature: bool,
) -> Result<Digest, StoreError> {
let head = Self::verify_chain(records, from)?;
for r in records {
match &r.attestation {
Some(a)
if verifier.verify(&a.key_id, &record_signing_input(r.hash), &a.signature) => {}
Some(a) => {
return Err(StoreError::Corrupt {
seq: r.seq(),
detail: AttestError::BadSignature {
seq: r.seq(),
key_id: a.key_id.clone(),
}
.to_string(),
});
}
None if require_signature => {
return Err(StoreError::Corrupt {
seq: r.seq(),
detail: AttestError::Unsigned { seq: r.seq() }.to_string(),
});
}
None => {}
}
}
Ok(head)
}
}
#[derive(Debug, Clone)]
pub struct Append {
pub run: RunId,
pub case: Option<CaseId>,
pub step: Option<StepId>,
pub phase: Phase,
pub effect_key: Option<EffectKey>,
pub kind: RecordKind,
}
impl Append {
pub fn new(run: RunId, kind: RecordKind) -> Self {
Self {
run,
case: None,
step: None,
phase: Phase::Forward,
effect_key: None,
kind,
}
}
#[must_use]
pub fn step(mut self, s: StepId) -> Self {
self.step = Some(s);
self
}
#[must_use]
pub fn phase(mut self, p: Phase) -> Self {
self.phase = p;
self
}
#[must_use]
pub fn case(mut self, c: CaseId) -> Self {
self.case = Some(c);
self
}
#[must_use]
pub fn effect(mut self, k: EffectKey) -> Self {
self.effect_key = Some(k);
self
}
#[must_use]
pub fn from_body(body: RecordBody) -> Self {
Self {
run: body.run,
case: body.case,
step: body.step,
phase: body.phase,
effect_key: body.effect_key,
kind: body.kind,
}
}
#[cfg(any(feature = "redb", feature = "postgres", test))]
pub(crate) fn into_body(self, seq: Seq, epoch: Epoch) -> RecordBody {
RecordBody {
seq,
run: self.run,
case: self.case,
step: self.step,
phase: self.phase,
epoch,
v: self.kind.version(),
effect_key: self.effect_key,
kind: self.kind,
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
fn body(seq: Seq, kind: RecordKind) -> RecordBody {
RecordBody {
seq,
run: RunId::generate(),
case: None,
step: None,
phase: Phase::Forward,
epoch: 1,
v: kind.version(),
effect_key: None,
kind,
}
}
#[test]
fn a_record_larger_than_the_limit_is_refused() {
let huge = "x".repeat(Record::MAX_RECORD_BYTES + 1);
let sealed = Record::seal(
body(
1,
RecordKind::RunAdmitted {
capability: huge,
governed_by: None,
input: json!(null),
input_label: crate::core::Label::trusted(),
policy_bundle: None,
canon: crate::core::canon::VERSION,
idempotency_key: None,
},
),
Digest::ZERO,
);
match sealed {
Err(StoreError::RecordTooLarge { bytes, limit }) => {
assert!(bytes > limit, "the error must report the real overage");
assert_eq!(limit, Record::MAX_RECORD_BYTES);
}
Err(other) => panic!("refused for the wrong reason: {other}"),
Ok(r) => panic!(
"a {}-byte record was accepted into an append-only chain",
r.raw().len()
),
}
}
#[test]
fn an_ordinary_record_is_nowhere_near_the_limit() {
let r = Record::seal(
body(
1,
RecordKind::RunAdmitted {
capability: "auditor@2.0.0".into(),
governed_by: None,
input: json!({ "ticket": "printer on fire" }),
input_label: crate::core::Label::trusted(),
policy_bundle: None,
canon: crate::core::canon::VERSION,
idempotency_key: None,
},
),
Digest::ZERO,
)
.expect("an ordinary record seals");
assert!(
r.raw().len() * 100 < Record::MAX_RECORD_BYTES,
"an ordinary record is {} bytes against a {}-byte ceiling; the \
ceiling is too close to normal traffic to be a safety net",
r.raw().len(),
Record::MAX_RECORD_BYTES
);
}
fn chain_of(n: u64) -> Vec<Record> {
let mut prev = Digest::ZERO;
let mut out = Vec::new();
for i in 1..=n {
let r = Record::seal(
body(
i,
RecordKind::StepStarted {
skill: format!("s{i}"),
},
),
prev,
)
.unwrap();
prev = r.hash;
out.push(r);
}
out
}
#[test]
fn sealing_is_deterministic() {
let b = body(1, RecordKind::StepStarted { skill: "x".into() });
let a = Record::seal(b.clone(), Digest::ZERO).unwrap();
let c = Record::seal(b, Digest::ZERO).unwrap();
assert_eq!(
a.hash, c.hash,
"same body + same prev must hash identically"
);
}
#[test]
fn valid_chain_verifies() {
let records = chain_of(5);
let head = Record::verify_chain(&records, Digest::ZERO).unwrap();
assert_eq!(head, records.last().unwrap().hash);
}
#[test]
fn tampered_payload_is_detected() {
let records = chain_of(3);
let mut tampered = records;
tampered[1].raw = b"{\"seq\":2,\"tampered\":true}".to_vec();
let err = Record::verify_chain(&tampered, Digest::ZERO).unwrap_err();
assert!(
matches!(err, StoreError::Corrupt { seq: 2, .. }),
"got {err:?}"
);
}
#[test]
fn deleted_record_is_detected_as_a_gap() {
let records = chain_of(4);
let mut with_hole = records;
with_hole.remove(2);
let err = Record::verify_chain(&with_hole, Digest::ZERO).unwrap_err();
assert!(matches!(err, StoreError::Corrupt { .. }));
}
#[test]
fn reordered_records_break_the_chain() {
let mut records = chain_of(4);
records.swap(1, 2);
assert!(Record::verify_chain(&records, Digest::ZERO).is_err());
}
#[test]
fn from_stored_rejects_a_hash_that_does_not_cover_the_bytes() {
let r = chain_of(1).pop().unwrap();
let err = Record::from_stored(b"{\"seq\":1}".to_vec(), Digest::ZERO, r.hash).unwrap_err();
assert!(matches!(err, StoreError::Corrupt { .. }));
}
#[test]
fn from_stored_roundtrips_a_genuine_record() {
let r = chain_of(1).pop().unwrap();
let back = Record::from_stored(r.raw().to_vec(), r.prev_hash, r.hash).unwrap();
assert_eq!(back.body, r.body);
}
#[test]
fn effect_records_carry_their_key() {
let key = EffectKey::from_hex(&Digest::of(b"k").to_hex()).unwrap();
let a = Append::new(
RunId::generate(),
RecordKind::EffectDone {
declared: crate::core::DeclaredOutput::untrusted(),
output: json!(1),
source: None,
spend: Spend::default(),
},
)
.effect(key);
let rec = Record::seal(a.into_body(1, 1), Digest::ZERO).unwrap();
assert_eq!(rec.effect_key(), Some(key));
}
}