use std::collections::BTreeMap;
use std::path::{Path, PathBuf};
use std::time::Duration;
use async_trait::async_trait;
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use crate::openspec::{
VerificationCompletionRole, VerificationDeclaration, VerificationExecutionClass,
};
pub const EVIDENCE_SCHEMA: &str = "conflux-verification-evidence-v2";
pub const EVIDENCE_AUTHORITY: &str = "conflux-runtime-executor";
pub const GATES_DIR: &str = "gates";
pub const LEGACY_EVIDENCE_DIR: &str = ".cflx/verification-evidence";
pub fn legacy_target_evidence(workspace: &Path) -> Option<PathBuf> {
let path = workspace.join(LEGACY_EVIDENCE_DIR);
std::fs::symlink_metadata(&path).ok().map(|_| path)
}
pub const DEFAULT_MIN_REUSE_SECONDS: i64 = 60;
const MAX_ARTIFACT_BYTES: usize = 4 * 1024 * 1024;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct ReusePolicy {
pub min_reuse_seconds: i64,
}
impl Default for ReusePolicy {
fn default() -> Self {
Self {
min_reuse_seconds: DEFAULT_MIN_REUSE_SECONDS,
}
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct ToolIdentity {
pub path: String,
pub executable_digest: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub version: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct VerificationEvidence {
pub schema: String,
pub authority: String,
pub verification_id: String,
pub commit_oid: String,
pub tree_oid: String,
#[serde(default)]
pub review_base_commit: String,
#[serde(default)]
pub review_range: String,
pub argv: Vec<String>,
pub cwd: String,
pub automation_path: String,
pub automation_blob_oid: String,
pub tool: ToolIdentity,
pub started_at: DateTime<Utc>,
pub ended_at: DateTime<Utc>,
pub exit_code: i32,
pub artifact_path: String,
pub artifact_digest: String,
pub clean_before: bool,
pub clean_after: bool,
}
impl VerificationEvidence {
pub fn elapsed_seconds(&self) -> Option<i64> {
let elapsed = self.ended_at.signed_duration_since(self.started_at);
let seconds = elapsed.num_seconds();
(seconds >= 0).then_some(seconds)
}
pub fn to_json(&self) -> serde_json::Value {
serde_json::json!({
"verification_id": self.verification_id,
"commit_oid": self.commit_oid,
"tree_oid": self.tree_oid,
"review_base_commit": self.review_base_commit,
"review_range": self.review_range,
"argv": self.argv,
"cwd": self.cwd,
"automation_path": self.automation_path,
"automation_blob_oid": self.automation_blob_oid,
"tool_path": self.tool.path,
"tool_digest": self.tool.executable_digest,
"tool_version": self.tool.version,
"exit_code": self.exit_code,
"artifact_path": self.artifact_path,
"artifact_digest": self.artifact_digest,
"elapsed_seconds": self.elapsed_seconds(),
})
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum EvidenceDefect {
Missing,
Unreadable(String),
Malformed(String),
UnknownSchema(String),
AgentAuthored(String),
}
impl EvidenceDefect {
pub fn code(&self) -> &'static str {
match self {
Self::Missing => "evidence_missing",
Self::Unreadable(_) => "evidence_unreadable",
Self::Malformed(_) => "evidence_malformed",
Self::UnknownSchema(_) => "evidence_unknown_schema",
Self::AgentAuthored(_) => "evidence_agent_authored",
}
}
pub fn detail(&self) -> String {
match self {
Self::Missing => "no runtime-authored evidence sidecar exists".to_string(),
Self::Unreadable(detail) => format!("evidence sidecar could not be read: {detail}"),
Self::Malformed(detail) => format!("evidence sidecar is malformed: {detail}"),
Self::UnknownSchema(found) => {
format!("evidence schema '{found}' is not '{EVIDENCE_SCHEMA}'")
}
Self::AgentAuthored(detail) => {
format!("evidence sidecar is not runtime-authored: {detail}")
}
}
}
}
fn is_full_length_object_id(value: &str) -> bool {
matches!(value.len(), 40 | 64) && value.chars().all(|c| c.is_ascii_hexdigit())
}
pub fn is_storable_verification_id(id: &str) -> bool {
crate::config::defaults::is_safe_path_component(id)
}
pub fn normalize_relative_cwd(raw: &str) -> Option<String> {
let trimmed = raw.trim();
if trimmed.is_empty() || trimmed == "." || trimmed == "./" {
return Some(".".to_string());
}
if trimmed.starts_with('/') || trimmed.contains('\\') {
return None;
}
let mut segments = Vec::new();
for segment in trimmed.split('/') {
match segment {
"" | "." => continue,
".." => return None,
other => segments.push(other),
}
}
if segments.is_empty() {
Some(".".to_string())
} else {
Some(segments.join("/"))
}
}
pub fn parse_evidence(bytes: &[u8]) -> Result<VerificationEvidence, EvidenceDefect> {
let value: serde_json::Value = serde_json::from_slice(bytes)
.map_err(|error| EvidenceDefect::Malformed(error.to_string()))?;
let schema = value
.get("schema")
.and_then(serde_json::Value::as_str)
.unwrap_or_default()
.to_string();
if schema != EVIDENCE_SCHEMA {
return Err(EvidenceDefect::UnknownSchema(schema));
}
let record: VerificationEvidence = serde_json::from_value(value)
.map_err(|error| EvidenceDefect::Malformed(error.to_string()))?;
if record.authority != EVIDENCE_AUTHORITY {
return Err(EvidenceDefect::AgentAuthored(format!(
"authority '{}' is not '{EVIDENCE_AUTHORITY}'",
record.authority
)));
}
if !is_storable_verification_id(&record.verification_id) {
return Err(EvidenceDefect::Malformed(format!(
"verification id '{}' is not storable",
record.verification_id
)));
}
for (field, value) in [
("commit_oid", &record.commit_oid),
("tree_oid", &record.tree_oid),
("automation_blob_oid", &record.automation_blob_oid),
("artifact_digest", &record.artifact_digest),
("tool.executable_digest", &record.tool.executable_digest),
] {
if !is_full_length_object_id(value) {
return Err(EvidenceDefect::Malformed(format!(
"{field} '{value}' is not a full-length object id"
)));
}
}
if !record.review_base_commit.is_empty()
&& !is_full_length_object_id(&record.review_base_commit)
{
return Err(EvidenceDefect::Malformed(format!(
"review_base_commit '{}' is not a full-length object id",
record.review_base_commit
)));
}
if record.argv.is_empty() || record.argv.iter().any(|arg| arg.trim().is_empty()) {
return Err(EvidenceDefect::Malformed(
"argv must be a non-empty array of non-empty arguments".to_string(),
));
}
if normalize_relative_cwd(&record.cwd).as_deref() != Some(record.cwd.as_str()) {
return Err(EvidenceDefect::Malformed(format!(
"cwd '{}' is not a normalized repository-relative path",
record.cwd
)));
}
if record.automation_path.trim().is_empty() || record.tool.path.trim().is_empty() {
return Err(EvidenceDefect::Malformed(
"automation path and tool path are required".to_string(),
));
}
if record.artifact_path != EvidenceStore::artifact_relative_path(&record.verification_id) {
return Err(EvidenceDefect::Malformed(format!(
"artifact path '{}' is not the store-relative '{}'",
record.artifact_path,
EvidenceStore::artifact_relative_path(&record.verification_id)
)));
}
if record.elapsed_seconds().is_none() {
return Err(EvidenceDefect::Malformed(
"ended_at precedes started_at".to_string(),
));
}
if record.exit_code != 0 {
return Err(EvidenceDefect::Malformed(format!(
"exit code {} is not success",
record.exit_code
)));
}
if !record.clean_before || !record.clean_after {
return Err(EvidenceDefect::Malformed(
"capture recorded a dirty index or worktree".to_string(),
));
}
Ok(record)
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum IneligibleReason {
MissingId,
UnstorableId(String),
NotRepositoryLocal,
NotChangeBlocking,
MissingCommand,
MissingAutomation,
ShellSyntax(String),
}
impl IneligibleReason {
pub fn code(&self) -> &'static str {
match self {
Self::MissingId => "declaration_missing_id",
Self::UnstorableId(_) => "declaration_unstorable_id",
Self::NotRepositoryLocal => "declaration_not_repository_local",
Self::NotChangeBlocking => "declaration_not_change_blocking",
Self::MissingCommand => "declaration_missing_command",
Self::MissingAutomation => "declaration_missing_automation",
Self::ShellSyntax(_) => "declaration_shell_syntax",
}
}
pub fn detail(&self) -> String {
match self {
Self::MissingId => "verification declaration has no id".to_string(),
Self::UnstorableId(id) => format!("verification id '{id}' is not storable"),
Self::NotRepositoryLocal => {
"only execution_class: repository-local evidence is reusable".to_string()
}
Self::NotChangeBlocking => {
"only completion_role: change-blocking evidence is reusable".to_string()
}
Self::MissingCommand => "declaration has no evidence/rerun command".to_string(),
Self::MissingAutomation => "declaration has no automation path".to_string(),
Self::ShellSyntax(command) => {
format!("command '{command}' is not a plain argv the runtime can supervise")
}
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct VerificationRequest {
pub verification_id: String,
pub argv: Vec<String>,
pub cwd: String,
pub automation_path: String,
}
const SHELL_METACHARACTERS: [char; 9] = ['|', '&', ';', '<', '>', '(', ')', '`', '$'];
pub fn eligible_request(
declaration: &VerificationDeclaration,
) -> Result<VerificationRequest, IneligibleReason> {
let id = declaration
.id
.as_deref()
.map(str::trim)
.filter(|id| !id.is_empty())
.ok_or(IneligibleReason::MissingId)?;
if !is_storable_verification_id(id) {
return Err(IneligibleReason::UnstorableId(id.to_string()));
}
let execution_class = declaration
.execution_class
.as_deref()
.and_then(VerificationExecutionClass::parse);
if !execution_class.is_some_and(VerificationExecutionClass::is_repository_local) {
return Err(IneligibleReason::NotRepositoryLocal);
}
let completion_role = declaration
.completion_role
.as_deref()
.and_then(VerificationCompletionRole::parse);
if completion_role != Some(VerificationCompletionRole::ChangeBlocking) {
return Err(IneligibleReason::NotChangeBlocking);
}
let command = declaration
.rerun
.as_deref()
.or(declaration.evidence.as_deref())
.map(str::trim)
.filter(|command| !command.is_empty())
.ok_or(IneligibleReason::MissingCommand)?;
if command.contains(SHELL_METACHARACTERS) {
return Err(IneligibleReason::ShellSyntax(command.to_string()));
}
let argv = shlex::split(command)
.filter(|argv| !argv.is_empty() && argv.iter().all(|arg| !arg.trim().is_empty()))
.ok_or_else(|| IneligibleReason::ShellSyntax(command.to_string()))?;
let automation_path = declaration
.automation
.as_deref()
.map(str::trim)
.filter(|path| !path.is_empty())
.ok_or(IneligibleReason::MissingAutomation)?;
if normalize_relative_cwd(automation_path).is_none() {
return Err(IneligibleReason::MissingAutomation);
}
Ok(VerificationRequest {
verification_id: id.to_string(),
argv,
cwd: ".".to_string(),
automation_path: automation_path.to_string(),
})
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct CurrentBindings {
pub commit_oid: String,
pub tree_oid: String,
pub review: Option<ReviewBinding>,
pub automation_blob_oid: String,
pub tool: ToolIdentity,
pub clean: bool,
pub artifact_digest: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum RerunReason {
Ineligible(IneligibleReason),
Defect(EvidenceDefect),
Mismatch {
field: &'static str,
recorded: String,
current: String,
},
DirtyWorktree(Vec<String>),
Unobservable(String),
BelowDurationThreshold {
elapsed: i64,
threshold: i64,
},
}
impl RerunReason {
pub fn code(&self) -> &'static str {
match self {
Self::Ineligible(reason) => reason.code(),
Self::Defect(defect) => defect.code(),
Self::Mismatch { .. } => "binding_mismatch",
Self::DirtyWorktree(_) => "worktree_dirty",
Self::Unobservable(_) => "state_unobservable",
Self::BelowDurationThreshold { .. } => "below_reuse_duration_threshold",
}
}
pub fn detail(&self) -> String {
match self {
Self::Ineligible(reason) => reason.detail(),
Self::Defect(defect) => defect.detail(),
Self::Mismatch {
field,
recorded,
current,
} => format!("{field} changed: evidence '{recorded}', current '{current}'"),
Self::DirtyWorktree(entries) => format!(
"worktree differs from the bound commit: {}",
entries.join(", ")
),
Self::Unobservable(detail) => {
format!("current repository state could not be proven: {detail}")
}
Self::BelowDurationThreshold { elapsed, threshold } => format!(
"measured {elapsed}s is below the {threshold}s reuse threshold; cheap commands rerun"
),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ReuseDecision {
Reuse {
verification_id: String,
evidence: Box<VerificationEvidence>,
},
Rerun {
verification_id: String,
reason: RerunReason,
},
}
impl ReuseDecision {
pub fn verification_id(&self) -> &str {
match self {
Self::Reuse {
verification_id, ..
}
| Self::Rerun {
verification_id, ..
} => verification_id,
}
}
pub fn is_reuse(&self) -> bool {
matches!(self, Self::Reuse { .. })
}
pub fn outcome(&self) -> &'static str {
if self.is_reuse() {
"reused"
} else {
"rerun"
}
}
pub fn to_json(&self) -> serde_json::Value {
match self {
Self::Reuse {
verification_id,
evidence,
} => serde_json::json!({
"verification_id": verification_id,
"outcome": "reused",
"evidence": evidence.to_json(),
}),
Self::Rerun {
verification_id,
reason,
} => serde_json::json!({
"verification_id": verification_id,
"outcome": "rerun",
"reason": reason.code(),
"detail": reason.detail(),
}),
}
}
pub fn summary(&self) -> String {
match self {
Self::Reuse {
verification_id,
evidence,
} => format!(
"{verification_id}: reused (commit {}, {}s, artifact {})",
&evidence.commit_oid[..12.min(evidence.commit_oid.len())],
evidence.elapsed_seconds().unwrap_or_default(),
evidence.artifact_path
),
Self::Rerun {
verification_id,
reason,
} => format!(
"{verification_id}: rerun ({}) {}",
reason.code(),
reason.detail()
),
}
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct VerificationReusePlan {
pub decisions: Vec<ReuseDecision>,
}
impl VerificationReusePlan {
pub fn is_empty(&self) -> bool {
self.decisions.is_empty()
}
pub fn reused_ids(&self) -> Vec<&str> {
self.decisions
.iter()
.filter(|decision| decision.is_reuse())
.map(ReuseDecision::verification_id)
.collect()
}
pub fn rerun_ids(&self) -> Vec<&str> {
self.decisions
.iter()
.filter(|decision| !decision.is_reuse())
.map(ReuseDecision::verification_id)
.collect()
}
pub fn to_json(&self) -> serde_json::Value {
serde_json::json!({
"verifications": self
.decisions
.iter()
.map(ReuseDecision::to_json)
.collect::<Vec<_>>(),
"reused": self.reused_ids(),
"rerun": self.rerun_ids(),
})
}
}
pub fn evaluate_reuse(
request: &VerificationRequest,
record: &VerificationEvidence,
current: &CurrentBindings,
policy: ReusePolicy,
) -> ReuseDecision {
let id = request.verification_id.clone();
let rerun = |reason: RerunReason| ReuseDecision::Rerun {
verification_id: id.clone(),
reason,
};
if record.verification_id != request.verification_id {
return rerun(RerunReason::Mismatch {
field: "verification_id",
recorded: record.verification_id.clone(),
current: request.verification_id.clone(),
});
}
if record.argv != request.argv {
return rerun(RerunReason::Mismatch {
field: "argv",
recorded: record.argv.join(" "),
current: request.argv.join(" "),
});
}
if record.cwd != request.cwd {
return rerun(RerunReason::Mismatch {
field: "cwd",
recorded: record.cwd.clone(),
current: request.cwd.clone(),
});
}
if record.automation_path != request.automation_path {
return rerun(RerunReason::Mismatch {
field: "automation_path",
recorded: record.automation_path.clone(),
current: request.automation_path.clone(),
});
}
if record.commit_oid != current.commit_oid {
return rerun(RerunReason::Mismatch {
field: "commit_oid",
recorded: record.commit_oid.clone(),
current: current.commit_oid.clone(),
});
}
if record.tree_oid != current.tree_oid {
return rerun(RerunReason::Mismatch {
field: "tree_oid",
recorded: record.tree_oid.clone(),
current: current.tree_oid.clone(),
});
}
if let Some(review) = current.review.as_ref() {
if record.review_base_commit != review.base_commit {
return rerun(RerunReason::Mismatch {
field: "review_base_commit",
recorded: record.review_base_commit.clone(),
current: review.base_commit.clone(),
});
}
if record.review_range != review.range {
return rerun(RerunReason::Mismatch {
field: "review_range",
recorded: record.review_range.clone(),
current: review.range.clone(),
});
}
}
if record.automation_blob_oid != current.automation_blob_oid {
return rerun(RerunReason::Mismatch {
field: "automation_blob_oid",
recorded: record.automation_blob_oid.clone(),
current: current.automation_blob_oid.clone(),
});
}
if record.tool.path != current.tool.path {
return rerun(RerunReason::Mismatch {
field: "tool_path",
recorded: record.tool.path.clone(),
current: current.tool.path.clone(),
});
}
if record.tool.executable_digest != current.tool.executable_digest {
return rerun(RerunReason::Mismatch {
field: "tool_digest",
recorded: record.tool.executable_digest.clone(),
current: current.tool.executable_digest.clone(),
});
}
if record.tool.version != current.tool.version {
return rerun(RerunReason::Mismatch {
field: "tool_version",
recorded: record.tool.version.clone().unwrap_or_default(),
current: current.tool.version.clone().unwrap_or_default(),
});
}
match current.artifact_digest.as_deref() {
None => {
return rerun(RerunReason::Defect(EvidenceDefect::Unreadable(format!(
"artifact '{}' is absent or unreadable",
record.artifact_path
))))
}
Some(digest) if digest != record.artifact_digest => {
return rerun(RerunReason::Mismatch {
field: "artifact_digest",
recorded: record.artifact_digest.clone(),
current: digest.to_string(),
})
}
Some(_) => {}
}
if !current.clean {
return rerun(RerunReason::DirtyWorktree(Vec::new()));
}
let elapsed = record.elapsed_seconds().unwrap_or_default();
if elapsed < policy.min_reuse_seconds {
return rerun(RerunReason::BelowDurationThreshold {
elapsed,
threshold: policy.min_reuse_seconds,
});
}
ReuseDecision::Reuse {
verification_id: id,
evidence: Box::new(record.clone()),
}
}
pub fn dirty_entries(porcelain_status: &str) -> Vec<String> {
porcelain_status
.lines()
.filter(|line| line.len() >= 4)
.map(|line| line.trim().to_string())
.collect()
}
#[async_trait]
pub trait RepositoryFacts: Send + Sync {
async fn head_commit(&self, workspace: &Path) -> Result<String, String>;
async fn resolve_revision(&self, workspace: &Path, revision: &str) -> Result<String, String> {
let _ = (workspace, revision);
Err(format!("revision '{revision}' cannot be resolved here"))
}
async fn head_tree(&self, workspace: &Path) -> Result<String, String>;
async fn tracked_blob_oid(&self, workspace: &Path, path: &str) -> Result<String, String>;
async fn hash_file(&self, workspace: &Path, path: &Path) -> Result<String, String>;
async fn porcelain_status(&self, workspace: &Path) -> Result<String, String>;
async fn resolve_tool(&self, workspace: &Path, program: &str) -> Result<ToolIdentity, String>;
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SupervisedOutcome {
pub exit_code: i32,
pub output: Vec<u8>,
}
const GATE_TERMINATION_GRACE_MS: u64 = 200;
const GATE_CLEANUP_TIMEOUT_MS: u64 = 5_000;
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum SupervisedRun {
Completed(SupervisedOutcome),
DeadlineExhausted(Box<crate::process_manager::ProcessGroupCleanupReport>),
}
#[async_trait]
pub trait CommandSupervisor: Send + Sync {
async fn run(&self, argv: &[String], cwd: &Path) -> Result<SupervisedOutcome, String>;
async fn run_bounded(
&self,
argv: &[String],
cwd: &Path,
budget: Option<Duration>,
) -> Result<SupervisedRun, String> {
let Some(budget) = budget else {
return self.run(argv, cwd).await.map(SupervisedRun::Completed);
};
match tokio::time::timeout(budget, self.run(argv, cwd)).await {
Ok(result) => result.map(SupervisedRun::Completed),
Err(_) => Ok(SupervisedRun::DeadlineExhausted(Box::new(
crate::process_manager::ProcessGroupCleanupReport::not_applicable(
"this supervisor owns no OS process, so cancelling its run is the whole \
termination",
),
))),
}
}
}
pub trait Clock: Send + Sync {
fn now(&self) -> DateTime<Utc>;
}
#[derive(Debug, Clone, Copy, Default)]
pub struct SystemClock;
impl Clock for SystemClock {
fn now(&self) -> DateTime<Utc> {
Utc::now()
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct CaptureRefusal {
pub verification_id: String,
pub code: &'static str,
pub detail: String,
}
impl CaptureRefusal {
fn new(verification_id: &str, code: &'static str, detail: impl Into<String>) -> Self {
Self {
verification_id: verification_id.to_string(),
code,
detail: detail.into(),
}
}
pub fn to_json(&self) -> serde_json::Value {
serde_json::json!({
"verification_id": self.verification_id,
"outcome": "not_captured",
"reason": self.code,
"detail": self.detail,
})
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum CaptureOutcome {
Captured {
evidence: Box<VerificationEvidence>,
artifact_location: PathBuf,
},
CommandFailed {
verification_id: String,
exit_code: i32,
artifact_path: String,
artifact_location: PathBuf,
},
Refused(CaptureRefusal),
DeadlineExhausted {
verification_id: String,
cleanup: Box<crate::process_manager::ProcessGroupCleanupReport>,
},
}
impl CaptureOutcome {
pub fn artifact_location(&self) -> Option<&Path> {
match self {
Self::Captured {
artifact_location, ..
}
| Self::CommandFailed {
artifact_location, ..
} => Some(artifact_location.as_path()),
Self::Refused(_) | Self::DeadlineExhausted { .. } => None,
}
}
pub fn to_json(&self) -> serde_json::Value {
match self {
Self::Captured {
evidence,
artifact_location,
} => serde_json::json!({
"verification_id": evidence.verification_id,
"outcome": "captured",
"evidence": evidence.to_json(),
"artifact_location": artifact_location.display().to_string(),
}),
Self::CommandFailed {
verification_id,
exit_code,
artifact_path,
artifact_location,
} => serde_json::json!({
"verification_id": verification_id,
"outcome": "command_failed",
"exit_code": exit_code,
"artifact_path": artifact_path,
"artifact_location": artifact_location.display().to_string(),
}),
Self::Refused(refusal) => refusal.to_json(),
Self::DeadlineExhausted {
verification_id,
cleanup,
} => serde_json::json!({
"verification_id": verification_id,
"outcome": "deadline_exhausted",
"cleanup_confirmed": cleanup.is_confirmed(),
"cleanup_diagnostics": cleanup.diagnostics(),
}),
}
}
}
#[derive(Debug, Clone)]
pub struct EvidenceStore {
root: PathBuf,
}
impl EvidenceStore {
pub fn new(root: impl Into<PathBuf>) -> Self {
Self { root: root.into() }
}
pub fn root(&self) -> &Path {
self.root.as_path()
}
pub fn directory(&self) -> PathBuf {
self.root.join(GATES_DIR)
}
pub fn envelope_relative_path(verification_id: &str) -> String {
format!("{GATES_DIR}/{verification_id}.json")
}
pub fn artifact_relative_path(verification_id: &str) -> String {
format!("{GATES_DIR}/{verification_id}.log")
}
pub fn artifact_path(&self, verification_id: &str) -> PathBuf {
self.root
.join(Self::artifact_relative_path(verification_id))
}
pub fn ensure_directory(&self) -> std::io::Result<()> {
let directory = self.directory();
std::fs::create_dir_all(&directory)?;
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
let _ = std::fs::set_permissions(&self.root, std::fs::Permissions::from_mode(0o700));
std::fs::set_permissions(&directory, std::fs::Permissions::from_mode(0o700))?;
}
Ok(())
}
pub fn load(&self, verification_id: &str) -> Result<VerificationEvidence, EvidenceDefect> {
if !is_storable_verification_id(verification_id) {
return Err(EvidenceDefect::Malformed(format!(
"verification id '{verification_id}' is not storable"
)));
}
let path = self
.root
.join(Self::envelope_relative_path(verification_id));
let metadata = match std::fs::metadata(&path) {
Ok(metadata) => metadata,
Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
return Err(EvidenceDefect::Missing)
}
Err(error) => return Err(EvidenceDefect::Unreadable(error.to_string())),
};
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
let mode = metadata.permissions().mode() & 0o777;
if mode != 0o400 {
return Err(EvidenceDefect::AgentAuthored(format!(
"envelope mode {mode:04o} is not the runtime-written 0400"
)));
}
}
#[cfg(not(unix))]
let _ = metadata;
let bytes =
std::fs::read(&path).map_err(|error| EvidenceDefect::Unreadable(error.to_string()))?;
parse_evidence(&bytes)
}
pub fn store(&self, evidence: &VerificationEvidence) -> std::io::Result<()> {
self.ensure_directory()?;
let directory = self.directory();
let final_path = self
.root
.join(Self::envelope_relative_path(&evidence.verification_id));
let temporary = directory.join(format!(
".{}.tmp-{}",
evidence.verification_id,
std::process::id()
));
let serialized = serde_json::to_vec_pretty(evidence)
.map_err(|error| std::io::Error::new(std::io::ErrorKind::InvalidData, error))?;
std::fs::write(&temporary, &serialized)?;
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
std::fs::set_permissions(&temporary, std::fs::Permissions::from_mode(0o400))?;
}
let _ = std::fs::remove_file(&final_path);
match std::fs::rename(&temporary, &final_path) {
Ok(()) => Ok(()),
Err(error) => {
let _ = std::fs::remove_file(&temporary);
Err(error)
}
}
}
fn store_artifact(&self, verification_id: &str, output: &[u8]) -> std::io::Result<PathBuf> {
self.ensure_directory()?;
let path = self.artifact_path(verification_id);
let bounded = if output.len() > MAX_ARTIFACT_BYTES {
&output[output.len() - MAX_ARTIFACT_BYTES..]
} else {
output
};
let _ = std::fs::remove_file(&path);
std::fs::write(&path, bounded)?;
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
std::fs::set_permissions(&path, std::fs::Permissions::from_mode(0o400))?;
}
Ok(path)
}
}
pub struct RuntimeVerificationExecutor<F, S, C> {
workspace: PathBuf,
store: EvidenceStore,
facts: F,
supervisor: S,
clock: C,
review: Option<ReviewBinding>,
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct ReviewBinding {
pub base_commit: String,
pub range: String,
}
impl<F, S, C> RuntimeVerificationExecutor<F, S, C>
where
F: RepositoryFacts,
S: CommandSupervisor,
C: Clock,
{
pub fn new(
workspace: impl Into<PathBuf>,
store: EvidenceStore,
facts: F,
supervisor: S,
clock: C,
) -> Self {
Self {
workspace: workspace.into(),
store,
facts,
supervisor,
clock,
review: None,
}
}
pub fn store(&self) -> &EvidenceStore {
&self.store
}
pub fn with_review_binding(mut self, review: ReviewBinding) -> Self {
self.review = Some(review);
self
}
async fn snapshot(&self, request: &VerificationRequest) -> Result<CurrentBindings, String> {
let workspace = self.workspace.as_path();
let commit_oid = self.facts.head_commit(workspace).await?;
let tree_oid = self.facts.head_tree(workspace).await?;
let automation_blob_oid = self
.facts
.tracked_blob_oid(workspace, &request.automation_path)
.await?;
let tool = self
.facts
.resolve_tool(workspace, request.argv.first().map_or("", String::as_str))
.await?;
let status = self.facts.porcelain_status(workspace).await?;
Ok(CurrentBindings {
commit_oid,
tree_oid,
review: self.review.clone(),
automation_blob_oid,
tool,
clean: dirty_entries(&status).is_empty(),
artifact_digest: None,
})
}
pub async fn capture(&self, request: &VerificationRequest) -> CaptureOutcome {
self.capture_bounded(request, None).await
}
pub async fn capture_bounded(
&self,
request: &VerificationRequest,
budget: Option<Duration>,
) -> CaptureOutcome {
let id = request.verification_id.as_str();
let store = &self.store;
if let Err(error) = store.ensure_directory() {
return CaptureOutcome::Refused(CaptureRefusal::new(
id,
"evidence_store_unavailable",
error.to_string(),
));
}
let before = match self.snapshot(request).await {
Ok(before) => before,
Err(error) => {
return CaptureOutcome::Refused(CaptureRefusal::new(
id,
"state_unobservable",
error,
))
}
};
if !before.clean {
return CaptureOutcome::Refused(CaptureRefusal::new(
id,
"worktree_dirty",
"the target worktree is dirty before execution".to_string(),
));
}
let cwd = if request.cwd == "." {
self.workspace.clone()
} else {
self.workspace.join(&request.cwd)
};
let started_at = self.clock.now();
let outcome = match self
.supervisor
.run_bounded(&request.argv, &cwd, budget)
.await
{
Ok(SupervisedRun::Completed(outcome)) => outcome,
Ok(SupervisedRun::DeadlineExhausted(cleanup)) => {
return CaptureOutcome::DeadlineExhausted {
verification_id: id.to_string(),
cleanup,
}
}
Err(error) => {
return CaptureOutcome::Refused(CaptureRefusal::new(
id,
"supervision_failed",
error,
))
}
};
let ended_at = self.clock.now();
let artifact_path = match store.store_artifact(id, &outcome.output) {
Ok(path) => path,
Err(error) => {
return CaptureOutcome::Refused(CaptureRefusal::new(
id,
"artifact_unwritable",
error.to_string(),
))
}
};
if outcome.exit_code != 0 {
return CaptureOutcome::CommandFailed {
verification_id: id.to_string(),
exit_code: outcome.exit_code,
artifact_path: EvidenceStore::artifact_relative_path(id),
artifact_location: artifact_path,
};
}
let artifact_digest = match self
.facts
.hash_file(self.workspace.as_path(), artifact_path.as_path())
.await
{
Ok(digest) if is_full_length_object_id(&digest) => digest,
Ok(digest) => {
return CaptureOutcome::Refused(CaptureRefusal::new(
id,
"artifact_digest_invalid",
format!("'{digest}' is not a full-length object id"),
))
}
Err(error) => {
return CaptureOutcome::Refused(CaptureRefusal::new(
id,
"artifact_digest_unavailable",
error,
))
}
};
let after = match self.snapshot(request).await {
Ok(after) => after,
Err(error) => {
return CaptureOutcome::Refused(CaptureRefusal::new(
id,
"state_unobservable",
error,
))
}
};
if !after.clean {
return CaptureOutcome::Refused(CaptureRefusal::new(
id,
"worktree_dirty",
"the target worktree became dirty while the command ran".to_string(),
));
}
for (field, recorded, current) in [
("commit_oid", &before.commit_oid, &after.commit_oid),
("tree_oid", &before.tree_oid, &after.tree_oid),
(
"automation_blob_oid",
&before.automation_blob_oid,
&after.automation_blob_oid,
),
(
"tool_digest",
&before.tool.executable_digest,
&after.tool.executable_digest,
),
] {
if recorded != current {
return CaptureOutcome::Refused(CaptureRefusal::new(
id,
"binding_moved_during_execution",
format!(
"{field} changed from '{recorded}' to '{current}' while the command ran"
),
));
}
}
let evidence = VerificationEvidence {
schema: EVIDENCE_SCHEMA.to_string(),
authority: EVIDENCE_AUTHORITY.to_string(),
verification_id: id.to_string(),
commit_oid: before.commit_oid.clone(),
tree_oid: before.tree_oid.clone(),
review_base_commit: before
.review
.as_ref()
.map(|review| review.base_commit.clone())
.unwrap_or_default(),
review_range: before
.review
.as_ref()
.map(|review| review.range.clone())
.unwrap_or_default(),
argv: request.argv.clone(),
cwd: request.cwd.clone(),
automation_path: request.automation_path.clone(),
automation_blob_oid: before.automation_blob_oid.clone(),
tool: before.tool.clone(),
started_at,
ended_at,
exit_code: outcome.exit_code,
artifact_path: EvidenceStore::artifact_relative_path(id),
artifact_digest,
clean_before: true,
clean_after: true,
};
if let Err(error) = self.store.store(&evidence) {
return CaptureOutcome::Refused(CaptureRefusal::new(
id,
"evidence_unwritable",
error.to_string(),
));
}
CaptureOutcome::Captured {
evidence: Box::new(evidence),
artifact_location: artifact_path,
}
}
}
pub async fn decide_reuse<F: RepositoryFacts>(
facts: &F,
workspace: &Path,
store: &EvidenceStore,
declaration: &VerificationDeclaration,
policy: ReusePolicy,
) -> Option<ReuseDecision> {
let declared_id = declaration
.id
.as_deref()
.map(str::trim)
.filter(|id| !id.is_empty())?;
if !declaration
.execution_class
.as_deref()
.and_then(VerificationExecutionClass::parse)
.is_some_and(VerificationExecutionClass::is_repository_local)
{
return None;
}
let request = match eligible_request(declaration) {
Ok(request) => request,
Err(reason) => {
return Some(ReuseDecision::Rerun {
verification_id: declared_id.to_string(),
reason: RerunReason::Ineligible(reason),
})
}
};
let record = match store.load(&request.verification_id) {
Ok(record) => record,
Err(defect) => {
return Some(ReuseDecision::Rerun {
verification_id: request.verification_id,
reason: RerunReason::Defect(defect),
})
}
};
let unobservable = |error: String| {
Some(ReuseDecision::Rerun {
verification_id: request.verification_id.clone(),
reason: RerunReason::Unobservable(error),
})
};
let commit_oid = match facts.head_commit(workspace).await {
Ok(value) => value,
Err(error) => return unobservable(error),
};
let tree_oid = match facts.head_tree(workspace).await {
Ok(value) => value,
Err(error) => return unobservable(error),
};
let automation_blob_oid = match facts
.tracked_blob_oid(workspace, &request.automation_path)
.await
{
Ok(value) => value,
Err(error) => return unobservable(error),
};
let tool = match facts
.resolve_tool(workspace, request.argv.first().map_or("", String::as_str))
.await
{
Ok(value) => value,
Err(error) => return unobservable(error),
};
let status = match facts.porcelain_status(workspace).await {
Ok(value) => value,
Err(error) => return unobservable(error),
};
let dirty = dirty_entries(&status);
if !dirty.is_empty() {
return Some(ReuseDecision::Rerun {
verification_id: request.verification_id,
reason: RerunReason::DirtyWorktree(dirty),
});
}
let artifact_digest = facts
.hash_file(workspace, &store.artifact_path(&request.verification_id))
.await
.ok()
.filter(|digest| is_full_length_object_id(digest));
let current = CurrentBindings {
commit_oid,
tree_oid,
review: None,
automation_blob_oid,
tool,
clean: true,
artifact_digest,
};
Some(evaluate_reuse(&request, &record, ¤t, policy))
}
pub async fn plan_reuse<F: RepositoryFacts>(
facts: &F,
workspace: &Path,
store: &EvidenceStore,
declarations: &[VerificationDeclaration],
policy: ReusePolicy,
) -> VerificationReusePlan {
let mut decisions = Vec::new();
let mut seen = BTreeMap::new();
for declaration in declarations {
if let Some(decision) = decide_reuse(facts, workspace, store, declaration, policy).await {
if seen
.insert(decision.verification_id().to_string(), ())
.is_none()
{
decisions.push(decision);
}
}
}
VerificationReusePlan { decisions }
}
pub struct GitRepositoryFacts;
#[async_trait]
impl RepositoryFacts for GitRepositoryFacts {
async fn head_commit(&self, workspace: &Path) -> Result<String, String> {
run_git(workspace, &["rev-parse", "HEAD"]).await
}
async fn head_tree(&self, workspace: &Path) -> Result<String, String> {
run_git(workspace, &["rev-parse", "HEAD^{tree}"]).await
}
async fn resolve_revision(&self, workspace: &Path, revision: &str) -> Result<String, String> {
run_git(workspace, &["rev-parse", &format!("{revision}^{{commit}}")]).await
}
async fn tracked_blob_oid(&self, workspace: &Path, path: &str) -> Result<String, String> {
run_git(workspace, &["rev-parse", &format!("HEAD:{path}")]).await
}
async fn hash_file(&self, workspace: &Path, path: &Path) -> Result<String, String> {
let path = path.to_string_lossy().to_string();
run_git(workspace, &["hash-object", "-t", "blob", "--", &path]).await
}
async fn porcelain_status(&self, workspace: &Path) -> Result<String, String> {
let argv = crate::vcs::git::commands::status_policy::read_only_status_argv(
crate::vcs::git::commands::status_policy::PORCELAIN_UNTRACKED_STATUS_ARGS,
);
run_git_raw(workspace, &argv).await
}
async fn resolve_tool(&self, workspace: &Path, program: &str) -> Result<ToolIdentity, String> {
if program.trim().is_empty() {
return Err("verification argv has no program".to_string());
}
let resolved = resolve_executable(program)
.ok_or_else(|| format!("executable '{program}' is not on PATH"))?;
let executable_digest = self.hash_file(workspace, resolved.as_path()).await?;
let version = tokio::process::Command::new(&resolved)
.arg(VERSION_FLAG)
.current_dir(workspace)
.stdin(std::process::Stdio::null())
.output()
.await
.ok()
.filter(|output| output.status.success())
.and_then(|output| usable_tool_version(&String::from_utf8_lossy(&output.stdout)));
Ok(ToolIdentity {
path: resolved.to_string_lossy().to_string(),
executable_digest,
version,
})
}
}
const VERSION_FLAG: &str = "--version";
fn usable_tool_version(raw: &str) -> Option<String> {
let version = raw.trim();
(!version.is_empty() && version != VERSION_FLAG).then(|| version.to_string())
}
fn resolve_executable(program: &str) -> Option<PathBuf> {
let candidate = Path::new(program);
if candidate.components().count() > 1 {
let absolute = if candidate.is_absolute() {
candidate.to_path_buf()
} else {
std::env::current_dir().ok()?.join(candidate)
};
return absolute.is_file().then_some(absolute);
}
let path = std::env::var_os("PATH")?;
std::env::split_paths(&path)
.map(|directory| directory.join(program))
.find(|candidate| candidate.is_file())
}
async fn run_git(workspace: &Path, args: &[&str]) -> Result<String, String> {
run_git_raw(workspace, args).await.map(|out| {
out.lines()
.next()
.map(str::trim)
.unwrap_or_default()
.to_string()
})
}
async fn run_git_raw(workspace: &Path, args: &[&str]) -> Result<String, String> {
let output = tokio::process::Command::new("git")
.args(args)
.current_dir(workspace)
.stdin(std::process::Stdio::null())
.output()
.await
.map_err(|error| format!("git {}: {error}", args.join(" ")))?;
if !output.status.success() {
return Err(format!(
"git {} failed: {}",
args.join(" "),
String::from_utf8_lossy(&output.stderr).trim()
));
}
Ok(String::from_utf8_lossy(&output.stdout).to_string())
}
pub struct DirectCommandSupervisor;
#[async_trait]
impl CommandSupervisor for DirectCommandSupervisor {
async fn run(&self, argv: &[String], cwd: &Path) -> Result<SupervisedOutcome, String> {
match self.run_bounded(argv, cwd, None).await? {
SupervisedRun::Completed(outcome) => Ok(outcome),
SupervisedRun::DeadlineExhausted(report) => Err(format!(
"unbounded supervision reported a deadline: {}",
report.diagnostics()
)),
}
}
async fn run_bounded(
&self,
argv: &[String],
cwd: &Path,
budget: Option<Duration>,
) -> Result<SupervisedRun, String> {
use tokio::io::AsyncReadExt;
let (program, arguments) = argv
.split_first()
.ok_or_else(|| "verification argv is empty".to_string())?;
let mut command = tokio::process::Command::new(program);
command
.args(arguments)
.current_dir(cwd)
.stdin(std::process::Stdio::null())
.stdout(std::process::Stdio::piped())
.stderr(std::process::Stdio::piped())
.kill_on_drop(true);
crate::process_manager::configure_process_group(&mut command);
let mut child = command
.spawn()
.map_err(|error| format!("{program}: {error}"))?;
let pgid = child.id().unwrap_or(0);
let stdout = child.stdout.take();
let stderr = child.stderr.take();
let stdout_task = tokio::spawn(async move {
let mut buffer = Vec::new();
if let Some(mut stream) = stdout {
let _ = stream.read_to_end(&mut buffer).await;
}
buffer
});
let stderr_task = tokio::spawn(async move {
let mut buffer = Vec::new();
if let Some(mut stream) = stderr {
let _ = stream.read_to_end(&mut buffer).await;
}
buffer
});
let status = match budget {
None => child.wait().await,
Some(budget) => match tokio::time::timeout(budget, child.wait()).await {
Ok(status) => status,
Err(_) => {
let _ = child.start_kill();
let _ = child.wait().await;
stdout_task.abort();
stderr_task.abort();
let report = crate::process_manager::cleanup_process_group_verified(
pgid,
GATE_TERMINATION_GRACE_MS,
GATE_CLEANUP_TIMEOUT_MS,
Some("acceptance-gate"),
None,
)
.await;
return Ok(SupervisedRun::DeadlineExhausted(Box::new(report)));
}
},
};
let status = status.map_err(|error| format!("{program}: {error}"))?;
let mut merged = stdout_task.await.unwrap_or_default();
merged.extend_from_slice(&stderr_task.await.unwrap_or_default());
Ok(SupervisedRun::Completed(SupervisedOutcome {
exit_code: status.code().unwrap_or(-1),
output: merged,
}))
}
}
#[cfg(test)]
#[path = "verification_evidence/tests.rs"]
mod tests;