//! 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, ) -> 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, ) -> 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 { context.validate()?; Ok(Self { context, slot, input, }) } } /// A callback handler runs outside the caller and broker queue task. pub type CallbackFuture = Pin> + Send>>; pub type CallbackHandler = Arc CallbackFuture + Send + Sync>; struct QueuedCallback { request: CallbackRequest, reply: oneshot::Sender>, } #[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, state: Arc, last_context: Arc>>, } impl CallbackBroker { pub fn new(capacity: usize, handler: CallbackHandler) -> Result { ensure!(capacity > 0, "callback broker capacity must be positive"); let (tx, mut rx) = mpsc::channel::(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 { 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 { self.last_context.lock().await.clone() } }