diff --git a/config/grid-agent.example.json b/config/grid-agent.example.json index 91ba2ef..78df48b 100644 --- a/config/grid-agent.example.json +++ b/config/grid-agent.example.json @@ -24,7 +24,19 @@ }, "storage_path": "data/grid-agent", "behavior": { - "heartbeat_seconds": 30 + "heartbeat_seconds": 30, + "settle_milliseconds": 2000, + "response_delay_min_milliseconds": 350, + "response_delay_max_milliseconds": 1200, + "attention_dwell_milliseconds": 4000, + "idle_interval_seconds": 120, + "action_timeout_seconds": 15, + "max_walk_duration_seconds": 10, + "stuck_timeout_seconds": 3, + "min_action_interval_milliseconds": 750, + "max_walk_distance_meters": 8, + "max_attention_distance_meters": 96, + "idle_look_enabled": true }, "reconnect": { "initial_delay_milliseconds": 1000, diff --git a/crates/metacrate-grid-agent/src/backend.rs b/crates/metacrate-grid-agent/src/backend.rs index 39ea7be..d43edb7 100644 --- a/crates/metacrate-grid-agent/src/backend.rs +++ b/crates/metacrate-grid-agent/src/backend.rs @@ -107,6 +107,14 @@ pub struct LibremetaverseWorldSnapshotSource { agent: Arc, } +#[cfg(feature = "live-grid")] +#[derive(Clone)] +pub struct LibremetaverseEmbodimentSink { + client: libremetaverse::GridClient, + agent: Arc, + generation: Arc>>, +} + #[cfg(feature = "live-grid")] impl fmt::Debug for LibremetaverseClientOwner { fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { @@ -130,6 +138,7 @@ pub struct LibremetaverseSessionBackend { password: crate::config::SecretString, interaction: Option, perception: Option, + behavior: Option, delivery_generation: Arc>>, } @@ -187,7 +196,7 @@ impl LibremetaverseClientOwner { &self, connection: crate::config::GridConnection, ) -> Result { - self.session_backend_inner(connection, None, None) + self.session_backend_inner(connection, None, None, None) } /// Creates a supervised live session whose generation owns exactly one @@ -201,7 +210,7 @@ impl LibremetaverseClientOwner { connection: crate::config::GridConnection, ingress: crate::interaction::InteractionIngress, ) -> Result { - self.session_backend_inner(connection, Some(ingress), None) + self.session_backend_inner(connection, Some(ingress), None, None) } /// Creates a session that generation-fences interaction and perception. @@ -215,7 +224,27 @@ impl LibremetaverseClientOwner { interaction: crate::interaction::InteractionIngress, perception: crate::perception::PerceptionIngress, ) -> Result { - self.session_backend_inner(connection, Some(interaction), Some(perception)) + self.session_backend_inner(connection, Some(interaction), Some(perception), None) + } + + /// Creates a session that generation-fences interaction, perception, and embodiment. + /// + /// # Errors + /// + /// Rejects a malformed avatar identity before any subscription or login. + pub fn session_backend_with_agent_services( + &self, + connection: crate::config::GridConnection, + interaction: crate::interaction::InteractionIngress, + perception: crate::perception::PerceptionIngress, + behavior: crate::behavior::BehaviorIngress, + ) -> Result { + self.session_backend_inner( + connection, + Some(interaction), + Some(perception), + Some(behavior), + ) } fn session_backend_inner( @@ -223,6 +252,7 @@ impl LibremetaverseClientOwner { connection: crate::config::GridConnection, interaction: Option, perception: Option, + behavior: Option, ) -> Result { let avatar = connection.avatar_name.trim(); let (first_name, last_name) = avatar @@ -242,6 +272,7 @@ impl LibremetaverseClientOwner { password: connection.password, interaction, perception, + behavior, delivery_generation: Arc::clone(&self.delivery_generation), }) } @@ -261,6 +292,15 @@ impl LibremetaverseClientOwner { agent: Arc::clone(&self.agent), } } + + #[must_use] + pub fn embodiment_sink(&self) -> LibremetaverseEmbodimentSink { + LibremetaverseEmbodimentSink { + client: self.client.clone(), + agent: Arc::clone(&self.agent), + generation: Arc::clone(&self.delivery_generation), + } + } } #[cfg(feature = "live-grid")] @@ -437,6 +477,222 @@ impl crate::perception::WorldSnapshotSource for LibremetaverseWorldSnapshotSourc } } +#[cfg(feature = "live-grid")] +impl LibremetaverseEmbodimentSink { + fn ensure_generation(&self, generation: u64) -> Result<(), crate::behavior::BehaviorError> { + let current = *self + .generation + .read() + .unwrap_or_else(std::sync::PoisonError::into_inner); + if current == Some(generation) && self.client.network().native_connected() { + Ok(()) + } else { + Err(crate::behavior::BehaviorError::NotReady) + } + } + + fn native_pose( + &self, + generation: u64, + ) -> Result { + self.ensure_generation(generation)?; + let simulator = self + .client + .network() + .native_current_sim() + .ok_or(crate::behavior::BehaviorError::NotReady)?; + let position = self.agent.sim_position(); + let rotation = self.agent.sim_rotation(); + let sin_yaw = 2.0 * f64::from(rotation.w * rotation.z + rotation.x * rotation.y); + let cos_yaw = 1.0 - 2.0 * f64::from(rotation.y * rotation.y + rotation.z * rotation.z); + Ok(crate::behavior::EmbodiedPose { + generation, + region_id: simulator.region_id, + position: native_position(position), + heading_degrees: sin_yaw.atan2(cos_yaw).to_degrees().rem_euclid(360.0), + sitting: self.agent.sitting_on() != 0, + }) + } +} + +#[cfg(feature = "live-grid")] +#[allow(clippy::cast_possible_truncation)] +impl crate::behavior::EmbodimentSink for LibremetaverseEmbodimentSink { + fn current_pose( + &self, + generation: u64, + cancellation: CancellationToken, + ) -> crate::behavior::EmbodimentFuture<'_, crate::behavior::EmbodiedPose> { + Box::pin(async move { + if cancellation.is_cancellation_requested() { + return Err(crate::behavior::BehaviorError::Cancelled); + } + self.native_pose(generation) + }) + } + + fn resolve_avatar( + &self, + generation: u64, + avatar_id: UUID, + cancellation: CancellationToken, + ) -> crate::behavior::EmbodimentFuture<'_, crate::perception::WorldPosition> { + Box::pin(async move { + if cancellation.is_cancellation_requested() { + return Err(crate::behavior::BehaviorError::Cancelled); + } + self.ensure_generation(generation)?; + let simulator = self + .client + .network() + .native_current_sim() + .ok_or(crate::behavior::BehaviorError::NotReady)?; + simulator + .objects_avatars + .read() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .values() + .find(|avatar| avatar.id == avatar_id) + .map(|avatar| native_position(avatar.position)) + .ok_or(crate::behavior::BehaviorError::TargetUnavailable) + }) + } + + fn face_point( + &self, + generation: u64, + point: crate::perception::WorldPosition, + cancellation: CancellationToken, + ) -> crate::behavior::EmbodimentFuture<'_, ()> { + Box::pin(async move { + if cancellation.is_cancellation_requested() { + return Err(crate::behavior::BehaviorError::Cancelled); + } + self.ensure_generation(generation)?; + self.agent + .movement + .turn_toward( + libremetaverse_types::Vector3 { + x: point.x as f32, + y: point.y as f32, + z: point.z as f32, + }, + Some(true), + ) + .map_err(|_| crate::behavior::BehaviorError::NativeOperation)?; + Ok(()) + }) + } + + fn begin_walk( + &self, + generation: u64, + point: crate::perception::WorldPosition, + cancellation: CancellationToken, + ) -> crate::behavior::EmbodimentFuture<'_, ()> { + Box::pin(async move { + if cancellation.is_cancellation_requested() { + return Err(crate::behavior::BehaviorError::Cancelled); + } + self.ensure_generation(generation)?; + self.agent + .auto_pilot_local( + point.x.round() as i32, + point.y.round() as i32, + point.z as f32, + ) + .map_err(|_| crate::behavior::BehaviorError::NativeOperation) + }) + } + + fn validate_walk_target( + &self, + generation: u64, + point: crate::perception::WorldPosition, + cancellation: CancellationToken, + ) -> crate::behavior::EmbodimentFuture<'_, ()> { + Box::pin(async move { + if cancellation.is_cancellation_requested() { + return Err(crate::behavior::BehaviorError::Cancelled); + } + self.ensure_generation(generation)?; + let simulator = self + .client + .network() + .native_current_sim() + .ok_or(crate::behavior::BehaviorError::NotReady)?; + let parcels = self.client.parcels(); + let current = parcels + .get_parcel_local_id(simulator.clone(), self.agent.sim_position()) + .map_err(|_| crate::behavior::BehaviorError::NativeOperation)?; + let target = parcels + .get_parcel_local_id( + simulator, + libremetaverse_types::Vector3 { + x: point.x as f32, + y: point.y as f32, + z: point.z as f32, + }, + ) + .map_err(|_| crate::behavior::BehaviorError::NativeOperation)?; + if current == 0 || current != target { + return Err(crate::behavior::BehaviorError::RegionBoundary); + } + Ok(()) + }) + } + + fn stop( + &self, + generation: u64, + cancellation: CancellationToken, + ) -> crate::behavior::EmbodimentFuture<'_, ()> { + Box::pin(async move { + if cancellation.is_cancellation_requested() { + return Err(crate::behavior::BehaviorError::Cancelled); + } + self.ensure_generation(generation)?; + self.agent + .auto_pilot_cancel() + .map(|_| ()) + .map_err(|_| crate::behavior::BehaviorError::NativeOperation) + }) + } + + fn sit( + &self, + generation: u64, + cancellation: CancellationToken, + ) -> crate::behavior::EmbodimentFuture<'_, ()> { + Box::pin(async move { + if cancellation.is_cancellation_requested() { + return Err(crate::behavior::BehaviorError::Cancelled); + } + self.ensure_generation(generation)?; + self.agent + .sit() + .map_err(|_| crate::behavior::BehaviorError::NativeOperation) + }) + } + + fn stand( + &self, + generation: u64, + cancellation: CancellationToken, + ) -> crate::behavior::EmbodimentFuture<'_, ()> { + Box::pin(async move { + if cancellation.is_cancellation_requested() { + return Err(crate::behavior::BehaviorError::Cancelled); + } + self.ensure_generation(generation)?; + self.agent + .stand() + .map(|_| ()) + .map_err(|_| crate::behavior::BehaviorError::NativeOperation) + }) + } +} + #[cfg(feature = "live-grid")] fn native_position(value: libremetaverse_types::Vector3) -> crate::perception::WorldPosition { crate::perception::WorldPosition { @@ -580,6 +836,7 @@ struct LibremetaverseSession { subscriptions: Vec, interaction: Option, perception: Option, + behavior: Option, delivery_generation: Arc>>, } @@ -616,16 +873,20 @@ impl crate::session::GridSession for LibremetaverseSession { cancellation: CancellationToken, ) -> crate::session::SessionFuture<'static, Result<(), crate::session::SessionFailure>> { Box::pin(async move { - *self - .delivery_generation - .write() - .unwrap_or_else(std::sync::PoisonError::into_inner) = None; if let Some(interaction) = &self.interaction { let _ = interaction.try_disconnected(); } if let Some(perception) = &self.perception { perception.disconnected(); } + if let Some(behavior) = &self.behavior { + behavior.disconnected(); + tokio::task::yield_now().await; + } + *self + .delivery_generation + .write() + .unwrap_or_else(std::sync::PoisonError::into_inner) = None; self.subscriptions.clear(); self.network .native_logout_async(Some(cancellation)) @@ -694,6 +955,7 @@ impl crate::session::GridSessionBackend for LibremetaverseSessionBackend { let disconnected_sender = sender.clone(); let disconnected_interaction = self.interaction.clone(); let disconnected_perception = self.perception.clone(); + let disconnected_behavior = self.behavior.clone(); let disconnected_generation = Arc::clone(&self.delivery_generation); let disconnected = self .network @@ -707,6 +969,9 @@ impl crate::session::GridSessionBackend for LibremetaverseSessionBackend { if let Some(perception) = &disconnected_perception { perception.disconnected(); } + if let Some(behavior) = &disconnected_behavior { + behavior.disconnected(); + } let kind = match event.reason() { libremetaverse::NetworkManagerDisconnectType::NetworkTimeout | libremetaverse::NetworkManagerDisconnectType::ClientInitiated => { @@ -727,6 +992,7 @@ impl crate::session::GridSessionBackend for LibremetaverseSessionBackend { let ready_sender = sender.clone(); let ready_interaction = self.interaction.clone(); let ready_perception = self.perception.clone(); + let ready_behavior = self.behavior.clone(); let ready_network = self.network.clone(); let ready_agent = Arc::clone(&self.agent); let ready_generation = Arc::clone(&self.delivery_generation); @@ -747,18 +1013,27 @@ impl crate::session::GridSessionBackend for LibremetaverseSessionBackend { { let _ = perception.connected(generation, simulator.region_id); } + if let (Some(behavior), Some(simulator)) = + (&ready_behavior, ready_network.native_current_sim()) + { + let _ = behavior.connected(generation, simulator.region_id); + } let _ = ready_sender.try_send(crate::session::SessionSignal::Ready); } })); let mut subscriptions = vec![disconnected, ready]; if let Some(perception) = &self.perception { let crossing_perception = perception.clone(); + let crossing_behavior = self.behavior.clone(); let crossing_network = self.network.clone(); subscriptions.push(self.network.native_subscribe_sim_changed(Arc::new( move |_| { if let Some(simulator) = crossing_network.native_current_sim() { let _ = crossing_perception.region_changed(generation, simulator.region_id); + if let Some(behavior) = &crossing_behavior { + let _ = behavior.region_changed(generation, simulator.region_id); + } } }, ))); @@ -806,6 +1081,11 @@ impl crate::session::GridSessionBackend for LibremetaverseSessionBackend { { let _ = perception.connected(generation, simulator.region_id); } + if let (Some(behavior), Some(simulator)) = + (&self.behavior, self.network.native_current_sim()) + { + let _ = behavior.connected(generation, simulator.region_id); + } let _ = sender.try_send(crate::session::SessionSignal::Ready); } let session: Box = Box::new(LibremetaverseSession { @@ -815,6 +1095,7 @@ impl crate::session::GridSessionBackend for LibremetaverseSessionBackend { subscriptions, interaction: self.interaction.clone(), perception: self.perception.clone(), + behavior: self.behavior.clone(), delivery_generation: Arc::clone(&self.delivery_generation), }); Ok(session) diff --git a/crates/metacrate-grid-agent/src/behavior.rs b/crates/metacrate-grid-agent/src/behavior.rs new file mode 100644 index 0000000..7c15697 --- /dev/null +++ b/crates/metacrate-grid-agent/src/behavior.rs @@ -0,0 +1,1423 @@ +//! Generation-fenced embodied attention and bounded presence behavior. + +#![allow(clippy::missing_errors_doc)] + +use crate::backend::{AuthorizedToolBackend, BackendError, BackendFuture}; +use crate::config::BehaviorSettings; +use crate::interaction::{PacerFuture, ResponsePacer}; +use crate::llm::{ToolDefinition, ToolSchema}; +use crate::perception::WorldPosition; +use crate::policy::{ + AllowedOrigins, ApprovalRule, AuthorizedAction, Capability, FixedCost, Idempotency, + OriginClass, PolicyError, PolicyReasonCode, PolicyTool, ResourceCost, ResourceEstimator, Risk, +}; +use crate::types::{BoundedText, MAX_BODY_BYTES, MAX_OBSERVABLE_DETAIL_BYTES, ToolCallOutcome}; +use libremetaverse_types::UUID; +use libremetaverse_types::compat::CancellationToken; +use serde::Serialize; +use serde_json::{Map, Value, json}; +use std::collections::{BTreeMap, BTreeSet}; +use std::error::Error; +use std::fmt; +use std::sync::Arc; +use std::sync::atomic::{AtomicU64, Ordering}; +use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH}; +use tokio::sync::{mpsc, oneshot, watch}; +use tokio::task::JoinHandle; + +pub const FACE_AVATAR_TOOL: &str = "behavior_face_avatar"; +pub const FACE_POINT_TOOL: &str = "behavior_face_point"; +pub const LOOK_AROUND_TOOL: &str = "behavior_look_around"; +pub const WALK_SHORT_TOOL: &str = "behavior_walk_short"; +pub const STOP_TOOL: &str = "behavior_stop"; +pub const SIT_TOOL: &str = "behavior_sit"; +pub const STAND_TOOL: &str = "behavior_stand"; +pub const CURRENT_POSE_TOOL: &str = "behavior_current_pose"; + +const WALK_POLL: Duration = Duration::from_millis(250); +const TURN_INTERVAL: Duration = Duration::from_millis(40); +const MAX_TURN_STEP_DEGREES: f64 = 30.0; +const ARRIVAL_METERS: f64 = 0.5; +const STUCK_PROGRESS_METERS: f64 = 0.05; +const MAX_TOOL_RESULT_BYTES: usize = 8 * 1024; + +pub type EmbodimentFuture<'a, T> = BackendFuture<'a, Result>; + +#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize)] +#[serde(rename_all = "snake_case")] +pub enum BehaviorMode { + Offline, + Settling, + Available, + Engaged, + Executing, + Roaming, + Paused, + Recovering, +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize)] +#[serde(rename_all = "snake_case")] +pub enum BehaviorOutcome { + Completed, + Cancelled, + TimedOut, + Stuck, + Rejected, + Preempted, +} + +#[derive(Clone, Debug, PartialEq)] +pub struct EmbodiedPose { + pub generation: u64, + pub region_id: UUID, + pub position: WorldPosition, + pub heading_degrees: f64, + pub sitting: bool, +} + +impl EmbodiedPose { + fn valid(&self, generation: u64, region_id: UUID) -> bool { + self.generation == generation + && self.region_id == region_id + && self.position.x.is_finite() + && self.position.y.is_finite() + && self.position.z.is_finite() + && self.heading_degrees.is_finite() + } +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub enum BehaviorError { + NotReady, + Paused, + EmergencyStopped, + InvalidArguments, + TargetUnavailable, + RegionBoundary, + RateLimited, + Cancelled, + TimedOut, + Stuck, + NativeOperation, + QueueClosed, +} + +impl fmt::Display for BehaviorError { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str(match self { + Self::NotReady => "embodied behavior is unavailable before session readiness", + Self::Paused => "embodied behavior is paused by the operator", + Self::EmergencyStopped => "embodied behavior emergency stop is active", + Self::InvalidArguments => "behavior tool arguments are invalid or outside bounds", + Self::TargetUnavailable => "the requested attention target is unavailable", + Self::RegionBoundary => "the requested action would leave the current region", + Self::RateLimited => "embodied action rate limit is active", + Self::Cancelled => "embodied action was cancelled", + Self::TimedOut => "embodied action timed out", + Self::Stuck => "bounded walk stopped after making no progress", + Self::NativeOperation => "native movement operation failed", + Self::QueueClosed => "embodied behavior controller is unavailable", + }) + } +} + +impl Error for BehaviorError {} + +/// High-level native boundary. No control flags, raw packets, teleport, flight, +/// animation, touch, following, or unrestricted autopilot are exposed here. +pub trait EmbodimentSink: Send + Sync + 'static { + fn current_pose( + &self, + generation: u64, + cancellation: CancellationToken, + ) -> EmbodimentFuture<'_, EmbodiedPose>; + fn resolve_avatar( + &self, + generation: u64, + avatar_id: UUID, + cancellation: CancellationToken, + ) -> EmbodimentFuture<'_, WorldPosition>; + fn face_point( + &self, + generation: u64, + point: WorldPosition, + cancellation: CancellationToken, + ) -> EmbodimentFuture<'_, ()>; + fn begin_walk( + &self, + generation: u64, + point: WorldPosition, + cancellation: CancellationToken, + ) -> EmbodimentFuture<'_, ()>; + fn validate_walk_target( + &self, + generation: u64, + point: WorldPosition, + cancellation: CancellationToken, + ) -> EmbodimentFuture<'_, ()>; + fn stop(&self, generation: u64, cancellation: CancellationToken) -> EmbodimentFuture<'_, ()>; + fn sit(&self, generation: u64, cancellation: CancellationToken) -> EmbodimentFuture<'_, ()>; + fn stand(&self, generation: u64, cancellation: CancellationToken) -> EmbodimentFuture<'_, ()>; +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub enum BehaviorTrigger { + SessionReady, + SessionDisconnected, + RegionChanged, + PublicResponse { + delivery_id: String, + avatar_id: UUID, + }, + AuthorizedTool { + authorization_id: u64, + tool: String, + }, + IdleTimer, + OperatorPause, + OperatorResume, + EmergencyStop, +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub enum BehaviorPolicyResult { + BuiltInAttention, + Authorized(u64), + InternalIdle, + OperatorOverride, +} + +#[derive(Clone, Debug, PartialEq)] +pub enum BehaviorObservation { + Transition { + from: BehaviorMode, + to: BehaviorMode, + trigger: BehaviorTrigger, + }, + Action { + generation: Option, + region_id: Option, + action: String, + trigger: BehaviorTrigger, + policy: BehaviorPolicyResult, + duration_millis: u64, + outcome: BehaviorOutcome, + }, +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +struct ReadyState { + generation: u64, + region_id: UUID, +} + +#[derive(Clone, Debug)] +enum BehaviorAction { + FaceAvatar(UUID), + FacePoint(WorldPosition), + LookAround, + WalkShort { + heading_degrees: f64, + distance_meters: f64, + }, + Stop, + Sit, + Stand, + CurrentPose, +} + +impl BehaviorAction { + fn name(&self) -> &'static str { + match self { + Self::FaceAvatar(_) => FACE_AVATAR_TOOL, + Self::FacePoint(_) => FACE_POINT_TOOL, + Self::LookAround => LOOK_AROUND_TOOL, + Self::WalkShort { .. } => WALK_SHORT_TOOL, + Self::Stop => STOP_TOOL, + Self::Sit => SIT_TOOL, + Self::Stand => STAND_TOOL, + Self::CurrentPose => CURRENT_POSE_TOOL, + } + } +} + +struct ActionRequest { + action: BehaviorAction, + trigger: BehaviorTrigger, + policy: BehaviorPolicyResult, + reply: oneshot::Sender>, +} + +enum Command { + Action(ActionRequest), + Attention { + delivery_id: String, + avatar_id: UUID, + reply: oneshot::Sender>, + }, + Roaming(bool), + Shutdown(oneshot::Sender<()>), +} + +#[derive(Clone)] +pub struct BehaviorIngress { + commands: mpsc::Sender, + ready: watch::Sender>, + paused: watch::Sender, + emergency: watch::Sender, +} + +impl fmt::Debug for BehaviorIngress { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("BehaviorIngress") + .finish_non_exhaustive() + } +} + +impl BehaviorIngress { + pub fn connected(&self, generation: u64, region_id: UUID) -> Result<(), BehaviorError> { + if generation == 0 || region_id == UUID::zero() { + return Err(BehaviorError::InvalidArguments); + } + self.ready + .send(Some(ReadyState { + generation, + region_id, + })) + .map_err(|_| BehaviorError::QueueClosed) + } + + pub fn region_changed(&self, generation: u64, region_id: UUID) -> Result<(), BehaviorError> { + self.connected(generation, region_id) + } + + pub fn disconnected(&self) { + let _ = self.ready.send(None); + } + + pub fn pause(&self) { + let _ = self.paused.send(true); + } + + pub fn resume(&self) { + let _ = self.paused.send(false); + } + + pub fn emergency_stop(&self) { + let _ = self.emergency.send(true); + } + + pub fn clear_emergency_stop(&self) { + let _ = self.emergency.send(false); + } + + /// Marks a separately authorized scheduler task as roaming. The controller + /// itself never invents a roaming route or exposes follow/wander behavior. + pub fn set_roaming(&self, roaming: bool) -> Result<(), BehaviorError> { + self.commands + .try_send(Command::Roaming(roaming)) + .map_err(|_| BehaviorError::QueueClosed) + } + + async fn request( + &self, + action: BehaviorAction, + trigger: BehaviorTrigger, + policy: BehaviorPolicyResult, + ) -> Result { + let (reply, receive) = oneshot::channel(); + self.commands + .send(Command::Action(ActionRequest { + action, + trigger, + policy, + reply, + })) + .await + .map_err(|_| BehaviorError::QueueClosed)?; + receive.await.map_err(|_| BehaviorError::QueueClosed)? + } + + pub async fn attention( + &self, + delivery_id: String, + avatar_id: UUID, + ) -> Result<(), BehaviorError> { + let (reply, receive) = oneshot::channel(); + self.commands + .send(Command::Attention { + delivery_id, + avatar_id, + reply, + }) + .await + .map_err(|_| BehaviorError::QueueClosed)?; + receive.await.map_err(|_| BehaviorError::QueueClosed)? + } +} + +impl ResponsePacer for BehaviorIngress { + fn prepare_public_response( + &self, + delivery_id: String, + avatar_id: UUID, + cancellation: CancellationToken, + ) -> PacerFuture<'_> { + Box::pin(async move { + tokio::select! { + () = cancellation.cancelled() => Err(crate::interaction::InteractionPacingError::Cancelled), + result = self.attention(delivery_id, avatar_id) => result.map_err(|_| crate::interaction::InteractionPacingError::Unavailable), + } + }) + } +} + +pub struct BehaviorHandle { + ingress: BehaviorIngress, + observations: mpsc::Receiver, + task: Option>, + shutdown_timeout: Duration, +} + +impl BehaviorHandle { + #[must_use] + pub fn ingress(&self) -> BehaviorIngress { + self.ingress.clone() + } + + pub async fn next_observation(&mut self) -> Option { + self.observations.recv().await + } + + pub async fn shutdown(mut self) -> Result<(), BehaviorError> { + let (reply, mut receive) = oneshot::channel(); + self.ingress + .commands + .send(Command::Shutdown(reply)) + .await + .map_err(|_| BehaviorError::QueueClosed)?; + tokio::time::timeout(self.shutdown_timeout, async { + loop { + tokio::select! { + result = &mut receive => return result.map_err(|_| BehaviorError::QueueClosed), + observation = self.observations.recv() => { + if observation.is_none() { + return receive.await.map_err(|_| BehaviorError::QueueClosed); + } + } + } + } + }) + .await + .map_err(|_| BehaviorError::TimedOut)??; + if let Some(task) = self.task.take() { + tokio::time::timeout(self.shutdown_timeout, task) + .await + .map_err(|_| BehaviorError::TimedOut)? + .map_err(|_| BehaviorError::QueueClosed)?; + } + Ok(()) + } +} + +pub struct BehaviorController { + settings: BehaviorSettings, + sink: Arc, + queue_capacity: usize, + observation_capacity: usize, + shutdown_timeout: Duration, + random: Arc, +} + +impl BehaviorController { + pub fn new( + settings: BehaviorSettings, + sink: Arc, + queue_capacity: usize, + observation_capacity: usize, + shutdown_timeout: Duration, + ) -> Result { + Self::with_random( + settings, + sink, + queue_capacity, + observation_capacity, + shutdown_timeout, + Arc::new(SystemBehaviorRandom::new()), + ) + } + + pub fn with_random( + settings: BehaviorSettings, + sink: Arc, + queue_capacity: usize, + observation_capacity: usize, + shutdown_timeout: Duration, + random: Arc, + ) -> Result { + if !settings.is_valid() { + return Err(BehaviorError::InvalidArguments); + } + if queue_capacity == 0 + || queue_capacity > 8_192 + || observation_capacity == 0 + || observation_capacity > 8_192 + || shutdown_timeout.is_zero() + || shutdown_timeout > Duration::from_mins(1) + { + return Err(BehaviorError::InvalidArguments); + } + Ok(Self { + settings, + sink, + queue_capacity, + observation_capacity, + shutdown_timeout, + random, + }) + } + + #[must_use] + pub fn start(self) -> BehaviorHandle { + let (commands, receiver) = mpsc::channel(self.queue_capacity); + let (ready, ready_rx) = watch::channel(None); + let (paused, paused_rx) = watch::channel(false); + let (emergency, emergency_rx) = watch::channel(false); + let (observations, observation_rx) = mpsc::channel(self.observation_capacity); + let ingress = BehaviorIngress { + commands, + ready, + paused, + emergency, + }; + let shutdown_timeout = self.shutdown_timeout; + let task = tokio::spawn(run_actor( + self, + receiver, + ready_rx, + paused_rx, + emergency_rx, + observations, + )); + BehaviorHandle { + ingress, + observations: observation_rx, + task: Some(task), + shutdown_timeout, + } + } +} + +pub trait BehaviorRandom: Send + Sync + 'static { + fn next_u64(&self) -> u64; +} + +struct SystemBehaviorRandom(AtomicU64); + +impl SystemBehaviorRandom { + fn new() -> Self { + let seed = SystemTime::now() + .duration_since(UNIX_EPOCH) + .map_or(1, |value| { + u64::try_from(value.as_nanos()).unwrap_or(u64::MAX) + }) + | 1; + Self(AtomicU64::new(seed)) + } +} + +impl BehaviorRandom for SystemBehaviorRandom { + fn next_u64(&self) -> u64 { + let mut current = self.0.load(Ordering::Relaxed); + loop { + let mut next = current; + next ^= next << 13; + next ^= next >> 7; + next ^= next << 17; + match self + .0 + .compare_exchange_weak(current, next, Ordering::Relaxed, Ordering::Relaxed) + { + Ok(_) => return next, + Err(actual) => current = actual, + } + } + } +} + +#[allow(clippy::too_many_lines)] // Lifecycle, preemption, idle, and commands share one ordered owner. +async fn run_actor( + controller: BehaviorController, + mut commands: mpsc::Receiver, + mut ready: watch::Receiver>, + mut paused: watch::Receiver, + mut emergency: watch::Receiver, + observations: mpsc::Sender, +) { + let mut mode = BehaviorMode::Offline; + let mut active_ready: Option = None; + let mut last_action: Option = None; + let idle = tokio::time::sleep(controller.settings.idle_interval); + tokio::pin!(idle); + loop { + tokio::select! { + changed = ready.changed() => { + if changed.is_err() { break; } + let next_ready = *ready.borrow(); + let trigger = match (active_ready, next_ready) { + (Some(previous), Some(next)) if previous.region_id != next.region_id => BehaviorTrigger::RegionChanged, + (_, Some(_)) => BehaviorTrigger::SessionReady, + _ => BehaviorTrigger::SessionDisconnected, + }; + if let Some(previous) = active_ready.filter(|previous| Some(*previous) != next_ready) { + let _ = controller.sink.stop(previous.generation, CancellationToken::default()).await; + } + active_ready = next_ready; + let next = if *paused.borrow() { BehaviorMode::Paused } else if next_ready.is_some() { BehaviorMode::Settling } else { BehaviorMode::Offline }; + transition(&observations, &mut mode, next, trigger).await; + if let Some(state) = next_ready { + tokio::time::sleep(controller.settings.settle_delay).await; + if ready.borrow().as_ref() == Some(&state) && !*paused.borrow() && !*emergency.borrow() { + transition(&observations, &mut mode, BehaviorMode::Available, BehaviorTrigger::SessionReady).await; + } + } + idle.as_mut().reset(tokio::time::Instant::now() + controller.settings.idle_interval); + } + changed = paused.changed() => { + if changed.is_err() { break; } + if *paused.borrow() { + let current_ready = *ready.borrow(); + if let Some(state) = current_ready { let _ = controller.sink.stop(state.generation, CancellationToken::default()).await; } + transition(&observations, &mut mode, BehaviorMode::Paused, BehaviorTrigger::OperatorPause).await; + } else { + let next = if ready.borrow().is_some() { BehaviorMode::Available } else { BehaviorMode::Offline }; + transition(&observations, &mut mode, next, BehaviorTrigger::OperatorResume).await; + } + } + changed = emergency.changed() => { + if changed.is_err() { break; } + if *emergency.borrow() { + let current_ready = *ready.borrow(); + if let Some(state) = current_ready { let _ = controller.sink.stop(state.generation, CancellationToken::default()).await; } + transition(&observations, &mut mode, BehaviorMode::Paused, BehaviorTrigger::EmergencyStop).await; + } else { + let next = if *paused.borrow() { + BehaviorMode::Paused + } else if ready.borrow().is_some() { + BehaviorMode::Available + } else { + BehaviorMode::Offline + }; + transition(&observations, &mut mode, next, BehaviorTrigger::OperatorResume).await; + } + } + () = &mut idle, if controller.settings.idle_look_enabled => { + if mode == BehaviorMode::Available { + let (reply, _) = oneshot::channel(); + let request = ActionRequest { action: BehaviorAction::LookAround, trigger: BehaviorTrigger::IdleTimer, policy: BehaviorPolicyResult::InternalIdle, reply }; + execute_request(&controller, request, &mut mode, &mut last_action, &ready, &paused, &emergency, &observations).await; + } + idle.as_mut().reset(tokio::time::Instant::now() + controller.settings.idle_interval); + } + command = commands.recv() => { + let Some(command) = command else { break; }; + match command { + Command::Shutdown(reply) => { + let current_ready = *ready.borrow(); + if let Some(state) = current_ready { let _ = controller.sink.stop(state.generation, CancellationToken::default()).await; } + let _ = reply.send(()); + break; + } + Command::Roaming(roaming) => { + let next = if *paused.borrow() { + BehaviorMode::Paused + } else if ready.borrow().is_none() { + BehaviorMode::Offline + } else if roaming { + BehaviorMode::Roaming + } else { + BehaviorMode::Available + }; + transition(&observations, &mut mode, next, BehaviorTrigger::IdleTimer).await; + } + Command::Action(request) => execute_request(&controller, request, &mut mode, &mut last_action, &ready, &paused, &emergency, &observations).await, + Command::Attention { delivery_id, avatar_id, reply } => { + let request = ActionRequest { + action: BehaviorAction::FaceAvatar(avatar_id), + trigger: BehaviorTrigger::PublicResponse { delivery_id, avatar_id }, + policy: BehaviorPolicyResult::BuiltInAttention, + reply: response_reply(reply), + }; + let delay = random_duration(&controller, controller.settings.response_delay_min, controller.settings.response_delay_max); + tokio::time::sleep(delay).await; + execute_request(&controller, request, &mut mode, &mut last_action, &ready, &paused, &emergency, &observations).await; + if mode == BehaviorMode::Engaged { + tokio::time::sleep(controller.settings.attention_dwell).await; + if !*paused.borrow() && ready.borrow().is_some() { transition(&observations, &mut mode, BehaviorMode::Available, BehaviorTrigger::IdleTimer).await; } + } + } + } + } + } + } +} + +fn response_reply( + reply: oneshot::Sender>, +) -> oneshot::Sender> { + let (tx, rx) = oneshot::channel(); + tokio::spawn(async move { + let result = rx + .await + .unwrap_or(Err(BehaviorError::QueueClosed)) + .map(|_| ()); + let _ = reply.send(result); + }); + tx +} + +async fn transition( + observations: &mpsc::Sender, + mode: &mut BehaviorMode, + to: BehaviorMode, + trigger: BehaviorTrigger, +) { + if *mode != to { + let from = *mode; + *mode = to; + let _ = observations + .send(BehaviorObservation::Transition { from, to, trigger }) + .await; + } +} + +#[allow(clippy::too_many_arguments, clippy::too_many_lines)] +async fn execute_request( + controller: &BehaviorController, + request: ActionRequest, + mode: &mut BehaviorMode, + last_action: &mut Option, + ready: &watch::Receiver>, + paused: &watch::Receiver, + emergency: &watch::Receiver, + observations: &mpsc::Sender, +) { + let started = Instant::now(); + let action_name = request.action.name().to_owned(); + let state = *ready.borrow(); + let mut result = if *emergency.borrow() { + Err(BehaviorError::EmergencyStopped) + } else if *paused.borrow() { + Err(BehaviorError::Paused) + } else if state.is_none() { + Err(BehaviorError::NotReady) + } else if last_action + .is_some_and(|last| last.elapsed() < controller.settings.min_action_interval) + && !matches!( + request.action, + BehaviorAction::Stop | BehaviorAction::CurrentPose + ) + { + Err(BehaviorError::RateLimited) + } else { + Ok(Value::Null) + }; + if result.is_ok() { + let state = state.expect("checked ready state"); + let target_mode = if matches!(request.policy, BehaviorPolicyResult::BuiltInAttention) { + BehaviorMode::Engaged + } else if matches!(request.policy, BehaviorPolicyResult::InternalIdle) { + BehaviorMode::Available + } else { + BehaviorMode::Executing + }; + transition(observations, mode, target_mode, request.trigger.clone()).await; + let mut action_ready = ready.clone(); + let mut action_paused = paused.clone(); + let mut action_emergency = emergency.clone(); + result = tokio::select! { + value = tokio::time::timeout(controller.settings.action_timeout, perform_action(controller, state, &request.action)) => value.unwrap_or(Err(BehaviorError::TimedOut)), + _ = action_ready.changed() => Err(BehaviorError::Cancelled), + _ = action_paused.changed() => Err(BehaviorError::Paused), + _ = action_emergency.changed() => Err(BehaviorError::EmergencyStopped), + }; + if result.is_err() && matches!(request.action, BehaviorAction::WalkShort { .. }) { + let _ = controller + .sink + .stop(state.generation, CancellationToken::default()) + .await; + } + *last_action = Some(Instant::now()); + if matches!( + result, + Err(BehaviorError::TimedOut | BehaviorError::Stuck | BehaviorError::NativeOperation) + ) { + transition( + observations, + mode, + BehaviorMode::Recovering, + request.trigger.clone(), + ) + .await; + } + if matches!(*mode, BehaviorMode::Executing | BehaviorMode::Recovering) { + transition( + observations, + mode, + BehaviorMode::Available, + request.trigger.clone(), + ) + .await; + } + if matches!( + result, + Err(BehaviorError::Paused | BehaviorError::EmergencyStopped) + ) { + transition( + observations, + mode, + BehaviorMode::Paused, + request.trigger.clone(), + ) + .await; + } else if matches!(result, Err(BehaviorError::Cancelled)) { + let lifecycle_mode = if ready.borrow().is_some() { + BehaviorMode::Settling + } else { + BehaviorMode::Offline + }; + transition(observations, mode, lifecycle_mode, request.trigger.clone()).await; + } else if result.is_err() + && matches!(request.policy, BehaviorPolicyResult::BuiltInAttention) + { + transition( + observations, + mode, + BehaviorMode::Available, + request.trigger.clone(), + ) + .await; + } + } + let outcome = match &result { + Ok(_) => BehaviorOutcome::Completed, + Err(BehaviorError::TimedOut) => BehaviorOutcome::TimedOut, + Err(BehaviorError::Stuck) => BehaviorOutcome::Stuck, + Err(BehaviorError::Cancelled) => BehaviorOutcome::Cancelled, + Err(BehaviorError::Paused | BehaviorError::EmergencyStopped) => BehaviorOutcome::Preempted, + Err(_) => BehaviorOutcome::Rejected, + }; + let _ = observations + .send(BehaviorObservation::Action { + generation: state.map(|value| value.generation), + region_id: state.map(|value| value.region_id), + action: action_name, + trigger: request.trigger, + policy: request.policy, + duration_millis: u64::try_from(started.elapsed().as_millis()).unwrap_or(u64::MAX), + outcome, + }) + .await; + let _ = request.reply.send(result); +} + +async fn perform_action( + controller: &BehaviorController, + state: ReadyState, + action: &BehaviorAction, +) -> Result { + let cancellation = CancellationToken::default(); + match action { + BehaviorAction::CurrentPose => pose_value( + &controller + .sink + .current_pose(state.generation, cancellation) + .await?, + state, + ), + BehaviorAction::FaceAvatar(avatar_id) => { + let point = controller + .sink + .resolve_avatar(state.generation, *avatar_id, cancellation.clone()) + .await?; + let pose = controller + .sink + .current_pose(state.generation, cancellation.clone()) + .await?; + validate_pose(&pose, state)?; + validate_point( + pose.position, + point, + controller.settings.max_attention_distance_meters, + )?; + smooth_face(controller, state, &pose, point, cancellation).await?; + Ok(json!({"status":"completed","target_avatar_id":avatar_id.to_string()})) + } + BehaviorAction::FacePoint(point) => { + let pose = controller + .sink + .current_pose(state.generation, cancellation.clone()) + .await?; + validate_pose(&pose, state)?; + validate_point( + pose.position, + *point, + controller.settings.max_attention_distance_meters, + )?; + smooth_face(controller, state, &pose, *point, cancellation).await?; + Ok(json!({"status":"completed"})) + } + BehaviorAction::LookAround => { + let pose = controller + .sink + .current_pose(state.generation, cancellation.clone()) + .await?; + validate_pose(&pose, state)?; + let offset = f64::from( + u32::try_from(controller.random.next_u64() % 121).expect("bounded random offset"), + ) - 60.0; + let heading = (pose.heading_degrees + offset).to_radians(); + let point = WorldPosition { + x: pose.position.x + heading.cos() * 4.0, + y: pose.position.y + heading.sin() * 4.0, + z: pose.position.z, + }; + smooth_face(controller, state, &pose, point, cancellation).await?; + Ok(json!({"status":"completed"})) + } + BehaviorAction::WalkShort { + heading_degrees, + distance_meters, + } => { + walk( + controller, + state, + *heading_degrees, + *distance_meters, + cancellation, + ) + .await + } + BehaviorAction::Stop => { + controller.sink.stop(state.generation, cancellation).await?; + Ok(json!({"status":"completed"})) + } + BehaviorAction::Sit => { + controller.sink.sit(state.generation, cancellation).await?; + Ok(json!({"status":"completed"})) + } + BehaviorAction::Stand => { + controller + .sink + .stand(state.generation, cancellation) + .await?; + Ok(json!({"status":"completed"})) + } + } +} + +async fn walk( + controller: &BehaviorController, + state: ReadyState, + heading_degrees: f64, + distance_meters: f64, + cancellation: CancellationToken, +) -> Result { + if !heading_degrees.is_finite() + || !distance_meters.is_finite() + || !(0.25..=f64::from(controller.settings.max_walk_distance_meters)) + .contains(&distance_meters) + { + return Err(BehaviorError::InvalidArguments); + } + let start = controller + .sink + .current_pose(state.generation, cancellation.clone()) + .await?; + validate_pose(&start, state)?; + let heading = heading_degrees.to_radians(); + let target = WorldPosition { + x: start.position.x + heading.cos() * distance_meters, + y: start.position.y + heading.sin() * distance_meters, + z: start.position.z, + }; + if !(0.5..=255.5).contains(&target.x) || !(0.5..=255.5).contains(&target.y) { + return Err(BehaviorError::RegionBoundary); + } + controller + .sink + .validate_walk_target(state.generation, target, cancellation.clone()) + .await?; + smooth_face(controller, state, &start, target, cancellation.clone()).await?; + controller + .sink + .begin_walk(state.generation, target, cancellation.clone()) + .await?; + let deadline = tokio::time::Instant::now() + controller.settings.max_walk_duration; + let mut last_progress = tokio::time::Instant::now(); + let mut prior = start.position; + let outcome = loop { + tokio::time::sleep(WALK_POLL).await; + if tokio::time::Instant::now() >= deadline { + break Err(BehaviorError::TimedOut); + } + let pose = match controller + .sink + .current_pose(state.generation, cancellation.clone()) + .await + { + Ok(pose) => pose, + Err(error) => break Err(error), + }; + if validate_pose(&pose, state).is_err() { + break Err(BehaviorError::RegionBoundary); + } + if distance(pose.position, target) <= ARRIVAL_METERS { + break Ok(pose); + } + if distance(pose.position, prior) >= STUCK_PROGRESS_METERS { + prior = pose.position; + last_progress = tokio::time::Instant::now(); + } else if last_progress.elapsed() >= controller.settings.stuck_timeout { + break Err(BehaviorError::Stuck); + } + }; + let stop_result = controller + .sink + .stop(state.generation, CancellationToken::default()) + .await; + match (outcome, stop_result) { + (Ok(pose), Ok(())) => pose_value(&pose, state), + (Err(error), _) | (_, Err(error)) => Err(error), + } +} + +async fn smooth_face( + controller: &BehaviorController, + state: ReadyState, + pose: &EmbodiedPose, + target: WorldPosition, + cancellation: CancellationToken, +) -> Result<(), BehaviorError> { + let desired = (target.y - pose.position.y) + .atan2(target.x - pose.position.x) + .to_degrees() + .rem_euclid(360.0); + let delta = (desired - pose.heading_degrees + 540.0).rem_euclid(360.0) - 180.0; + let steps = (1..=6) + .find(|step| delta.abs() <= MAX_TURN_STEP_DEGREES * f64::from(*step)) + .unwrap_or(6); + let horizontal_distance = ((target.x - pose.position.x).powi(2) + + (target.y - pose.position.y).powi(2)) + .sqrt() + .clamp(1.0, 4.0); + for step in 1..=steps { + if cancellation.is_cancellation_requested() { + return Err(BehaviorError::Cancelled); + } + let point = if step == steps { + target + } else { + let fraction = f64::from(step) / f64::from(steps); + let heading = (pose.heading_degrees + delta * fraction).to_radians(); + WorldPosition { + x: pose.position.x + heading.cos() * horizontal_distance, + y: pose.position.y + heading.sin() * horizontal_distance, + z: pose.position.z + (target.z - pose.position.z) * fraction, + } + }; + controller + .sink + .face_point(state.generation, point, cancellation.clone()) + .await?; + if step != steps { + tokio::time::sleep(TURN_INTERVAL).await; + } + } + Ok(()) +} + +fn pose_value(pose: &EmbodiedPose, state: ReadyState) -> Result { + validate_pose(pose, state)?; + Ok( + json!({"status":"observed","generation":pose.generation,"region_id":pose.region_id.to_string(),"position":pose.position,"heading_degrees":pose.heading_degrees,"sitting":pose.sitting}), + ) +} + +fn validate_pose(pose: &EmbodiedPose, state: ReadyState) -> Result<(), BehaviorError> { + pose.valid(state.generation, state.region_id) + .then_some(()) + .ok_or(BehaviorError::RegionBoundary) +} + +fn validate_point( + origin: WorldPosition, + point: WorldPosition, + maximum: u32, +) -> Result<(), BehaviorError> { + if !point.x.is_finite() + || !point.y.is_finite() + || !point.z.is_finite() + || distance(origin, point) > f64::from(maximum) + || !(0.0..=256.0).contains(&point.x) + || !(0.0..=256.0).contains(&point.y) + { + return Err(BehaviorError::InvalidArguments); + } + Ok(()) +} + +fn distance(left: WorldPosition, right: WorldPosition) -> f64 { + ((left.x - right.x).powi(2) + (left.y - right.y).powi(2) + (left.z - right.z).powi(2)).sqrt() +} + +fn random_duration( + controller: &BehaviorController, + minimum: Duration, + maximum: Duration, +) -> Duration { + let range = maximum.saturating_sub(minimum).as_millis(); + if range == 0 { + return minimum; + } + let offset = u128::from(controller.random.next_u64()) % (range + 1); + minimum + Duration::from_millis(u64::try_from(offset).unwrap_or(u64::MAX)) +} + +#[derive(Clone)] +pub struct BehaviorBackend { + ingress: BehaviorIngress, +} + +struct WalkCost { + maximum_meters: u32, +} + +impl ResourceEstimator for WalkCost { + #[allow(clippy::cast_possible_truncation, clippy::cast_sign_loss)] + fn estimate(&self, arguments: &Value) -> Result { + let distance = arguments + .get("distance_meters") + .and_then(Value::as_f64) + .filter(|value| { + value.is_finite() && (0.25..=f64::from(self.maximum_meters)).contains(value) + }) + .ok_or(PolicyReasonCode::InvalidArguments)?; + Ok(ResourceCost { + tool_calls: 1, + movement_millimeters: (distance * 1_000.0).ceil() as u64, + ..ResourceCost::default() + }) + } +} + +impl BehaviorBackend { + #[must_use] + pub const fn new(ingress: BehaviorIngress) -> Self { + Self { ingress } + } +} + +impl AuthorizedToolBackend for BehaviorBackend { + fn apply( + &self, + action: AuthorizedAction, + cancellation: CancellationToken, + ) -> BackendFuture<'_, Result> { + Box::pin(async move { + let call_id = action.call().call_id.clone(); + let authorization_id = action.authorization_id(); + let parsed = parse_action( + action.call().name.as_str(), + action.call().arguments_json.as_str(), + ); + let result = match parsed { + Ok(behavior_action) => tokio::select! { + () = cancellation.cancelled() => Err(BehaviorError::Cancelled), + value = self.ingress.request(behavior_action, BehaviorTrigger::AuthorizedTool { authorization_id, tool: action.call().name.as_str().to_owned() }, BehaviorPolicyResult::Authorized(authorization_id)) => value, + }, + Err(error) => Err(error), + }; + Ok(match result { + Ok(value) => { + let serialized = + serde_json::to_string(&value).map_err(|_| BackendError::Operation { + operation: "serialize behavior result", + })?; + if serialized.len() > MAX_TOOL_RESULT_BYTES { + return Err(BackendError::Operation { + operation: "bounded behavior result", + }); + } + ToolCallOutcome::Completed { + call_id, + result: BoundedText::::new("behavior.result", serialized) + .map_err(|_| BackendError::Operation { + operation: "bounded behavior result", + })?, + } + } + Err(error) => ToolCallOutcome::Rejected { + call_id, + reason: BoundedText::::new( + "behavior.rejection", + error.to_string(), + ) + .map_err(|_| BackendError::Operation { + operation: "bounded behavior rejection", + })?, + }, + }) + }) + } +} + +fn parse_action(name: &str, json_arguments: &str) -> Result { + let value: Value = + serde_json::from_str(json_arguments).map_err(|_| BehaviorError::InvalidArguments)?; + let args = value.as_object().ok_or(BehaviorError::InvalidArguments)?; + match name { + FACE_AVATAR_TOOL => Ok(BehaviorAction::FaceAvatar(parse_uuid(args, "avatar_id")?)), + FACE_POINT_TOOL => Ok(BehaviorAction::FacePoint(parse_point(args)?)), + LOOK_AROUND_TOOL => { + require_empty(args)?; + Ok(BehaviorAction::LookAround) + } + WALK_SHORT_TOOL => Ok(BehaviorAction::WalkShort { + heading_degrees: number(args, "heading_degrees")?, + distance_meters: number(args, "distance_meters")?, + }), + STOP_TOOL => { + require_empty(args)?; + Ok(BehaviorAction::Stop) + } + SIT_TOOL => { + require_empty(args)?; + Ok(BehaviorAction::Sit) + } + STAND_TOOL => { + require_empty(args)?; + Ok(BehaviorAction::Stand) + } + CURRENT_POSE_TOOL => { + require_empty(args)?; + Ok(BehaviorAction::CurrentPose) + } + _ => Err(BehaviorError::InvalidArguments), + } +} + +fn parse_uuid(args: &Map, name: &str) -> Result { + let text = args + .get(name) + .and_then(Value::as_str) + .ok_or(BehaviorError::InvalidArguments)?; + let id = UUID::parse(text.to_owned()).map_err(|_| BehaviorError::InvalidArguments)?; + (id != UUID::zero()) + .then_some(id) + .ok_or(BehaviorError::InvalidArguments) +} + +fn parse_point(args: &Map) -> Result { + if args.len() != 3 { + return Err(BehaviorError::InvalidArguments); + } + Ok(WorldPosition { + x: number(args, "x")?, + y: number(args, "y")?, + z: number(args, "z")?, + }) +} + +fn number(args: &Map, name: &str) -> Result { + args.get(name) + .and_then(Value::as_f64) + .filter(|value| value.is_finite()) + .ok_or(BehaviorError::InvalidArguments) +} + +fn require_empty(args: &Map) -> Result<(), BehaviorError> { + args.is_empty() + .then_some(()) + .ok_or(BehaviorError::InvalidArguments) +} + +#[allow(clippy::too_many_lines)] // One registry table keeps the eight schemas auditable together. +pub fn behavior_policy_tools(settings: &BehaviorSettings) -> Result, PolicyError> { + let origins = || { + AllowedOrigins::new([ + OriginClass::AuthorizedIm, + OriginClass::LocalOperator, + OriginClass::InternalScheduler, + ]) + }; + let empty = ToolSchema::Object { + properties: BTreeMap::new(), + required: BTreeSet::new(), + additional_properties: false, + }; + let object = |properties: &[(&str, ToolSchema)], required: &[&str]| ToolSchema::Object { + properties: properties + .iter() + .map(|(name, schema)| ((*name).to_owned(), schema.clone())) + .collect(), + required: required.iter().map(|name| (*name).to_owned()).collect(), + additional_properties: false, + }; + let specs = [ + ( + FACE_AVATAR_TOOL, + "Turn attention toward a currently visible avatar.", + object(&[("avatar_id", ToolSchema::String)], &["avatar_id"]), + true, + 0, + ), + ( + FACE_POINT_TOOL, + "Turn attention toward a nearby point in the current region.", + object( + &[ + ("x", ToolSchema::Number), + ("y", ToolSchema::Number), + ("z", ToolSchema::Number), + ], + &["x", "y", "z"], + ), + true, + 0, + ), + ( + LOOK_AROUND_TOOL, + "Make one small bounded attention shift.", + empty.clone(), + true, + 0, + ), + ( + WALK_SHORT_TOOL, + "Walk a short bounded distance on a heading, then stop.", + object( + &[ + ("heading_degrees", ToolSchema::Number), + ("distance_meters", ToolSchema::Number), + ], + &["heading_degrees", "distance_meters"], + ), + true, + u64::from(settings.max_walk_distance_meters) * 1_000, + ), + ( + STOP_TOOL, + "Immediately stop current bounded movement.", + empty.clone(), + true, + 0, + ), + ( + SIT_TOOL, + "Sit on the ground when supported.", + empty.clone(), + true, + 0, + ), + ( + STAND_TOOL, + "Stand from the current seated pose.", + empty.clone(), + true, + 0, + ), + ( + CURRENT_POSE_TOOL, + "Read current embodied pose and session provenance.", + empty, + false, + 0, + ), + ]; + specs + .into_iter() + .map(|(name, description, schema, mutating, movement_mm)| { + let cost = ResourceCost { + tool_calls: 1, + movement_millimeters: movement_mm, + ..ResourceCost::default() + }; + let estimator: Arc = if name == WALK_SHORT_TOOL { + Arc::new(WalkCost { + maximum_meters: settings.max_walk_distance_meters, + }) + } else { + Arc::new(FixedCost(cost)) + }; + PolicyTool::new( + ToolDefinition { + name: BoundedText::new("behavior.tool.name", name)?, + description: BoundedText::new("behavior.tool.description", description)?, + schema, + mutating, + }, + if mutating { + Capability::Movement + } else { + Capability::Informational + }, + if mutating { + Risk::Movement + } else { + Risk::ReadOnly + }, + origins()?, + cost, + Idempotency::Idempotent, + ApprovalRule::Never, + true, + estimator, + ) + }) + .collect() +} + +/// Exact-name router used when one policy gateway serves perception and behavior. +pub struct AuthorizedBackendRouter { + routes: BTreeMap>, +} + +impl AuthorizedBackendRouter { + pub fn new( + routes: impl IntoIterator)>, + ) -> Result { + let mut mapped = BTreeMap::new(); + for (name, backend) in routes { + if name.is_empty() || mapped.insert(name, backend).is_some() { + return Err(BackendError::Configuration { + component: "authorized backend routes", + }); + } + } + if mapped.is_empty() { + return Err(BackendError::Configuration { + component: "authorized backend routes", + }); + } + Ok(Self { routes: mapped }) + } +} + +impl AuthorizedToolBackend for AuthorizedBackendRouter { + fn apply( + &self, + action: AuthorizedAction, + cancellation: CancellationToken, + ) -> BackendFuture<'_, Result> { + Box::pin(async move { + let backend = self + .routes + .get(action.call().name.as_str()) + .ok_or(BackendError::RejectedMutation)?; + backend.apply(action, cancellation).await + }) + } +} diff --git a/crates/metacrate-grid-agent/src/behavior_tests.rs b/crates/metacrate-grid-agent/src/behavior_tests.rs new file mode 100644 index 0000000..adc784a --- /dev/null +++ b/crates/metacrate-grid-agent/src/behavior_tests.rs @@ -0,0 +1,562 @@ +use crate::backend::AuthorizedToolBackend; +use crate::behavior::*; +use crate::config::BehaviorSettings; +use crate::perception::WorldPosition; +use crate::policy::{ + ActionOrigin, MemoryPolicyAudit, PolicyAuditSink, PolicyGateway, PolicyLimits, + PolicyRequestContext, +}; +use crate::types::{ProposedToolCall, ToolCallOutcome}; +use libremetaverse_types::UUID; +use libremetaverse_types::compat::CancellationToken; +use serde_json::{Value, json}; +use std::collections::{BTreeMap, BTreeSet}; +use std::sync::{Arc, Mutex}; +use std::time::Duration; + +fn uuid(number: u64) -> UUID { + UUID::new_with_string(format!("00000000-0000-4000-8000-{number:012x}")).expect("UUID") +} + +fn settings() -> BehaviorSettings { + BehaviorSettings { + heartbeat: Duration::from_secs(30), + settle_delay: Duration::from_millis(10), + response_delay_min: Duration::from_millis(20), + response_delay_max: Duration::from_millis(20), + attention_dwell: Duration::from_millis(30), + idle_interval: Duration::from_secs(15), + action_timeout: Duration::from_secs(5), + max_walk_duration: Duration::from_secs(4), + stuck_timeout: Duration::from_secs(1), + min_action_interval: Duration::ZERO, + max_walk_distance_meters: 8, + max_attention_distance_meters: 96, + idle_look_enabled: true, + } +} + +#[derive(Clone)] +struct FakeSink { + state: Arc>, +} + +struct FakeState { + pose: EmbodiedPose, + avatars: BTreeMap, + calls: Vec, + walk_target: Option, + stuck: bool, + walk_allowed: bool, +} + +impl FakeSink { + fn new(generation: u64, region_id: UUID) -> Self { + Self { + state: Arc::new(Mutex::new(FakeState { + pose: EmbodiedPose { + generation, + region_id, + position: WorldPosition { + x: 128.0, + y: 128.0, + z: 24.0, + }, + heading_degrees: 0.0, + sitting: false, + }, + avatars: BTreeMap::from([( + uuid(2), + WorldPosition { + x: 132.0, + y: 128.0, + z: 24.0, + }, + )]), + calls: Vec::new(), + walk_target: None, + stuck: false, + walk_allowed: true, + })), + } + } + + fn calls(&self) -> Vec { + self.state.lock().expect("state").calls.clone() + } + fn set_region(&self, region_id: UUID) { + self.state.lock().expect("state").pose.region_id = region_id; + } + fn set_stuck(&self, stuck: bool) { + self.state.lock().expect("state").stuck = stuck; + } + fn set_walk_allowed(&self, allowed: bool) { + self.state.lock().expect("state").walk_allowed = allowed; + } +} + +impl EmbodimentSink for FakeSink { + fn current_pose( + &self, + generation: u64, + _: CancellationToken, + ) -> EmbodimentFuture<'_, EmbodiedPose> { + Box::pin(async move { + let mut state = self.state.lock().expect("state"); + if state.pose.generation != generation { + return Err(BehaviorError::NotReady); + } + if let Some(target) = state.walk_target + && !state.stuck + { + state.pose.position = target; + } + Ok(state.pose.clone()) + }) + } + fn resolve_avatar( + &self, + generation: u64, + id: UUID, + _: CancellationToken, + ) -> EmbodimentFuture<'_, WorldPosition> { + Box::pin(async move { + let state = self.state.lock().expect("state"); + if state.pose.generation != generation { + return Err(BehaviorError::NotReady); + } + state + .avatars + .get(&id) + .copied() + .ok_or(BehaviorError::TargetUnavailable) + }) + } + fn face_point( + &self, + _: u64, + _: WorldPosition, + _: CancellationToken, + ) -> EmbodimentFuture<'_, ()> { + Box::pin(async move { + self.state.lock().expect("state").calls.push("face".into()); + Ok(()) + }) + } + fn begin_walk( + &self, + _: u64, + point: WorldPosition, + _: CancellationToken, + ) -> EmbodimentFuture<'_, ()> { + Box::pin(async move { + let mut state = self.state.lock().expect("state"); + state.calls.push("walk".into()); + state.walk_target = Some(point); + Ok(()) + }) + } + fn validate_walk_target( + &self, + _: u64, + _: WorldPosition, + _: CancellationToken, + ) -> EmbodimentFuture<'_, ()> { + Box::pin(async move { + self.state + .lock() + .expect("state") + .walk_allowed + .then_some(()) + .ok_or(BehaviorError::RegionBoundary) + }) + } + fn stop(&self, _: u64, _: CancellationToken) -> EmbodimentFuture<'_, ()> { + Box::pin(async move { + let mut state = self.state.lock().expect("state"); + state.calls.push("stop".into()); + state.walk_target = None; + Ok(()) + }) + } + fn sit(&self, _: u64, _: CancellationToken) -> EmbodimentFuture<'_, ()> { + Box::pin(async move { + let mut state = self.state.lock().expect("state"); + state.calls.push("sit".into()); + state.pose.sitting = true; + Ok(()) + }) + } + fn stand(&self, _: u64, _: CancellationToken) -> EmbodimentFuture<'_, ()> { + Box::pin(async move { + let mut state = self.state.lock().expect("state"); + state.calls.push("stand".into()); + state.pose.sitting = false; + Ok(()) + }) + } +} + +struct FixedRandom(u64); +impl BehaviorRandom for FixedRandom { + fn next_u64(&self) -> u64 { + self.0 + } +} + +fn authorize( + settings: &BehaviorSettings, + name: &str, + arguments: &Value, +) -> crate::policy::AuthorizedAction { + let avatar = uuid(900); + let audit: Arc = Arc::new(MemoryPolicyAudit::new(64).expect("audit")); + let gateway = PolicyGateway::new( + BTreeSet::from([avatar]), + behavior_policy_tools(settings).expect("tools"), + PolicyLimits::default(), + audit, + ) + .expect("gateway"); + let context = PolicyRequestContext::new( + ActionOrigin::instant_message(avatar), + "behavior-session", + "behavior-correlation", + ) + .expect("context"); + let call = ProposedToolCall::new("behavior-call", name, arguments.to_string()).expect("call"); + gateway + .evaluate(&context, &call, arguments, None, 100) + .expect("evaluation") + .into_authorization() + .expect("authorization") +} + +async fn run_tool( + backend: &BehaviorBackend, + settings: &BehaviorSettings, + name: &str, + args: Value, +) -> ToolCallOutcome { + backend + .apply( + authorize(settings, name, &args), + CancellationToken::default(), + ) + .await + .expect("backend") +} + +#[test] +fn policy_exposes_only_bounded_high_level_actions() { + let tools = behavior_policy_tools(&settings()).expect("tools"); + assert_eq!(tools.len(), 8); + let names = tools + .iter() + .map(|tool| tool.definition.name.as_str()) + .collect::>(); + assert_eq!( + names, + [ + FACE_AVATAR_TOOL, + FACE_POINT_TOOL, + LOOK_AROUND_TOOL, + WALK_SHORT_TOOL, + STOP_TOOL, + SIT_TOOL, + STAND_TOOL, + CURRENT_POSE_TOOL + ] + ); + assert!(names.iter().all(|name| !name.contains("teleport") + && !name.contains("control") + && !name.contains("follow"))); +} + +#[tokio::test(start_paused = true)] +async fn public_attention_waits_turns_and_returns_to_available() { + let region = uuid(10); + let sink = Arc::new(FakeSink::new(1, region)); + let handle = BehaviorController::with_random( + settings(), + sink.clone(), + 16, + 32, + Duration::from_secs(1), + Arc::new(FixedRandom(0)), + ) + .expect("controller") + .start(); + let ingress = handle.ingress(); + ingress.connected(1, region).expect("ready"); + tokio::time::advance(Duration::from_millis(11)).await; + tokio::task::yield_now().await; + let pacing = tokio::spawn(async move { ingress.attention("delivery-1".into(), uuid(2)).await }); + tokio::time::advance(Duration::from_millis(21)).await; + tokio::task::yield_now().await; + assert!(pacing.await.expect("task").is_ok()); + assert_eq!(sink.calls(), vec!["face"]); + tokio::time::advance(Duration::from_millis(31)).await; + handle.shutdown().await.expect("shutdown"); +} + +#[tokio::test(start_paused = true)] +async fn quarter_turn_uses_bounded_intermediate_camera_updates() { + let region = uuid(10); + let sink = Arc::new(FakeSink::new(1, region)); + let mut config = settings(); + config.settle_delay = Duration::ZERO; + let handle = + BehaviorController::new(config.clone(), sink.clone(), 16, 32, Duration::from_secs(1)) + .expect("controller") + .start(); + let ingress = handle.ingress(); + ingress.connected(1, region).expect("ready"); + tokio::task::yield_now().await; + let backend = BehaviorBackend::new(ingress); + let turn = tokio::spawn({ + let config = config.clone(); + async move { + run_tool( + &backend, + &config, + FACE_POINT_TOOL, + json!({"x":128.0,"y":132.0,"z":24.0}), + ) + .await + } + }); + tokio::task::yield_now().await; + tokio::time::advance(Duration::from_millis(100)).await; + tokio::task::yield_now().await; + assert!(matches!( + turn.await.expect("turn"), + ToolCallOutcome::Completed { .. } + )); + assert_eq!(sink.calls(), vec!["face", "face", "face"]); + handle.shutdown().await.expect("shutdown"); +} + +#[tokio::test(start_paused = true)] +async fn bounded_walk_stops_and_region_change_rejects_stale_pose() { + let region = uuid(10); + let sink = Arc::new(FakeSink::new(1, region)); + let mut config = settings(); + config.settle_delay = Duration::ZERO; + let handle = + BehaviorController::new(config.clone(), sink.clone(), 16, 32, Duration::from_secs(1)) + .expect("controller") + .start(); + let ingress = handle.ingress(); + ingress.connected(1, region).expect("ready"); + tokio::task::yield_now().await; + let backend = BehaviorBackend::new(ingress.clone()); + let task = tokio::spawn({ + let backend = backend.clone(); + let config = config.clone(); + async move { + run_tool( + &backend, + &config, + WALK_SHORT_TOOL, + json!({"heading_degrees":0.0,"distance_meters":2.0}), + ) + .await + } + }); + tokio::time::advance(Duration::from_millis(300)).await; + tokio::task::yield_now().await; + assert!(matches!( + task.await.expect("task"), + ToolCallOutcome::Completed { .. } + )); + assert_eq!(sink.calls(), vec!["face", "walk", "stop"]); + let new_region = uuid(11); + sink.set_region(new_region); + ingress.region_changed(1, new_region).expect("cross"); + tokio::task::yield_now().await; + let pose = run_tool(&backend, &config, CURRENT_POSE_TOOL, json!({})).await; + assert!(matches!(pose, ToolCallOutcome::Completed { .. })); + handle.shutdown().await.expect("shutdown"); +} + +#[tokio::test(start_paused = true)] +async fn parcel_boundary_rejects_walk_before_any_movement_update() { + let region = uuid(10); + let sink = Arc::new(FakeSink::new(1, region)); + sink.set_walk_allowed(false); + let mut config = settings(); + config.settle_delay = Duration::ZERO; + let handle = + BehaviorController::new(config.clone(), sink.clone(), 16, 32, Duration::from_secs(1)) + .expect("controller") + .start(); + let ingress = handle.ingress(); + ingress.connected(1, region).expect("ready"); + tokio::task::yield_now().await; + let backend = BehaviorBackend::new(ingress); + assert!(matches!( + run_tool( + &backend, + &config, + WALK_SHORT_TOOL, + json!({"heading_degrees":0.0,"distance_meters":2.0}), + ) + .await, + ToolCallOutcome::Rejected { .. } + )); + assert_eq!(sink.calls(), vec!["stop"]); + handle.shutdown().await.expect("shutdown"); +} + +#[tokio::test(start_paused = true)] +async fn pause_preempts_stuck_walk_and_blocks_all_motion() { + let region = uuid(10); + let sink = Arc::new(FakeSink::new(1, region)); + sink.set_stuck(true); + let mut config = settings(); + config.settle_delay = Duration::ZERO; + let handle = + BehaviorController::new(config.clone(), sink.clone(), 16, 32, Duration::from_secs(1)) + .expect("controller") + .start(); + let ingress = handle.ingress(); + ingress.connected(1, region).expect("ready"); + tokio::task::yield_now().await; + let backend = BehaviorBackend::new(ingress.clone()); + let task = tokio::spawn({ + let backend = backend.clone(); + let config = config.clone(); + async move { + run_tool( + &backend, + &config, + WALK_SHORT_TOOL, + json!({"heading_degrees":90.0,"distance_meters":2.0}), + ) + .await + } + }); + tokio::time::advance(Duration::from_millis(300)).await; + ingress.pause(); + tokio::task::yield_now().await; + assert!(matches!( + task.await.expect("task"), + ToolCallOutcome::Rejected { .. } + )); + let before = sink.calls().len(); + assert!(matches!( + run_tool(&backend, &config, SIT_TOOL, json!({})).await, + ToolCallOutcome::Rejected { .. } + )); + assert_eq!(sink.calls().len(), before); + handle.shutdown().await.expect("shutdown"); +} + +#[tokio::test(start_paused = true)] +async fn offline_and_emergency_stop_emit_zero_motion() { + let region = uuid(10); + let sink = Arc::new(FakeSink::new(1, region)); + let mut config = settings(); + config.settle_delay = Duration::ZERO; + let handle = + BehaviorController::new(config.clone(), sink.clone(), 16, 32, Duration::from_secs(1)) + .expect("controller") + .start(); + let ingress = handle.ingress(); + let backend = BehaviorBackend::new(ingress.clone()); + assert!(matches!( + run_tool(&backend, &config, SIT_TOOL, json!({})).await, + ToolCallOutcome::Rejected { .. } + )); + ingress.connected(1, region).expect("ready"); + tokio::task::yield_now().await; + ingress.emergency_stop(); + tokio::task::yield_now().await; + let after_stop = sink.calls().len(); + assert!(matches!( + run_tool(&backend, &config, STAND_TOOL, json!({})).await, + ToolCallOutcome::Rejected { .. } + )); + assert_eq!(sink.calls().len(), after_stop); + handle.shutdown().await.expect("shutdown"); +} + +#[tokio::test(start_paused = true)] +async fn seeded_idle_is_one_low_frequency_look_without_movement_or_chat() { + let region = uuid(10); + let sink = Arc::new(FakeSink::new(1, region)); + let mut config = settings(); + config.settle_delay = Duration::ZERO; + let handle = BehaviorController::with_random( + config, + sink.clone(), + 16, + 32, + Duration::from_secs(1), + Arc::new(FixedRandom(30)), + ) + .expect("controller") + .start(); + let ingress = handle.ingress(); + ingress.connected(1, region).expect("ready"); + tokio::task::yield_now().await; + tokio::time::advance(Duration::from_secs(14)).await; + tokio::task::yield_now().await; + assert!(sink.calls().is_empty()); + tokio::time::advance(Duration::from_secs(2)).await; + tokio::task::yield_now().await; + assert_eq!(sink.calls(), vec!["face"]); + handle.shutdown().await.expect("shutdown"); +} + +#[tokio::test(start_paused = true)] +async fn region_change_preempts_walk_and_disappearance_cancels_attention() { + let region = uuid(10); + let sink = Arc::new(FakeSink::new(1, region)); + sink.set_stuck(true); + let mut config = settings(); + config.settle_delay = Duration::ZERO; + config.response_delay_min = Duration::ZERO; + config.response_delay_max = Duration::ZERO; + let handle = + BehaviorController::new(config.clone(), sink.clone(), 16, 32, Duration::from_secs(1)) + .expect("controller") + .start(); + let ingress = handle.ingress(); + ingress.connected(1, region).expect("ready"); + tokio::task::yield_now().await; + let backend = BehaviorBackend::new(ingress.clone()); + let task = tokio::spawn({ + let backend = backend.clone(); + let config = config.clone(); + async move { + run_tool( + &backend, + &config, + WALK_SHORT_TOOL, + json!({"heading_degrees":180.0,"distance_meters":2.0}), + ) + .await + } + }); + tokio::time::advance(Duration::from_millis(300)).await; + let next = uuid(11); + sink.set_region(next); + ingress.region_changed(1, next).expect("cross"); + tokio::task::yield_now().await; + assert!(matches!( + task.await.expect("walk"), + ToolCallOutcome::Rejected { .. } + )); + let missing = tokio::spawn({ + let ingress = ingress.clone(); + async move { ingress.attention("missing".into(), uuid(999)).await } + }); + tokio::task::yield_now().await; + assert_eq!( + missing.await.expect("attention"), + Err(BehaviorError::TargetUnavailable) + ); + handle.shutdown().await.expect("shutdown"); +} diff --git a/crates/metacrate-grid-agent/src/config.rs b/crates/metacrate-grid-agent/src/config.rs index 3696483..9d8d55b 100644 --- a/crates/metacrate-grid-agent/src/config.rs +++ b/crates/metacrate-grid-agent/src/config.rs @@ -200,6 +200,43 @@ impl Default for Limits { #[derive(Clone, Debug, Eq, PartialEq)] pub struct BehaviorSettings { pub heartbeat: Duration, + pub settle_delay: Duration, + pub response_delay_min: Duration, + pub response_delay_max: Duration, + pub attention_dwell: Duration, + pub idle_interval: Duration, + pub action_timeout: Duration, + pub max_walk_duration: Duration, + pub stuck_timeout: Duration, + pub min_action_interval: Duration, + pub max_walk_distance_meters: u32, + pub max_attention_distance_meters: u32, + pub idle_look_enabled: bool, +} + +impl BehaviorSettings { + pub(crate) fn is_valid(&self) -> bool { + if self.heartbeat.is_zero() + || self.settle_delay > Duration::from_secs(30) + || self.response_delay_min > self.response_delay_max + || self.response_delay_max > Duration::from_secs(30) + || self.attention_dwell > Duration::from_mins(1) + || self.idle_interval < Duration::from_secs(15) + || self.idle_interval > Duration::from_hours(1) + || self.action_timeout.is_zero() + || self.action_timeout > Duration::from_mins(1) + || self.max_walk_duration.is_zero() + || self.max_walk_duration > self.action_timeout + || self.stuck_timeout.is_zero() + || self.stuck_timeout > self.max_walk_duration + || self.min_action_interval > Duration::from_secs(30) + || !(1..=20).contains(&self.max_walk_distance_meters) + || !(1..=256).contains(&self.max_attention_distance_meters) + { + return false; + } + true + } } /// Conversation-memory settings. Persistence is opt-in and uses @@ -260,6 +297,9 @@ impl AgentConfig { self.interaction .validate() .map_err(|_| ConfigError::InvalidInteraction)?; + if !self.behavior.is_valid() { + return Err(ConfigError::InvalidBehavior); + } if self.mode != OperatingMode::OfflineFake && self.grid.is_none() { return Err(ConfigError::Missing { field: "grid", @@ -435,6 +475,7 @@ pub enum ConfigError { InvalidReconnect, InvalidConversationMemory, InvalidInteraction, + InvalidBehavior, } impl fmt::Display for ConfigError { @@ -489,6 +530,7 @@ impl fmt::Display for ConfigError { formatter.write_str("invalid conversation-memory bounds") } Self::InvalidInteraction => formatter.write_str("invalid interaction bounds"), + Self::InvalidBehavior => formatter.write_str("invalid embodied behavior bounds"), } } } @@ -564,6 +606,18 @@ struct RawLimits { #[serde(default, deny_unknown_fields)] struct RawBehavior { heartbeat_seconds: Option, + settle_milliseconds: Option, + response_delay_min_milliseconds: Option, + response_delay_max_milliseconds: Option, + attention_dwell_milliseconds: Option, + idle_interval_seconds: Option, + action_timeout_seconds: Option, + max_walk_duration_seconds: Option, + stuck_timeout_seconds: Option, + min_action_interval_milliseconds: Option, + max_walk_distance_meters: Option, + max_attention_distance_meters: Option, + idle_look_enabled: Option, } #[derive(Clone, Default, Deserialize)] @@ -896,6 +950,48 @@ fn resolve( 1, 300, )?, + settle_delay: Duration::from_millis(raw.behavior.settle_milliseconds.unwrap_or(2_000)), + response_delay_min: Duration::from_millis( + raw.behavior.response_delay_min_milliseconds.unwrap_or(350), + ), + response_delay_max: Duration::from_millis( + raw.behavior + .response_delay_max_milliseconds + .unwrap_or(1_200), + ), + attention_dwell: Duration::from_millis( + raw.behavior.attention_dwell_milliseconds.unwrap_or(4_000), + ), + idle_interval: checked_duration( + "behavior.idle_interval_seconds", + raw.behavior.idle_interval_seconds.unwrap_or(120), + 15, + 3_600, + )?, + action_timeout: checked_duration( + "behavior.action_timeout_seconds", + raw.behavior.action_timeout_seconds.unwrap_or(15), + 1, + 60, + )?, + max_walk_duration: checked_duration( + "behavior.max_walk_duration_seconds", + raw.behavior.max_walk_duration_seconds.unwrap_or(10), + 1, + 60, + )?, + stuck_timeout: checked_duration( + "behavior.stuck_timeout_seconds", + raw.behavior.stuck_timeout_seconds.unwrap_or(3), + 1, + 60, + )?, + min_action_interval: Duration::from_millis( + raw.behavior.min_action_interval_milliseconds.unwrap_or(750), + ), + max_walk_distance_meters: raw.behavior.max_walk_distance_meters.unwrap_or(8), + max_attention_distance_meters: raw.behavior.max_attention_distance_meters.unwrap_or(96), + idle_look_enabled: raw.behavior.idle_look_enabled.unwrap_or(true), }, reconnect, conversation, diff --git a/crates/metacrate-grid-agent/src/interaction.rs b/crates/metacrate-grid-agent/src/interaction.rs index ccf9570..2458b35 100644 --- a/crates/metacrate-grid-agent/src/interaction.rs +++ b/crates/metacrate-grid-agent/src/interaction.rs @@ -270,6 +270,26 @@ pub trait InteractionResponder: Send + Sync + 'static { ) -> ResponderFuture<'_>; } +pub type PacerFuture<'a> = + Pin> + Send + 'a>>; + +/// Optional embodied response hook. It may turn toward a nearby speaker and +/// wait a bounded, human-scale delay before public delivery. +pub trait ResponsePacer: Send + Sync + 'static { + fn prepare_public_response( + &self, + delivery_id: String, + avatar_id: UUID, + cancellation: CancellationToken, + ) -> PacerFuture<'_>; +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum InteractionPacingError { + Cancelled, + Unavailable, +} + #[derive(Clone, Copy, Debug, Eq, PartialEq)] pub enum InteractionModelError { Cancelled, @@ -450,6 +470,7 @@ pub struct InteractionCoordinator { conversation: Arc, responder: Arc, sink: Arc, + pacer: Option>, input_capacity: usize, observation_capacity: usize, shutdown_timeout: Duration, @@ -499,12 +520,19 @@ impl InteractionCoordinator { conversation, responder, sink, + pacer: None, input_capacity, observation_capacity, shutdown_timeout, }) } + #[must_use] + pub fn with_response_pacer(mut self, pacer: Arc) -> Self { + self.pacer = Some(pacer); + self + } + #[must_use] pub fn start(self) -> InteractionHandle { let (commands, command_receiver) = mpsc::channel(self.input_capacity); @@ -947,6 +975,7 @@ fn schedule_ready( Arc::clone(&coordinator.conversation), Arc::clone(&coordinator.responder), Arc::clone(&coordinator.sink), + coordinator.pacer.clone(), Arc::clone(limiter), observations.clone(), completions.clone(), @@ -966,6 +995,7 @@ async fn process_batch( conversation: Arc, responder: Arc, sink: Arc, + pacer: Option>, limiter: Arc, observations: mpsc::Sender, completions: mpsc::Sender, @@ -1001,6 +1031,7 @@ async fn process_batch( &sink, &limiter, &cancellation, + pacer.as_ref(), "I could not retain that message safely.", ) .await; @@ -1031,6 +1062,7 @@ async fn process_batch( &sink, &limiter, &cancellation, + pacer.as_ref(), "For safety, commands must be requested through an authorized direct message.", ) .await; @@ -1101,6 +1133,15 @@ async fn process_batch( return; }; let parts = split_utf8(&safe, settings.grid_chunk_bytes); + if public_engaged && let Some(pacer) = &pacer { + let _ = pacer + .prepare_public_response( + delivery_id.as_str().to_owned(), + trigger.sender_id, + cancellation.clone(), + ) + .await; + } let delivered = deliver_parts( generation, trigger, @@ -1173,6 +1214,7 @@ async fn process_batch( &sink, &limiter, &cancellation, + pacer.as_ref(), "I could not answer that just now.", ) .await; @@ -1209,6 +1251,7 @@ async fn fixed_failure( sink: &Arc, limiter: &Arc, cancellation: &CancellationToken, + pacer: Option<&Arc>, text: &str, ) -> usize { let useful = trigger.channel == InteractionChannel::DirectIm @@ -1218,6 +1261,18 @@ async fn fixed_failure( let Some(session_id) = session_id.filter(|_| useful) else { return 0; }; + if trigger.channel == InteractionChannel::PublicChat + && trigger.nearby + && let Some(pacer) = pacer + { + let _ = pacer + .prepare_public_response( + trigger.delivery_id.as_str().to_owned(), + trigger.sender_id, + cancellation.clone(), + ) + .await; + } let parts = split_utf8(text, settings.grid_chunk_bytes); deliver_parts( generation, diff --git a/crates/metacrate-grid-agent/src/interaction_tests.rs b/crates/metacrate-grid-agent/src/interaction_tests.rs index 25e1e0c..f8427a6 100644 --- a/crates/metacrate-grid-agent/src/interaction_tests.rs +++ b/crates/metacrate-grid-agent/src/interaction_tests.rs @@ -165,6 +165,28 @@ impl InteractionSink for FakeSink { } } +struct FakePacer { + called: AtomicBool, + delay: Duration, +} + +impl ResponsePacer for FakePacer { + fn prepare_public_response( + &self, + _delivery_id: String, + _avatar_id: UUID, + cancellation: libremetaverse_types::compat::CancellationToken, + ) -> PacerFuture<'_> { + self.called.store(true, Ordering::Release); + Box::pin(async move { + tokio::select! { + () = cancellation.cancelled() => Err(InteractionPacingError::Cancelled), + () = tokio::time::sleep(self.delay) => Ok(()), + } + }) + } +} + fn coordinator( interaction_settings: InteractionSettings, authorized: BTreeSet, @@ -235,6 +257,38 @@ fn request_text(request: &ResponseRequest) -> String { .join("\n") } +#[tokio::test(start_paused = true)] +async fn public_response_pacer_runs_before_visible_delivery() { + let responder = Arc::new(FakeResponder::with_response("Hello there")); + let sink = Arc::new(FakeSink::default()); + let pacer = Arc::new(FakePacer { + called: AtomicBool::new(false), + delay: Duration::from_millis(500), + }); + let mut handle = coordinator(settings(), BTreeSet::new(), responder, sink.clone()) + .with_response_pacer(pacer.clone()) + .start(); + handle.connected(1).expect("connect"); + handle + .submit(inbound( + "paced", + 1, + InteractionChannel::PublicChat, + "hello metacrate", + )) + .await + .expect("submit"); + settle_debounce().await; + assert!(pacer.called.load(Ordering::Acquire)); + assert!(sink.messages().is_empty()); + tokio::time::advance(Duration::from_millis(501)).await; + for _ in 0..4 { + tokio::task::yield_now().await; + } + assert_eq!(sink.messages().len(), 1); + handle.shutdown().await.expect("shutdown"); +} + #[tokio::test(start_paused = true)] #[allow(clippy::too_many_lines)] async fn mentions_aliases_greetings_ambient_and_public_routes_are_correct() { diff --git a/crates/metacrate-grid-agent/src/lib.rs b/crates/metacrate-grid-agent/src/lib.rs index 997d03f..56bef74 100644 --- a/crates/metacrate-grid-agent/src/lib.rs +++ b/crates/metacrate-grid-agent/src/lib.rs @@ -5,6 +5,7 @@ //! subprocess, provider-SDK, or platform-specific dependency. pub mod backend; +pub mod behavior; pub mod config; pub mod conversation; pub mod interaction; @@ -16,6 +17,8 @@ pub mod session; pub mod tool_loop; pub mod types; +#[cfg(test)] +mod behavior_tests; #[cfg(test)] mod conversation_tests; #[cfg(test)] @@ -33,8 +36,15 @@ pub use backend::{ }; #[cfg(feature = "live-grid")] pub use backend::{ - LibremetaverseClientOwner, LibremetaverseInteractionSink, LibremetaverseSessionBackend, - LibremetaverseWorldSnapshotSource, + LibremetaverseClientOwner, LibremetaverseEmbodimentSink, LibremetaverseInteractionSink, + LibremetaverseSessionBackend, LibremetaverseWorldSnapshotSource, +}; +pub use behavior::{ + AuthorizedBackendRouter, BehaviorBackend, BehaviorController, BehaviorError, BehaviorHandle, + BehaviorIngress, BehaviorMode, BehaviorObservation, BehaviorOutcome, BehaviorPolicyResult, + BehaviorRandom, BehaviorTrigger, CURRENT_POSE_TOOL, EmbodiedPose, EmbodimentFuture, + EmbodimentSink, FACE_AVATAR_TOOL, FACE_POINT_TOOL, LOOK_AROUND_TOOL, SIT_TOOL, STAND_TOOL, + STOP_TOOL, WALK_SHORT_TOOL, behavior_policy_tools, }; pub use config::{ AgentConfig, BehaviorSettings, ConfigError, ConfigLoader, ConversationSettings, EndpointUrl, @@ -51,9 +61,10 @@ pub use interaction::{ DeliveryFuture, DeliveryOutcome, ImDialogKind, InboundInteraction, InboundSource, InteractionChannel, InteractionCoordinator, InteractionDeliveryError, InteractionError, InteractionHandle, InteractionIngress, InteractionIntent, InteractionModelError, - InteractionObservation, InteractionOrigin, InteractionResponder, InteractionSettings, - InteractionSink, OutboundInteraction, PolicyLlmResponder, ResponderFuture, ResponseRequest, - SuppressionReason, VisibleResponse, split_utf8, + InteractionObservation, InteractionOrigin, InteractionPacingError, InteractionResponder, + InteractionSettings, InteractionSink, OutboundInteraction, PacerFuture, PolicyLlmResponder, + ResponderFuture, ResponsePacer, ResponseRequest, SuppressionReason, VisibleResponse, + split_utf8, }; pub use llm::{ Completion, CompletionMessage, ContentPart, ImageDetail, LlmClient, LlmError, diff --git a/crates/metacrate-grid-agent/src/main.rs b/crates/metacrate-grid-agent/src/main.rs index f3a4f20..e61b4d0 100644 --- a/crates/metacrate-grid-agent/src/main.rs +++ b/crates/metacrate-grid-agent/src/main.rs @@ -156,14 +156,16 @@ async fn run_live( })?; let owner = LibremetaverseClientOwner::new()?; let mut live = start_live_interactions(&config, &owner)?; - let backend = match owner.session_backend_with_services( + let backend = match owner.session_backend_with_agent_services( connection, live.interaction.ingress(), live.perception.clone(), + live.behavior.ingress(), ) { Ok(backend) => backend, Err(error) => { live.interaction.shutdown().await?; + live.behavior.shutdown().await?; return Err(error.into()); } }; @@ -212,14 +214,18 @@ async fn run_live( if let Err(error) = readiness { let session_result = handle.shutdown().await; let interaction_result = live.interaction.shutdown().await; + let behavior_result = live.behavior.shutdown().await; session_result?; interaction_result?; + behavior_result?; return Err(error.into()); } let session_result = handle.shutdown().await; let interaction_result = live.interaction.shutdown().await; + let behavior_result = live.behavior.shutdown().await; session_result?; interaction_result?; + behavior_result?; println!("grid agent completed one supervised login/logout cycle"); return Ok(()); } @@ -254,12 +260,18 @@ async fn run_live( let Some(event) = event else { break; }; println!("grid perception event={event:?}"); } + event = live.behavior.next_observation() => { + let Some(event) = event else { break; }; + println!("grid behavior event={event:?}"); + } } } let session_result = handle.shutdown().await; let interaction_result = live.interaction.shutdown().await; + let behavior_result = live.behavior.shutdown().await; session_result?; interaction_result?; + behavior_result?; if let Some(error) = signal_error { return Err(error.into()); } @@ -273,6 +285,7 @@ struct LiveInteractions { perception: metacrate_grid_agent::PerceptionIngress, perception_observations: tokio::sync::mpsc::Receiver, + behavior: metacrate_grid_agent::BehaviorHandle, } #[cfg(feature = "live-grid")] @@ -281,9 +294,10 @@ fn start_live_interactions( owner: &metacrate_grid_agent::LibremetaverseClientOwner, ) -> Result> { use metacrate_grid_agent::{ - AuthorizedToolBackend, ConversationStore, InteractionCoordinator, LlmClient, - LlmTransportLimits, MemoryPolicyAudit, PerceptionBackend, PolicyGateway, PolicyLimits, - PolicyLlmResponder, ToolLoopLimits, perception_policy_tools, + AuthorizedBackendRouter, AuthorizedToolBackend, BehaviorBackend, BehaviorController, + ConversationStore, InteractionCoordinator, LlmClient, LlmTransportLimits, + MemoryPolicyAudit, PerceptionBackend, PolicyGateway, PolicyLimits, PolicyLlmResponder, + ToolLoopLimits, behavior_policy_tools, perception_policy_tools, }; let transport_limits = LlmTransportLimits { @@ -304,11 +318,34 @@ fn start_live_interactions( let perception_observations = perception .take_observations() .ok_or_else(|| CliError("perception observation receiver already claimed".into()))?; + let behavior = BehaviorController::new( + config.behavior.clone(), + Arc::new(owner.embodiment_sink()), + config.limits.control_queue, + config.limits.observable_queue, + config.timeouts.shutdown, + )? + .start(); + let behavior_ingress = behavior.ingress(); let client = Arc::new(LlmClient::new(config.llm.clone(), transport_limits)?); let audit = Arc::new(MemoryPolicyAudit::new(config.limits.observable_queue)?); + let mut tools = perception_policy_tools()?; + tools.extend(behavior_policy_tools(&config.behavior)?); + let routes = tools + .iter() + .map(|tool| tool.definition.name.as_str().to_owned()) + .map(|name| { + let backend: Arc = if name.starts_with("behavior_") { + Arc::new(BehaviorBackend::new(behavior_ingress.clone())) + } else { + perception.clone() + }; + (name, backend) + }) + .collect::>(); let gateway = Arc::new(PolicyGateway::new( config.authorized_avatar_uuids.clone(), - perception_policy_tools()?, + tools, PolicyLimits::default(), audit, )?); @@ -325,15 +362,17 @@ fn start_live_interactions( .duration_since(UNIX_EPOCH) .map_or(0, |duration| duration.as_secs()) }); - let perception_backend: Arc = perception; + let routed_backend: Arc = + Arc::new(AuthorizedBackendRouter::new(routes)?); let responder = Arc::new(PolicyLlmResponder::new( client, gateway, - perception_backend, + routed_backend, loop_limits, now, )?); let sink = Arc::new(owner.interaction_sink()); + let pacer: Arc = Arc::new(behavior_ingress); let interaction = InteractionCoordinator::new( config.interaction.clone(), libremetaverse_types::UUID::zero(), @@ -345,10 +384,12 @@ fn start_live_interactions( config.limits.observable_queue, config.timeouts.shutdown, )? + .with_response_pacer(pacer) .start(); Ok(LiveInteractions { interaction, perception: perception_ingress, perception_observations, + behavior, }) } diff --git a/crates/metacrate-grid-agent/tests/dependency_policy.rs b/crates/metacrate-grid-agent/tests/dependency_policy.rs index 883d5a0..5553852 100644 --- a/crates/metacrate-grid-agent/tests/dependency_policy.rs +++ b/crates/metacrate-grid-agent/tests/dependency_policy.rs @@ -43,10 +43,10 @@ fn package_has_only_reviewed_rust_dependencies_and_no_build_script() { #[test] fn runtime_source_has_no_subprocess_or_native_abi_escape_hatch() { let source = PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("src"); - let mut files = Vec::with_capacity(18); + let mut files = Vec::with_capacity(20); collect_rust_files(&source, &mut files); assert!( - files.len() <= 18, + files.len() <= 20, "source-file count needs a reviewed bound update" ); for path in files { @@ -116,6 +116,37 @@ fn live_session_adapter_reuses_native_lifecycle_and_messaging_managers() { assert!(!backend.contains("reqwest")); } +#[test] +fn embodied_adapter_uses_only_reviewed_high_level_native_movement() { + let root = PathBuf::from(env!("CARGO_MANIFEST_DIR")); + let backend = fs::read_to_string(root.join("src/backend.rs")).expect("backend source"); + for required in [ + ".movement\n .turn_toward(", + ".auto_pilot_local(", + ".auto_pilot_cancel()", + ".get_parcel_local_id(", + ".sit()", + ".stand()", + ] { + assert!( + backend.contains(required), + "missing native embodiment mapping {required}" + ); + } + let behavior = fs::read_to_string(root.join("src/behavior.rs")).expect("behavior source"); + for forbidden in [ + "Teleport(", + "TouchObject(", + "set_agent_controls(", + "FollowAvatar(", + ] { + assert!( + !behavior.contains(forbidden), + "model-facing behavior contains forbidden primitive {forbidden}" + ); + } +} + fn collect_rust_files(directory: &Path, output: &mut Vec) { for entry in fs::read_dir(directory).expect("read source directory") { let path = entry.expect("source entry").path(); diff --git a/docs/grid-agent-behavior.md b/docs/grid-agent-behavior.md new file mode 100644 index 0000000..f092775 --- /dev/null +++ b/docs/grid-agent-behavior.md @@ -0,0 +1,19 @@ +# Grid-agent embodied behavior + +The embodied controller is the only route from an LLM-authorized action to avatar movement. Its public surface is deliberately high level: face a visible avatar or nearby point, make one small look shift, walk a short distance on a heading, stop, sit, stand, and read the current pose. It exposes no control flags, raw packets, follow/wander primitive, unrestricted autopilot, flight, teleport, touch, or arbitrary animation. + +## State and preemption + +The observable modes are `offline`, `settling`, `available`, `engaged`, `executing`, `roaming`, `paused`, and `recovering`. Session readiness enters a bounded settling period. Public-response pacing enters `engaged`; authorized actions enter `executing`; timeout, stuck, and native failures pass through `recovering`. Roaming is only a marker for a separately policy-authorized scheduler task—the controller does not invent routes. Operator pause and the global emergency stop preempt movement and send a stop request. Disconnects and region changes cancel in-flight work and invalidate every old target. + +Every transition records its trigger. Every action record contains generation and region provenance, action name, trigger, policy authorization class or ID, duration, and outcome without including message text or tool arguments. + +## Safety envelope + +- Only authenticated IMs, the authenticated local operator, and authorized scheduler grants may invoke behavior tools. Public speaker attention is a fixed built-in response behavior, not a public tool authorization. +- Face targets must be present in the current avatar cache or be a finite nearby point in the current 256 m region. +- A short walk is limited by configured distance and duration. Before moving, the native adapter verifies that the target's cached 4 m parcel-map cell has the same nonzero parcel ID as the current position. The controller then faces the target, starts one local movement request, polls pose, and always cancels movement on arrival, timeout, stuck detection, pause, emergency stop, disconnect, or region change. +- Actions have a global timeout and minimum interval. Idle behavior is one low-frequency look shift; it never moves continuously and never emits chat. +- Response delay, attention dwell, idle interval, timeouts, distance bounds, and rate limits are validated configuration. Random pacing is injectable so paused-time tests are deterministic. + +The live adapter reuses the single native client/manager graph. Turns use `AgentMovement::turn_toward`; bounded walking uses a local autopilot target that is always paired with `auto_pilot_cancel`; sit and stand use the existing agent manager. Native state reads and actions are generation-fenced before use.