use std::sync::atomic::{AtomicU32, Ordering};
use std::sync::Arc;
use std::time::{Duration, SystemTime, UNIX_EPOCH};
use axum::{
body::Bytes,
extract::State,
http::{header::CONTENT_TYPE, HeaderMap, StatusCode},
Json,
};
use crate::core::policy::{PolicyMatch, ResidentBundle};
use crate::daemon::identity::{self, IdentitySignals};
use crate::daemon::AppState;
use crate::envelope::{
new_event_id, AgentType, EventEnvelope, HookEventType, Verdict, VerdictResponse,
VerdictResponseContext, VerdictResponseContextBody, VerdictResponseContextHeadline,
};
use crate::privacy;
use crate::zone_eval::{
AgentContext as ZoneAgentContext, Decision as ZoneDecision, Event as ZoneEvent, EventEnv,
SessionFacts,
};
const CT_CLOUDEVENTS_SINGLE: &str = "application/cloudevents+json";
const CT_CLOUDEVENTS_BATCH: &str = "application/cloudevents-batch+json";
const CT_JSON: &str = "application/json";
struct HookActivityGuard<'a> {
in_flight: &'a AtomicU32,
}
impl<'a> HookActivityGuard<'a> {
fn new(state: &'a AppState) -> Self {
let now = SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|d| d.as_secs())
.unwrap_or(0);
state.last_hook_at_unix_secs.store(now, Ordering::Relaxed);
state.hooks_in_flight.fetch_add(1, Ordering::AcqRel);
Self {
in_flight: &state.hooks_in_flight,
}
}
}
impl Drop for HookActivityGuard<'_> {
fn drop(&mut self) {
self.in_flight.fetch_sub(1, Ordering::AcqRel);
}
}
pub async fn ingest_cloudevent(
State(state): State<Arc<AppState>>,
headers: HeaderMap,
body: Bytes,
) -> (StatusCode, HeaderMap, Json<VerdictResponse>) {
let start = std::time::Instant::now();
let _activity_guard = HookActivityGuard::new(&state);
let ct = headers
.get(CONTENT_TYPE)
.and_then(|v| v.to_str().ok())
.unwrap_or("");
let is_batch = ct.starts_with(CT_CLOUDEVENTS_BATCH);
let is_single = ct.starts_with(CT_CLOUDEVENTS_SINGLE) || ct.starts_with(CT_JSON);
if !is_batch && !is_single {
return reject(
StatusCode::UNSUPPORTED_MEDIA_TYPE,
"unsupported Content-Type; expected application/cloudevents+json or application/cloudevents-batch+json",
start,
);
}
let envelopes: Vec<EventEnvelope> = if is_batch {
match serde_json::from_slice::<Vec<EventEnvelope>>(&body) {
Ok(v) if !v.is_empty() => v,
Ok(_) => {
return reject(
StatusCode::BAD_REQUEST,
"cloudevents-batch body must not be empty",
start,
);
}
Err(e) => {
return reject(
StatusCode::BAD_REQUEST,
&format!("invalid cloudevents batch body: {e}"),
start,
);
}
}
} else {
match serde_json::from_slice::<EventEnvelope>(&body) {
Ok(env) => vec![env],
Err(e) => {
return reject(
StatusCode::BAD_REQUEST,
&format!("invalid cloudevent body: {e}"),
start,
);
}
}
};
let mut worst: Option<BatchOutcome> = None;
let mut session_ref: Option<String> = None;
let mut response_headers = HeaderMap::new();
for envelope in envelopes {
if is_doctor_probe(envelope.data.as_ref()) {
continue;
}
let (
ev_type,
ev_id,
ev_subject,
dedup_hit,
policy_match,
zone_decision,
hold_wait,
delivery,
) = process_envelope(state.clone(), envelope, start, ProcessMode::Live).await;
if dedup_hit {
response_headers.insert(
"x-openlatch-dedup",
"true".parse().expect("static header value is valid"),
);
}
if session_ref.is_none() {
session_ref = ev_subject.filter(|s| !s.is_empty());
}
join_most_restrictive(
&mut worst,
BatchOutcome {
verdict: verdict_for(&ev_type, zone_decision.as_ref(), policy_match.as_ref()),
event_type: ev_type,
event_id: ev_id,
policy_match,
zone_decision,
hold_wait,
delivery,
},
);
}
let outcome = worst.unwrap_or_else(|| BatchOutcome {
verdict: Verdict::Allow,
event_type: HookEventType::Unknown(String::new()),
event_id: new_event_id(),
policy_match: None,
zone_decision: None,
hold_wait: None,
delivery: crate::daemon::optimize::Delivery::default(),
});
let latency_ms = start.elapsed().as_secs_f64() * 1000.0;
let mut response = match outcome.verdict {
Verdict::Block => {
let m = outcome.policy_match.as_ref();
VerdictResponse::block(
outcome.event_id,
latency_ms,
m.map(|m| m.reason.clone()).or_else(|| {
outcome
.zone_decision
.as_ref()
.map(|decision| decision.reason.clone())
}),
m.map(|m| m.severity.to_string()),
m.map(|m| m.rule_id.clone()),
)
}
Verdict::Ask => {
let mut response = VerdictResponse::ask(outcome.event_id, latency_ms);
response.reason = outcome
.zone_decision
.as_ref()
.map(|d| d.reason.clone())
.filter(|r| !r.is_empty());
response
}
Verdict::Optimize => VerdictResponse::optimize(outcome.event_id, latency_ms),
Verdict::Allow => VerdictResponse::allow(outcome.event_id, latency_ms),
};
if let Some((token, timeout_ms)) = outcome.hold_wait {
response.schema_version = "1.2".to_string();
response.hold_token = Some(token);
response.hold_timeout_ms = Some(timeout_ms);
response.defer = true;
}
if outcome.delivery.updated_input.is_some()
|| outcome.delivery.additional_context.is_some()
|| outcome.delivery.system_message.is_some()
{
response.schema_version = "1.2".to_string();
response.updated_input = outcome.delivery.updated_input.unwrap_or_default();
response.additional_context = outcome.delivery.additional_context;
response.system_message = outcome.delivery.system_message;
}
if !matches!(outcome.verdict, Verdict::Block) && !state.pending_alerts.is_empty() {
if let Some(subject) = session_ref.as_deref() {
if let Some(alert) = state.pending_alerts.pop_for_session(subject) {
let surface = match outcome.event_type {
HookEventType::PreToolUse => "pretooluse_reason",
HookEventType::SessionStart => "session_start_context",
_ => "other",
};
crate::telemetry::capture_global(
crate::telemetry::Event::config_alert_surfaced_in_session(
&alert.severity,
surface,
),
);
attach_alert_context(&mut response, &alert);
}
}
}
(StatusCode::OK, response_headers, Json(response))
}
pub(crate) fn is_doctor_probe(data: Option<&serde_json::Value>) -> bool {
data.and_then(|d| d.get("tool_name"))
.and_then(serde_json::Value::as_str)
== Some("OpenlatchDoctorProbe")
}
struct BatchOutcome {
verdict: Verdict,
event_type: HookEventType,
event_id: String,
policy_match: Option<PolicyMatch>,
zone_decision: Option<ZoneDecision>,
hold_wait: Option<(String, u64)>,
delivery: crate::daemon::optimize::Delivery,
}
fn verdict_rank(verdict: Verdict) -> u8 {
match verdict {
Verdict::Allow => 0,
Verdict::Optimize => 1,
Verdict::Ask => 2,
Verdict::Block => 3,
}
}
fn join_most_restrictive(worst: &mut Option<BatchOutcome>, candidate: BatchOutcome) {
match worst {
Some(current) if verdict_rank(candidate.verdict) <= verdict_rank(current.verdict) => {}
_ => *worst = Some(candidate),
}
}
fn verdict_for(
event_type: &HookEventType,
zone_decision: Option<&ZoneDecision>,
policy_match: Option<&PolicyMatch>,
) -> Verdict {
if policy_match.is_some_and(|m| !m.shadow) {
return Verdict::Block;
}
if let Some(decision) = zone_decision {
return decision.verdict;
}
match event_type {
HookEventType::Stop => Verdict::Allow,
_ => Verdict::Allow,
}
}
fn attach_alert_context(
response: &mut VerdictResponse,
alert: &crate::daemon::config_monitor::PendingAlert,
) {
let headline = clamp_str(&alert.headline, 120);
let body = clamp_str(&alert.body, 500);
let Some(ctx) = build_context(&headline, &body) else {
return;
};
response.context = Some(ctx);
response.reason = Some(headline);
response.schema_version = "1.1".to_string();
}
fn build_context(headline: &str, body: &str) -> Option<VerdictResponseContext> {
let headline = VerdictResponseContextHeadline::try_from(headline.to_string()).ok()?;
let body = VerdictResponseContextBody::try_from(body.to_string()).ok()?;
Some(VerdictResponseContext {
body,
evidence: Vec::new(),
headline,
})
}
fn clamp_str(s: &str, max: usize) -> String {
if s.chars().count() <= max {
return s.to_string();
}
s.chars().take(max.saturating_sub(1)).chain(['…']).collect()
}
#[derive(Clone, Copy, PartialEq, Eq)]
pub(crate) enum ProcessMode {
Live,
Replay,
}
pub(crate) async fn process_envelope_for_replay(
state: Arc<AppState>,
envelope: EventEnvelope,
start: std::time::Instant,
) -> (HookEventType, String, Option<String>, bool) {
let (ev_type, ev_id, subject, dedup_hit, _policy_match, _zone_decision, _hold_wait, _delivery) =
process_envelope(state, envelope, start, ProcessMode::Replay).await;
(ev_type, ev_id, subject, dedup_hit)
}
const REPLAY_YIELD_FILL_RATIO_INV: usize = 4;
const REPLAY_BACKPRESSURE_TICK: Duration = Duration::from_millis(50);
async fn forward_replay_event(
tx: &tokio::sync::mpsc::Sender<crate::cloud::CloudEvent>,
mut cloud_event: crate::cloud::CloudEvent,
) {
use tokio::sync::mpsc::error::TrySendError;
loop {
let max_cap = tx.max_capacity();
let free = tx.capacity();
if free.saturating_mul(REPLAY_YIELD_FILL_RATIO_INV) < max_cap {
if tx.is_closed() {
tracing::warn!("cloud channel closed during replay — worker has exited");
return;
}
tokio::time::sleep(REPLAY_BACKPRESSURE_TICK).await;
continue;
}
match tx.try_send(cloud_event) {
Ok(()) => return,
Err(TrySendError::Full(ev)) => {
cloud_event = ev;
tokio::time::sleep(REPLAY_BACKPRESSURE_TICK).await;
}
Err(TrySendError::Closed(_)) => {
tracing::warn!("cloud channel closed during replay — worker has exited");
return;
}
}
}
}
async fn process_envelope(
state: Arc<AppState>,
mut envelope: EventEnvelope,
start: std::time::Instant,
mode: ProcessMode,
) -> (
HookEventType,
String,
Option<String>,
bool,
Option<PolicyMatch>,
Option<ZoneDecision>,
Option<(String, u64)>,
crate::daemon::optimize::Delivery,
) {
let event_type = envelope.type_.clone();
let event_id = envelope.id.clone();
let policy_bundle: Option<Arc<Option<ResidentBundle>>> = match mode {
ProcessMode::Live => state.policy.as_ref().map(|p| p.handle.load_full()),
ProcessMode::Replay => None,
};
let resident = policy_bundle.as_deref().and_then(|b| b.as_ref());
let policy_match = resident.and_then(|bundle| {
crate::core::policy::evaluate(bundle, event_type.as_str(), envelope.data.as_ref())
});
let zone_evaluation = resident
.and_then(|bundle| bundle.zone.as_ref().map(|zone| (bundle, zone)))
.and_then(|(bundle, zone)| {
let event = map_zone_event(&envelope, bundle, state.config.agent_id.as_deref())?;
let now = std::time::Instant::now();
let now_ms = chrono::Utc::now().timestamp_millis();
let policy = state.policy.as_ref()?;
let outcome =
policy
.sessions
.evaluate(&event.session_id, zone.state_layout, now, |prior| {
let outcome = crate::daemon::optimize::evaluate_and_dispatch(
zone,
&event,
envelope.source.as_str(),
&policy.dispatch_sessions,
&state.config.log_dir,
prior,
now,
now_ms,
);
let state = outcome.state_to_commit.clone();
(outcome, state)
});
let mut decision = outcome.decision;
if degrade_unapplied_rewrite(&mut decision) {
decision.verdict = Verdict::Ask;
}
Some((decision, event, outcome.delivery))
});
let mut zone_decision = zone_evaluation
.as_ref()
.map(|(decision, _, _)| decision.clone());
let delivery = zone_evaluation
.as_ref()
.map(|(_, _, delivery)| delivery.clone())
.unwrap_or_default();
let mut hold_wait = None;
if let (Some(policy), Some(bundle), Some((_, event, _)), Some(decision), Some(request)) = (
state.policy.as_ref(),
resident,
zone_evaluation.as_ref(),
zone_decision.as_ref(),
zone_decision
.as_ref()
.and_then(|decision| decision.hold.clone()),
) {
if !event.tool_use_id.is_empty() {
let host_timeout = bundle
.hold_config
.as_ref()
.and_then(|hold| hold.host_timeout_s)
.unwrap_or(900)
.max(31) as u64;
let max_pending = bundle
.hold_config
.as_ref()
.and_then(|hold| hold.max_pending_holds)
.unwrap_or(16)
.max(1) as usize;
let correlation_id = envelope
.data
.as_ref()
.and_then(|data| data.get("correlation_id"))
.and_then(serde_json::Value::as_str)
.unwrap_or(&event_id)
.to_string();
let artifact = decision.artifact_id.as_deref().and_then(|artifact_id| {
bundle.zone.as_ref().and_then(|zone| {
zone.artifacts
.iter()
.find(|artifact| artifact.artifact_id() == Some(artifact_id))
})
});
let registration = crate::daemon::hold_queue::HoldRegistration {
tool_use_id: event.tool_use_id.clone(),
atom_id: decision.atom_id.clone(),
policy_id: artifact.and_then(|artifact| artifact.envelope.policy_id.clone()),
policy_public_id: decision.policy_public_id.clone(),
effects: decision
.effects
.iter()
.filter_map(|effect| serde_json::to_value(effect).ok())
.collect(),
};
match policy.holds.register(
registration,
request,
event.clone(),
correlation_id,
envelope.time.to_rfc3339(),
envelope.subject.clone(),
Duration::from_secs(host_timeout),
max_pending,
) {
crate::daemon::hold_queue::RegisterResult::Pending { timeout } => {
hold_wait = Some((
event.tool_use_id.clone(),
timeout.as_millis().min(u64::MAX as u128) as u64,
));
}
crate::daemon::hold_queue::RegisterResult::Full(resolution) => {
if let Some(decision) = zone_decision.as_mut() {
decision.verdict = resolution.verdict;
decision.hold = None;
}
}
}
}
}
if let (Some(m), Some(bundle)) = (policy_match.as_ref(), resident) {
tracing::info!(
target: "policy",
rule_id = %m.rule_id,
mode = %m.mode,
verdict = %verdict_for(&event_type, zone_decision.as_ref(), Some(m)),
shadow = m.shadow,
revision = bundle.revision,
"policy decision"
);
}
if let (Some(install_id), Some(session_id)) = (
state.config.agent_id.as_deref(),
envelope.subject.as_deref(),
) {
let data = envelope.data.as_ref();
let str_field = |k: &str| {
data.and_then(|d| d.get(k))
.and_then(|v| v.as_str())
.filter(|s| !s.is_empty())
};
let signals = crate::model_relay::session::SessionSignals {
tool_use_id: str_field("tool_use_id"),
prompt: str_field("prompt"),
};
state.registry.upsert_signals(
install_id,
install_id,
envelope.source.as_str(),
session_id,
&signals,
);
}
if let Some(provider) = envelope
.data
.as_ref()
.and_then(|d| d.get("model"))
.and_then(|m| m.get("provider"))
.and_then(|p| p.as_str())
.map(str::trim)
.filter(|p| !p.is_empty() && *p != "unknown")
{
state.registry.note_provider(
envelope.source.as_str(),
provider,
crate::hooks::provider_endpoints::now_unix(),
);
}
let event_cwd = envelope
.data
.as_ref()
.and_then(|d| d.get("cwd"))
.and_then(|v| v.as_str());
let (osuser, identity) = identity_stamp(
mode,
resident.is_some_and(|b| b.capture_identity_signals),
envelope.subject.as_deref(),
event_cwd,
identity::os_user,
identity::observe_session,
);
envelope.osuser = osuser;
envelope.gitemail = identity.git_email;
envelope.provideracct = identity.provider_account;
let config_monitor_tx = state.config_monitor.request_tx();
if matches!(event_type, HookEventType::SessionStart) {
if let Err(e) = crate::daemon::config_monitor::enrich_session_start(
&mut envelope.data,
&state.content_hash_cache,
config_monitor_tx.as_ref(),
)
.await
{
tracing::warn!(
code = crate::error::ERR_INVENTORY_ENRICH_FAILED,
error = ?e,
"session_start enrichment failed; forwarding un-enriched"
);
}
} else if let Some(cwd) = event_cwd {
crate::daemon::config_monitor::try_register_project_from_cwd(
cwd,
config_monitor_tx.as_ref(),
)
.await;
}
let raw_subject = envelope.subject.clone();
let subject = raw_subject.clone().unwrap_or_else(|| "unknown".into());
let type_str = envelope.type_.as_str().to_string();
let data_for_dedup = envelope.data.clone().unwrap_or_default();
let is_duplicate = state
.dedup
.check_and_insert(&subject, &type_str, &data_for_dedup);
if is_duplicate {
tracing::debug!(
code = crate::error::ERR_EVENT_DEDUPED,
session_id = %subject,
event_type = %type_str,
"Event deduplicated within TTL window"
);
return (
event_type,
event_id,
raw_subject,
true,
policy_match,
zone_decision,
hold_wait,
delivery,
);
}
if matches!(envelope.source, AgentType::Unknown(_)) {
crate::telemetry::capture_hook_source_unknown(envelope.source.as_str());
}
if matches!(envelope.type_, HookEventType::Unknown(_)) {
crate::telemetry::capture_hook_type_unknown(envelope.type_.as_str());
}
if let Some(data) = envelope.data.as_mut() {
privacy::filter_event_with(data, &state.privacy_filter);
}
if envelope.os.is_none() {
envelope.os = Some(crate::envelope::os_string().to_string());
}
if envelope.arch.is_none() {
envelope.arch = Some(crate::envelope::arch_string().to_string());
}
envelope.clientversion = Some(env!("OPENLATCH_VERSION").to_string());
if envelope.datacontenttype.is_none() {
envelope.datacontenttype = Some("application/json".to_string());
}
if envelope.localipv4.is_none() {
envelope.localipv4 = state.local_ipv4;
}
if envelope.localipv6.is_none() {
envelope.localipv6 = state.local_ipv6;
}
if envelope.publicipv4.is_none() {
envelope.publicipv4 = state.public_ipv4;
}
if envelope.publicipv6.is_none() {
envelope.publicipv6 = state.public_ipv6;
}
let mut envelope_value = serde_json::to_value(&envelope).unwrap_or(serde_json::Value::Null);
let verdict_for_stored =
verdict_for(&event_type, zone_decision.as_ref(), policy_match.as_ref());
let latency_snapshot_ms = start.elapsed().as_millis() as u64;
if let Some(obj) = envelope_value.as_object_mut() {
obj.insert(
"olverdict".to_string(),
serde_json::json!(verdict_for_stored.to_string()),
);
obj.insert(
"ollatencyms".to_string(),
serde_json::json!(latency_snapshot_ms),
);
{
use crate::core::envelope::normalize;
let data = obj.get("data").cloned().unwrap_or(serde_json::Value::Null);
if let Some(prompt) = normalize::prompt_of(&data) {
obj.insert("olprompt".to_string(), serde_json::json!(prompt));
}
let (provider, slug) = normalize::model_of(&data);
if let Some(provider) = provider {
obj.insert("olmodelprovider".to_string(), serde_json::json!(provider));
}
if let Some(slug) = slug {
obj.insert("olmodelslug".to_string(), serde_json::json!(slug));
}
}
if mode == ProcessMode::Live {
stamp_egress_extensions(obj, egress_stamp(&state.egress.snapshot()));
}
if let Some(policy) = state.policy.as_ref() {
if mode == ProcessMode::Live {
stamp_policy_extensions(
obj,
policy_match.as_ref(),
resident,
!policy.last_fetch_ok.load(Ordering::Relaxed),
SystemTime::now(),
);
stamp_record_v2_if_winner(
obj,
zone_decision.as_ref(),
policy_match.as_ref(),
resident,
envelope.source.as_str(),
&event_type,
hold_wait.is_some(),
Some(&delivery),
);
}
}
}
let envelope_json = serde_json::to_string(&envelope_value).unwrap_or_default();
state.event_logger.log(envelope_json);
state.event_counter.fetch_add(1, Ordering::Relaxed);
if let Some(tx) = &state.cloud_tx {
let agent_id = state.config.agent_id.as_deref().unwrap_or_default();
if agent_id.is_empty() {
tracing::warn!(
code = "OL-1200",
"cloud forward skipped: [daemon].agent_id is empty in config.toml — run 'openlatch init' to populate it"
);
} else {
let cloud_event = crate::cloud::CloudEvent {
envelope: match mode {
ProcessMode::Live => envelope_value,
ProcessMode::Replay => serde_json::to_value(&envelope).unwrap_or_default(),
},
agent_id: agent_id.to_string(),
};
match mode {
ProcessMode::Live => {
let max_cap = tx.max_capacity();
match tx.try_send(cloud_event) {
Ok(()) => {}
Err(tokio::sync::mpsc::error::TrySendError::Full(_)) => {
if let Some(cs) = state.cloud_state.as_ref() {
cs.record_live_drop();
}
tracing::warn!(code = "OL-1200", "cloud channel full - event dropped");
}
Err(tokio::sync::mpsc::error::TrySendError::Closed(_)) => {
tracing::warn!("cloud channel closed - worker has exited");
}
}
if let Some(cs) = state.cloud_state.as_ref() {
if max_cap > 0 {
let free = tx.capacity();
let depth = max_cap.saturating_sub(free);
let pct = (depth as u64).saturating_mul(100) / max_cap as u64;
cs.note_channel_depth_pct(pct);
}
}
}
ProcessMode::Replay => {
forward_replay_event(tx, cloud_event).await;
}
}
}
}
if matches!(event_type, HookEventType::SessionEnd) {
if let (Some(policy), Some(subject)) = (state.policy.as_ref(), raw_subject.as_deref()) {
policy.sessions.remove(subject);
policy.dispatch_sessions.remove(subject);
}
}
(
event_type,
event_id,
raw_subject,
false,
policy_match,
zone_decision,
hold_wait,
delivery,
)
}
fn degrade_unapplied_rewrite(decision: &mut ZoneDecision) -> bool {
if decision.verdict == Verdict::Optimize && decision.rewrite.is_some() {
decision.verdict = Verdict::Ask;
true
} else {
false
}
}
fn map_zone_event(
envelope: &EventEnvelope,
bundle: &ResidentBundle,
agent_id: Option<&str>,
) -> Option<ZoneEvent> {
let data = envelope.data.as_ref()?;
let object = data.as_object()?;
let tool_input = object
.get("tool_input")
.cloned()
.or_else(|| object.get("parameters").cloned())
.unwrap_or(serde_json::Value::Null);
let tool_result = object
.get("tool_result")
.cloned()
.or_else(|| object.get("result").cloned())
.or_else(|| object.get("tool_response").cloned())
.or_else(|| {
matches!(envelope.type_, HookEventType::PostToolUseFailure).then(|| {
serde_json::json!({
"is_error": true,
"error": object.get("error").cloned().unwrap_or(serde_json::Value::Null),
"is_interrupt": object.get("is_interrupt").cloned().unwrap_or(serde_json::Value::Bool(false)),
})
})
});
let cwd = object
.get("cwd")
.and_then(serde_json::Value::as_str)
.map(str::to_string);
let environment = bundle
.zone
.as_ref()
.and_then(|z| z.meta.as_ref())
.and_then(|m| m.environment.as_ref())
.map(ToString::to_string);
Some(ZoneEvent {
event_type: match envelope.type_ {
HookEventType::SessionEnd => "stop".to_string(),
_ => envelope.type_.as_str().to_string(),
},
session_id: envelope.subject.clone().unwrap_or_default(),
tool_use_id: object
.get("tool_use_id")
.and_then(serde_json::Value::as_str)
.unwrap_or_default()
.to_string(),
tool_name: object
.get("tool_name")
.and_then(serde_json::Value::as_str)
.or_else(|| object.get("toolName").and_then(serde_json::Value::as_str))
.unwrap_or_default()
.to_string(),
tool_input,
tool_result,
agent: Some(ZoneAgentContext {
agent_id: envelope
.agentid
.clone()
.or_else(|| agent_id.map(str::to_string)),
agent_row_id: None,
agent_type: bundle
.zone
.as_ref()
.and_then(|zone| zone.meta.as_ref())
.and_then(|meta| meta.agent_type.clone())
.map(|category| category.to_string()),
environment,
function: bundle.agent_function.as_ref().map(ToString::to_string),
principal: None,
}),
binding: serde_json::json!({"kind": format!("{}_hook", envelope.source.as_str().replace('-', "_")), "mode": "command"}),
env: Some(EventEnv {
home: dirs::home_dir().map(|path| path.to_string_lossy().into_owned()),
cwd,
path_dirs: std::env::var_os("PATH")
.map(|p| {
std::env::split_paths(&p)
.map(|p| p.to_string_lossy().into_owned())
.collect()
})
.unwrap_or_default(),
additional_dirs: Vec::new(),
}),
session: object
.get("session")
.cloned()
.and_then(|v| serde_json::from_value::<SessionFacts>(v).ok()),
spend_delta: crate::daemon::spend::from_payload(data),
})
}
fn stamp_record_v2(
obj: &mut serde_json::Map<String, serde_json::Value>,
decision: &ZoneDecision,
bundle: Option<&ResidentBundle>,
source: &str,
event_type: &HookEventType,
deferred: bool,
delivery: Option<&crate::daemon::optimize::Delivery>,
) {
let scalar = |obj: &mut serde_json::Map<String, serde_json::Value>,
key: &str,
value: Option<serde_json::Value>| {
if let Some(value) = value {
obj.insert(key.to_string(), value);
}
};
scalar(
obj,
"olatomid",
decision.atom_id.clone().map(serde_json::Value::String),
);
scalar(
obj,
"olpolicyruleid",
decision.artifact_id.clone().map(serde_json::Value::String),
);
let artifact = decision.artifact_id.as_deref().and_then(|artifact_id| {
bundle
.and_then(|resident| resident.zone.as_ref())
.and_then(|zone| {
zone.artifacts
.iter()
.find(|artifact| artifact.artifact_id() == Some(artifact_id))
})
});
scalar(
obj,
"olpolicyid",
artifact
.and_then(|artifact| artifact.envelope.policy_id.clone())
.map(serde_json::Value::String),
);
scalar(
obj,
"oldimension",
decision
.dimension
.as_ref()
.map(|v| serde_json::json!(v.to_string())),
);
scalar(
obj,
"olmode",
decision
.mode
.as_ref()
.map(|v| serde_json::json!(v.to_string())),
);
scalar(obj, "oltier", decision.tier.map(serde_json::Value::from));
scalar(
obj,
"ollayer",
artifact
.and_then(|artifact| artifact.envelope.zone_layer.as_ref())
.and_then(|layer| layer.node_id.clone())
.map(serde_json::Value::String),
);
scalar(
obj,
"olspechash",
artifact
.and_then(|artifact| artifact.envelope.spec_hash.clone())
.map(serde_json::Value::String),
);
obj.insert(
"olbinding".into(),
serde_json::json!(format!("{}_hook", source.replace('-', "_"))),
);
let (enforced, result) = delivery
.and_then(|delivery| delivery.enforced.zip(delivery.result))
.unwrap_or_else(|| delivery_outcome(source, event_type, decision, deferred));
obj.insert("olenforced".into(), serde_json::json!(i64::from(enforced)));
obj.insert("olresult".into(), serde_json::json!(result));
if let Some(optimize) = &decision.optimize {
obj.insert(
"oloptimize".into(),
serde_json::to_value(optimize).unwrap_or_default(),
);
if let Some(selected) = optimize.actual.as_ref().or(optimize.monitor.as_ref()) {
obj.insert(
"ollever".into(),
serde_json::json!(selected.lever.to_string()),
);
}
}
if delivery.is_some_and(|delivery| delivery.second_optimize_degraded) {
obj.insert(
"oloptimizedegraded".into(),
serde_json::json!("second_optimize"),
);
}
obj.insert(
"oleffects".into(),
serde_json::Value::String(
serde_json::to_string(&decision.effects).unwrap_or_else(|_| "[]".into()),
),
);
obj.insert(
"olinconclusive".into(),
serde_json::Value::String(
serde_json::to_string(&decision.inconclusive_facts).unwrap_or_else(|_| "[]".into()),
),
);
obj.insert(
"olundecided".into(),
serde_json::json!(i64::from(decision.undecided)),
);
if let Some(rewrite) = &decision.rewrite {
obj.insert(
"ollever".into(),
serde_json::json!(rewrite.lever.to_string()),
);
}
if let Some(shadow) = decision.would_have_verdict {
obj.insert(
"olverdictshadow".into(),
serde_json::json!(shadow.to_string()),
);
}
}
fn delivery_outcome(
source: &str,
event_type: &HookEventType,
decision: &ZoneDecision,
deferred: bool,
) -> (bool, &'static str) {
let enforcing = decision
.mode
.as_ref()
.is_some_and(|mode| mode.0 == "enforce");
if deferred {
return (enforcing, "deferred");
}
if decision.would_have_verdict.is_some() && !enforcing {
return (false, "flagged");
}
if !enforcing || decision.artifact_id.is_none() {
return (false, result_for(decision.verdict));
}
let pre_action = matches!(event_type, HookEventType::PreToolUse);
let supported_agent = matches!(source, "claude-code" | "codex-cli");
match decision.verdict {
Verdict::Allow => (true, "allowed"),
Verdict::Ask => (pre_action && supported_agent, "flagged"),
Verdict::Block if pre_action && supported_agent => (true, "blocked"),
Verdict::Block => (false, "allowed"),
Verdict::Optimize => (false, "flagged"),
}
}
pub async fn wait_for_hold(
State(state): State<Arc<AppState>>,
axum::extract::Path(tool_use_id): axum::extract::Path<String>,
) -> (StatusCode, Json<VerdictResponse>) {
let Some(policy) = state.policy.as_ref() else {
return (
StatusCode::NOT_FOUND,
Json(VerdictResponse::allow(new_event_id(), 0.0)),
);
};
let Some(mut resolution) = policy.holds.wait(&tool_use_id).await else {
return (
StatusCode::NOT_FOUND,
Json(VerdictResponse::allow(new_event_id(), 0.0)),
);
};
if resolution.answer == crate::daemon::hold_queue::HoldAnswer::Retired {
let resident = policy.handle.load_full();
if let Some(zone) = resident
.as_ref()
.as_ref()
.and_then(|bundle| bundle.zone.as_ref())
{
let now = std::time::Instant::now();
let decision = policy.sessions.evaluate(
&resolution.event.session_id,
zone.state_layout,
now,
|prior| {
crate::zone_eval::evaluate(
zone,
&resolution.event,
prior,
chrono::Utc::now().timestamp_millis(),
)
},
);
resolution.verdict = decision.verdict;
resolution.result = result_for(decision.verdict);
}
}
emit_hold_resolution(&state, &tool_use_id, &resolution);
let mut response = response_for_verdict(
resolution.verdict,
new_event_id(),
0.0,
Some("OpenLatch hold resolved with a refusal".to_string()),
);
response.schema_version = "1.2".to_string();
(StatusCode::OK, Json(response))
}
fn response_for_verdict(
verdict: Verdict,
event_id: String,
latency_ms: f64,
block_reason: Option<String>,
) -> VerdictResponse {
match verdict {
Verdict::Allow => VerdictResponse::allow(event_id, latency_ms),
Verdict::Ask => VerdictResponse::ask(event_id, latency_ms),
Verdict::Optimize => VerdictResponse::optimize(event_id, latency_ms),
Verdict::Block => VerdictResponse::block(event_id, latency_ms, block_reason, None, None),
}
}
fn result_for(verdict: Verdict) -> &'static str {
match verdict {
Verdict::Allow => "allowed",
Verdict::Ask => "flagged",
Verdict::Block => "blocked",
Verdict::Optimize => "rewritten",
}
}
fn emit_hold_resolution(
state: &AppState,
tool_use_id: &str,
resolution: &crate::daemon::hold_queue::HoldResolution,
) {
let envelope = hold_resolution_envelope(tool_use_id, resolution);
state
.event_logger
.log(serde_json::to_string(&envelope).unwrap_or_default());
let Some(tx) = state.cloud_tx.as_ref() else {
return;
};
let agent_id = state.config.agent_id.as_deref().unwrap_or_default();
if agent_id.is_empty() {
tracing::warn!(
code = "OL-1200",
"hold resolution upload skipped: [daemon].agent_id is empty in config.toml"
);
return;
}
try_enqueue_hold_resolution(tx, state.cloud_state.as_ref(), agent_id, envelope);
}
fn try_enqueue_hold_resolution(
tx: &tokio::sync::mpsc::Sender<crate::cloud::CloudEvent>,
cloud_state: Option<&crate::cloud::CloudState>,
agent_id: &str,
envelope: serde_json::Value,
) {
match tx.try_send(crate::cloud::CloudEvent {
envelope,
agent_id: agent_id.to_string(),
}) {
Ok(()) => {}
Err(tokio::sync::mpsc::error::TrySendError::Full(_)) => {
if let Some(cloud_state) = cloud_state {
cloud_state.record_live_drop();
}
tracing::warn!(
code = "OL-1200",
"cloud channel full - hold resolution dropped"
);
}
Err(tokio::sync::mpsc::error::TrySendError::Closed(_)) => {
tracing::warn!("cloud channel closed - hold resolution not uploaded");
}
}
}
fn hold_resolution_envelope(
tool_use_id: &str,
resolution: &crate::daemon::hold_queue::HoldResolution,
) -> serde_json::Value {
let mut envelope = serde_json::json!({
"specversion": "1.0",
"id": new_event_id(),
"source": "openlatch-client",
"type": "ai.openlatch.decision.resolved",
"time": crate::daemon::hold_queue::resolution_time(&resolution.original_time),
"data": {
"tool_use_id": tool_use_id,
"correlation_id": resolution.correlation_id,
},
"olverdict": resolution.verdict.to_string(),
"olresult": resolution.result,
"olenforced": 1,
"olhostprompted": 0,
});
if let Some(subject) = resolution.subject.as_ref() {
envelope["subject"] = serde_json::Value::String(subject.clone());
}
envelope
}
#[allow(clippy::too_many_arguments)]
fn stamp_record_v2_if_winner(
obj: &mut serde_json::Map<String, serde_json::Value>,
zone_decision: Option<&ZoneDecision>,
policy_match: Option<&PolicyMatch>,
bundle: Option<&ResidentBundle>,
source: &str,
event_type: &HookEventType,
deferred: bool,
delivery: Option<&crate::daemon::optimize::Delivery>,
) {
let Some(decision) =
zone_decision.filter(|_| !policy_match.is_some_and(|matched| !matched.shadow))
else {
return;
};
stamp_record_v2(
obj, decision, bundle, source, event_type, deferred, delivery,
);
}
fn stamp_policy_extensions(
obj: &mut serde_json::Map<String, serde_json::Value>,
policy_match: Option<&PolicyMatch>,
bundle: Option<&ResidentBundle>,
offline: bool,
now: SystemTime,
) {
if let Some(m) = policy_match {
if m.shadow {
obj.insert(
"olverdictshadow".to_string(),
serde_json::json!(Verdict::Block.to_string()),
);
}
obj.insert(
"olpolicyruleid".to_string(),
serde_json::json!(m.rule_id.clone()),
);
}
if let Some(b) = bundle {
obj.insert(
"olpolicybundlerev".to_string(),
serde_json::json!(b.revision),
);
obj.insert(
"olpolicybundleage".to_string(),
serde_json::json!(b.age_seconds(now)),
);
}
obj.insert("olpolicyoffline".to_string(), serde_json::json!(offline));
}
fn identity_stamp(
mode: ProcessMode,
capture: bool,
subject: Option<&str>,
cwd: Option<&str>,
resolve_os_user: impl FnOnce() -> Option<String>,
observe_session: impl FnOnce(&str, Option<&str>) -> IdentitySignals,
) -> (Option<String>, IdentitySignals) {
if mode != ProcessMode::Live || !capture {
return (None, IdentitySignals::default());
}
let signals = match subject {
Some(session_id) => observe_session(session_id, cwd),
None => IdentitySignals::default(),
};
(resolve_os_user(), signals)
}
#[derive(Debug, Default, PartialEq)]
struct EgressStamp {
proxytype: Option<String>,
proxysource: Option<String>,
proxyintercepted: Option<bool>,
}
fn egress_stamp(snapshot: &crate::egress::EgressSnapshot) -> EgressStamp {
use crate::egress::{ProxySource, ProxyType};
let proxytype = match snapshot.proxy_type() {
ProxyType::Http => "http",
ProxyType::Https => "https",
ProxyType::Socks5 => "socks5",
ProxyType::Pac => "pac",
ProxyType::Direct => return EgressStamp::default(),
};
EgressStamp {
proxytype: Some(proxytype.to_string()),
proxysource: snapshot.source.map(|source| {
match source {
ProxySource::Manual => "manual",
ProxySource::Env => "env",
ProxySource::Windows => "windows",
ProxySource::Macos => "macos",
ProxySource::Gnome => "gnome",
ProxySource::Pac => "pac",
ProxySource::Wpad => "wpad",
}
.to_string()
}),
proxyintercepted: snapshot.tls_intercepted,
}
}
fn stamp_egress_extensions(
obj: &mut serde_json::Map<String, serde_json::Value>,
stamp: EgressStamp,
) {
if let Some(proxytype) = stamp.proxytype {
obj.insert("proxytype".to_string(), serde_json::json!(proxytype));
}
if let Some(proxysource) = stamp.proxysource {
obj.insert("proxysource".to_string(), serde_json::json!(proxysource));
}
if let Some(intercepted) = stamp.proxyintercepted {
obj.insert(
"proxyintercepted".to_string(),
serde_json::json!(intercepted),
);
}
}
fn reject(
status: StatusCode,
message: &str,
start: std::time::Instant,
) -> (StatusCode, HeaderMap, Json<VerdictResponse>) {
tracing::warn!(%message, http_status = %status, "cloudevents ingest rejected");
let latency_ms = start.elapsed().as_secs_f64() * 1000.0;
let response = VerdictResponse::allow(new_event_id(), latency_ms);
(status, HeaderMap::new(), Json(response))
}
pub async fn health(State(state): State<Arc<AppState>>) -> Json<serde_json::Value> {
let uptime_secs = state.started_at.elapsed().as_secs();
let degraded = state.health.degraded_names();
if !degraded.is_empty() {
tracing::debug!(
code = crate::error::ERR_SUBSYSTEM_DEGRADED,
subsystems = %degraded.join(","),
"/health reporting degraded"
);
}
let egress_failed = state.egress.status() == crate::egress::EgressStatus::Failed;
Json(serde_json::json!({
"status": health_status(!degraded.is_empty(), egress_failed),
"version": env!("OPENLATCH_VERSION"),
"pid": std::process::id(),
"uptime_secs": uptime_secs,
"subsystems": state.health.subsystems_json(),
}))
}
pub async fn metrics(State(state): State<Arc<AppState>>) -> Json<serde_json::Value> {
let events = state.event_counter.load(Ordering::Relaxed);
let uptime_secs = state.started_at.elapsed().as_secs();
let update_available = state.get_available_update();
let (
cloud_status,
cloud_forwarded_count,
cloud_last_sync_secs,
cloud_drop_count,
cloud_api_url,
cloud_consecutive_live_drops,
cloud_emergency_mode,
cloud_channel_high_water_ms,
) = if let Some(ref cs) = state.cloud_state {
let forwarded = cs.forwarded_count();
let last_sync = cs.last_sync_secs();
let drops = cs.drop_count();
let live_drops = cs.consecutive_live_drops();
let emergency = cs.is_emergency_mode();
let high_water_ms = cs.channel_high_water_window_ms();
let status = if cs.is_auth_error() {
"auth_error"
} else if cs.is_no_credential() {
"no_credential"
} else if emergency {
"emergency"
} else if cs.consecutive_drops() > 0 {
"network_error"
} else {
"connected"
};
(
status,
forwarded,
last_sync,
drops,
state.config.cloud.api_url.as_str(),
live_drops,
emergency,
high_water_ms,
)
} else {
("not_configured", 0u64, 0u64, 0u64, "", 0u64, false, 0u64)
};
let (outbox_pending_bytes, outbox_pending_count) = match state.outbox.as_ref() {
Some(o) => (o.byte_size(), o.pending_count()),
None => (0u64, 0u64),
};
let fallback_pending_count = fallback_log_line_count();
let policy = policy_metrics(state.policy.as_ref(), SystemTime::now());
let state_evicted = state
.policy
.as_ref()
.map(|policy| policy.sessions.evicted_count())
.unwrap_or(0);
let pending_holds = state
.policy
.as_ref()
.map(|policy| policy.holds.pending_count())
.unwrap_or(0);
let egress = egress_metrics(&state.egress);
let (cloud_channel_depth, cloud_channel_size_max) = match state.cloud_tx.as_ref() {
Some(tx) => {
let max = tx.max_capacity();
let free = tx.capacity();
(max.saturating_sub(free) as u64, max as u64)
}
None => (0u64, 0u64),
};
Json(serde_json::json!({
"events_processed": events,
"uptime_secs": uptime_secs,
"update_available": update_available,
"cloud_status": cloud_status,
"cloud_forwarded_count": cloud_forwarded_count,
"cloud_last_sync_secs": cloud_last_sync_secs,
"cloud_drop_count": cloud_drop_count,
"cloud_api_url": cloud_api_url,
"cloud_channel_depth": cloud_channel_depth,
"cloud_channel_size_max": cloud_channel_size_max,
"cloud_consecutive_live_drops": cloud_consecutive_live_drops,
"cloud_emergency_mode": cloud_emergency_mode,
"cloud_channel_high_water_ms": cloud_channel_high_water_ms,
"cloud_channel_high_water_pct_trip": crate::cloud::worker::HIGH_WATER_TRIP_PCT,
"outbox_pending_bytes": outbox_pending_bytes,
"outbox_pending_count": outbox_pending_count,
"fallback_pending_count": fallback_pending_count,
"policy_enabled": policy.enabled,
"policy_has_bundle": policy.has_bundle,
"policy_schema_version": policy.schema_version,
"policy_bundle_digest": policy.bundle_digest,
"policy_agent_function": policy.agent_function,
"policy_revision": policy.revision,
"policy_bundle_age_seconds": policy.bundle_age_seconds,
"policy_last_poll_ok_secs": policy.last_poll_ok_secs,
"policy_floor_blocked": policy.floor_blocked,
"policy_installed_version": policy.installed_version,
"policy_min_client_version": policy.min_client_version,
"policy_omitted_kinds": policy.omitted_kinds,
"policy_skipped_items": policy.skipped_items,
"policy_state_evicted": state_evicted,
"policy_pending_holds": pending_holds,
"egress_up": egress.up,
"egress_status": egress.status,
"egress_consecutive_failures": egress.consecutive_failures,
"proxy_in_use": egress.proxy_in_use,
"subsystem_restarts_total": state.health.total_restarts(),
"subsystem_degraded_count": state.health.degraded_count(),
}))
}
fn health_status(subsystem_degraded: bool, egress_failed: bool) -> &'static str {
if subsystem_degraded || egress_failed {
"degraded"
} else {
"ok"
}
}
struct EgressMetrics {
up: bool,
status: &'static str,
consecutive_failures: u32,
proxy_in_use: bool,
}
fn egress_metrics(state: &crate::egress::EgressState) -> EgressMetrics {
let status = state.status();
EgressMetrics {
up: status != crate::egress::EgressStatus::Failed,
status: status.as_str(),
consecutive_failures: state.consecutive_failures(),
proxy_in_use: state.snapshot().proxy_in_use,
}
}
struct PolicyMetrics {
enabled: bool,
has_bundle: bool,
schema_version: serde_json::Value,
bundle_digest: serde_json::Value,
revision: serde_json::Value,
bundle_age_seconds: serde_json::Value,
last_poll_ok_secs: serde_json::Value,
agent_function: serde_json::Value,
floor_blocked: bool,
installed_version: serde_json::Value,
min_client_version: serde_json::Value,
omitted_kinds: serde_json::Value,
skipped_items: usize,
}
fn policy_metrics(policy: Option<&crate::daemon::PolicyRuntime>, now: SystemTime) -> PolicyMetrics {
let Some(policy) = policy else {
return PolicyMetrics {
enabled: false,
has_bundle: false,
schema_version: serde_json::Value::Null,
bundle_digest: serde_json::Value::Null,
revision: serde_json::Value::Null,
bundle_age_seconds: serde_json::Value::Null,
last_poll_ok_secs: serde_json::Value::Null,
agent_function: serde_json::Value::Null,
floor_blocked: false,
installed_version: serde_json::Value::Null,
min_client_version: serde_json::Value::Null,
omitted_kinds: serde_json::json!([]),
skipped_items: 0,
};
};
let resident = policy.handle.load_full();
let (
has_bundle,
schema_version,
bundle_digest,
revision,
bundle_age_seconds,
agent_function,
omitted_kinds,
skipped_items,
) = match resident.as_ref() {
Some(b) => (
true,
serde_json::json!(b.schema_version),
serde_json::json!(b.digest),
serde_json::json!(b.revision),
serde_json::json!(b.age_seconds(now)),
serde_json::json!(b.agent_function),
serde_json::json!(b.omitted_kinds),
b.zone.as_ref().map_or(0, |zone| zone.skipped.len()),
),
None => (
false,
serde_json::Value::Null,
serde_json::Value::Null,
serde_json::Value::Null,
serde_json::Value::Null,
serde_json::Value::Null,
serde_json::json!([]),
0,
),
};
let floor = policy.floor.read().ok().and_then(|guard| guard.clone());
let last_poll_ok_at = policy.last_poll_ok_at.load(Ordering::Relaxed);
let last_poll_ok_secs = if last_poll_ok_at > 0 {
let now_unix = now
.duration_since(UNIX_EPOCH)
.map(|d| d.as_secs() as i64)
.unwrap_or(0);
serde_json::json!(now_unix.saturating_sub(last_poll_ok_at).max(0))
} else {
serde_json::Value::Null
};
PolicyMetrics {
enabled: true,
has_bundle,
schema_version,
bundle_digest,
revision,
bundle_age_seconds,
last_poll_ok_secs,
agent_function,
floor_blocked: floor.is_some(),
installed_version: floor.as_ref().map_or(serde_json::Value::Null, |state| {
serde_json::json!(state.installed_version)
}),
min_client_version: floor.as_ref().map_or(serde_json::Value::Null, |state| {
serde_json::json!(state.minimum_version)
}),
omitted_kinds,
skipped_items,
}
}
fn fallback_log_line_count() -> u64 {
use std::io::{BufRead, Seek, SeekFrom};
let log_dir = crate::config::openlatch_dir().join("logs");
let path = log_dir.join("fallback.jsonl");
let offset_path = log_dir.join("fallback.jsonl.offset");
let offset: u64 = std::fs::read_to_string(&offset_path)
.ok()
.and_then(|s| s.trim().parse().ok())
.unwrap_or(0);
let Ok(file) = std::fs::File::open(&path) else {
return 0;
};
let mut reader = std::io::BufReader::new(file);
if offset > 0 && reader.seek(SeekFrom::Start(offset)).is_err() {
return 0;
}
reader
.lines()
.map_while(Result::ok)
.filter(|l| !l.trim().is_empty())
.count() as u64
}
pub async fn shutdown_handler(State(state): State<Arc<AppState>>) -> StatusCode {
let mut tx = state.shutdown_tx.lock().await;
if let Some(sender) = tx.take() {
let _ = sender.send(());
StatusCode::OK
} else {
StatusCode::GONE
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::core::policy::test_support::{resident, rule};
use crate::generated::types::{PolicyRuleMode, PolicyRuleSeverity};
fn approved_hold_resolution() -> crate::daemon::hold_queue::HoldResolution {
crate::daemon::hold_queue::HoldResolution {
answer: crate::daemon::hold_queue::HoldAnswer::Approved,
verdict: Verdict::Allow,
result: "held_approved",
event: crate::zone_eval::Event::default(),
correlation_id: "corr-original".into(),
original_time: "2999-09-11T12:34:56.789Z".into(),
subject: Some("session-original".into()),
}
}
#[test]
fn timeout_hold_resolution_preserves_correlation_and_reports_no_host_prompt() {
let mut resolution = approved_hold_resolution();
resolution.answer = crate::daemon::hold_queue::HoldAnswer::Timeout;
resolution.verdict = Verdict::Block;
resolution.result = "held_timeout";
let envelope = hold_resolution_envelope("tool-original", &resolution);
let object = envelope.as_object().expect("CloudEvent object");
let mut keys: Vec<_> = object.keys().map(String::as_str).collect();
keys.sort_unstable();
assert_eq!(
keys,
[
"data",
"id",
"olenforced",
"olhostprompted",
"olresult",
"olverdict",
"source",
"specversion",
"subject",
"time",
"type",
]
);
assert_eq!(
envelope["data"],
serde_json::json!({
"tool_use_id": "tool-original",
"correlation_id": "corr-original",
})
);
assert_eq!(envelope["subject"], "session-original");
assert_eq!(envelope["time"], "2999-09-11T12:34:56.789Z");
assert_eq!(envelope["olverdict"], "block");
assert_eq!(envelope["olresult"], "held_timeout");
assert_eq!(envelope["olenforced"], 1);
assert_eq!(envelope["olhostprompted"], 0);
}
#[tokio::test]
async fn hold_resolution_uses_the_existing_cloud_worker_channel() {
let (tx, mut rx) = tokio::sync::mpsc::channel(1);
let envelope = hold_resolution_envelope("tool-original", &approved_hold_resolution());
try_enqueue_hold_resolution(&tx, None, "agent-isolated", envelope.clone());
let forwarded = rx.recv().await.expect("resolution queued");
assert_eq!(forwarded.agent_id, "agent-isolated");
assert_eq!(forwarded.envelope, envelope);
}
#[test]
fn doctor_probe_events_are_not_ingested() {
let probe = serde_json::json!({
"hook_event_name": "PreToolUse",
"tool_name": "OpenlatchDoctorProbe",
"tool_input": {},
});
assert!(
is_doctor_probe(Some(&probe)),
"the probe must be skipped before process_envelope, or every doctor run \
fabricates a tool call in the customer's audit trail"
);
let real = serde_json::json!({
"hook_event_name": "PreToolUse",
"tool_name": "Bash",
"tool_input": { "command": "echo OpenlatchDoctorProbe" },
});
assert!(!is_doctor_probe(Some(&real)));
assert!(!is_doctor_probe(Some(&serde_json::json!({}))));
assert!(
!is_doctor_probe(None),
"an envelope with no data at all is ordinary traffic, not a probe"
);
assert!(is_doctor_probe(Some(&serde_json::json!({
"tool_name": crate::hooks::bindings::codex_cli::DOCTOR_PROBE_TOOL_NAME,
}))));
}
#[test]
fn every_tool_shape_reaches_zone_mapping_with_full_input() {
let bundle = resident(vec![], true);
for (tool_name, tool_input) in [
("Read", serde_json::json!({"file_path": "/tmp/a"})),
("Bash", serde_json::json!({"command": ["echo a", "echo b"]})),
(
"mcp__storage__delete",
serde_json::json!({"Bucket": "prod", "Keys": ["a", "b"]}),
),
] {
let envelope: EventEnvelope = serde_json::from_value(serde_json::json!({
"specversion": "1.0",
"id": "evt-map",
"source": "claude-code",
"type": "pre_tool_use",
"time": "2026-09-10T12:00:00Z",
"subject": "session-map",
"data": {
"tool_use_id": "tool-map",
"tool_name": tool_name,
"tool_input": tool_input,
"cwd": "/work",
}
}))
.expect("valid envelope");
let mapped =
map_zone_event(&envelope, &bundle, Some("agent-map")).expect("mapped event");
assert_eq!(mapped.tool_name, tool_name);
assert_eq!(mapped.tool_input, tool_input, "tool={tool_name}");
assert_eq!(mapped.session_id, "session-map");
assert_eq!(mapped.tool_use_id, "tool-map");
assert_eq!(mapped.agent.unwrap().agent_id.as_deref(), Some("agent-map"));
}
}
#[test]
fn a_cline_payload_reaches_the_zone_evaluator() {
let bundle = resident(vec![], true);
let envelope: EventEnvelope = serde_json::from_value(serde_json::json!({
"specversion": "1.0",
"id": "evt-cline",
"source": "cline",
"type": "pre_tool_use",
"time": "2026-09-15T12:00:00Z",
"subject": "session-cline",
"data": {
"tool_use_id": "tool-cline",
"toolName": "run_commands",
"parameters": { "command": "rm -rf /tmp" },
"cwd": "/work",
}
}))
.expect("valid envelope");
let mapped = map_zone_event(&envelope, &bundle, Some("agent-cline")).expect("mapped event");
assert_eq!(mapped.tool_name, "run_commands");
assert_eq!(
mapped.tool_input,
serde_json::json!({ "command": "rm -rf /tmp" })
);
assert_eq!(mapped.session_id, "session-cline");
assert_eq!(mapped.tool_use_id, "tool-cline");
let data = envelope.data.as_ref().expect("data survives mapping");
assert!(data.get("tool_name").is_none());
assert!(data.get("tool_input").is_none());
assert_eq!(
data.get("toolName").and_then(serde_json::Value::as_str),
Some("run_commands")
);
}
#[test]
fn the_native_zone_spelling_wins_over_the_cline_fallback() {
let bundle = resident(vec![], true);
let envelope: EventEnvelope = serde_json::from_value(serde_json::json!({
"specversion": "1.0", "id": "evt-both", "source": "claude-code",
"type": "pre_tool_use", "time": "2026-09-15T12:00:00Z",
"data": {
"tool_name": "Bash",
"toolName": "run_commands",
"tool_input": { "command": "native" },
"parameters": { "command": "fallback" },
}
}))
.expect("valid envelope");
let mapped = map_zone_event(&envelope, &bundle, None).expect("mapped event");
assert_eq!(mapped.tool_name, "Bash");
assert_eq!(
mapped.tool_input,
serde_json::json!({ "command": "native" })
);
}
#[test]
fn claude_native_post_payloads_reach_result_and_error_evaluation() {
let bundle = resident(vec![], true);
let success: EventEnvelope = serde_json::from_value(serde_json::json!({
"specversion":"1.0","id":"evt-success","source":"claude-code",
"type":"post_tool_use","time":"2026-09-10T12:00:00Z","subject":"session-map",
"data":{"tool_use_id":"tool-1","tool_name":"Bash","tool_input":{"command":"cargo test"},"tool_response":"native stdout"}
})).expect("valid native PostToolUse");
let mapped = map_zone_event(&success, &bundle, None).expect("mapped success");
assert_eq!(mapped.tool_result, Some(serde_json::json!("native stdout")));
let failure: EventEnvelope = serde_json::from_value(serde_json::json!({
"specversion":"1.0","id":"evt-failure","source":"claude-code",
"type":"post_tool_use_failure","time":"2026-09-10T12:00:00Z","subject":"session-map",
"data":{"tool_use_id":"tool-2","tool_name":"Bash","tool_input":{"command":"cargo test"},"error":"exit 1","is_interrupt":false}
})).expect("valid native PostToolUseFailure");
let mapped = map_zone_event(&failure, &bundle, None).expect("mapped failure");
assert_eq!(mapped.event_type, "post_tool_use_failure");
assert_eq!(
mapped.tool_result.as_ref().and_then(|v| v.get("is_error")),
Some(&serde_json::json!(true))
);
assert_eq!(
mapped.tool_result.as_ref().and_then(|v| v.get("error")),
Some(&serde_json::json!("exit 1"))
);
}
#[test]
fn zone_agent_type_is_the_bundle_category_not_the_agent_platform() {
let mut bundle = resident(vec![], true);
bundle.zone = Some(
crate::zone_eval::load(serde_json::json!({
"schema_version": 2,
"organization_id": "org",
"revision": 1,
"built_at": "2026-09-10T00:00:00Z",
"enforcement_enabled": true,
"signature": null,
"artifacts": [],
"meta": {"agent_type": "coding"}
}))
.expect("zone bundle"),
);
let envelope: EventEnvelope = serde_json::from_value(serde_json::json!({
"specversion": "1.0", "id": "evt", "source": "claude-code",
"type": "pre_tool_use", "time": "2026-09-10T12:00:00Z",
"data": {"tool_name": "Read", "tool_input": {}}
}))
.expect("envelope");
let event = map_zone_event(&envelope, &bundle, Some("agent")).expect("mapped");
assert_eq!(event.agent.unwrap().agent_type.as_deref(), Some("coding"));
}
#[test]
fn a_block_with_an_unselected_rewrite_stays_blocked() {
let mut decision = ZoneDecision {
verdict: Verdict::Block,
rewrite: Some(crate::zone_eval::Rewrite {
lever: crate::generated::types::Lever("steer".into()),
artifact_id: Some("rewrite".into()),
steer_instruction: Some("try another way".into()),
}),
..ZoneDecision::default()
};
assert!(!degrade_unapplied_rewrite(&mut decision));
assert_eq!(decision.verdict, Verdict::Block);
}
#[test]
fn record_v2_names_the_artifact_shadow_and_actual_delivery() {
let mut document: serde_json::Value = serde_json::from_str(include_str!(
"../../tools/policy-seed/fixtures/schema-2.json"
))
.expect("schema-2 fixture");
document["artifacts"][0]["policy_id"] = serde_json::json!("policy-internal");
document["artifacts"][0]["spec_hash"] = serde_json::json!("sha256:spec");
document["artifacts"][0]["zone_layer"] =
serde_json::json!({"node_id": "unit-root", "path": ["unit-root"]});
let (_, bundle) = crate::core::policy::project_document(document).expect("bundle");
let mut monitor = ZoneDecision {
artifact_id: Some("seed-shell".into()),
atom_id: Some("seed-shell-atom".into()),
would_have_verdict: Some(Verdict::Block),
..ZoneDecision::default()
};
let mut obj = serde_json::Map::new();
stamp_record_v2(
&mut obj,
&monitor,
Some(&bundle),
"claude-code",
&HookEventType::PreToolUse,
false,
None,
);
assert_eq!(obj["olpolicyruleid"], "seed-shell");
assert_eq!(obj["olpolicyid"], "policy-internal");
assert_eq!(obj["ollayer"], "unit-root");
assert_eq!(obj["olspechash"], "sha256:spec");
assert_eq!(obj["olbinding"], "claude_code_hook");
assert_eq!(obj["olverdictshadow"], "block");
assert_eq!(obj["olenforced"], 0);
assert_eq!(obj["olresult"], "flagged");
monitor.verdict = Verdict::Block;
monitor.mode = Some(crate::generated::types::PolicyMode("enforce".into()));
monitor.would_have_verdict = None;
let mut post = serde_json::Map::new();
stamp_record_v2(
&mut post,
&monitor,
None,
"claude-code",
&HookEventType::PostToolUse,
false,
None,
);
assert_eq!(post["olenforced"], 0);
assert_eq!(post["olresult"], "allowed");
let mut pre = serde_json::Map::new();
stamp_record_v2(
&mut pre,
&monitor,
None,
"claude-code",
&HookEventType::PreToolUse,
false,
None,
);
assert_eq!(pre["olenforced"], 1);
assert_eq!(pre["olresult"], "blocked");
let final_block = crate::daemon::optimize::Delivery {
enforced: Some(true),
result: Some("blocked"),
..crate::daemon::optimize::Delivery::default()
};
let mut final_record = serde_json::Map::new();
stamp_record_v2(
&mut final_record,
&monitor,
None,
"claude-code",
&HookEventType::PreToolUse,
false,
Some(&final_block),
);
assert_eq!(final_record["olpolicyruleid"], "seed-shell");
assert_eq!(final_record["olenforced"], 1);
assert_eq!(final_record["olresult"], "blocked");
}
#[test]
fn evaluated_monitor_ask_and_block_reach_record_v2_egress() {
for shadow_verdict in ["ask", "block"] {
let mut document: serde_json::Value = serde_json::from_str(include_str!(
"../../tools/policy-seed/fixtures/schema-2.json"
))
.expect("schema-2 fixture");
document["artifacts"][0]["policy_id"] = serde_json::json!("policy-monitor");
document["artifacts"][0]["policy_public_id"] = serde_json::json!("POL-MONITOR");
document["artifacts"][0]["mode"] = serde_json::json!("monitor");
document["artifacts"][0]["body"]["verdict"] = serde_json::json!(shadow_verdict);
let (_, bundle) = crate::core::policy::project_document(document).expect("bundle");
let envelope: EventEnvelope = serde_json::from_value(serde_json::json!({
"specversion": "1.0", "id": format!("evt-monitor-{shadow_verdict}"),
"source": "claude-code", "type": "pre_tool_use",
"time": "2026-09-10T12:00:00Z", "subject": "session-monitor",
"data": {"tool_use_id": "tool-monitor", "tool_name": "Bash",
"tool_input": {"command": "curl https://example.test"}}
}))
.expect("valid envelope");
let event = map_zone_event(&envelope, &bundle, None).expect("mapped event");
let (decision, _) = crate::zone_eval::evaluate(
bundle.zone.as_ref().expect("zone bundle"),
&event,
None,
1_757_505_600_000,
);
assert_eq!(decision.verdict, Verdict::Allow);
assert_eq!(
decision
.would_have_verdict
.map(|v| v.to_string())
.as_deref(),
Some(shadow_verdict)
);
let mut obj = serde_json::Map::new();
stamp_record_v2_if_winner(
&mut obj,
Some(&decision),
None,
Some(&bundle),
"claude-code",
&HookEventType::PreToolUse,
false,
None,
);
assert_eq!(obj["olatomid"], "seed-shell-atom");
assert_eq!(obj["olpolicyruleid"], "seed-shell");
assert_eq!(obj["olpolicyid"], "policy-monitor");
assert_eq!(obj["olverdictshadow"], shadow_verdict);
assert_eq!(obj["olenforced"], 0);
assert_eq!(obj["olresult"], "flagged");
}
}
#[test]
fn enforcing_legacy_match_keeps_its_provenance_over_a_zone_decision() {
let bundle = resident(
vec![rule("OL-LEGACY-WINNER", "*rm*", PolicyRuleMode::Enforce)],
true,
);
let legacy = match_of("OL-LEGACY-WINNER", false);
let zone = ZoneDecision {
verdict: Verdict::Allow,
artifact_id: Some("zone-loser".into()),
atom_id: Some("zone-atom".into()),
..ZoneDecision::default()
};
let mut obj = serde_json::Map::new();
stamp_policy_extensions(
&mut obj,
Some(&legacy),
Some(&bundle),
false,
bundle.built_at,
);
stamp_record_v2_if_winner(
&mut obj,
Some(&zone),
Some(&legacy),
Some(&bundle),
"claude-code",
&HookEventType::PreToolUse,
false,
None,
);
assert_eq!(obj["olpolicyruleid"], "OL-LEGACY-WINNER");
assert!(obj.get("olatomid").is_none());
assert!(obj.get("olresult").is_none());
assert!(obj.get("olenforced").is_none());
}
fn match_of(rule_id: &str, shadow: bool) -> PolicyMatch {
PolicyMatch {
rule_id: rule_id.to_string(),
reason: format!("{rule_id} says no"),
severity: PolicyRuleSeverity::High,
mode: if shadow {
PolicyRuleMode::Observe
} else {
PolicyRuleMode::Enforce
},
shadow,
}
}
fn outcome(verdict: Verdict, id: &str) -> BatchOutcome {
BatchOutcome {
verdict,
event_type: match verdict {
Verdict::Ask => HookEventType::Stop,
_ => HookEventType::PreToolUse,
},
event_id: id.to_string(),
policy_match: match verdict {
Verdict::Block => Some(match_of("OL-CMD-ENF", false)),
_ => None,
},
zone_decision: None,
hold_wait: None,
delivery: crate::daemon::optimize::Delivery::default(),
}
}
#[test]
fn a_non_shadow_match_blocks_and_a_shadow_match_does_not() {
assert_eq!(
verdict_for(
&HookEventType::PreToolUse,
None,
Some(&match_of("OL-CMD-A", false))
),
Verdict::Block
);
assert_eq!(
verdict_for(
&HookEventType::PreToolUse,
None,
Some(&match_of("OL-CMD-A", true))
),
Verdict::Allow,
"observe never blocks"
);
assert_eq!(
verdict_for(&HookEventType::PreToolUse, None, None),
Verdict::Allow
);
assert_eq!(
verdict_for(&HookEventType::Stop, None, None),
Verdict::Allow
);
}
#[test]
fn block_outranks_the_stop_arm() {
assert_eq!(
verdict_for(
&HookEventType::Stop,
None,
Some(&match_of("OL-CMD-A", false))
),
Verdict::Block
);
}
#[test]
fn the_fold_is_most_restrictive_wins() {
for (a, b) in [
(Verdict::Allow, Verdict::Optimize),
(Verdict::Allow, Verdict::Ask),
(Verdict::Allow, Verdict::Block),
(Verdict::Optimize, Verdict::Ask),
(Verdict::Optimize, Verdict::Block),
(Verdict::Ask, Verdict::Block),
] {
for (first, second) in [(a, b), (b, a)] {
let mut worst = None;
join_most_restrictive(&mut worst, outcome(first, "evt_1"));
join_most_restrictive(&mut worst, outcome(second, "evt_2"));
let got = worst.expect("folded");
assert_eq!(
verdict_rank(got.verdict),
verdict_rank(a).max(verdict_rank(b)),
"{first:?} then {second:?} must fold to the more restrictive"
);
}
}
}
#[test]
fn the_fold_is_order_independent() {
let orders = [
[Verdict::Block, Verdict::Ask, Verdict::Allow],
[Verdict::Allow, Verdict::Block, Verdict::Ask],
[Verdict::Ask, Verdict::Allow, Verdict::Block],
];
for order in orders {
let mut worst = None;
for (i, v) in order.iter().enumerate() {
join_most_restrictive(&mut worst, outcome(*v, &format!("evt_{i}")));
}
assert_eq!(worst.expect("folded").verdict, Verdict::Block, "{order:?}");
}
}
#[test]
fn equal_rank_keeps_the_first_envelope() {
let mut worst = None;
join_most_restrictive(&mut worst, outcome(Verdict::Allow, "evt_1"));
join_most_restrictive(&mut worst, outcome(Verdict::Allow, "evt_2"));
assert_eq!(worst.expect("folded").event_id, "evt_1");
}
#[test]
fn the_deciding_envelope_carries_its_own_match() {
let mut worst = None;
join_most_restrictive(&mut worst, outcome(Verdict::Allow, "evt_1"));
join_most_restrictive(&mut worst, outcome(Verdict::Block, "evt_2"));
let got = worst.expect("folded");
assert_eq!(got.event_id, "evt_2");
assert_eq!(
got.policy_match.expect("block carries its match").rule_id,
"OL-CMD-ENF"
);
}
fn stamped(
policy_match: Option<&PolicyMatch>,
bundle: Option<&ResidentBundle>,
offline: bool,
now: SystemTime,
) -> serde_json::Map<String, serde_json::Value> {
let mut obj = serde_json::Map::new();
stamp_policy_extensions(&mut obj, policy_match, bundle, offline, now);
obj
}
#[test]
fn a_deny_stamps_rule_id_revision_age_and_offline_but_no_shadow() {
let bundle = resident(
vec![rule("OL-CMD-ENF", "*rm*", PolicyRuleMode::Enforce)],
true,
);
let now = bundle.built_at + Duration::from_secs(90);
let obj = stamped(
Some(&match_of("OL-CMD-ENF", false)),
Some(&bundle),
false,
now,
);
assert_eq!(obj.get("olpolicyruleid").unwrap(), "OL-CMD-ENF");
assert_eq!(obj.get("olpolicybundlerev").unwrap(), 42);
assert_eq!(obj.get("olpolicybundleage").unwrap(), 90);
assert_eq!(obj.get("olpolicyoffline").unwrap(), false);
assert!(
obj.get("olverdictshadow").is_none(),
"an enforced deny is not a shadow verdict"
);
}
#[test]
fn an_observe_match_stamps_the_shadow_verdict() {
let bundle = resident(vec![], true);
let obj = stamped(
Some(&match_of("OL-CMD-OBS", true)),
Some(&bundle),
false,
bundle.built_at,
);
assert_eq!(obj.get("olverdictshadow").unwrap(), "block");
assert_eq!(obj.get("olpolicyruleid").unwrap(), "OL-CMD-OBS");
}
#[test]
fn no_match_still_reports_the_bundle_revision() {
let bundle = resident(vec![], true);
let obj = stamped(None, Some(&bundle), false, bundle.built_at);
assert_eq!(obj.get("olpolicybundlerev").unwrap(), 42);
assert!(obj.get("olpolicyruleid").is_none());
assert!(obj.get("olverdictshadow").is_none());
}
#[test]
fn with_no_bundle_the_bundle_fields_are_omitted_and_offline_still_reported() {
let obj = stamped(None, None, true, SystemTime::now());
assert!(obj.get("olpolicybundlerev").is_none());
assert!(obj.get("olpolicybundleage").is_none());
assert_eq!(obj.get("olpolicyoffline").unwrap(), true);
}
#[test]
fn bundle_age_is_clamped_at_zero_when_the_local_clock_is_behind() {
let bundle = resident(vec![], true);
let behind = bundle.built_at - Duration::from_secs(3_600);
let obj = stamped(None, Some(&bundle), false, behind);
assert_eq!(obj.get("olpolicybundleage").unwrap(), 0);
}
fn runtime(
bundle: Option<ResidentBundle>,
last_poll_ok_at: i64,
) -> crate::daemon::PolicyRuntime {
crate::daemon::PolicyRuntime {
handle: crate::core::policy::new_handle(bundle),
last_fetch_ok: Arc::new(std::sync::atomic::AtomicBool::new(true)),
last_poll_ok_at: Arc::new(std::sync::atomic::AtomicI64::new(last_poll_ok_at)),
floor: Arc::new(std::sync::RwLock::new(None)),
sessions: Arc::new(crate::daemon::session_state::SessionStateManager::default()),
dispatch_sessions: Arc::new(
crate::daemon::session_state::DispatchStateManager::default(),
),
holds: Arc::new(crate::daemon::hold_queue::HoldQueue::default()),
reinforcement: Arc::new(crate::daemon::reinforcement::ReinforcementStore::default()),
}
}
#[test]
fn health_degrades_on_a_failed_egress_even_with_every_subsystem_running() {
assert_eq!(health_status(false, false), "ok");
assert_eq!(
health_status(false, true),
"degraded",
"a daemon that cannot get a byte out is not healthy"
);
assert_eq!(health_status(true, false), "degraded");
assert_eq!(health_status(true, true), "degraded");
}
fn proxied_egress() -> crate::egress::EgressConfig {
let toml = crate::egress::ProxyToml {
mode: Some("manual".into()),
url: Some("http://alice-proxy.corp:8080".into()),
..Default::default()
};
struct NoEnv;
impl crate::egress::EnvSource for NoEnv {
fn var(&self, _key: &str) -> Option<String> {
None
}
}
crate::egress::EgressConfig::resolve(Some(&toml), &NoEnv, 7443, 7444).expect("resolve")
}
#[test]
fn egress_metrics_report_a_healthy_direct_host_as_up() {
let m = egress_metrics(&crate::egress::EgressState::new(
&crate::egress::EgressConfig::direct(),
));
assert!(m.up);
assert_eq!(m.status, "ok");
assert_eq!(m.consecutive_failures, 0);
assert!(!m.proxy_in_use);
}
#[test]
fn egress_metrics_go_down_only_at_the_threshold() {
let state = crate::egress::EgressState::new(&crate::egress::EgressConfig::direct());
state.record_failure(crate::error::ERR_PROXY_UNREACHABLE, "refused");
let one = egress_metrics(&state);
assert!(one.up, "one blip is not an outage");
assert_eq!(one.consecutive_failures, 1);
state.record_failure(crate::error::ERR_PROXY_UNREACHABLE, "refused");
let two = egress_metrics(&state);
assert!(!two.up);
assert_eq!(two.status, "failed");
assert_eq!(two.consecutive_failures, 2);
}
#[test]
fn the_unauthenticated_scrape_never_carries_the_proxy_url() {
let cfg = proxied_egress();
let state = crate::egress::EgressState::new(&cfg);
let m = egress_metrics(&state);
assert!(m.proxy_in_use);
let rendered = serde_json::json!({
"egress_up": m.up,
"egress_status": m.status,
"egress_consecutive_failures": m.consecutive_failures,
"proxy_in_use": m.proxy_in_use,
})
.to_string();
assert!(
!rendered.contains("alice-proxy.corp"),
"the proxy host leaked into an unauthenticated surface: {rendered}"
);
assert!(!rendered.contains("8080"));
}
fn stamped_envelope(snapshot: &crate::egress::EgressSnapshot) -> serde_json::Value {
let envelope = EventEnvelope {
specversion: "1.0".into(),
id: new_event_id(),
source: AgentType::ClaudeCode,
type_: HookEventType::PreToolUse,
time: chrono::Utc::now(),
datacontenttype: None,
subject: Some("sess_abc".into()),
data: None,
os: None,
arch: None,
localipv4: None,
localipv6: None,
publicipv4: None,
publicipv6: None,
clientversion: None,
agentversion: None,
agentid: None,
osuser: None,
gitemail: None,
provideracct: None,
wireformat: None,
};
let mut wire = serde_json::to_value(&envelope).expect("serialise envelope");
stamp_egress_extensions(
wire.as_object_mut().expect("an envelope is a JSON object"),
egress_stamp(snapshot),
);
wire
}
fn snapshot_of(cfg: &crate::egress::EgressConfig) -> crate::egress::EgressSnapshot {
crate::egress::EgressSnapshot::from_config(cfg)
}
#[test]
fn a_direct_host_stamps_none_of_the_three() {
let wire = stamped_envelope(&snapshot_of(&crate::egress::EgressConfig::direct()));
for attribute in ["proxytype", "proxysource", "proxyintercepted"] {
assert!(
wire.get(attribute).is_none(),
"{attribute} must be ABSENT on a direct host, not null: {wire}"
);
}
}
#[test]
fn a_proxied_host_stamps_type_and_source_but_not_interception() {
let wire = stamped_envelope(&snapshot_of(&proxied_egress()));
assert_eq!(wire["proxytype"], "http");
assert!(
wire.get("proxyintercepted").is_none(),
"absence must mean 'not yet known', never 'not intercepted': {wire}"
);
let rendered = wire.to_string();
assert!(!rendered.contains("alice-proxy.corp"), "{rendered}");
assert!(!rendered.contains("8080"), "{rendered}");
}
#[test]
fn an_observed_handshake_stamps_the_interception_verdict() {
let mut snapshot = snapshot_of(&proxied_egress());
snapshot.source = Some(crate::egress::ProxySource::Windows);
snapshot.tls_intercepted = Some(true);
let wire = stamped_envelope(&snapshot);
assert_eq!(wire["proxytype"], "http");
assert_eq!(wire["proxysource"], "windows");
assert_eq!(wire["proxyintercepted"], true);
snapshot.tls_intercepted = Some(false);
assert_eq!(stamped_envelope(&snapshot)["proxyintercepted"], false);
}
#[test]
fn a_pac_route_stamps_pac_not_the_returned_scheme() {
let mut snapshot = snapshot_of(&proxied_egress());
snapshot.source = Some(crate::egress::ProxySource::Pac);
let wire = stamped_envelope(&snapshot);
assert_eq!(wire["proxytype"], "pac");
assert_eq!(wire["proxysource"], "pac");
}
#[test]
fn a_proxy_with_no_recorded_source_still_stamps_its_type() {
let wire = stamped_envelope(&snapshot_of(&proxied_egress()));
assert_eq!(wire["proxytype"], "http");
assert!(wire.get("proxysource").is_none());
}
#[test]
fn metrics_report_policy_disabled_as_such() {
let m = policy_metrics(None, SystemTime::now());
assert!(!m.enabled);
assert!(!m.has_bundle);
assert!(m.schema_version.is_null());
assert!(m.bundle_digest.is_null());
assert!(m.revision.is_null());
assert!(m.bundle_age_seconds.is_null());
assert!(m.last_poll_ok_secs.is_null());
}
#[test]
fn metrics_distinguish_no_bundle_from_a_resident_one() {
let cold = policy_metrics(Some(&runtime(None, 0)), SystemTime::now());
assert!(cold.enabled);
assert!(!cold.has_bundle, "never fetched: allow, marked no-bundle");
assert!(cold.schema_version.is_null());
assert!(cold.bundle_digest.is_null());
assert!(cold.revision.is_null());
assert!(cold.last_poll_ok_secs.is_null(), "0 means never, not now");
let bundle = resident(vec![], true);
let now = bundle.built_at + Duration::from_secs(120);
let warm = policy_metrics(Some(&runtime(Some(bundle), 1_784_624_400)), now);
assert!(warm.has_bundle, "last-good resident: keep enforcing");
assert_eq!(warm.schema_version, 1);
assert!(warm.bundle_digest.is_null());
assert_eq!(warm.revision, 42);
assert_eq!(warm.bundle_age_seconds, 120);
}
#[test]
fn metrics_report_the_resident_agent_function() {
let mut bundle = resident(vec![], true);
bundle.agent_function = Some(crate::generated::types::AgentFunction::ItOps);
let m = policy_metrics(Some(&runtime(Some(bundle), 1_000)), SystemTime::now());
assert_eq!(
m.agent_function, "it_ops",
"the wire spelling, not the variant"
);
let contextless = policy_metrics(
Some(&runtime(Some(resident(vec![], true)), 1_000)),
SystemTime::now(),
);
assert!(contextless.agent_function.is_null());
assert!(policy_metrics(Some(&runtime(None, 0)), SystemTime::now())
.agent_function
.is_null());
}
#[test]
fn metrics_bundle_age_is_clamped_and_poll_age_is_seconds_since_the_last_ok() {
let bundle = resident(vec![], true);
let behind = bundle.built_at - Duration::from_secs(60);
let m = policy_metrics(Some(&runtime(Some(bundle), 1_000)), behind);
assert_eq!(m.bundle_age_seconds, 0, "clock skew never goes negative");
let expected = behind
.duration_since(UNIX_EPOCH)
.expect("after epoch")
.as_secs() as i64
- 1_000;
assert_eq!(m.last_poll_ok_secs, expected);
}
fn queued_alert() -> crate::daemon::config_monitor::PendingAlert {
crate::daemon::config_monitor::PendingAlert {
alert_id: "alert_1".to_string(),
session_ref_id: "sess_abc".to_string(),
severity: "high".to_string(),
headline: "Configuration alert pending".to_string(),
body: "MCP server 'evil' was added.".to_string(),
source_id: "claude-code:mcp:deadbeef".to_string(),
created_at: "2026-07-21T09:00:00Z".to_string(),
delivered: false,
}
}
#[test]
fn attaching_an_alert_to_a_block_would_overwrite_the_rules_reason() {
let m = match_of("OL-CMD-ENF", false);
let mut response = VerdictResponse::block(
"evt_1".to_string(),
1.0,
Some(m.reason.clone()),
Some(m.severity.to_string()),
Some(m.rule_id.clone()),
);
assert_eq!(response.reason.as_deref(), Some("OL-CMD-ENF says no"));
assert!(response.context.is_none());
attach_alert_context(&mut response, &queued_alert());
assert_eq!(
response.reason.as_deref(),
Some("Configuration alert pending"),
"attach_alert_context overwrites `reason` unconditionally"
);
assert!(
response.context.is_some(),
"and the translator prefers context over reason on a block"
);
}
#[test]
fn a_guarded_block_renders_the_rules_reason_to_the_agent() {
let m = match_of("OL-CMD-ENF", false);
let response = VerdictResponse::block(
"evt_1".to_string(),
1.0,
Some(m.reason.clone()),
Some(m.severity.to_string()),
Some(m.rule_id.clone()),
);
let out = crate::hook_output::translate(
"claude-code",
"pre_tool_use",
&crate::hook_output::Verdict {
decision: &response.verdict.to_string(),
reason: response.reason.as_deref(),
context: None,
},
);
assert_eq!(
out["hookSpecificOutput"]["permissionDecision"], "deny",
"a policy block must reach the agent as a native deny"
);
assert_eq!(
out["hookSpecificOutput"]["permissionDecisionReason"],
"OL-CMD-ENF says no"
);
}
#[test]
fn extension_names_are_lowercase_alphanumeric() {
let bundle = resident(vec![], true);
let obj = stamped(
Some(&match_of("OL-CMD-A", true)),
Some(&bundle),
true,
bundle.built_at,
);
assert_eq!(
obj.len(),
5,
"five of the six; olverdict is stamped by the caller"
);
for name in obj.keys() {
assert!(
name.chars()
.all(|c| c.is_ascii_lowercase() || c.is_ascii_digit()),
"illegal CloudEvents extension name: {name}"
);
}
}
fn fake_signals() -> IdentitySignals {
IdentitySignals {
git_email: Some("dev@example.com".to_string()),
provider_account: Some("dev@anthropic.example".to_string()),
}
}
fn stamp(
mode: ProcessMode,
capture: bool,
subject: Option<&str>,
) -> (Option<String>, IdentitySignals) {
identity_stamp(
mode,
capture,
subject,
Some("/repo"),
|| Some("injected-os-user".to_string()),
|_, _| fake_signals(),
)
}
#[test]
fn capture_on_with_a_session_stamps_all_three() {
let (osuser, signals) = stamp(ProcessMode::Live, true, Some("sess_a"));
assert_eq!(osuser.as_deref(), Some("injected-os-user"));
assert_eq!(signals, fake_signals());
}
#[test]
fn capture_off_or_no_bundle_stamps_nothing() {
let (osuser, signals) = stamp(ProcessMode::Live, false, Some("sess_a"));
assert_eq!(osuser, None);
assert_eq!(signals, IdentitySignals::default());
}
#[test]
fn replay_is_never_identity_stamped_even_with_capture_on() {
let (osuser, signals) = stamp(ProcessMode::Replay, true, Some("sess_a"));
assert_eq!(osuser, None);
assert_eq!(signals, IdentitySignals::default());
}
#[test]
fn a_sessionless_envelope_carries_osuser_only() {
let (osuser, signals) = identity_stamp(
ProcessMode::Live,
true,
None,
Some("/repo"),
|| Some("injected-os-user".to_string()),
|_, _| panic!("no subject means no session lookup"),
);
assert_eq!(osuser.as_deref(), Some("injected-os-user"));
assert_eq!(signals, IdentitySignals::default());
}
#[test]
fn the_events_cwd_reaches_the_resolver() {
let (_, signals) = identity_stamp(
ProcessMode::Live,
true,
Some("sess_a"),
Some("/repo/sub"),
|| None,
|session_id, cwd| {
assert_eq!(session_id, "sess_a");
IdentitySignals {
git_email: cwd.map(str::to_owned),
provider_account: None,
}
},
);
assert_eq!(signals.git_email.as_deref(), Some("/repo/sub"));
}
}