ive harnessed the harness
Something went wrong. Try again.
8.4 kB · 242 lines
Rust
at main
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243//! 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<String>,}
pub struct Policy { pub release: Release, pub worker: Worker, callbacks: CallbackBroker,}impl Policy { pub async fn start( config: &WorkerConfig, session: SessionId, release: Release, ) -> Result<Arc<Self>> { 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<Value> { 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<CallbackContext> { self.callbacks.last_context().await }
pub fn callback_stats(&self) -> crate::remote::callback::CallbackStats { self.callbacks.stats() }
fn decode_review(result: Result<Value>) -> 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<AttentionDecision> { 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)]}) }}