use std::collections::{BTreeMap, BTreeSet};
use std::io::BufRead;
use serde::Serialize;
use serde_json::Value;
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,
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 mut runs: Vec<(RunId, Vec<RecordBody>)> = Vec::new();
let mut unreadable = Vec::new();
let mut header = false;
for (index, line) in input.lines().enumerate() {
let line = line?;
if line.trim().is_empty() {
continue;
}
let value: Value = serde_json::from_str(&line)
.map_err(|e| CheckError::NotAnExport(format!("line {} is not JSON: {e}", index + 1)))?;
match value.get("kind").and_then(Value::as_str) {
Some("agentplane.export") => {
let version = value.get("version").and_then(Value::as_u64);
if version != Some(u64::from(crate::export::FORMAT_VERSION)) {
return Err(CheckError::NotAnExport(format!(
"the export is at format version {version:?}, and this build reads {}",
crate::export::FORMAT_VERSION
)));
}
header = true;
}
_ if !header => {
return Err(CheckError::NotAnExport(
"the first line is not an agentplane export header".into(),
));
}
Some("agentplane.export.run") => {
let run = value
.get("run")
.cloned()
.and_then(|r| serde_json::from_value::<RunId>(r).ok())
.ok_or_else(|| {
CheckError::NotAnExport(format!("line {} names no run", index + 1))
})?;
runs.push((run, Vec::new()));
}
Some("agentplane.export.end") => {
if let Some(list) = value.get("unreadable").and_then(Value::as_array) {
unreadable.extend(list.iter().filter_map(|u| {
u.get("run")
.cloned()
.and_then(|r| serde_json::from_value::<RunId>(r).ok())
}));
}
}
Some(_) => {}
None => {
let raw = value.get("raw").and_then(Value::as_str).ok_or_else(|| {
CheckError::NotAnExport(format!("line {} carries no wire bytes", index + 1))
})?;
let body: RecordBody = serde_json::from_str(raw).map_err(|e| {
CheckError::NotAnExport(format!(
"line {} holds a record this build does not read: {e}",
index + 1
))
})?;
let Some((run, records)) = runs.last_mut() else {
return Err(CheckError::NotAnExport(format!(
"line {} is a record before any run block",
index + 1
)));
};
if body.run != *run {
return Err(CheckError::NotAnExport(format!(
"line {} belongs to run {}, filed under {run}",
index + 1,
body.run
)));
}
records.push(body);
}
}
}
if !header {
return Err(CheckError::NotAnExport("the input is empty".into()));
}
Ok((runs, unreadable))
}