//! Policies return proposals. Validation and failure disposition belong here. use crate::{ behavior::Release, protocol::*, remote::callback::{CallbackBroker, CallbackContext, CallbackRequest}, worker::{no_host_operations, WorkCommand, Worker, WorkerConfig}, }; use anyhow::{bail, ensure, Context, Result}; use serde::{Deserialize, Serialize}; use serde_json::{json, Value}; use std::{sync::Arc, time::Duration}; #[derive(Clone, Debug, Serialize, Deserialize)] #[serde(tag = "action", rename_all = "snake_case", deny_unknown_fields)] pub enum DeliveryDecision { Allow, Hold { reason: String }, } #[derive(Clone, Debug, Serialize, Deserialize)] #[serde(tag = "action", rename_all = "snake_case", deny_unknown_fields)] pub enum AttentionDecision { Wake { priority: Priority }, Defer { seconds: f64 }, Ignore { reason: String }, } #[derive(Clone, Copy, Debug, Serialize, Deserialize)] #[serde(rename_all = "snake_case")] pub enum Priority { Operator, Directed, Ambient, } #[derive(Clone, Debug)] pub struct ReviewedDelivery { pub decision: DeliveryDecision, pub fault: Option, } pub struct Policy { pub release: Release, pub worker: Worker, callbacks: CallbackBroker, } impl Policy { pub async fn start( config: &WorkerConfig, session: SessionId, release: Release, ) -> Result> { let worker = Worker::spawn(config, session, Role::Hooks, Some(&release)).await?; let revision = release.id.clone(); let hook_worker = worker.clone(); let callbacks = CallbackBroker::new( 32, Arc::new(move |callback| { let hook_worker = hook_worker.clone(); let revision = revision.clone(); Box::pin(async move { ensure!( callback.context.policy_revision == revision, "callback policy revision is not active" ); let outcome = hook_worker .call( WorkCommand::Invoke { id: OperationId::fresh(), slot: callback.slot, input: callback.input, }, Duration::from_secs(2), no_host_operations(), ) .await?; match outcome { Outcome::Hook { value } => Ok(value), Outcome::Failed { error } => bail!("hook failed: {error}"), _ => bail!("invalid policy result"), } }) }), )?; Ok(Arc::new(Self { release, worker, callbacks, })) } pub async fn validate_synthetic(&self) -> Result<()> { self.attention(1, "operator", 0) .await .context("candidate attention hook synthetic validation failed")?; let test_effect = EffectId::fresh(); let test_review = self.review(&test_effect, "candidate synthetic test").await; if let Some(fault) = test_review.fault { bail!("candidate delivery hook synthetic validation failed: {fault}"); } Ok(()) } async fn invoke(&self, slot: HookSlot, input: Value) -> Result { let outcome = self .worker .call( WorkCommand::Invoke { id: OperationId::fresh(), slot, input, }, Duration::from_secs(2), no_host_operations(), ) .await?; match outcome { Outcome::Hook { value } => Ok(value), Outcome::Failed { error } => bail!("hook failed: {error}"), _ => bail!("invalid policy result"), } } /// Unavailable required review means hold, never accidental approval. pub async fn review(&self, effect: &EffectId, text: &str) -> ReviewedDelivery { let value = self .invoke( HookSlot::Delivery, json!({"effect_id": effect, "destination":"operator", "text":text}), ) .await; Self::decode_review(value) } /// Reviews an effect through the independent callback lane. The context is /// immutable for the originating execution and policy release. pub async fn review_with_context( &self, context: CallbackContext, effect: &EffectId, text: &str, ) -> ReviewedDelivery { let result = async { let request = CallbackRequest::new( context, HookSlot::Delivery, json!({"effect_id": effect, "destination":"operator", "text":text}), )?; self.callbacks.submit(request, Duration::from_secs(2)).await } .await; Self::decode_review(result) } pub async fn last_callback_context(&self) -> Option { self.callbacks.last_context().await } pub fn callback_stats(&self) -> crate::remote::callback::CallbackStats { self.callbacks.stats() } fn decode_review(result: Result) -> ReviewedDelivery { let result = result.and_then(|value| { let decision: DeliveryDecision = serde_json::from_value(value)?; if let DeliveryDecision::Hold { reason } = &decision { ensure!( !reason.trim().is_empty() && reason.len() <= 2048, "invalid hold reason" ); } Ok(decision) }); match result { Ok(decision) => ReviewedDelivery { decision, fault: None, }, Err(error) => ReviewedDelivery { decision: DeliveryDecision::Hold { reason: "required delivery policy unavailable or invalid".into(), }, fault: Some(format!("{error:#}")), }, } } /// This slice exposes and validates attention policy but does not yet own a /// model scheduler. The caller must not acknowledge input merely for a plan. pub async fn attention( &self, event_id: i64, source: &str, received_ms: i64, ) -> Result { let decision: AttentionDecision = serde_json::from_value( self.invoke( HookSlot::Attention, json!({"event_id":event_id,"source":source,"received_ms":received_ms}), ) .await?, )?; match &decision { AttentionDecision::Defer { seconds } => ensure!( seconds.is_finite() && *seconds > 0.0 && *seconds <= 600.0, "invalid defer duration" ), AttentionDecision::Ignore { reason } => ensure!( !reason.trim().is_empty() && reason.len() <= 2048, "invalid ignore reason" ), AttentionDecision::Wake { priority: Priority::Operator, } => ensure!( source == "operator", "external input cannot request operator authority" ), _ => (), } if source == "operator" { ensure!( matches!( decision, AttentionDecision::Wake { priority: Priority::Operator } ), "operator attention cannot be suppressed" ); } Ok(decision) } pub fn describe(&self, slot: HookSlot) -> Value { json!({"slot":slot,"kind":"decision","placement":"hook_worker","revision":self.release.id, "handler":self.release.hooks.get(&slot),"deadline_ms":2000, "failure":match slot { HookSlot::Delivery => "hold_effect", HookSlot::Attention => "retain_pending_input" }, "host_effects":false}) } pub fn list(&self) -> Value { json!({"revision":self.release.id,"slots":[self.describe(HookSlot::Attention),self.describe(HookSlot::Delivery)]}) } }