#[cfg(test)]
pub use molo_agent::ReActAgent;
pub use molo_agent::{agent, tool};
#[cfg(test)]
pub use molo_core::ToolCall;
pub use molo_core::{effect, message, observability, provider, run};
use crate::agent::{
AgentAction, AgentError, AgentKernel, ModelObservation, ModelRequest, Observation,
};
use crate::effect::{
DisplayOutput, EffectKind, EffectObservation, EffectOutput, EffectRequest, EffectStatus,
RiskLevel,
};
use crate::provider::{Provider, ProviderError, ProviderRequestContext};
use crate::run::{Artifact, RunContext, RunMetadata, RunOutput, RunRequest};
use async_trait::async_trait;
use futures::stream::{FuturesUnordered, StreamExt};
use serde::{Deserialize, Serialize};
use std::collections::{BTreeMap, HashMap};
use std::fmt;
use std::sync::{Arc, Mutex};
use std::time::Duration;
pub use crate::observability::RedactionRecord;
#[derive(Debug)]
pub struct HarnessRuntime<P, H> {
provider: P,
harness: H,
config: HarnessRuntimeConfig,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(default)]
#[non_exhaustive]
pub struct HarnessRuntimeConfig {
pub(crate) max_agent_steps: usize,
pub(crate) max_effect_batch_concurrency: usize,
pub(crate) fail_fast_effect_batches: bool,
}
impl Default for HarnessRuntimeConfig {
fn default() -> Self {
Self {
max_agent_steps: 256,
max_effect_batch_concurrency: 1,
fail_fast_effect_batches: false,
}
}
}
impl HarnessRuntimeConfig {
pub fn new() -> Self {
Self::default()
}
pub fn max_agent_steps(&self) -> usize {
self.max_agent_steps
}
pub fn with_max_agent_steps(mut self, max_agent_steps: usize) -> Self {
self.max_agent_steps = max_agent_steps;
self
}
pub fn max_effect_batch_concurrency(&self) -> usize {
self.max_effect_batch_concurrency
}
pub fn with_max_effect_batch_concurrency(
mut self,
max_effect_batch_concurrency: usize,
) -> Self {
self.max_effect_batch_concurrency = max_effect_batch_concurrency;
self
}
pub fn fail_fast_effect_batches(&self) -> bool {
self.fail_fast_effect_batches
}
pub fn with_fail_fast_effect_batches(mut self, fail_fast_effect_batches: bool) -> Self {
self.fail_fast_effect_batches = fail_fast_effect_batches;
self
}
}
impl<P, H> HarnessRuntime<P, H>
where
P: Provider,
H: Harness,
{
pub fn new(provider: P, harness: H) -> Self {
Self {
provider,
harness,
config: HarnessRuntimeConfig::default(),
}
}
pub fn with_config(mut self, config: HarnessRuntimeConfig) -> Self {
self.config = config;
self
}
pub async fn run<K>(
&self,
kernel: &mut K,
request: RunRequest,
context: RunContext,
) -> Result<RunOutput, HarnessRuntimeError>
where
K: AgentKernel,
{
check_run_context(&context)?;
let mut action = kernel.start(request, &context).await?;
for _ in 0..self.config.max_agent_steps {
check_run_context(&context)?;
let observation = match action {
AgentAction::Respond { output } => return Ok(output),
AgentAction::RequestModel { request } => {
Observation::Model(self.execute_model_request(request, &context).await?)
}
AgentAction::RequestEffect { request } => {
let observation = self.harness.execute(request, &context).await?;
Observation::Effect(observation)
}
AgentAction::RequestEffects { requests } => {
let observations = self.execute_effect_batch(requests, &context).await?;
Observation::Effects(observations)
}
_ => {
return Err(
AgentError::InvalidStep("unsupported agent action".to_string()).into(),
);
}
};
action = kernel.observe(observation, &context).await?;
}
Err(HarnessRuntimeError::TooManyAgentSteps(
self.config.max_agent_steps,
))
}
async fn execute_effect_batch(
&self,
requests: Vec<EffectRequest>,
context: &RunContext,
) -> Result<Vec<EffectObservation>, HarnessError> {
check_run_context(context)?;
if requests.is_empty() {
return Ok(Vec::new());
}
let concurrency = self.config.max_effect_batch_concurrency.max(1);
if concurrency == 1 {
let mut observations = Vec::with_capacity(requests.len());
for request in requests {
let effect_id = request.id.clone();
match self.harness.execute(request, context).await {
Ok(observation) => observations.push(observation),
Err(error) if self.config.fail_fast_effect_batches => return Err(error),
Err(error) if is_terminal_context_error(&error) => return Err(error),
Err(error) => {
observations.push(observation_from_harness_error(effect_id, error))
}
}
}
return Ok(observations);
}
let request_count = requests.len();
let mut pending = requests.into_iter().enumerate();
let mut in_flight = FuturesUnordered::new();
let mut observations: Vec<Option<EffectObservation>> = std::iter::repeat_with(|| None)
.take(request_count)
.collect();
for _ in 0..concurrency {
let Some((index, request)) = pending.next() else {
break;
};
in_flight.push(execute_indexed_effect(
&self.harness,
index,
request,
context,
));
}
while let Some((index, effect_id, result)) = in_flight.next().await {
match result {
Ok(observation) => observations[index] = Some(observation),
Err(error) if self.config.fail_fast_effect_batches => return Err(error),
Err(error) if is_terminal_context_error(&error) => return Err(error),
Err(error) => {
observations[index] = Some(observation_from_harness_error(effect_id, error))
}
}
if let Some((next_index, request)) = pending.next() {
in_flight.push(execute_indexed_effect(
&self.harness,
next_index,
request,
context,
));
}
}
Ok(observations
.into_iter()
.map(|observation| {
observation.expect("batch scheduler must fill every requested observation")
})
.collect())
}
async fn execute_model_request(
&self,
request: ModelRequest,
context: &RunContext,
) -> Result<ModelObservation, ProviderError> {
let request_id = request.id;
let provider_context = ProviderRequestContext::from_run_context(&request_id, context);
let response = self
.provider
.chat_with_context(request.chat, &provider_context)
.await?;
Ok(ModelObservation::new(request_id, response))
}
}
async fn execute_indexed_effect<'a, H>(
harness: &'a H,
index: usize,
request: EffectRequest,
context: &'a RunContext,
) -> (usize, String, Result<EffectObservation, HarnessError>)
where
H: Harness,
{
let effect_id = request.id.clone();
let result = harness.execute(request, context).await;
(index, effect_id, result)
}
#[async_trait]
pub trait Harness: Send + Sync {
async fn execute(
&self,
request: EffectRequest,
context: &RunContext,
) -> Result<EffectObservation, HarnessError>;
async fn execute_batch(
&self,
requests: Vec<EffectRequest>,
context: &RunContext,
) -> Result<Vec<EffectObservation>, HarnessError> {
let mut observations = Vec::with_capacity(requests.len());
for request in requests {
observations.push(self.execute(request, context).await?);
}
Ok(observations)
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ClassifiedEffect {
pub request: EffectRequest,
pub requested_risk: RiskLevel,
pub effective_risk: RiskLevel,
pub reasons: Vec<String>,
pub metadata: RunMetadata,
}
#[async_trait]
pub trait RiskClassifier: Send + Sync {
async fn classify(
&self,
request: EffectRequest,
context: &RunContext,
) -> Result<ClassifiedEffect, HarnessError>;
}
#[derive(Debug, Default, Clone, Copy)]
pub struct DefaultRiskClassifier;
#[async_trait]
impl RiskClassifier for DefaultRiskClassifier {
async fn classify(
&self,
request: EffectRequest,
_context: &RunContext,
) -> Result<ClassifiedEffect, HarnessError> {
validate_effect_request(&request)?;
let requested_risk = request.risk;
let kind_floor = risk_floor_for_kind(&request.kind);
let mut effective_risk = max_risk(requested_risk, kind_floor);
let mut reasons = vec![format!("kind floor: {:?}", kind_floor)];
let payload_text = request.payload.to_string().to_ascii_lowercase();
let description = request.description.to_ascii_lowercase();
let combined = format!("{description} {payload_text}");
if contains_critical_pattern(&combined) {
effective_risk = max_risk(effective_risk, RiskLevel::Critical);
reasons.push("critical payload pattern".to_string());
} else if contains_high_pattern(&combined) {
effective_risk = max_risk(effective_risk, RiskLevel::High);
reasons.push("high-risk payload pattern".to_string());
}
Ok(ClassifiedEffect {
request,
requested_risk,
effective_risk,
reasons,
metadata: RunMetadata::new(),
})
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub enum PolicyDecision {
Allow,
Deny {
reason: String,
},
RequireApproval {
reason: String,
},
}
#[async_trait]
pub trait PolicyEngine: Send + Sync {
async fn evaluate(
&self,
effect: &ClassifiedEffect,
context: &RunContext,
) -> Result<PolicyDecision, HarnessError>;
}
#[derive(Debug, Default, Clone, Copy)]
pub struct DefaultPolicyEngine;
#[async_trait]
impl PolicyEngine for DefaultPolicyEngine {
async fn evaluate(
&self,
effect: &ClassifiedEffect,
_context: &RunContext,
) -> Result<PolicyDecision, HarnessError> {
Ok(match effect.effective_risk {
RiskLevel::Low | RiskLevel::Medium => PolicyDecision::Allow,
RiskLevel::High => PolicyDecision::RequireApproval {
reason: "high-risk effect requires approval".to_string(),
},
RiskLevel::Critical => PolicyDecision::Deny {
reason: "critical-risk effect denied by default policy".to_string(),
},
_ => PolicyDecision::RequireApproval {
reason: "unknown-risk effect requires approval".to_string(),
},
})
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ApprovalRequest {
pub run_id: String,
pub effect_id: String,
pub kind: EffectKind,
pub description: String,
pub risk: RiskLevel,
pub reason: String,
pub payload_summary: String,
pub sandbox: SandboxPolicy,
pub network: NetworkPolicy,
pub metadata: RunMetadata,
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub enum ApprovalDecision {
AllowOnce,
AllowForSession,
Deny {
reason: String,
},
}
#[async_trait]
pub trait ApprovalBroker: Send + Sync {
async fn approve(
&self,
request: ApprovalRequest,
context: &RunContext,
) -> Result<ApprovalDecision, ApprovalError>;
}
#[derive(Debug, Default, Clone, Copy)]
pub struct AlwaysAllowApprovalBroker;
#[async_trait]
impl ApprovalBroker for AlwaysAllowApprovalBroker {
async fn approve(
&self,
_request: ApprovalRequest,
_context: &RunContext,
) -> Result<ApprovalDecision, ApprovalError> {
Ok(ApprovalDecision::AllowOnce)
}
}
#[derive(Debug, Default, Clone, Copy)]
pub struct AlwaysDenyApprovalBroker;
#[async_trait]
impl ApprovalBroker for AlwaysDenyApprovalBroker {
async fn approve(
&self,
_request: ApprovalRequest,
_context: &RunContext,
) -> Result<ApprovalDecision, ApprovalError> {
Ok(ApprovalDecision::Deny {
reason: "denied by approval broker".to_string(),
})
}
}
#[derive(Debug, Clone)]
pub struct StaticApprovalBroker {
decision: ApprovalDecision,
}
impl StaticApprovalBroker {
pub fn new(decision: ApprovalDecision) -> Self {
Self { decision }
}
}
#[async_trait]
impl ApprovalBroker for StaticApprovalBroker {
async fn approve(
&self,
_request: ApprovalRequest,
_context: &RunContext,
) -> Result<ApprovalDecision, ApprovalError> {
Ok(self.decision.clone())
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[non_exhaustive]
pub enum SandboxPolicy {
ReadOnly,
WorkspaceWrite,
FullAccess,
Custom(String),
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[non_exhaustive]
pub enum NetworkPolicy {
Deny,
AllowListed(Vec<String>),
AllowAll,
Custom(String),
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(default)]
#[non_exhaustive]
pub struct OutputLimit {
pub(crate) model_bytes: usize,
pub(crate) display_bytes: usize,
pub(crate) debug_bytes: usize,
}
impl Default for OutputLimit {
fn default() -> Self {
Self {
model_bytes: 64 * 1024,
display_bytes: 256 * 1024,
debug_bytes: 16 * 1024,
}
}
}
impl OutputLimit {
pub fn new(model_bytes: usize, display_bytes: usize, debug_bytes: usize) -> Self {
Self {
model_bytes,
display_bytes,
debug_bytes,
}
}
pub fn model_bytes(&self) -> usize {
self.model_bytes
}
pub fn with_model_bytes(mut self, model_bytes: usize) -> Self {
self.model_bytes = model_bytes;
self
}
pub fn display_bytes(&self) -> usize {
self.display_bytes
}
pub fn with_display_bytes(mut self, display_bytes: usize) -> Self {
self.display_bytes = display_bytes;
self
}
pub fn debug_bytes(&self) -> usize {
self.debug_bytes
}
pub fn with_debug_bytes(mut self, debug_bytes: usize) -> Self {
self.debug_bytes = debug_bytes;
self
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(default)]
#[non_exhaustive]
pub struct ExecutionPolicy {
pub(crate) sandbox: SandboxPolicy,
pub(crate) network: NetworkPolicy,
pub(crate) timeout: Option<Duration>,
pub(crate) output_limit: OutputLimit,
}
impl Default for ExecutionPolicy {
fn default() -> Self {
Self {
sandbox: SandboxPolicy::ReadOnly,
network: NetworkPolicy::Deny,
timeout: Some(Duration::from_secs(30)),
output_limit: OutputLimit::default(),
}
}
}
impl ExecutionPolicy {
pub fn new(sandbox: SandboxPolicy, network: NetworkPolicy) -> Self {
Self {
sandbox,
network,
..Self::default()
}
}
pub fn sandbox(&self) -> &SandboxPolicy {
&self.sandbox
}
pub fn with_sandbox(mut self, sandbox: SandboxPolicy) -> Self {
self.sandbox = sandbox;
self
}
pub fn network(&self) -> &NetworkPolicy {
&self.network
}
pub fn with_network(mut self, network: NetworkPolicy) -> Self {
self.network = network;
self
}
pub fn timeout(&self) -> Option<Duration> {
self.timeout
}
pub fn with_timeout(mut self, timeout: Option<Duration>) -> Self {
self.timeout = timeout;
self
}
pub fn output_limit(&self) -> &OutputLimit {
&self.output_limit
}
pub fn with_output_limit(mut self, output_limit: OutputLimit) -> Self {
self.output_limit = output_limit;
self
}
}
#[async_trait]
pub trait EffectExecutor: Send + Sync {
async fn execute(
&self,
request: &EffectRequest,
policy: &ExecutionPolicy,
context: &RunContext,
) -> Result<RawEffectOutput, ExecutionError>;
}
#[derive(Debug, Clone, PartialEq)]
pub struct RawEffectOutput {
pub observation_for_model: String,
pub display: Option<DisplayOutput>,
pub artifacts: Vec<Artifact>,
pub metadata: RunMetadata,
pub debug: Option<String>,
}
impl RawEffectOutput {
pub fn text(observation_for_model: impl Into<String>) -> Self {
Self {
observation_for_model: observation_for_model.into(),
display: None,
artifacts: Vec::new(),
metadata: RunMetadata::new(),
debug: None,
}
}
pub fn with_display(mut self, display: DisplayOutput) -> Self {
self.display = Some(display);
self
}
pub fn with_artifacts(mut self, artifacts: Vec<Artifact>) -> Self {
self.artifacts = artifacts;
self
}
pub fn with_metadata(mut self, metadata: RunMetadata) -> Self {
self.metadata = metadata;
self
}
pub fn with_debug(mut self, debug: impl Into<String>) -> Self {
self.debug = Some(debug.into());
self
}
}
#[derive(Debug, Default, Clone, Copy)]
pub struct NoopEffectExecutor;
#[async_trait]
impl EffectExecutor for NoopEffectExecutor {
async fn execute(
&self,
request: &EffectRequest,
_policy: &ExecutionPolicy,
_context: &RunContext,
) -> Result<RawEffectOutput, ExecutionError> {
Err(ExecutionError::Unsupported(format!(
"effect kind {:?} is not supported by NoopEffectExecutor",
request.kind
)))
}
}
#[derive(Debug, Clone, Default)]
pub struct StaticEffectExecutor {
outputs: BTreeMap<String, Result<RawEffectOutput, ExecutionError>>,
}
impl StaticEffectExecutor {
pub fn new() -> Self {
Self::default()
}
pub fn with_output(mut self, effect_id: impl Into<String>, output: RawEffectOutput) -> Self {
self.outputs.insert(effect_id.into(), Ok(output));
self
}
pub fn with_error(mut self, effect_id: impl Into<String>, error: ExecutionError) -> Self {
self.outputs.insert(effect_id.into(), Err(error));
self
}
}
#[async_trait]
impl EffectExecutor for StaticEffectExecutor {
async fn execute(
&self,
request: &EffectRequest,
_policy: &ExecutionPolicy,
_context: &RunContext,
) -> Result<RawEffectOutput, ExecutionError> {
match self.outputs.get(&request.id) {
Some(Ok(output)) => Ok(output.clone()),
Some(Err(error)) => Err(error.clone()),
None => Err(ExecutionError::Unsupported(format!(
"no static output for effect {}",
request.id
))),
}
}
}
#[derive(Default, Clone)]
pub struct RouterEffectExecutor {
routes: HashMap<EffectKindKey, Arc<dyn EffectExecutor>>,
}
impl fmt::Debug for RouterEffectExecutor {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("RouterEffectExecutor")
.field("routes", &self.routes.keys().collect::<Vec<_>>())
.finish()
}
}
impl RouterEffectExecutor {
pub fn new() -> Self {
Self::default()
}
pub fn route(mut self, kind: EffectKind, executor: impl EffectExecutor + 'static) -> Self {
self.routes
.insert(EffectKindKey::from(kind), Arc::new(executor));
self
}
}
#[async_trait]
impl EffectExecutor for RouterEffectExecutor {
async fn execute(
&self,
request: &EffectRequest,
policy: &ExecutionPolicy,
context: &RunContext,
) -> Result<RawEffectOutput, ExecutionError> {
let key = EffectKindKey::from(request.kind.clone());
let Some(executor) = self.routes.get(&key) else {
return Err(ExecutionError::Unsupported(format!(
"no executor registered for effect kind {:?}",
request.kind
)));
};
executor.execute(request, policy, context).await
}
}
#[derive(Debug, Clone, PartialEq)]
pub struct LimitedOutput {
pub output: EffectOutput,
pub truncated: bool,
pub redactions: Vec<RedactionRecord>,
pub debug: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct RedactedText {
pub text: String,
pub redactions: Vec<RedactionRecord>,
}
pub trait Redactor: Send + Sync {
fn redact_model_text(&self, text: &str) -> RedactedText;
fn redact_display_text(&self, text: &str) -> RedactedText;
fn redact_debug_text(&self, text: &str) -> RedactedText;
fn redact_metadata(&self, metadata: RunMetadata) -> RunMetadata;
}
#[derive(Debug, Default, Clone, Copy)]
pub struct NoopRedactor;
impl Redactor for NoopRedactor {
fn redact_model_text(&self, text: &str) -> RedactedText {
RedactedText {
text: text.to_string(),
redactions: Vec::new(),
}
}
fn redact_display_text(&self, text: &str) -> RedactedText {
RedactedText {
text: text.to_string(),
redactions: Vec::new(),
}
}
fn redact_debug_text(&self, text: &str) -> RedactedText {
RedactedText {
text: text.to_string(),
redactions: Vec::new(),
}
}
fn redact_metadata(&self, metadata: RunMetadata) -> RunMetadata {
metadata
}
}
#[derive(Debug, Clone)]
pub struct PatternRedactor {
patterns: Vec<String>,
replacement: String,
}
impl PatternRedactor {
pub fn new(patterns: impl IntoIterator<Item = impl Into<String>>) -> Self {
Self {
patterns: patterns.into_iter().map(Into::into).collect(),
replacement: "[REDACTED]".to_string(),
}
}
pub fn with_replacement(mut self, replacement: impl Into<String>) -> Self {
self.replacement = replacement.into();
self
}
fn redact_field(&self, field: &str, text: &str) -> RedactedText {
let mut redacted = text.to_string();
let mut records = Vec::new();
for pattern in &self.patterns {
if pattern.is_empty() || !redacted.contains(pattern) {
continue;
}
redacted = redacted.replace(pattern, &self.replacement);
records.push(RedactionRecord {
field: field.to_string(),
reason: "pattern match".to_string(),
});
}
RedactedText {
text: redacted,
redactions: records,
}
}
}
impl Redactor for PatternRedactor {
fn redact_model_text(&self, text: &str) -> RedactedText {
self.redact_field("model", text)
}
fn redact_display_text(&self, text: &str) -> RedactedText {
self.redact_field("display", text)
}
fn redact_debug_text(&self, text: &str) -> RedactedText {
self.redact_field("debug", text)
}
fn redact_metadata(&self, mut metadata: RunMetadata) -> RunMetadata {
for value in metadata.values_mut() {
redact_json_value(value, &self.patterns, &self.replacement);
}
metadata
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[non_exhaustive]
pub enum AuditEvent {
EffectRequested {
effect_id: String,
kind: EffectKind,
description: String,
risk: RiskLevel,
},
EffectClassified {
effect_id: String,
requested_risk: RiskLevel,
effective_risk: RiskLevel,
reasons: Vec<String>,
},
PolicyDecided {
effect_id: String,
decision: String,
},
ApprovalRequested {
effect_id: String,
reason: String,
},
ApprovalDecided {
effect_id: String,
decision: String,
},
EffectStarted {
effect_id: String,
policy: ExecutionPolicySummary,
},
EffectCompleted {
effect_id: String,
truncated: bool,
},
EffectDenied {
effect_id: String,
reason: String,
},
EffectFailed {
effect_id: String,
reason: String,
},
EffectTimedOut {
effect_id: String,
reason: String,
},
EffectCancelled {
effect_id: String,
reason: String,
},
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ExecutionPolicySummary {
pub sandbox: String,
pub network: String,
pub timeout_ms: Option<u64>,
}
#[async_trait]
pub trait AuditSink: Send + Sync {
async fn record(&self, event: AuditEvent, context: &RunContext) -> Result<(), AuditError>;
}
#[derive(Debug, Default, Clone, Copy)]
pub struct NoopAuditSink;
#[async_trait]
impl AuditSink for NoopAuditSink {
async fn record(&self, _event: AuditEvent, _context: &RunContext) -> Result<(), AuditError> {
Ok(())
}
}
#[derive(Debug, Default, Clone)]
pub struct VecAuditSink {
events: Arc<Mutex<Vec<AuditEvent>>>,
}
impl VecAuditSink {
pub fn new() -> Self {
Self::default()
}
pub fn events(&self) -> Vec<AuditEvent> {
self.events
.lock()
.expect("VecAuditSink lock poisoned")
.clone()
}
}
#[async_trait]
impl AuditSink for VecAuditSink {
async fn record(&self, event: AuditEvent, _context: &RunContext) -> Result<(), AuditError> {
self.events
.lock()
.expect("VecAuditSink lock poisoned")
.push(event);
Ok(())
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[non_exhaustive]
pub enum TranscriptRecord {
RunStarted {
run_id: String,
request: RunRequest,
},
AgentAction {
run_id: String,
action: AgentActionSummary,
},
EffectRequested {
run_id: String,
effect_id: String,
kind: EffectKind,
description: String,
risk: RiskLevel,
},
EffectClassified {
run_id: String,
effect_id: String,
requested_risk: RiskLevel,
effective_risk: RiskLevel,
reasons: Vec<String>,
},
PolicyDecided {
run_id: String,
effect_id: String,
decision: String,
},
ApprovalRequested {
run_id: String,
effect_id: String,
reason: String,
},
ApprovalDecided {
run_id: String,
effect_id: String,
decision: String,
},
ExecutorStarted {
run_id: String,
effect_id: String,
policy: ExecutionPolicySummary,
},
ExecutorCompleted {
run_id: String,
effect_id: String,
status: EffectStatus,
},
ModelObservation {
run_id: String,
request_id: String,
summary: ModelSummary,
},
EffectObservation {
run_id: String,
effect_id: String,
status: EffectStatus,
},
RunCompleted {
run_id: String,
output: RunOutput,
},
RunFailed {
run_id: String,
error: String,
},
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[non_exhaustive]
pub enum AgentActionSummary {
Respond,
RequestModel {
request_id: String,
},
RequestEffect {
effect_id: String,
},
RequestEffects {
effect_ids: Vec<String>,
},
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ModelSummary {
pub content_bytes: usize,
pub tool_calls: usize,
}
#[async_trait]
pub trait TranscriptStore: Send + Sync {
async fn append(
&self,
record: TranscriptRecord,
context: &RunContext,
) -> Result<(), TranscriptError>;
}
#[derive(Debug, Default, Clone, Copy)]
pub struct NoopTranscriptStore;
#[async_trait]
impl TranscriptStore for NoopTranscriptStore {
async fn append(
&self,
_record: TranscriptRecord,
_context: &RunContext,
) -> Result<(), TranscriptError> {
Ok(())
}
}
#[derive(Debug, Default, Clone)]
pub struct VecTranscriptStore {
records: Arc<Mutex<Vec<TranscriptRecord>>>,
}
impl VecTranscriptStore {
pub fn new() -> Self {
Self::default()
}
pub fn records(&self) -> Vec<TranscriptRecord> {
self.records
.lock()
.expect("VecTranscriptStore lock poisoned")
.clone()
}
}
#[async_trait]
impl TranscriptStore for VecTranscriptStore {
async fn append(
&self,
record: TranscriptRecord,
_context: &RunContext,
) -> Result<(), TranscriptError> {
self.records
.lock()
.expect("VecTranscriptStore lock poisoned")
.push(record);
Ok(())
}
}
pub struct BasicHarness<E, P, A, S, T> {
executor: E,
policy: P,
approval: A,
audit: S,
transcript: T,
classifier: DefaultRiskClassifier,
redactor: Arc<dyn Redactor>,
session_approvals: Mutex<Vec<SessionApproval>>,
config: HarnessConfig,
}
impl<E, P, A, S, T> fmt::Debug for BasicHarness<E, P, A, S, T>
where
E: fmt::Debug,
P: fmt::Debug,
A: fmt::Debug,
S: fmt::Debug,
T: fmt::Debug,
{
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("BasicHarness")
.field("executor", &self.executor)
.field("policy", &self.policy)
.field("approval", &self.approval)
.field("audit", &self.audit)
.field("transcript", &self.transcript)
.field("classifier", &self.classifier)
.field("redactor", &"dyn Redactor")
.field("session_approvals", &self.session_approvals)
.field("config", &self.config)
.finish()
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct SessionApproval {
run_id: String,
kind: EffectKindKey,
risk_ceiling: RiskLevel,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(default)]
#[non_exhaustive]
pub struct HarnessConfig {
pub(crate) default_sandbox: SandboxPolicy,
pub(crate) default_network: NetworkPolicy,
pub(crate) default_timeout: Duration,
pub(crate) output_limit: OutputLimit,
pub(crate) fail_closed_on_audit_error: bool,
pub(crate) fail_closed_on_transcript_error: bool,
}
impl Default for HarnessConfig {
fn default() -> Self {
Self {
default_sandbox: SandboxPolicy::ReadOnly,
default_network: NetworkPolicy::Deny,
default_timeout: Duration::from_secs(30),
output_limit: OutputLimit::default(),
fail_closed_on_audit_error: true,
fail_closed_on_transcript_error: false,
}
}
}
impl HarnessConfig {
pub fn new() -> Self {
Self::default()
}
pub fn default_sandbox(&self) -> &SandboxPolicy {
&self.default_sandbox
}
pub fn with_default_sandbox(mut self, default_sandbox: SandboxPolicy) -> Self {
self.default_sandbox = default_sandbox;
self
}
pub fn default_network(&self) -> &NetworkPolicy {
&self.default_network
}
pub fn with_default_network(mut self, default_network: NetworkPolicy) -> Self {
self.default_network = default_network;
self
}
pub fn default_timeout(&self) -> Duration {
self.default_timeout
}
pub fn with_default_timeout(mut self, default_timeout: Duration) -> Self {
self.default_timeout = default_timeout;
self
}
pub fn output_limit(&self) -> &OutputLimit {
&self.output_limit
}
pub fn with_output_limit(mut self, output_limit: OutputLimit) -> Self {
self.output_limit = output_limit;
self
}
pub fn fail_closed_on_audit_error(&self) -> bool {
self.fail_closed_on_audit_error
}
pub fn with_fail_closed_on_audit_error(mut self, fail_closed_on_audit_error: bool) -> Self {
self.fail_closed_on_audit_error = fail_closed_on_audit_error;
self
}
pub fn fail_closed_on_transcript_error(&self) -> bool {
self.fail_closed_on_transcript_error
}
pub fn with_fail_closed_on_transcript_error(
mut self,
fail_closed_on_transcript_error: bool,
) -> Self {
self.fail_closed_on_transcript_error = fail_closed_on_transcript_error;
self
}
}
impl
BasicHarness<
NoopEffectExecutor,
DefaultPolicyEngine,
AlwaysDenyApprovalBroker,
NoopAuditSink,
NoopTranscriptStore,
>
{
pub fn noop() -> Self {
Self::new(
NoopEffectExecutor,
DefaultPolicyEngine,
AlwaysDenyApprovalBroker,
NoopAuditSink,
NoopTranscriptStore,
)
}
}
impl<E, P, A, S, T> BasicHarness<E, P, A, S, T>
where
E: EffectExecutor,
P: PolicyEngine,
A: ApprovalBroker,
S: AuditSink,
T: TranscriptStore,
{
pub fn new(executor: E, policy: P, approval: A, audit: S, transcript: T) -> Self {
Self {
executor,
policy,
approval,
audit,
transcript,
classifier: DefaultRiskClassifier,
redactor: Arc::new(NoopRedactor),
session_approvals: Mutex::new(Vec::new()),
config: HarnessConfig::default(),
}
}
pub fn with_config(mut self, config: HarnessConfig) -> Self {
self.config = config;
self
}
pub fn with_noop_redactor(mut self) -> Self {
self.redactor = Arc::new(NoopRedactor);
self
}
pub fn with_redactor(mut self, redactor: impl Redactor + 'static) -> Self {
self.redactor = Arc::new(redactor);
self
}
async fn audit(&self, event: AuditEvent, context: &RunContext) -> Result<(), HarnessError> {
match self.audit.record(event, context).await {
Ok(()) => Ok(()),
Err(error) if self.config.fail_closed_on_audit_error => Err(error.into()),
Err(_) => Ok(()),
}
}
async fn transcript(
&self,
record: TranscriptRecord,
context: &RunContext,
) -> Result<(), HarnessError> {
match self.transcript.append(record, context).await {
Ok(()) => Ok(()),
Err(error) if self.config.fail_closed_on_transcript_error => Err(error.into()),
Err(_) => Ok(()),
}
}
}
#[async_trait]
impl<E, P, A, S, T> Harness for BasicHarness<E, P, A, S, T>
where
E: EffectExecutor,
P: PolicyEngine,
A: ApprovalBroker,
S: AuditSink,
T: TranscriptStore,
{
async fn execute(
&self,
request: EffectRequest,
context: &RunContext,
) -> Result<EffectObservation, HarnessError> {
check_run_context(context)?;
validate_effect_request(&request)?;
self.audit(
AuditEvent::EffectRequested {
effect_id: request.id.clone(),
kind: request.kind.clone(),
description: request.description.clone(),
risk: request.risk,
},
context,
)
.await?;
self.transcript(
TranscriptRecord::EffectRequested {
run_id: context.run_id.clone(),
effect_id: request.id.clone(),
kind: request.kind.clone(),
description: request.description.clone(),
risk: request.risk,
},
context,
)
.await?;
let classified = self.classifier.classify(request, context).await?;
self.audit(
AuditEvent::EffectClassified {
effect_id: classified.request.id.clone(),
requested_risk: classified.requested_risk,
effective_risk: classified.effective_risk,
reasons: classified.reasons.clone(),
},
context,
)
.await?;
self.transcript(
TranscriptRecord::EffectClassified {
run_id: context.run_id.clone(),
effect_id: classified.request.id.clone(),
requested_risk: classified.requested_risk,
effective_risk: classified.effective_risk,
reasons: classified.reasons.clone(),
},
context,
)
.await?;
let execution_policy = self.execution_policy(&classified, context);
let decision = self.policy.evaluate(&classified, context).await?;
let decision_summary = policy_decision_summary(&decision);
self.audit(
AuditEvent::PolicyDecided {
effect_id: classified.request.id.clone(),
decision: decision_summary.clone(),
},
context,
)
.await?;
self.transcript(
TranscriptRecord::PolicyDecided {
run_id: context.run_id.clone(),
effect_id: classified.request.id.clone(),
decision: decision_summary,
},
context,
)
.await?;
match decision {
PolicyDecision::Allow => {}
PolicyDecision::Deny { reason } => {
return self
.denied_observation(classified.request, reason, context)
.await;
}
PolicyDecision::RequireApproval { reason } => {
if self.is_session_approved(&classified, context) {
self.audit(
AuditEvent::ApprovalDecided {
effect_id: classified.request.id.clone(),
decision: "allow for session".to_string(),
},
context,
)
.await?;
self.transcript(
TranscriptRecord::ApprovalDecided {
run_id: context.run_id.clone(),
effect_id: classified.request.id.clone(),
decision: "allow for session".to_string(),
},
context,
)
.await?;
} else {
let approval_request = ApprovalRequest {
run_id: context.run_id.clone(),
effect_id: classified.request.id.clone(),
kind: classified.request.kind.clone(),
description: classified.request.description.clone(),
risk: classified.effective_risk,
reason: reason.clone(),
payload_summary: payload_summary(&classified.request.payload),
sandbox: execution_policy.sandbox.clone(),
network: execution_policy.network.clone(),
metadata: classified.request.metadata.clone(),
};
self.audit(
AuditEvent::ApprovalRequested {
effect_id: classified.request.id.clone(),
reason: reason.clone(),
},
context,
)
.await?;
self.transcript(
TranscriptRecord::ApprovalRequested {
run_id: context.run_id.clone(),
effect_id: classified.request.id.clone(),
reason,
},
context,
)
.await?;
let approval = self.approval.approve(approval_request, context).await?;
let approval_summary = approval_decision_summary(&approval);
self.audit(
AuditEvent::ApprovalDecided {
effect_id: classified.request.id.clone(),
decision: approval_summary.clone(),
},
context,
)
.await?;
self.transcript(
TranscriptRecord::ApprovalDecided {
run_id: context.run_id.clone(),
effect_id: classified.request.id.clone(),
decision: approval_summary,
},
context,
)
.await?;
if let ApprovalDecision::Deny { reason } = approval {
return self
.denied_observation(classified.request, reason, context)
.await;
}
if matches!(approval, ApprovalDecision::AllowForSession) {
self.remember_session_approval(&classified, context);
}
}
}
}
self.audit(
AuditEvent::EffectStarted {
effect_id: classified.request.id.clone(),
policy: ExecutionPolicySummary::from_policy(&execution_policy),
},
context,
)
.await?;
self.transcript(
TranscriptRecord::ExecutorStarted {
run_id: context.run_id.clone(),
effect_id: classified.request.id.clone(),
policy: ExecutionPolicySummary::from_policy(&execution_policy),
},
context,
)
.await?;
let effect_id = classified.request.id.clone();
let execution = run_executor_with_context(
&self.executor,
&classified.request,
&execution_policy,
context,
)
.await;
match execution {
Ok(raw) => {
let limited =
limit_and_redact(raw, &self.config.output_limit, self.redactor.as_ref());
self.audit(
AuditEvent::EffectCompleted {
effect_id: effect_id.clone(),
truncated: limited.truncated,
},
context,
)
.await?;
self.transcript(
TranscriptRecord::ExecutorCompleted {
run_id: context.run_id.clone(),
effect_id: effect_id.clone(),
status: EffectStatus::Succeeded,
},
context,
)
.await?;
let mut observation_metadata = RunMetadata::new();
observation_metadata.insert(
"truncated".to_string(),
serde_json::json!(limited.truncated),
);
observation_metadata.insert(
"redactions_applied".to_string(),
serde_json::json!(limited.redactions.len()),
);
let observation = EffectObservation {
effect_id: effect_id.clone(),
status: EffectStatus::Succeeded,
output: limited.output,
metadata: observation_metadata,
};
self.transcript(
TranscriptRecord::EffectObservation {
run_id: context.run_id.clone(),
effect_id,
status: observation.status.clone(),
},
context,
)
.await?;
Ok(observation)
}
Err(ExecutionError::TimedOut(reason)) => {
self.audit(
AuditEvent::EffectTimedOut {
effect_id: effect_id.clone(),
reason: reason.clone(),
},
context,
)
.await?;
self.transcript(
TranscriptRecord::ExecutorCompleted {
run_id: context.run_id.clone(),
effect_id: effect_id.clone(),
status: EffectStatus::TimedOut,
},
context,
)
.await?;
self.terminal_observation(
effect_id,
EffectStatus::TimedOut,
format!("effect timed out: {reason}"),
context,
)
.await
}
Err(ExecutionError::Denied(reason)) => {
self.audit(
AuditEvent::EffectDenied {
effect_id: effect_id.clone(),
reason: reason.clone(),
},
context,
)
.await?;
self.transcript(
TranscriptRecord::ExecutorCompleted {
run_id: context.run_id.clone(),
effect_id: effect_id.clone(),
status: EffectStatus::Denied,
},
context,
)
.await?;
self.terminal_observation(
effect_id,
EffectStatus::Denied,
format!("effect denied: {reason}"),
context,
)
.await
}
Err(ExecutionError::Cancelled(reason)) => {
self.audit(
AuditEvent::EffectCancelled {
effect_id: effect_id.clone(),
reason: reason.clone(),
},
context,
)
.await?;
self.transcript(
TranscriptRecord::ExecutorCompleted {
run_id: context.run_id.clone(),
effect_id: effect_id.clone(),
status: EffectStatus::Cancelled,
},
context,
)
.await?;
self.terminal_observation(
effect_id,
EffectStatus::Cancelled,
format!("effect cancelled: {reason}"),
context,
)
.await
}
Err(error) => {
let reason = error.to_string();
self.audit(
AuditEvent::EffectFailed {
effect_id: effect_id.clone(),
reason: reason.clone(),
},
context,
)
.await?;
self.transcript(
TranscriptRecord::ExecutorCompleted {
run_id: context.run_id.clone(),
effect_id: effect_id.clone(),
status: EffectStatus::Failed,
},
context,
)
.await?;
self.terminal_observation(
effect_id,
EffectStatus::Failed,
format!("effect failed: {reason}"),
context,
)
.await
}
}
}
}
impl<E, P, A, S, T> BasicHarness<E, P, A, S, T>
where
E: EffectExecutor,
P: PolicyEngine,
A: ApprovalBroker,
S: AuditSink,
T: TranscriptStore,
{
fn execution_policy(
&self,
classified: &ClassifiedEffect,
context: &RunContext,
) -> ExecutionPolicy {
let mut timeout = Some(self.config.default_timeout);
if let Some(request_timeout) = classified.request.timeout {
timeout = Some(timeout.map_or(request_timeout, |default| default.min(request_timeout)));
}
if let Some(remaining) = context.remaining() {
timeout = Some(timeout.map_or(remaining, |current| current.min(remaining)));
}
ExecutionPolicy {
sandbox: self.config.default_sandbox.clone(),
network: self.config.default_network.clone(),
timeout,
output_limit: self.config.output_limit.clone(),
}
}
fn is_session_approved(&self, classified: &ClassifiedEffect, context: &RunContext) -> bool {
let kind = EffectKindKey::from(classified.request.kind.clone());
let approvals = self
.session_approvals
.lock()
.expect("BasicHarness session approval lock poisoned");
approvals.iter().any(|approval| {
approval.run_id == context.run_id
&& approval.kind == kind
&& risk_rank(classified.effective_risk) <= risk_rank(approval.risk_ceiling)
})
}
fn remember_session_approval(&self, classified: &ClassifiedEffect, context: &RunContext) {
let approval = SessionApproval {
run_id: context.run_id.clone(),
kind: EffectKindKey::from(classified.request.kind.clone()),
risk_ceiling: classified.effective_risk,
};
let mut approvals = self
.session_approvals
.lock()
.expect("BasicHarness session approval lock poisoned");
if !approvals.contains(&approval) {
approvals.push(approval);
}
}
async fn denied_observation(
&self,
request: EffectRequest,
reason: String,
context: &RunContext,
) -> Result<EffectObservation, HarnessError>
where
S: AuditSink,
T: TranscriptStore,
{
self.audit(
AuditEvent::EffectDenied {
effect_id: request.id.clone(),
reason: reason.clone(),
},
context,
)
.await?;
self.terminal_observation(
request.id,
EffectStatus::Denied,
format!("effect denied: {reason}"),
context,
)
.await
}
async fn terminal_observation(
&self,
effect_id: String,
status: EffectStatus,
observation_for_model: String,
context: &RunContext,
) -> Result<EffectObservation, HarnessError>
where
T: TranscriptStore,
{
let observation = EffectObservation {
effect_id: effect_id.clone(),
status,
output: EffectOutput::text(observation_for_model),
metadata: RunMetadata::new(),
};
self.transcript(
TranscriptRecord::EffectObservation {
run_id: context.run_id.clone(),
effect_id,
status: observation.status.clone(),
},
context,
)
.await?;
Ok(observation)
}
}
#[derive(Debug, thiserror::Error)]
#[non_exhaustive]
pub enum HarnessError {
#[error("invalid effect request: {0}")]
InvalidRequest(String),
#[error("policy error: {0}")]
Policy(String),
#[error("approval error: {0}")]
Approval(#[from] ApprovalError),
#[error("execution error: {0}")]
Execution(#[from] ExecutionError),
#[error("audit error: {0}")]
Audit(#[from] AuditError),
#[error("transcript error: {0}")]
Transcript(#[from] TranscriptError),
#[error("run cancelled")]
Cancelled,
#[error("run deadline exceeded")]
DeadlineExceeded,
}
#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
#[non_exhaustive]
pub enum ApprovalError {
#[error("{0}")]
Broker(String),
}
#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
#[non_exhaustive]
pub enum ExecutionError {
#[error("unsupported effect: {0}")]
Unsupported(String),
#[error("failed: {0}")]
Failed(String),
#[error("denied: {0}")]
Denied(String),
#[error("timed out: {0}")]
TimedOut(String),
#[error("cancelled: {0}")]
Cancelled(String),
}
#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
#[non_exhaustive]
pub enum AuditError {
#[error("{0}")]
Sink(String),
}
#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
#[non_exhaustive]
pub enum TranscriptError {
#[error("{0}")]
Store(String),
}
#[derive(Debug, thiserror::Error)]
#[non_exhaustive]
pub enum HarnessRuntimeError {
#[error("agent error: {0}")]
Agent(#[from] AgentError),
#[error("provider error: {0}")]
Provider(#[from] ProviderError),
#[error("harness error: {0}")]
Harness(#[from] HarnessError),
#[error("too many agent steps: {0}")]
TooManyAgentSteps(usize),
}
impl ExecutionPolicySummary {
fn from_policy(policy: &ExecutionPolicy) -> Self {
Self {
sandbox: format!("{:?}", policy.sandbox),
network: format!("{:?}", policy.network),
timeout_ms: policy
.timeout
.map(|timeout| timeout.as_millis().min(u128::from(u64::MAX)) as u64),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
enum EffectKindKey {
ReadFile,
WriteFile,
ApplyPatch,
Search,
ExecuteCommand,
Git,
Network,
Browser,
Mcp,
Custom(String),
}
impl From<EffectKind> for EffectKindKey {
fn from(kind: EffectKind) -> Self {
match kind {
EffectKind::ReadFile => Self::ReadFile,
EffectKind::WriteFile => Self::WriteFile,
EffectKind::ApplyPatch => Self::ApplyPatch,
EffectKind::Search => Self::Search,
EffectKind::ExecuteCommand => Self::ExecuteCommand,
EffectKind::Git => Self::Git,
EffectKind::Network => Self::Network,
EffectKind::Browser => Self::Browser,
EffectKind::Mcp => Self::Mcp,
EffectKind::Custom(kind) => Self::Custom(kind),
_ => Self::Custom("unknown".to_string()),
}
}
}
fn validate_effect_request(request: &EffectRequest) -> Result<(), HarnessError> {
if request.id.trim().is_empty() {
return Err(HarnessError::InvalidRequest(
"effect id must not be empty".to_string(),
));
}
if request.description.trim().is_empty() {
return Err(HarnessError::InvalidRequest(
"effect description must not be empty".to_string(),
));
}
Ok(())
}
fn check_run_context(context: &RunContext) -> Result<(), HarnessError> {
if context.is_cancelled() {
Err(HarnessError::Cancelled)
} else if context.is_expired() {
Err(HarnessError::DeadlineExceeded)
} else {
Ok(())
}
}
fn is_terminal_context_error(error: &HarnessError) -> bool {
matches!(
error,
HarnessError::Cancelled | HarnessError::DeadlineExceeded
)
}
fn observation_from_harness_error(effect_id: String, error: HarnessError) -> EffectObservation {
let mut metadata = RunMetadata::new();
metadata.insert("harness_error".to_string(), serde_json::json!(true));
EffectObservation {
effect_id,
status: EffectStatus::Failed,
output: EffectOutput::text(format!("effect failed: {error}")),
metadata,
}
}
async fn run_executor_with_context<E>(
executor: &E,
request: &EffectRequest,
policy: &ExecutionPolicy,
context: &RunContext,
) -> Result<RawEffectOutput, ExecutionError>
where
E: EffectExecutor,
{
if context.is_cancelled() {
return Err(ExecutionError::Cancelled("run cancelled".to_string()));
}
if context.is_expired() {
return Err(ExecutionError::TimedOut(
"run deadline exceeded".to_string(),
));
}
match policy.timeout {
Some(timeout) if timeout.is_zero() => {
Err(ExecutionError::TimedOut("timeout elapsed".to_string()))
}
Some(timeout) => {
tokio::select! {
_ = context.cancellation.cancelled() => {
Err(ExecutionError::Cancelled("run cancelled".to_string()))
}
_ = tokio::time::sleep(timeout) => {
Err(ExecutionError::TimedOut(format!("exceeded {timeout:?}")))
}
output = executor.execute(request, policy, context) => output,
}
}
None => {
tokio::select! {
_ = context.cancellation.cancelled() => {
Err(ExecutionError::Cancelled("run cancelled".to_string()))
}
output = executor.execute(request, policy, context) => output,
}
}
}
}
fn limit_and_redact(
raw: RawEffectOutput,
limit: &OutputLimit,
redactor: &(impl Redactor + ?Sized),
) -> LimitedOutput {
let model_redacted = redactor.redact_model_text(&raw.observation_for_model);
let (mut model_text, model_truncated) = truncate_with_marker(
model_redacted.text,
limit.model_bytes,
"\n[output truncated]",
);
let mut redactions = model_redacted.redactions;
let mut truncated = model_truncated;
if model_truncated && !model_text.contains("[output truncated]") {
model_text.push_str("\n[output truncated]");
}
let display = raw.display.map(|display| {
let display_redacted = redactor.redact_display_text(&display.content);
redactions.extend(display_redacted.redactions);
let (content, was_truncated) =
truncate_with_marker(display_redacted.text, limit.display_bytes, "\n[truncated]");
truncated |= was_truncated;
DisplayOutput {
content,
metadata: redactor.redact_metadata(display.metadata),
..display
}
});
let debug = raw.debug.map(|debug| {
let debug_redacted = redactor.redact_debug_text(&debug);
redactions.extend(debug_redacted.redactions);
let (debug, was_truncated) =
truncate_with_marker(debug_redacted.text, limit.debug_bytes, "\n[truncated]");
truncated |= was_truncated;
debug
});
LimitedOutput {
output: EffectOutput::text(model_text)
.with_artifacts(raw.artifacts)
.with_metadata(redactor.redact_metadata(raw.metadata)),
truncated,
redactions,
debug,
}
.with_display(display)
}
impl LimitedOutput {
fn with_display(mut self, display: Option<DisplayOutput>) -> Self {
self.output.display = display;
self
}
}
fn truncate_with_marker(mut text: String, limit: usize, marker: &str) -> (String, bool) {
if text.len() <= limit {
return (text, false);
}
if limit == 0 {
return (String::new(), true);
}
let marker_len = marker.len().min(limit);
let keep = limit.saturating_sub(marker_len);
let mut end = keep.min(text.len());
while !text.is_char_boundary(end) {
end -= 1;
}
text.truncate(end);
if marker_len == marker.len() {
text.push_str(marker);
}
(text, true)
}
fn risk_floor_for_kind(kind: &EffectKind) -> RiskLevel {
match kind {
EffectKind::ReadFile | EffectKind::Search => RiskLevel::Low,
EffectKind::WriteFile
| EffectKind::ApplyPatch
| EffectKind::ExecuteCommand
| EffectKind::Git
| EffectKind::Network
| EffectKind::Browser
| EffectKind::Mcp
| EffectKind::Custom(_) => RiskLevel::Medium,
_ => RiskLevel::Medium,
}
}
fn max_risk(left: RiskLevel, right: RiskLevel) -> RiskLevel {
if risk_rank(left) >= risk_rank(right) {
left
} else {
right
}
}
fn risk_rank(risk: RiskLevel) -> u8 {
match risk {
RiskLevel::Low => 0,
RiskLevel::Medium => 1,
RiskLevel::High => 2,
RiskLevel::Critical => 3,
_ => 3,
}
}
fn contains_high_pattern(text: &str) -> bool {
["sudo", "force push", "--force", " outside workspace"]
.iter()
.any(|pattern| text.contains(pattern))
}
fn contains_critical_pattern(text: &str) -> bool {
[
"rm -rf /",
"mkfs",
"dd if=",
"shutdown",
"reboot",
"chmod -r 777 /",
]
.iter()
.any(|pattern| text.contains(pattern))
}
fn payload_summary(payload: &serde_json::Value) -> String {
payload.to_string().chars().take(512).collect()
}
fn policy_decision_summary(decision: &PolicyDecision) -> String {
match decision {
PolicyDecision::Allow => "allow".to_string(),
PolicyDecision::Deny { reason } => format!("deny: {reason}"),
PolicyDecision::RequireApproval { reason } => format!("require approval: {reason}"),
}
}
fn approval_decision_summary(decision: &ApprovalDecision) -> String {
match decision {
ApprovalDecision::AllowOnce => "allow once".to_string(),
ApprovalDecision::AllowForSession => "allow for session".to_string(),
ApprovalDecision::Deny { reason } => format!("deny: {reason}"),
}
}
fn redact_json_value(value: &mut serde_json::Value, patterns: &[String], replacement: &str) {
match value {
serde_json::Value::String(text) => {
for pattern in patterns {
if !pattern.is_empty() && text.contains(pattern) {
*text = text.replace(pattern, replacement);
}
}
}
serde_json::Value::Array(values) => {
for value in values {
redact_json_value(value, patterns, replacement);
}
}
serde_json::Value::Object(map) => {
for value in map.values_mut() {
redact_json_value(value, patterns, replacement);
}
}
serde_json::Value::Null | serde_json::Value::Bool(_) | serde_json::Value::Number(_) => {}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::agent::{AgentKernel, ModelRequest};
use crate::message::Message;
use crate::provider::{FakeProvider, FakeReply, ProviderError};
use crate::tool::{Tool, ToolContext, ToolError, ToolRegistry, ToolResult, ToolSchema};
use serde_json::json;
struct SingleModelKernel;
#[async_trait]
impl AgentKernel for SingleModelKernel {
async fn start(
&mut self,
_request: RunRequest,
_context: &RunContext,
) -> Result<AgentAction, AgentError> {
Ok(AgentAction::RequestModel {
request: ModelRequest::new("model-1", Default::default()),
})
}
async fn observe(
&mut self,
observation: Observation,
context: &RunContext,
) -> Result<AgentAction, AgentError> {
let Observation::Model(observation) = observation else {
return Err(AgentError::InvalidStep(
"expected model observation".to_string(),
));
};
let answer = match &observation.response.message {
Message::Assistant { content, .. } => content.clone(),
_ => String::new(),
};
Ok(AgentAction::Respond {
output: RunOutput {
run_id: context.run_id.clone(),
answer,
summary: Default::default(),
final_message: observation.response.message,
artifacts: Vec::new(),
metadata: RunMetadata::new(),
},
})
}
}
struct EffectKernel {
observed: bool,
}
#[async_trait]
impl AgentKernel for EffectKernel {
async fn start(
&mut self,
_request: RunRequest,
_context: &RunContext,
) -> Result<AgentAction, AgentError> {
Ok(AgentAction::RequestEffect {
request: EffectRequest::new(EffectKind::ReadFile, "read config", json!({}))
.with_id("effect-1"),
})
}
async fn observe(
&mut self,
observation: Observation,
context: &RunContext,
) -> Result<AgentAction, AgentError> {
let Observation::Effect(observation) = observation else {
return Err(AgentError::InvalidStep(
"expected effect observation".to_string(),
));
};
self.observed = true;
Ok(AgentAction::Respond {
output: RunOutput {
run_id: context.run_id.clone(),
answer: observation.output.observation_for_model,
summary: Default::default(),
final_message: Message::assistant("done"),
artifacts: Vec::new(),
metadata: RunMetadata::new(),
},
})
}
}
#[tokio::test]
async fn runtime_drives_model_request() {
let provider = FakeProvider::new([FakeReply::Text("hello".to_string())]);
let harness = BasicHarness::noop();
let runtime = HarnessRuntime::new(provider, harness);
let output = runtime
.run(
&mut SingleModelKernel,
RunRequest::text("hi"),
RunContext::new("run-model"),
)
.await
.unwrap();
assert_eq!(output.answer, "hello");
}
#[tokio::test]
async fn runtime_drives_effect_request() {
let provider = FakeProvider::new([]);
let executor = StaticEffectExecutor::new()
.with_output("effect-1", RawEffectOutput::text("file contents"));
let audit = VecAuditSink::new();
let harness = BasicHarness::new(
executor,
DefaultPolicyEngine,
AlwaysAllowApprovalBroker,
audit.clone(),
NoopTranscriptStore,
);
let runtime = HarnessRuntime::new(provider, harness);
let output = runtime
.run(
&mut EffectKernel { observed: false },
RunRequest::text("read"),
RunContext::new("run-effect"),
)
.await
.unwrap();
assert_eq!(output.answer, "file contents");
assert!(audit
.events()
.iter()
.any(|event| matches!(event, AuditEvent::EffectCompleted { effect_id, .. } if effect_id == "effect-1")));
}
#[tokio::test]
async fn runtime_returns_provider_error() {
let provider = FakeProvider::new([FakeReply::Error(ProviderError::Api {
status: 500,
code: None,
message: "provider down".to_string(),
})]);
let harness = BasicHarness::noop();
let runtime = HarnessRuntime::new(provider, harness);
let err = runtime
.run(
&mut SingleModelKernel,
RunRequest::text("hi"),
RunContext::new("run-provider-error"),
)
.await
.unwrap_err();
assert!(matches!(err, HarnessRuntimeError::Provider(_)));
}
struct LoopKernel;
#[async_trait]
impl AgentKernel for LoopKernel {
async fn start(
&mut self,
_request: RunRequest,
_context: &RunContext,
) -> Result<AgentAction, AgentError> {
Ok(AgentAction::RequestModel {
request: ModelRequest::new("model-1", Default::default()),
})
}
async fn observe(
&mut self,
_observation: Observation,
_context: &RunContext,
) -> Result<AgentAction, AgentError> {
Ok(AgentAction::RequestModel {
request: ModelRequest::new("model-loop", Default::default()),
})
}
}
#[tokio::test]
async fn runtime_limits_agent_steps() {
let provider = FakeProvider::new([
FakeReply::Text("one".to_string()),
FakeReply::Text("two".to_string()),
FakeReply::Text("three".to_string()),
]);
let harness = BasicHarness::noop();
let runtime = HarnessRuntime::new(provider, harness).with_config(HarnessRuntimeConfig {
max_agent_steps: 1,
..Default::default()
});
let err = runtime
.run(
&mut LoopKernel,
RunRequest::text("hi"),
RunContext::new("run-step-limit"),
)
.await
.unwrap_err();
assert!(matches!(err, HarnessRuntimeError::TooManyAgentSteps(1)));
}
#[tokio::test]
async fn react_kernel_runs_effect_through_runtime() {
let provider = FakeProvider::new([
FakeReply::ToolCalls {
content: String::new(),
calls: vec![crate::ToolCall {
id: "call-1".to_string(),
name: "read_file".to_string(),
arguments: "{}".to_string(),
}],
},
FakeReply::Text("done after observation".to_string()),
]);
let executor =
StaticEffectExecutor::new().with_output("effect-1", RawEffectOutput::text("observed"));
let harness = BasicHarness::new(
executor,
DefaultPolicyEngine,
AlwaysAllowApprovalBroker,
NoopAuditSink,
NoopTranscriptStore,
);
let runtime = HarnessRuntime::new(provider, harness);
let mut registry = ToolRegistry::new();
registry.register(EffectTool);
let mut kernel = crate::ReActAgent::kernel(registry, "");
let output = runtime
.run(
&mut kernel,
RunRequest::text("read"),
RunContext::new("run-react-runtime"),
)
.await
.unwrap();
assert_eq!(output.answer, "done after observation");
}
#[tokio::test]
async fn basic_harness_denies_critical_effect() {
let harness = BasicHarness::new(
StaticEffectExecutor::new(),
DefaultPolicyEngine,
AlwaysAllowApprovalBroker,
NoopAuditSink,
NoopTranscriptStore,
);
let observation = harness
.execute(
EffectRequest::new(
EffectKind::ExecuteCommand,
"run rm -rf /",
json!({"cmd": "rm -rf /"}),
)
.with_id("critical"),
&RunContext::new("run-deny"),
)
.await
.unwrap();
assert_eq!(observation.status, EffectStatus::Denied);
assert!(
observation
.output
.observation_for_model
.contains("effect denied")
);
}
#[tokio::test]
async fn basic_harness_batch_can_mix_denied_and_succeeded() {
let harness = BasicHarness::new(
StaticEffectExecutor::new()
.with_output("read-ok", RawEffectOutput::text("read succeeded")),
DefaultPolicyEngine,
AlwaysAllowApprovalBroker,
NoopAuditSink,
NoopTranscriptStore,
);
let observations = harness
.execute_batch(
vec![
EffectRequest::new(EffectKind::ReadFile, "read", json!({})).with_id("read-ok"),
EffectRequest::new(
EffectKind::ExecuteCommand,
"run rm -rf /",
json!({"cmd": "rm -rf /"}),
)
.with_id("deny-critical"),
],
&RunContext::new("run-batch"),
)
.await
.unwrap();
assert_eq!(observations[0].status, EffectStatus::Succeeded);
assert_eq!(observations[1].status, EffectStatus::Denied);
}
#[derive(Debug, Clone)]
struct ConcurrencyHarness {
in_flight: Arc<std::sync::atomic::AtomicUsize>,
max_seen: Arc<std::sync::atomic::AtomicUsize>,
}
#[async_trait]
impl Harness for ConcurrencyHarness {
async fn execute(
&self,
request: EffectRequest,
_context: &RunContext,
) -> Result<EffectObservation, HarnessError> {
let current = self
.in_flight
.fetch_add(1, std::sync::atomic::Ordering::SeqCst)
+ 1;
self.max_seen
.fetch_max(current, std::sync::atomic::Ordering::SeqCst);
let delay = if request.id == "effect-1" { 30 } else { 5 };
tokio::time::sleep(Duration::from_millis(delay)).await;
self.in_flight
.fetch_sub(1, std::sync::atomic::Ordering::SeqCst);
Ok(EffectObservation::succeeded(
request.id.clone(),
format!("done {}", request.id),
))
}
}
#[tokio::test]
async fn runtime_batch_honors_concurrency_limit_and_request_order() {
let in_flight = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let max_seen = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let harness = ConcurrencyHarness {
in_flight,
max_seen: max_seen.clone(),
};
let runtime = HarnessRuntime::new(FakeProvider::new([]), harness)
.with_config(HarnessRuntimeConfig::default().with_max_effect_batch_concurrency(2));
let observations = runtime
.execute_effect_batch(
vec![
EffectRequest::new(EffectKind::ReadFile, "one", json!({})).with_id("effect-1"),
EffectRequest::new(EffectKind::ReadFile, "two", json!({})).with_id("effect-2"),
EffectRequest::new(EffectKind::ReadFile, "three", json!({}))
.with_id("effect-3"),
],
&RunContext::new("run-batch-runtime"),
)
.await
.unwrap();
assert_eq!(
observations
.iter()
.map(|observation| observation.effect_id.as_str())
.collect::<Vec<_>>(),
["effect-1", "effect-2", "effect-3"]
);
assert_eq!(max_seen.load(std::sync::atomic::Ordering::SeqCst), 2);
}
#[derive(Debug, Clone)]
struct ErroringHarness {
seen: Arc<Mutex<Vec<String>>>,
}
#[async_trait]
impl Harness for ErroringHarness {
async fn execute(
&self,
request: EffectRequest,
_context: &RunContext,
) -> Result<EffectObservation, HarnessError> {
self.seen
.lock()
.expect("seen lock poisoned")
.push(request.id.clone());
if request.id == "bad" {
return Err(HarnessError::Execution(ExecutionError::Failed(
"boom".to_string(),
)));
}
Ok(EffectObservation::succeeded(request.id, "ok"))
}
}
#[tokio::test]
async fn runtime_batch_fail_fast_stops_after_first_harness_error() {
let seen = Arc::new(Mutex::new(Vec::new()));
let runtime = HarnessRuntime::new(
FakeProvider::new([]),
ErroringHarness { seen: seen.clone() },
)
.with_config(
HarnessRuntimeConfig::default()
.with_max_effect_batch_concurrency(1)
.with_fail_fast_effect_batches(true),
);
let error = runtime
.execute_effect_batch(
vec![
EffectRequest::new(EffectKind::ReadFile, "bad", json!({})).with_id("bad"),
EffectRequest::new(EffectKind::ReadFile, "later", json!({})).with_id("later"),
],
&RunContext::new("run-batch-fail-fast"),
)
.await
.unwrap_err();
assert!(matches!(error, HarnessError::Execution(_)));
assert_eq!(seen.lock().expect("seen lock poisoned").as_slice(), ["bad"]);
}
#[tokio::test]
async fn runtime_batch_fail_open_collects_error_observation_and_continues() {
let runtime = HarnessRuntime::new(
FakeProvider::new([]),
ErroringHarness {
seen: Arc::new(Mutex::new(Vec::new())),
},
)
.with_config(
HarnessRuntimeConfig::default()
.with_max_effect_batch_concurrency(1)
.with_fail_fast_effect_batches(false),
);
let observations = runtime
.execute_effect_batch(
vec![
EffectRequest::new(EffectKind::ReadFile, "bad", json!({})).with_id("bad"),
EffectRequest::new(EffectKind::ReadFile, "later", json!({})).with_id("later"),
],
&RunContext::new("run-batch-fail-open"),
)
.await
.unwrap();
assert_eq!(observations[0].effect_id, "bad");
assert_eq!(observations[0].status, EffectStatus::Failed);
assert_eq!(observations[1].status, EffectStatus::Succeeded);
}
#[tokio::test]
async fn approval_deny_becomes_denied_observation() {
let harness = BasicHarness::new(
StaticEffectExecutor::new(),
DefaultPolicyEngine,
AlwaysDenyApprovalBroker,
NoopAuditSink,
NoopTranscriptStore,
);
let observation = harness
.execute(
EffectRequest::new(
EffectKind::ExecuteCommand,
"run high-risk command",
json!({}),
)
.with_id("approval-deny")
.with_risk(RiskLevel::High),
&RunContext::new("run-approval-deny"),
)
.await
.unwrap();
assert_eq!(observation.status, EffectStatus::Denied);
assert!(
observation
.output
.observation_for_model
.contains("approval broker")
);
}
#[derive(Debug, Clone)]
struct CountingApprovalBroker {
calls: Arc<Mutex<usize>>,
decision: ApprovalDecision,
}
#[async_trait]
impl ApprovalBroker for CountingApprovalBroker {
async fn approve(
&self,
_request: ApprovalRequest,
_context: &RunContext,
) -> Result<ApprovalDecision, ApprovalError> {
*self.calls.lock().expect("approval counter lock poisoned") += 1;
Ok(self.decision.clone())
}
}
#[tokio::test]
async fn allow_for_session_skips_repeated_approval_for_same_run_and_kind() {
let calls = Arc::new(Mutex::new(0));
let broker = CountingApprovalBroker {
calls: calls.clone(),
decision: ApprovalDecision::AllowForSession,
};
let harness = BasicHarness::new(
StaticEffectExecutor::new()
.with_output("effect-1", RawEffectOutput::text("first"))
.with_output("effect-2", RawEffectOutput::text("second")),
DefaultPolicyEngine,
broker,
NoopAuditSink,
NoopTranscriptStore,
);
let context = RunContext::new("run-session-approval");
let first = harness
.execute(
EffectRequest::new(EffectKind::ExecuteCommand, "run command", json!({}))
.with_id("effect-1")
.with_risk(RiskLevel::High),
&context,
)
.await
.unwrap();
let second = harness
.execute(
EffectRequest::new(EffectKind::ExecuteCommand, "run another command", json!({}))
.with_id("effect-2")
.with_risk(RiskLevel::High),
&context,
)
.await
.unwrap();
assert_eq!(first.status, EffectStatus::Succeeded);
assert_eq!(second.status, EffectStatus::Succeeded);
assert_eq!(*calls.lock().expect("approval counter lock poisoned"), 1);
}
#[tokio::test]
async fn executor_failure_is_effect_observation() {
let harness = BasicHarness::new(
StaticEffectExecutor::new()
.with_error("effect-1", ExecutionError::Failed("boom".to_string())),
DefaultPolicyEngine,
AlwaysAllowApprovalBroker,
NoopAuditSink,
NoopTranscriptStore,
);
let observation = harness
.execute(
EffectRequest::new(EffectKind::ReadFile, "read", json!({})).with_id("effect-1"),
&RunContext::new("run-fail"),
)
.await
.unwrap();
assert_eq!(observation.status, EffectStatus::Failed);
assert!(observation.output.observation_for_model.contains("boom"));
}
#[tokio::test]
async fn executor_denial_is_denied_observation() {
let harness = BasicHarness::new(
StaticEffectExecutor::new().with_error(
"effect-1",
ExecutionError::Denied("capability mismatch".to_string()),
),
DefaultPolicyEngine,
AlwaysAllowApprovalBroker,
NoopAuditSink,
NoopTranscriptStore,
);
let observation = harness
.execute(
EffectRequest::new(EffectKind::ReadFile, "read", json!({})).with_id("effect-1"),
&RunContext::new("run-denied-by-executor"),
)
.await
.unwrap();
assert_eq!(observation.status, EffectStatus::Denied);
assert!(
observation
.output
.observation_for_model
.contains("capability mismatch")
);
}
#[derive(Debug, Default, Clone, Copy)]
struct SlowExecutor;
#[async_trait]
impl EffectExecutor for SlowExecutor {
async fn execute(
&self,
_request: &EffectRequest,
_policy: &ExecutionPolicy,
_context: &RunContext,
) -> Result<RawEffectOutput, ExecutionError> {
tokio::time::sleep(Duration::from_millis(50)).await;
Ok(RawEffectOutput::text("late"))
}
}
#[tokio::test]
async fn executor_timeout_becomes_timed_out_observation() {
let harness = BasicHarness::new(
SlowExecutor,
DefaultPolicyEngine,
AlwaysAllowApprovalBroker,
NoopAuditSink,
NoopTranscriptStore,
)
.with_config(HarnessConfig {
default_timeout: Duration::from_millis(1),
..Default::default()
});
let observation = harness
.execute(
EffectRequest::new(EffectKind::ReadFile, "read slowly", json!({}))
.with_id("slow-effect"),
&RunContext::new("run-timeout"),
)
.await
.unwrap();
assert_eq!(observation.status, EffectStatus::TimedOut);
}
#[tokio::test]
async fn output_is_limited_and_redacted() {
let raw =
RawEffectOutput::text("secret-token-1234567890 plus a long tail that must be cut")
.with_debug("debug secret-token");
let limit = OutputLimit {
model_bytes: 32,
display_bytes: 20,
debug_bytes: 20,
};
let redactor = PatternRedactor::new(["secret-token"]);
let output = limit_and_redact(raw, &limit, &redactor);
assert!(output.truncated);
assert!(!output.output.observation_for_model.contains("secret-token"));
assert!(output.output.observation_for_model.contains("[REDACTED]"));
}
#[tokio::test]
async fn basic_harness_uses_configured_redactor() {
let harness = BasicHarness::new(
StaticEffectExecutor::new()
.with_output("effect-1", RawEffectOutput::text("secret-token in output")),
DefaultPolicyEngine,
AlwaysAllowApprovalBroker,
NoopAuditSink,
NoopTranscriptStore,
)
.with_redactor(PatternRedactor::new(["secret-token"]));
let observation = harness
.execute(
EffectRequest::new(EffectKind::ReadFile, "read", json!({})).with_id("effect-1"),
&RunContext::new("run-redactor"),
)
.await
.unwrap();
assert_eq!(observation.status, EffectStatus::Succeeded);
assert!(
!observation
.output
.observation_for_model
.contains("secret-token")
);
assert!(
observation
.output
.observation_for_model
.contains("[REDACTED]")
);
assert_eq!(
observation
.metadata
.get("redactions_applied")
.and_then(serde_json::Value::as_u64),
Some(1)
);
}
#[tokio::test]
async fn redacted_effect_output_is_used_for_observation_audit_and_transcript() {
let audit = VecAuditSink::new();
let transcript = VecTranscriptStore::new();
let mut metadata = RunMetadata::new();
metadata.insert(
"token".to_string(),
serde_json::Value::String("secret-token".to_string()),
);
let harness = BasicHarness::new(
StaticEffectExecutor::new().with_output(
"effect-1",
RawEffectOutput::text("model sees secret-token")
.with_display(DisplayOutput::new(
crate::effect::DisplayFormat::PlainText,
"display sees secret-token",
))
.with_metadata(metadata)
.with_debug("debug sees secret-token"),
),
DefaultPolicyEngine,
AlwaysAllowApprovalBroker,
audit.clone(),
transcript.clone(),
)
.with_redactor(PatternRedactor::new(["secret-token"]));
let observation = harness
.execute(
EffectRequest::new(EffectKind::ReadFile, "read", json!({})).with_id("effect-1"),
&RunContext::new("run-redaction-records"),
)
.await
.unwrap();
assert_eq!(observation.status, EffectStatus::Succeeded);
let observation_json = serde_json::to_string(&observation).unwrap();
assert!(!observation_json.contains("secret-token"));
assert!(observation_json.contains("[REDACTED]"));
assert!(
observation
.metadata
.get("redactions_applied")
.and_then(serde_json::Value::as_u64)
.unwrap_or_default()
> 0
);
let audit_json = serde_json::to_string(&audit.events()).unwrap();
assert!(!audit_json.contains("secret-token"));
let transcript_json = serde_json::to_string(&transcript.records()).unwrap();
assert!(!transcript_json.contains("secret-token"));
}
#[tokio::test]
async fn transcript_records_effect_lifecycle() {
let transcript = VecTranscriptStore::new();
let harness = BasicHarness::new(
StaticEffectExecutor::new().with_output("effect-1", RawEffectOutput::text("ok")),
DefaultPolicyEngine,
AlwaysAllowApprovalBroker,
NoopAuditSink,
transcript.clone(),
);
let observation = harness
.execute(
EffectRequest::new(EffectKind::ReadFile, "read config", json!({}))
.with_id("effect-1"),
&RunContext::new("run-transcript-lifecycle"),
)
.await
.unwrap();
assert_eq!(observation.status, EffectStatus::Succeeded);
let records = transcript.records();
assert!(records.iter().any(
|record| matches!(record, TranscriptRecord::EffectRequested { effect_id, .. } if effect_id == "effect-1")
));
assert!(records.iter().any(
|record| matches!(record, TranscriptRecord::EffectClassified { effect_id, .. } if effect_id == "effect-1")
));
assert!(records.iter().any(
|record| matches!(record, TranscriptRecord::PolicyDecided { effect_id, .. } if effect_id == "effect-1")
));
assert!(records.iter().any(
|record| matches!(record, TranscriptRecord::ExecutorStarted { effect_id, .. } if effect_id == "effect-1")
));
assert!(records.iter().any(
|record| matches!(record, TranscriptRecord::ExecutorCompleted { effect_id, status, .. } if effect_id == "effect-1" && *status == EffectStatus::Succeeded)
));
assert!(records.iter().any(
|record| matches!(record, TranscriptRecord::EffectObservation { effect_id, status, .. } if effect_id == "effect-1" && *status == EffectStatus::Succeeded)
));
}
#[derive(Debug, Default, Clone, Copy)]
struct FailingAuditSink;
#[async_trait]
impl AuditSink for FailingAuditSink {
async fn record(
&self,
_event: AuditEvent,
_context: &RunContext,
) -> Result<(), AuditError> {
Err(AuditError::Sink("audit offline".to_string()))
}
}
#[tokio::test]
async fn audit_failure_fails_closed_before_execution() {
let harness = BasicHarness::new(
StaticEffectExecutor::new()
.with_output("effect-1", RawEffectOutput::text("should not execute")),
DefaultPolicyEngine,
AlwaysAllowApprovalBroker,
FailingAuditSink,
NoopTranscriptStore,
);
let err = harness
.execute(
EffectRequest::new(EffectKind::ReadFile, "read", json!({})).with_id("effect-1"),
&RunContext::new("run-audit-fail"),
)
.await
.unwrap_err();
assert!(matches!(err, HarnessError::Audit(_)));
}
#[derive(Debug, Default, Clone, Copy)]
struct FailingTranscriptStore;
#[async_trait]
impl TranscriptStore for FailingTranscriptStore {
async fn append(
&self,
_record: TranscriptRecord,
_context: &RunContext,
) -> Result<(), TranscriptError> {
Err(TranscriptError::Store("transcript offline".to_string()))
}
}
#[tokio::test]
async fn transcript_failure_defaults_to_fail_open() {
let harness = BasicHarness::new(
StaticEffectExecutor::new().with_output("effect-1", RawEffectOutput::text("ok")),
DefaultPolicyEngine,
AlwaysAllowApprovalBroker,
NoopAuditSink,
FailingTranscriptStore,
);
let observation = harness
.execute(
EffectRequest::new(EffectKind::ReadFile, "read", json!({})).with_id("effect-1"),
&RunContext::new("run-transcript-open"),
)
.await
.unwrap();
assert_eq!(observation.status, EffectStatus::Succeeded);
}
#[tokio::test]
async fn transcript_failure_can_fail_closed() {
let harness = BasicHarness::new(
StaticEffectExecutor::new().with_output("effect-1", RawEffectOutput::text("ok")),
DefaultPolicyEngine,
AlwaysAllowApprovalBroker,
NoopAuditSink,
FailingTranscriptStore,
)
.with_config(HarnessConfig {
fail_closed_on_transcript_error: true,
..Default::default()
});
let err = harness
.execute(
EffectRequest::new(EffectKind::ReadFile, "read", json!({})).with_id("effect-1"),
&RunContext::new("run-transcript-closed"),
)
.await
.unwrap_err();
assert!(matches!(err, HarnessError::Transcript(_)));
}
#[tokio::test]
async fn invalid_effect_request_returns_harness_error() {
let harness = BasicHarness::noop();
let err = harness
.execute(
EffectRequest::new(EffectKind::ReadFile, "read", json!({})).with_id(""),
&RunContext::new("run-invalid"),
)
.await
.unwrap_err();
assert!(matches!(err, HarnessError::InvalidRequest(_)));
}
struct EffectTool;
#[async_trait]
impl Tool for EffectTool {
fn schema(&self) -> ToolSchema {
ToolSchema::new("read_file", "Read a file", json!({}))
}
async fn call(
&self,
_arguments: serde_json::Value,
_context: ToolContext<'_>,
) -> Result<ToolResult, ToolError> {
Ok(ToolResult::Effect(
EffectRequest::new(EffectKind::ReadFile, "read", json!({})).with_id("effect-1"),
))
}
}
#[tokio::test]
async fn registry_effect_tool_stays_external_to_harness() {
let mut registry = ToolRegistry::new();
registry.register(EffectTool);
assert_eq!(registry.names(), ["read_file"]);
}
}