ive harnessed the harness
Something went wrong. Try again.
6.6 kB · 201 lines
Rust
at main
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202//! A bounded, task-backed callback lane for host policy work.//!//! Remote workbench control remains owned by the worker actor. Callback policy//! is enqueued on this independent lane, so a remote cell never waits while a//! transport or session lock is held.use crate::protocol::{Generation, HookSlot, OperationId, SessionId, TargetStamp};use anyhow::{ensure, Context, Result};use serde::{Deserialize, Serialize};use serde_json::Value;use std::{ future::Future, pin::Pin, sync::{ atomic::{AtomicUsize, Ordering}, Arc, }, time::{Duration, Instant},};use tokio::sync::{mpsc, oneshot};
/// The placement selected for the originating workbench generation.#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]#[serde(tag = "kind", rename_all = "snake_case", deny_unknown_fields)]pub enum CallbackTarget { Local, Remote { stamp: TargetStamp },}
/// Immutable execution identity carried with every policy callback.#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]#[serde(deny_unknown_fields)]pub struct CallbackContext { pub session_id: SessionId, pub execution_id: OperationId, pub target: CallbackTarget, pub generation: Generation, pub policy_revision: String,}
impl CallbackContext { pub fn local( session_id: SessionId, execution_id: OperationId, generation: Generation, policy_revision: impl Into<String>, ) -> Self { Self { session_id, execution_id, target: CallbackTarget::Local, generation, policy_revision: policy_revision.into(), } }
pub fn remote( session_id: SessionId, execution_id: OperationId, stamp: TargetStamp, generation: Generation, policy_revision: impl Into<String>, ) -> Self { Self { session_id, execution_id, target: CallbackTarget::Remote { stamp }, generation, policy_revision: policy_revision.into(), } }
fn validate(&self) -> Result<()> { ensure!( !self.policy_revision.trim().is_empty(), "callback policy revision is empty" ); if let CallbackTarget::Remote { stamp } = &self.target { ensure!( !stamp.target_id.as_str().trim().is_empty() && !stamp.namespace_id.as_str().trim().is_empty() && !stamp.namespace_generation.as_str().trim().is_empty(), "callback target binding is incomplete" ); } Ok(()) }}
/// One host callback request, including the source worker's complete pin.#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]#[serde(deny_unknown_fields)]pub struct CallbackRequest { pub context: CallbackContext, pub slot: HookSlot, pub input: Value,}
impl CallbackRequest { pub fn new(context: CallbackContext, slot: HookSlot, input: Value) -> Result<Self> { context.validate()?; Ok(Self { context, slot, input, }) }}
/// A callback handler runs outside the caller and broker queue task.pub type CallbackFuture = Pin<Box<dyn Future<Output = Result<Value>> + Send>>;pub type CallbackHandler = Arc<dyn Fn(CallbackRequest) -> CallbackFuture + Send + Sync>;
struct QueuedCallback { request: CallbackRequest, reply: oneshot::Sender<Result<Value>>,}
#[derive(Default)]struct BrokerState { submitted: AtomicUsize, serviced: AtomicUsize,}
#[derive(Clone, Debug, Default, PartialEq, Eq)]pub struct CallbackStats { pub submitted: usize, pub serviced: usize,}
/// A bounded callback broker. It has no session or database lock and can be/// awaited from the remote worker's host-request path without re-entering it.#[derive(Clone)]pub struct CallbackBroker { tx: mpsc::Sender<QueuedCallback>, state: Arc<BrokerState>, last_context: Arc<tokio::sync::Mutex<Option<CallbackContext>>>,}
impl CallbackBroker { pub fn new(capacity: usize, handler: CallbackHandler) -> Result<Self> { ensure!(capacity > 0, "callback broker capacity must be positive"); let (tx, mut rx) = mpsc::channel::<QueuedCallback>(capacity); let state = Arc::new(BrokerState::default()); let last_context = Arc::new(tokio::sync::Mutex::new(None)); let worker_state = state.clone(); let worker_context = last_context.clone(); tokio::spawn(async move { while let Some(queued) = rx.recv().await { if queued.reply.is_closed() { continue; } worker_state.serviced.fetch_add(1, Ordering::Relaxed); *worker_context.lock().await = Some(queued.request.context.clone()); let handler = handler.clone(); tokio::spawn(async move { let result = handler(queued.request).await; let _ = queued.reply.send(result); }); } }); Ok(Self { tx, state, last_context, }) }
/// Enqueue and await a callback without holding caller-owned state. pub async fn submit(&self, request: CallbackRequest, budget: Duration) -> Result<Value> { request.context.validate()?; ensure!( budget > Duration::ZERO, "callback deadline must be positive" ); self.state.submitted.fetch_add(1, Ordering::Relaxed); let (reply, receive) = oneshot::channel(); let deadline = Instant::now() + budget; let queued = QueuedCallback { request, reply }; let remaining = deadline.saturating_duration_since(Instant::now()); tokio::time::timeout(remaining, self.tx.send(queued)) .await .context("callback broker enqueue timed out")? .context("callback broker is closed")?; let remaining = deadline.saturating_duration_since(Instant::now()); tokio::time::timeout(remaining, receive) .await .context("callback handler timed out")? .context("callback handler dropped its reply")? }
pub fn stats(&self) -> CallbackStats { CallbackStats { submitted: self.state.submitted.load(Ordering::Relaxed), serviced: self.state.serviced.load(Ordering::Relaxed), } }
pub async fn last_context(&self) -> Option<CallbackContext> { self.last_context.lock().await.clone() }}