use std::collections::{BTreeMap, BTreeSet};
use std::io::BufRead;
use serde::Serialize;
use crate::core::{
ACTIONS, Delegation, Digest, EffectKey, GroupOutcome, PolicyBundleIdentity, PolicyDecision,
PolicyEngine, RunId, StepId,
};
use crate::journal::{AgentIdentity, RecordBody, RecordKind, payload};
use super::requests::{self, Acting, GatedRequest};
const WAIT_KINDS: &[&str] = &["timer.sleep", "event.await"];
pub const OUTSIDE_EXPORT: &str = "a refused admission leaves no journal, and the served \
surfaces' gates journal neither their roles nor their callers; \
none of these is in an export, and none is evaluated";
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum TenantSource {
Supplied,
Default,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum Unevaluable {
Sealed,
Erased,
RequestNotJournaled,
GateSkipped,
GateIndistinguishable,
AgentNotJournaled,
ChainUnreadable,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub struct NotEvaluable {
#[serde(skip_serializing_if = "Option::is_none")]
pub step: Option<StepId>,
#[serde(skip_serializing_if = "Option::is_none")]
pub effect_key: Option<EffectKey>,
pub action: String,
pub resource: String,
pub reason: Unevaluable,
}
#[derive(Debug, Clone, PartialEq)]
pub struct Rebuilt {
pub step: Option<StepId>,
pub effect_key: Option<EffectKey>,
pub request: GatedRequest,
}
#[derive(Debug, Clone, PartialEq)]
pub struct RunRequests {
pub run: RunId,
pub recorded_bundle: Option<PolicyBundleIdentity>,
pub requests: Vec<Rebuilt>,
pub not_evaluable: Vec<NotEvaluable>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum Mode {
Recorded,
Mismatch,
Ungoverned,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub struct Finding {
#[serde(skip_serializing_if = "Option::is_none")]
pub step: Option<StepId>,
#[serde(skip_serializing_if = "Option::is_none")]
pub effect_key: Option<EffectKey>,
pub action: String,
pub resource: String,
pub reason: String,
pub malformed: bool,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize)]
pub struct Diff {
pub newly_denied: Vec<Finding>,
pub malformed_under_candidate: Vec<Finding>,
}
impl Diff {
fn is_empty(&self) -> bool {
self.newly_denied.is_empty() && self.malformed_under_candidate.is_empty()
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub struct RunReport {
pub run: RunId,
pub mode: Mode,
#[serde(skip_serializing_if = "Option::is_none")]
pub recorded_bundle: Option<Digest>,
pub evaluated: usize,
pub findings: Vec<Finding>,
#[serde(skip_serializing_if = "Option::is_none")]
pub diff: Option<Diff>,
pub not_evaluable: Vec<NotEvaluable>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub struct Tenant {
pub value: String,
pub source: TenantSource,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub struct Report {
pub bundle: Digest,
#[serde(skip_serializing_if = "Option::is_none")]
pub candidate: Option<Digest>,
pub tenant: Tenant,
pub runs: Vec<RunReport>,
pub unreadable: Vec<RunId>,
pub outside_export: &'static str,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Verdict {
Clean,
Findings,
Partial,
}
impl Report {
#[must_use]
pub fn evaluated(&self) -> usize {
self.runs.iter().map(|r| r.evaluated).sum()
}
#[must_use]
pub fn verdict(&self) -> Verdict {
let disagrees = self.runs.iter().any(|r| {
r.mode == Mode::Mismatch
|| !r.findings.is_empty()
|| r.diff.as_ref().is_some_and(|d| !d.is_empty())
});
if disagrees {
return Verdict::Findings;
}
if self.evaluated() == 0 {
return Verdict::Partial;
}
Verdict::Clean
}
}
#[derive(Debug, thiserror::Error)]
pub enum CheckError {
#[error("reading the export failed: {0}")]
Io(#[from] std::io::Error),
#[error("{0}")]
NotAnExport(String),
#[error("opening a sealed payload failed: {0}")]
Keys(String),
}
#[derive(Debug, Clone, Copy)]
pub struct Check<'a> {
tenant: &'a str,
source: TenantSource,
#[cfg(feature = "keyring")]
keys: Option<&'a dyn crate::keyring::KeyRing>,
}
impl<'a> Check<'a> {
#[must_use]
pub const fn new(tenant: &'a str, source: TenantSource) -> Self {
Self {
tenant,
source,
#[cfg(feature = "keyring")]
keys: None,
}
}
#[cfg(feature = "keyring")]
#[must_use]
pub const fn with_keys(mut self, keys: &'a dyn crate::keyring::KeyRing) -> Self {
self.keys = Some(keys);
self
}
pub async fn run<R: BufRead>(
&self,
input: R,
bundle: &dyn PolicyEngine,
candidate: Option<&dyn PolicyEngine>,
) -> Result<Report, CheckError> {
let (rebuilt, unreadable) = self.read(input).await?;
let digest = bundle.digest();
let runs = rebuilt
.into_iter()
.map(|run| judge(run, bundle, digest, candidate))
.collect();
Ok(Report {
bundle: digest,
candidate: candidate.map(PolicyEngine::digest),
tenant: Tenant {
value: self.tenant.to_owned(),
source: self.source,
},
runs,
unreadable,
outside_export: OUTSIDE_EXPORT,
})
}
pub async fn rebuild<R: BufRead>(&self, input: R) -> Result<Vec<RunRequests>, CheckError> {
Ok(self.read(input).await?.0)
}
async fn read<R: BufRead>(
&self,
input: R,
) -> Result<(Vec<RunRequests>, Vec<RunId>), CheckError> {
let (runs, unreadable) = read_export(input)?;
let mut out = Vec::with_capacity(runs.len());
for (run, records) in runs {
out.push(self.rebuild_run(run, records).await?);
}
Ok((out, unreadable))
}
#[cfg_attr(
not(feature = "keyring"),
allow(clippy::unused_async, clippy::unused_async_trait_impl)
)]
async fn open(
&self,
body: &RecordBody,
kind: &mut RecordKind,
) -> Result<Option<Unevaluable>, CheckError> {
#[cfg(feature = "keyring")]
if let Some(keys) = self.keys {
let opened =
crate::keyring::open_payloads(keys, self.tenant, body.run, body.effect_key, kind)
.await
.map_err(|e| CheckError::Keys(e.to_string()))?;
if opened.erased > 0 {
return Ok(Some(Unevaluable::Erased));
}
}
let _ = body;
let sealed = payload::payloads(kind)
.into_iter()
.any(|field| match field {
payload::SealedField::Value(v) => payload::is_sealed(v),
payload::SealedField::Text(t) => payload::is_sealed_text(t),
});
Ok(sealed.then_some(Unevaluable::Sealed))
}
#[allow(clippy::too_many_lines)]
async fn rebuild_run(
&self,
run: RunId,
records: Vec<RecordBody>,
) -> Result<RunRequests, CheckError> {
let mut out = RunRequests {
run,
recorded_bundle: None,
requests: Vec::new(),
not_evaluable: Vec::new(),
};
let mut admitted: Option<(RecordBody, String, Option<AgentIdentity>)> = None;
let mut chain_links = None;
let mut skills = BTreeSet::new();
for body in &records {
match &body.kind {
RecordKind::RunAdmitted {
capability,
governed_by,
policy_bundle,
..
} if admitted.is_none() => {
out.recorded_bundle = policy_bundle.as_deref().cloned();
admitted = Some((
body.clone(),
capability.clone(),
governed_by.as_deref().cloned(),
));
}
RecordKind::IdentityBound { chain } if chain_links.is_none() => {
chain_links = Some(chain.clone());
}
RecordKind::StepStarted { skill } if body.phase.is_forward() => {
skills.insert(skill.clone());
}
_ => {}
}
}
let Some((admission, capability, governed_by)) = admitted else {
return Ok(out);
};
let Ok(chain) = chain_links.map(Delegation::rehydrate).transpose() else {
out.not_evaluable.push(NotEvaluable {
step: None,
effect_key: None,
action: crate::core::ACTION_ADMIT.to_owned(),
resource: capability,
reason: Unevaluable::ChainUnreadable,
});
return Ok(out);
};
let acting = Acting {
tenant: self.tenant,
capability: &capability,
agent: governed_by.as_ref(),
chain: chain.as_ref(),
};
let one_skill = skills.len() <= 1;
let mut seen = BTreeSet::new();
let mut kind = admission.kind.clone();
match self.open(&admission, &mut kind).await? {
Some(reason) => out.not_evaluable.push(NotEvaluable {
step: None,
effect_key: None,
action: crate::core::ACTION_ADMIT.to_owned(),
resource: capability.clone(),
reason,
}),
None => {
if let RecordKind::RunAdmitted { input, .. } = &kind {
push(
&mut out,
&mut seen,
None,
None,
requests::admission(&acting, input),
);
}
}
}
let mut groups: BTreeMap<StepId, Vec<RecordBody>> = BTreeMap::new();
let mut settled: Vec<RecordBody> = Vec::new();
for body in records {
let step = body.step;
match &body.kind {
RecordKind::GroupOpened { .. } if body.phase.is_forward() => {
if let Some(step) = step {
groups.insert(step, Vec::new());
}
}
RecordKind::GroupSettled { outcome, .. } if body.phase.is_forward() => {
let members = step.and_then(|s| groups.remove(&s)).unwrap_or_default();
if *outcome == GroupOutcome::Committed {
settled.extend(members);
} else {
for member in members {
out.not_evaluable.push(unevaluable_effect(
&member,
Unevaluable::GateIndistinguishable,
));
}
}
}
RecordKind::EffectStarted { descriptor, .. } => {
if !body.phase.is_forward() {
out.not_evaluable
.push(unevaluable_effect(&body, Unevaluable::GateSkipped));
} else if WAIT_KINDS.contains(&descriptor.kind.as_str()) {
out.not_evaluable.push(unevaluable_effect(
&body,
Unevaluable::GateIndistinguishable,
));
} else if let Some(members) = step.and_then(|s| groups.get_mut(&s)) {
members.push(body);
} else {
settled.push(body);
}
}
RecordKind::Released { .. } => settled.push(body),
RecordKind::PolicyDenied {
action, resource, ..
} if ACTIONS.contains(&action.as_str()) => {
out.not_evaluable.push(NotEvaluable {
step,
effect_key: body.effect_key,
action: action.clone(),
resource: resource.clone(),
reason: Unevaluable::RequestNotJournaled,
});
}
_ => {}
}
}
for member in groups.into_values().flatten() {
out.not_evaluable.push(unevaluable_effect(
&member,
Unevaluable::GateIndistinguishable,
));
}
for body in settled {
let Some(step) = body.step else { continue };
if !one_skill {
out.not_evaluable
.push(unevaluable_effect(&body, Unevaluable::AgentNotJournaled));
continue;
}
let mut kind = body.kind.clone();
if let Some(reason) = self.open(&body, &mut kind).await? {
out.not_evaluable.push(unevaluable_effect(&body, reason));
continue;
}
let request = match &kind {
RecordKind::EffectStarted {
descriptor,
mutates,
outbound_label,
..
} => requests::effect(
&acting,
run,
step,
&descriptor.kind,
&descriptor.args,
*mutates,
outbound_label.as_ref(),
),
RecordKind::Released { release, label, .. } => {
requests::release(&acting, run, step, release, label)
}
_ => continue,
};
push(&mut out, &mut seen, Some(step), body.effect_key, request);
}
Ok(out)
}
}
fn push(
out: &mut RunRequests,
seen: &mut BTreeSet<Vec<u8>>,
step: Option<StepId>,
effect_key: Option<EffectKey>,
request: GatedRequest,
) {
let identity = crate::core::canon::value_bytes(&serde_json::json!([
request.principal_kind.entity_type(),
request.principal,
request.action,
request.resource,
request.context,
]));
if seen.insert(identity) {
out.requests.push(Rebuilt {
step,
effect_key,
request,
});
}
}
fn unevaluable_effect(body: &RecordBody, reason: Unevaluable) -> NotEvaluable {
let (action, resource) = match &body.kind {
RecordKind::Released { .. } => (crate::core::ACTION_RELEASE, requests::RELEASE_RESOURCE),
RecordKind::EffectStarted { descriptor, .. } => {
(crate::core::ACTION_PERFORM, descriptor.kind.as_str())
}
other => (crate::core::ACTION_PERFORM, other.kind_str()),
};
NotEvaluable {
step: body.step,
effect_key: body.effect_key,
action: action.to_owned(),
resource: resource.to_owned(),
reason,
}
}
fn judge(
run: RunRequests,
bundle: &dyn PolicyEngine,
digest: Digest,
candidate: Option<&dyn PolicyEngine>,
) -> RunReport {
let recorded = run
.recorded_bundle
.as_ref()
.map(PolicyBundleIdentity::digest);
let mode = match recorded {
None => Mode::Ungoverned,
Some(recorded) if recorded == digest => Mode::Recorded,
Some(_) => Mode::Mismatch,
};
let mut report = RunReport {
run: run.run,
mode,
recorded_bundle: recorded,
evaluated: 0,
findings: Vec::new(),
diff: None,
not_evaluable: Vec::new(),
};
if mode == Mode::Ungoverned {
return report;
}
report.not_evaluable = run.not_evaluable;
if mode == Mode::Recorded {
for rebuilt in &run.requests {
report.evaluated += 1;
if let Some(finding) =
disagreement(rebuilt, &bundle.authorize(&rebuilt.request.as_request()))
{
report.findings.push(finding);
}
}
}
if let Some(candidate) = candidate {
let mut diff = Diff::default();
for rebuilt in &run.requests {
let decision = candidate.authorize(&rebuilt.request.as_request());
if let Some(finding) = disagreement(rebuilt, &decision) {
if finding.malformed {
diff.malformed_under_candidate.push(finding);
} else {
diff.newly_denied.push(finding);
}
}
}
report.diff = Some(diff);
}
report
}
fn disagreement(rebuilt: &Rebuilt, decision: &PolicyDecision) -> Option<Finding> {
let reason = decision.reason()?;
Some(Finding {
step: rebuilt.step,
effect_key: rebuilt.effect_key,
action: rebuilt.request.action.to_owned(),
resource: rebuilt.request.resource.clone(),
reason: reason.to_owned(),
malformed: decision.is_malformed(),
})
}
#[allow(clippy::type_complexity)]
fn read_export<R: BufRead>(
input: R,
) -> Result<(Vec<(RunId, Vec<RecordBody>)>, Vec<RunId>), CheckError> {
let read = crate::export::read_runs(input).map_err(|e| match e {
crate::export::ReadError::Io(e) => CheckError::Io(e),
crate::export::ReadError::NotAnExport(e) => CheckError::NotAnExport(e),
})?;
Ok((read.runs, read.unreadable))
}