From 08696cd2feb190a58605d4da3265155ac49e5a6f Mon Sep 17 00:00:00 2001 From: Chili Palmer Date: Tue, 18 Aug 2026 11:08:50 +0200 Subject: [PATCH] Implement authorized landmark roaming (#130) --- crates/metacrate-grid-agent/src/backend.rs | 8 + .../src/landmark_tests.rs | 557 ++++++ crates/metacrate-grid-agent/src/landmarks.rs | 1738 +++++++++++++++++ crates/metacrate-grid-agent/src/lib.rs | 15 + crates/metacrate-grid-agent/src/main.rs | 58 +- .../tests/dependency_policy.rs | 4 +- docs/grid-agent-landmarks.md | 40 + 7 files changed, 2413 insertions(+), 7 deletions(-) create mode 100644 crates/metacrate-grid-agent/src/landmark_tests.rs create mode 100644 crates/metacrate-grid-agent/src/landmarks.rs create mode 100644 docs/grid-agent-landmarks.md diff --git a/crates/metacrate-grid-agent/src/backend.rs b/crates/metacrate-grid-agent/src/backend.rs index dcb3f91..50a4df9 100644 --- a/crates/metacrate-grid-agent/src/backend.rs +++ b/crates/metacrate-grid-agent/src/backend.rs @@ -187,6 +187,14 @@ impl LibremetaverseClientOwner { &self.client } + /// Returns the one native agent manager owned by this composition root. + /// `MetaCrate` adapters clone the `Arc`; they never construct a parallel + /// compatibility manager graph. + #[must_use] + pub fn agent(&self) -> Arc { + Arc::clone(&self.agent) + } + /// Creates the supervised live-session adapter without logging in. /// /// # Errors diff --git a/crates/metacrate-grid-agent/src/landmark_tests.rs b/crates/metacrate-grid-agent/src/landmark_tests.rs new file mode 100644 index 0000000..5f24a88 --- /dev/null +++ b/crates/metacrate-grid-agent/src/landmark_tests.rs @@ -0,0 +1,557 @@ +use crate::ProposedToolCall; +use crate::landmarks::*; +use crate::policy::{ + ActionOrigin, MemoryPolicyAudit, PolicyGateway, PolicyLimits, PolicyRequestContext, +}; +use libremetaverse_types::compat::CancellationTokenSource; +use libremetaverse_types::{AssetType, UUID, compat::CancellationToken}; +use std::collections::BTreeSet; +use std::sync::{Arc, Mutex}; +use std::time::{Duration, SystemTime, UNIX_EPOCH}; + +fn id(value: u128) -> UUID { + UUID::new_with_string(format!("{value:032x}")).expect("uuid") +} + +#[derive(Default)] +struct FakeGrid { + accepted: Mutex>, + declined: Mutex>, + teleports: Mutex>, + valid: Mutex, + succeed: Mutex, +} + +impl LandmarkGrid for FakeGrid { + fn accept_offer(&self, offer_id: &str, _: CancellationToken) -> LandmarkFuture<'_, ()> { + let offer_id = offer_id.to_owned(); + Box::pin(async move { + self.accepted.lock().expect("accepted").push(offer_id); + Ok(()) + }) + } + + fn decline_offer(&self, offer_id: &str, _: CancellationToken) -> LandmarkFuture<'_, ()> { + let offer_id = offer_id.to_owned(); + Box::pin(async move { + self.declined.lock().expect("declined").push(offer_id); + Ok(()) + }) + } + + fn verify_landmark( + &self, + _: UUID, + _: UUID, + _: u64, + _: CancellationToken, + ) -> LandmarkFuture<'_, bool> { + Box::pin(async move { Ok(*self.valid.lock().expect("valid")) }) + } + + fn teleport_landmark(&self, asset_id: UUID, _: CancellationToken) -> LandmarkFuture<'_, bool> { + Box::pin(async move { + self.teleports.lock().expect("teleports").push(asset_id); + Ok(*self.succeed.lock().expect("succeed")) + }) + } +} + +#[derive(Debug)] +struct FixedRandom; +impl RoamingRandom for FixedRandom { + fn index(&self, upper_exclusive: usize) -> usize { + upper_exclusive.saturating_sub(1) + } + fn interval_seconds(&self, minimum: u64, _: u64) -> u64 { + minimum + } +} + +fn landmark(inventory: u128, asset: u128, owner: u128, name: &str) -> OfferedInventoryNode { + OfferedInventoryNode { + inventory_id: id(inventory), + owner_id: id(owner), + name: name.into(), + kind: OfferedInventoryKind::Landmark { + asset_id: id(asset), + permissions_fingerprint: 0xCAFE, + }, + children: Vec::new(), + } +} + +fn offer(sender: u128, root: OfferedInventoryNode, offer_id: &str) -> LandmarkOffer { + LandmarkOffer { + offer_id: offer_id.into(), + sender_id: id(sender), + root, + received_unix_millis: 1_000, + } +} + +fn service(grid: Arc, authorized: &[u128]) -> LandmarkService { + *grid.valid.lock().expect("valid") = true; + *grid.succeed.lock().expect("succeed") = true; + LandmarkService::new( + grid, + Arc::new(FixedRandom), + authorized.iter().copied().map(id).collect(), + LandmarkLimits { + teleport_cooldown: Duration::ZERO, + ..LandmarkLimits::default() + }, + None, + ) + .expect("service") +} + +#[tokio::test] +async fn authorized_nested_landmarks_accept_once_and_unauthorized_or_bad_types_decline() { + let grid = Arc::new(FakeGrid::default()); + let service = service(grid.clone(), &[10]); + let nested = OfferedInventoryNode { + inventory_id: id(100), + owner_id: id(10), + name: "Trips".into(), + kind: OfferedInventoryKind::Folder, + children: vec![ + landmark(101, 201, 10, "Alpha"), + landmark(102, 202, 10, "Beta 世界"), + ], + }; + assert_eq!( + service + .ingest_offer( + offer(10, nested.clone(), "offer-1"), + CancellationToken::default() + ) + .await, + Ok(OfferDecision::Accepted) + ); + assert_eq!(service.entries().len(), 2); + assert_eq!( + service + .ingest_offer(offer(10, nested, "offer-1"), CancellationToken::default()) + .await, + Ok(OfferDecision::Duplicate) + ); + assert_eq!(grid.accepted.lock().expect("accepted").len(), 1); + + assert_eq!( + service + .ingest_offer( + offer(11, landmark(103, 203, 11, "No"), "offer-2"), + CancellationToken::default(), + ) + .await, + Err(LandmarkError::Unauthorized) + ); + let bad = OfferedInventoryNode { + inventory_id: id(104), + owner_id: id(10), + name: "Money".into(), + kind: OfferedInventoryKind::Other(AssetType::CallingCard), + children: Vec::new(), + }; + assert_eq!( + service + .ingest_offer(offer(10, bad, "offer-3"), CancellationToken::default()) + .await, + Err(LandmarkError::UnsupportedInventory) + ); + assert_eq!(grid.declined.lock().expect("declined").len(), 2); +} + +#[tokio::test] +async fn cycles_duplicate_assets_and_ambiguous_names_fail_closed() { + let grid = Arc::new(FakeGrid::default()); + let service = service(grid, &[20]); + let duplicate = OfferedInventoryNode { + inventory_id: id(300), + owner_id: id(20), + name: "Dupes".into(), + kind: OfferedInventoryKind::Folder, + children: vec![ + landmark(301, 401, 20, "Same"), + landmark(302, 401, 20, "Same"), + ], + }; + assert_eq!( + service + .ingest_offer(offer(20, duplicate, "dup"), CancellationToken::default()) + .await, + Err(LandmarkError::InvalidOffer) + ); + + let folder = OfferedInventoryNode { + inventory_id: id(310), + owner_id: id(20), + name: "Names".into(), + kind: OfferedInventoryKind::Folder, + children: vec![ + landmark(311, 411, 20, "Same"), + landmark(312, 412, 20, "Same"), + ], + }; + service + .ingest_offer(offer(20, folder, "names"), CancellationToken::default()) + .await + .expect("accepted"); + assert_eq!( + service.select("Same"), + Err(LandmarkError::AmbiguousSelection) + ); + assert!(service.select(&format!("lm-{}", id(311))).is_ok()); +} + +#[tokio::test] +async fn teleport_revalidates_stable_inventory_and_never_uses_message_destination() { + let grid = Arc::new(FakeGrid::default()); + let service = service(grid.clone(), &[30]); + service + .ingest_offer( + offer(30, landmark(501, 601, 30, "Home"), "home"), + CancellationToken::default(), + ) + .await + .expect("offer"); + let receipt = service + .teleport( + id(30), + "Home", + "command-1", + TeleportTrigger::Command, + CancellationToken::default(), + ) + .await + .expect("teleport"); + assert!(receipt.succeeded); + assert_eq!(service.observations().len(), 2); + assert_eq!( + service.observations()[1].kind, + LandmarkObservationKind::TeleportSucceeded + ); + assert_eq!( + grid.teleports.lock().expect("teleports").as_slice(), + &[id(601)] + ); + assert_eq!( + service + .teleport( + id(31), + &format!("lm-{}", id(501)), + "spoof", + TeleportTrigger::Command, + CancellationToken::default(), + ) + .await, + Err(LandmarkError::Unauthorized) + ); + *grid.valid.lock().expect("valid") = false; + assert_eq!( + service + .teleport( + id(30), + "Home", + "stale", + TeleportTrigger::Command, + CancellationToken::default(), + ) + .await, + Err(LandmarkError::StaleLandmark) + ); + assert_eq!(grid.teleports.lock().expect("teleports").len(), 1); +} + +#[tokio::test] +async fn cancellation_failure_and_cooldown_are_bounded_without_duplicate_teleports() { + let grid = Arc::new(FakeGrid::default()); + *grid.valid.lock().expect("valid") = true; + *grid.succeed.lock().expect("succeed") = false; + let service = LandmarkService::new( + grid.clone(), + Arc::new(FixedRandom), + BTreeSet::from([id(35)]), + LandmarkLimits { + teleport_cooldown: Duration::from_secs(30), + ..LandmarkLimits::default() + }, + None, + ) + .expect("service"); + service + .ingest_offer( + offer(35, landmark(551, 651, 35, "Failure"), "failure"), + CancellationToken::default(), + ) + .await + .expect("offer"); + let cancelled = CancellationTokenSource::new(); + cancelled.cancel(); + assert_eq!( + service + .teleport( + id(35), + "Failure", + "cancelled", + TeleportTrigger::Command, + cancelled.token(), + ) + .await, + Err(LandmarkError::Cancelled) + ); + assert!(grid.teleports.lock().expect("teleports").is_empty()); + assert_eq!( + service + .teleport( + id(35), + "Failure", + "failed", + TeleportTrigger::Command, + CancellationToken::default(), + ) + .await, + Err(LandmarkError::TeleportFailed) + ); + // A failed native attempt still starts the cooldown and cannot be replayed. + assert_eq!( + service + .teleport( + id(35), + "Failure", + "duplicate", + TeleportTrigger::Command, + CancellationToken::default(), + ) + .await, + Err(LandmarkError::Cooldown) + ); + assert_eq!(grid.teleports.lock().expect("teleports").len(), 1); +} + +#[tokio::test] +async fn roaming_persists_principal_skips_missed_runs_avoids_repeat_and_obeys_pause() { + let grid = Arc::new(FakeGrid::default()); + let service = service(grid, &[40]); + for (offer_id, inventory, asset, name) in [("a", 701, 801, "A"), ("b", 702, 802, "B")] { + service + .ingest_offer( + offer(40, landmark(inventory, asset, 40, name), offer_id), + CancellationToken::default(), + ) + .await + .expect("offer"); + } + let schedule = service + .upsert_schedule( + id(40), + "daily", + Duration::from_mins(5), + Duration::from_mins(10), + true, + 1_000, + ) + .expect("schedule"); + assert_eq!(schedule.authorizing_avatar, id(40).to_string()); + assert_eq!( + service.due_roam_selection( + "daily", + schedule.next_run_unix_millis, + RoamingPause { + conversation: true, + ..RoamingPause::default() + }, + ), + Err(LandmarkError::Paused) + ); + let first = service + .due_roam_selection( + "daily", + schedule.next_run_unix_millis, + RoamingPause::default(), + ) + .expect("due") + .expect("selection"); + let next = service.schedules()[0].next_run_unix_millis; + assert!(next > schedule.next_run_unix_millis); + let second = service + .due_roam_selection("daily", next, RoamingPause::default()) + .expect("next due") + .expect("next selection"); + assert_ne!(first.1, second.1); + assert_eq!(first.0, id(40)); +} + +#[tokio::test] +async fn catalog_and_schedule_survive_restart_without_replaying_a_missed_run() { + let directory = std::env::temp_dir().join(format!("metacrate-landmarks-{}", id(999))); + std::fs::create_dir_all(&directory).expect("directory"); + let path = directory.join("catalog.json"); + let grid = Arc::new(FakeGrid::default()); + *grid.valid.lock().expect("valid") = true; + *grid.succeed.lock().expect("succeed") = true; + { + let service = LandmarkService::new( + grid.clone(), + Arc::new(FixedRandom), + BTreeSet::from([id(50)]), + LandmarkLimits { + teleport_cooldown: Duration::ZERO, + ..LandmarkLimits::default() + }, + Some(path.clone()), + ) + .expect("service"); + service + .ingest_offer( + offer(50, landmark(901, 902, 50, "Persisted"), "persist"), + CancellationToken::default(), + ) + .await + .expect("offer"); + service + .upsert_schedule( + id(50), + "restart", + Duration::from_mins(5), + Duration::from_mins(5), + true, + 10_000, + ) + .expect("schedule"); + } + let restored = LandmarkService::new( + grid, + Arc::new(FixedRandom), + BTreeSet::from([id(50)]), + LandmarkLimits { + teleport_cooldown: Duration::ZERO, + ..LandmarkLimits::default() + }, + Some(path.clone()), + ) + .expect("restore"); + assert_eq!(restored.entries().len(), 1); + let before = restored.schedules()[0].next_run_unix_millis; + let selection = restored + .due_roam_selection("restart", before + 3_600_000, RoamingPause::default()) + .expect("due") + .expect("selection"); + assert_eq!(selection.0, id(50)); + assert!(restored.schedules()[0].next_run_unix_millis > before + 3_600_000); + let _ = std::fs::remove_file(path); + let _ = std::fs::remove_dir(directory); +} + +#[tokio::test(start_paused = true)] +async fn owned_runner_executes_due_schedule_once_and_cancels_cleanly() { + let grid = Arc::new(FakeGrid::default()); + let service = Arc::new(service(grid.clone(), &[55])); + service + .ingest_offer( + offer(55, landmark(951, 952, 55, "Runner"), "runner-offer"), + CancellationToken::default(), + ) + .await + .expect("offer"); + let now = u64::try_from( + SystemTime::now() + .duration_since(UNIX_EPOCH) + .expect("clock") + .as_millis(), + ) + .expect("millis"); + service + .upsert_schedule( + id(55), + "runner", + Duration::from_mins(5), + Duration::from_mins(5), + true, + now.saturating_sub(301_000), + ) + .expect("schedule"); + let runner = service.start_roaming_runner(RoamingPause::default(), None); + tokio::time::advance(Duration::from_secs(1)).await; + for _ in 0..8 { + tokio::task::yield_now().await; + } + assert_eq!( + grid.teleports.lock().expect("teleports").as_slice(), + &[id(952)] + ); + runner.shutdown().await; + tokio::time::advance(Duration::from_hours(1)).await; + assert_eq!(grid.teleports.lock().expect("teleports").len(), 1); +} + +#[test] +fn unsafe_bounds_and_empty_authorization_fail_before_side_effects() { + let grid = Arc::new(FakeGrid::default()); + let result = LandmarkService::new( + grid, + Arc::new(FixedRandom), + BTreeSet::new(), + LandmarkLimits { + max_folder_depth: 0, + ..LandmarkLimits::default() + }, + None, + ); + assert!(matches!(result, Err(LandmarkError::UnsafeLimits))); +} + +#[test] +fn public_and_unauthorized_im_cannot_select_teleport_or_change_roaming() { + let tools = landmark_policy_tools(LandmarkLimits::default()).expect("tools"); + let gateway = PolicyGateway::new( + BTreeSet::from([id(60)]), + tools, + PolicyLimits::default(), + Arc::new(MemoryPolicyAudit::new(64).expect("audit")), + ) + .expect("gateway"); + for (origin, name, arguments) in [ + ( + ActionOrigin::public_chat(id(60)), + LANDMARK_TELEPORT_TOOL, + serde_json::json!({"selector":"lm-safe"}), + ), + ( + ActionOrigin::instant_message(id(61)), + LANDMARK_SCHEDULE_TOOL, + serde_json::json!({ + "schedule_id":"no", + "minimum_interval_seconds":300, + "maximum_interval_seconds":600, + "enabled":true + }), + ), + ] { + let encoded = arguments.to_string(); + let context = PolicyRequestContext::new(origin, "session", "correlation").expect("context"); + let call = ProposedToolCall::new("call", name, encoded).expect("call"); + let decision = gateway + .evaluate(&context, &call, &arguments, None, 1) + .expect("decision"); + assert!(decision.into_authorization().is_none()); + } + + let arguments = serde_json::json!({"selector":"lm-safe"}); + let context = PolicyRequestContext::new( + ActionOrigin::instant_message(id(60)), + "session", + "authorized", + ) + .expect("context"); + let call = + ProposedToolCall::new("call", LANDMARK_TELEPORT_TOOL, arguments.to_string()).expect("call"); + assert!( + gateway + .evaluate(&context, &call, &arguments, None, 1) + .expect("decision") + .into_authorization() + .is_some() + ); +} diff --git a/crates/metacrate-grid-agent/src/landmarks.rs b/crates/metacrate-grid-agent/src/landmarks.rs new file mode 100644 index 0000000..b7bdd0f --- /dev/null +++ b/crates/metacrate-grid-agent/src/landmarks.rs @@ -0,0 +1,1738 @@ +//! Authorized landmark intake, catalog, teleport, and roaming. + +#![allow(clippy::missing_errors_doc)] + +use crate::backend::{AuthorizedToolBackend, BackendError, BackendFuture}; +use crate::llm::{ToolDefinition, ToolSchema}; +use crate::policy::{ + AllowedOrigins, ApprovalRule, AuthorizedAction, Capability, FixedCost, Idempotency, + OriginClass, PolicyError, PolicyTool, ResourceCost, Risk, +}; +use crate::types::{ + BoundedText, MAX_BODY_BYTES, MAX_IDENTIFIER_BYTES, MAX_OBSERVABLE_DETAIL_BYTES, ToolCallOutcome, +}; +use libremetaverse_types::{AssetType, UUID, compat::CancellationToken}; +use serde::{Deserialize, Serialize}; +use serde_json::json; +use std::collections::{BTreeMap, BTreeSet, VecDeque}; +use std::fmt; +use std::future::Future; +use std::path::{Path, PathBuf}; +use std::pin::Pin; +use std::sync::{Arc, Mutex}; +use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH}; + +pub const MAX_LANDMARK_NAME_BYTES: usize = 256; +pub const LANDMARK_LIST_TOOL: &str = "landmark_catalog_list"; +pub const LANDMARK_TELEPORT_TOOL: &str = "landmark_teleport"; +pub const LANDMARK_SCHEDULE_TOOL: &str = "landmark_roaming_schedule"; +pub const LANDMARK_STATUS_TOOL: &str = "landmark_roaming_status"; + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub struct LandmarkLimits { + pub max_catalog_entries: usize, + pub max_folder_depth: usize, + pub max_folder_nodes: usize, + pub max_asset_fetches: usize, + pub teleport_timeout: Duration, + pub teleport_cooldown: Duration, + pub minimum_roam_interval: Duration, + pub maximum_roam_interval: Duration, +} + +impl Default for LandmarkLimits { + fn default() -> Self { + Self { + max_catalog_entries: 512, + max_folder_depth: 8, + max_folder_nodes: 2_048, + max_asset_fetches: 512, + teleport_timeout: Duration::from_mins(1), + teleport_cooldown: Duration::from_secs(30), + minimum_roam_interval: Duration::from_mins(5), + maximum_roam_interval: Duration::from_hours(24), + } + } +} + +impl LandmarkLimits { + fn valid(self) -> bool { + (1..=4_096).contains(&self.max_catalog_entries) + && (1..=32).contains(&self.max_folder_depth) + && self.max_folder_nodes >= self.max_catalog_entries + && self.max_folder_nodes <= 16_384 + && self.max_asset_fetches <= self.max_folder_nodes + && !self.teleport_timeout.is_zero() + && self.teleport_timeout <= Duration::from_mins(5) + && self.teleport_cooldown <= Duration::from_hours(1) + && self.minimum_roam_interval >= Duration::from_secs(30) + && self.maximum_roam_interval >= self.minimum_roam_interval + && self.maximum_roam_interval <= Duration::from_hours(168) + } +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum LandmarkValidation { + Valid, + Stale, + PermissionChanged, + Failed, +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum LandmarkOutcome { + NeverUsed, + Teleported, + Failed, + Cancelled, +} + +#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)] +pub struct LandmarkEntry { + pub stable_id: String, + pub source_sender: String, + pub inventory_id: String, + pub asset_id: String, + pub permissions_fingerprint: u64, + pub display_name: String, + pub received_unix_millis: u64, + pub validation: LandmarkValidation, + pub last_outcome: LandmarkOutcome, +} + +impl LandmarkEntry { + fn inventory_uuid(&self) -> Result { + UUID::new_with_string(self.inventory_id.clone()).map_err(|_| LandmarkError::CorruptCatalog) + } +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub enum OfferedInventoryKind { + Landmark { + asset_id: UUID, + permissions_fingerprint: u64, + }, + Folder, + Link { + target_id: UUID, + }, + Other(AssetType), +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct OfferedInventoryNode { + pub inventory_id: UUID, + pub owner_id: UUID, + pub name: String, + pub kind: OfferedInventoryKind, + pub children: Vec, +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct LandmarkOffer { + pub offer_id: String, + pub sender_id: UUID, + pub root: OfferedInventoryNode, + pub received_unix_millis: u64, +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum OfferDecision { + Accepted, + Declined, + Duplicate, +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum TeleportTrigger { + Command, + Schedule, +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct TeleportReceipt { + pub correlation_id: String, + pub stable_id: String, + pub trigger: TeleportTrigger, + pub succeeded: bool, +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum LandmarkObservationKind { + Selected, + TeleportSucceeded, + TeleportFailed, +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct LandmarkObservation { + pub correlation_id: String, + pub stable_id: String, + pub trigger: TeleportTrigger, + pub kind: LandmarkObservationKind, +} + +#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)] +pub struct RoamingSchedule { + pub schedule_id: String, + pub authorizing_avatar: String, + pub minimum_interval_seconds: u64, + pub maximum_interval_seconds: u64, + pub enabled: bool, + pub next_run_unix_millis: u64, + pub last_stable_id: Option, + pub expires_unix_millis: u64, + pub remaining_runs: u64, +} + +#[allow(clippy::struct_excessive_bools)] +#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)] +pub struct RoamingPause { + pub conversation: bool, + pub operator: bool, + pub degraded: bool, + pub build: bool, + pub visual_capture: bool, +} + +impl RoamingPause { + #[must_use] + pub const fn active(self) -> bool { + self.conversation || self.operator || self.degraded || self.build || self.visual_capture + } +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub enum LandmarkError { + UnsafeLimits, + Unauthorized, + InvalidOffer, + DuplicateOffer, + UnsupportedInventory, + TraversalLimit, + CatalogFull, + AmbiguousSelection, + UnknownSelection, + StaleLandmark, + Busy, + Cooldown, + Cancelled, + TeleportFailed, + Persistence, + CorruptCatalog, + InvalidSchedule, + Paused, +} + +impl fmt::Display for LandmarkError { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str(match self { + Self::UnsafeLimits => "unsafe landmark workflow limits", + Self::Unauthorized => "landmark workflow requires an authorized principal", + Self::InvalidOffer => "inventory offer is invalid or stale", + Self::DuplicateOffer => "inventory offer was already handled", + Self::UnsupportedInventory => "offer contains non-landmark inventory", + Self::TraversalLimit => "landmark folder traversal limit reached", + Self::CatalogFull => "landmark catalog is full", + Self::AmbiguousSelection => "landmark name is ambiguous; select by stable ID", + Self::UnknownSelection => "landmark selection was not found", + Self::StaleLandmark => "landmark changed after catalog validation", + Self::Busy => "another teleport is active", + Self::Cooldown => "teleport cooldown is active", + Self::Cancelled => "landmark operation was cancelled", + Self::TeleportFailed => "landmark teleport failed", + Self::Persistence => "landmark catalog persistence failed", + Self::CorruptCatalog => "landmark catalog is corrupt", + Self::InvalidSchedule => "roaming schedule is invalid", + Self::Paused => "roaming is paused", + }) + } +} + +impl std::error::Error for LandmarkError {} + +pub type LandmarkFuture<'a, T> = + Pin> + Send + 'a>>; + +pub trait LandmarkGrid: Send + Sync + 'static { + fn accept_offer( + &self, + offer_id: &str, + cancellation: CancellationToken, + ) -> LandmarkFuture<'_, ()>; + fn decline_offer( + &self, + offer_id: &str, + cancellation: CancellationToken, + ) -> LandmarkFuture<'_, ()>; + fn verify_landmark( + &self, + inventory_id: UUID, + asset_id: UUID, + permissions_fingerprint: u64, + cancellation: CancellationToken, + ) -> LandmarkFuture<'_, bool>; + fn teleport_landmark( + &self, + asset_id: UUID, + cancellation: CancellationToken, + ) -> LandmarkFuture<'_, bool>; +} + +pub trait RoamingRandom: Send + Sync + 'static { + fn index(&self, upper_exclusive: usize) -> usize; + fn interval_seconds(&self, minimum: u64, maximum: u64) -> u64; +} + +#[derive(Debug)] +pub struct SystemRoamingRandom(Mutex); + +impl Default for SystemRoamingRandom { + fn default() -> Self { + let duration = SystemTime::now() + .duration_since(UNIX_EPOCH) + .unwrap_or_default(); + let seed = duration.as_secs() ^ u64::from(duration.subsec_nanos()); + Self(Mutex::new(seed ^ 0x9e37_79b9_7f4a_7c15)) + } +} + +impl RoamingRandom for SystemRoamingRandom { + fn index(&self, upper_exclusive: usize) -> usize { + if upper_exclusive == 0 { + return 0; + } + usize::try_from(self.next() % upper_exclusive as u64).unwrap_or(0) + } + + fn interval_seconds(&self, minimum: u64, maximum: u64) -> u64 { + minimum.saturating_add(self.next() % maximum.saturating_sub(minimum).saturating_add(1)) + } +} + +impl SystemRoamingRandom { + fn next(&self) -> u64 { + let mut state = self + .0 + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + *state ^= *state << 13; + *state ^= *state >> 7; + *state ^= *state << 17; + *state + } +} + +#[derive(Default, Serialize, Deserialize)] +struct PersistedLandmarks { + version: u32, + entries: Vec, + schedules: Vec, +} + +struct LandmarkState { + entries: BTreeMap, + schedules: BTreeMap, + handled_offers: BTreeSet, + last_teleport: Option, +} + +pub struct LandmarkService { + grid: Arc, + random: Arc, + authorized: BTreeSet, + limits: LandmarkLimits, + path: Option, + state: Mutex, + observations: Mutex>, + teleport: Arc, +} + +pub struct LandmarkRoamingHandle { + pause: tokio::sync::watch::Sender, + cancellation: libremetaverse_types::compat::CancellationTokenSource, + worker: Option>, +} + +impl fmt::Debug for LandmarkRoamingHandle { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("LandmarkRoamingHandle") + .field("pause", &*self.pause.borrow()) + .field( + "worker_finished", + &self + .worker + .as_ref() + .is_none_or(tokio::task::JoinHandle::is_finished), + ) + .finish_non_exhaustive() + } +} + +impl LandmarkRoamingHandle { + pub fn set_pause(&self, pause: RoamingPause) { + self.pause.send_replace(pause); + } + + pub fn update_pause(&self, update: impl FnOnce(&mut RoamingPause)) { + self.pause.send_modify(update); + } + + pub async fn shutdown(mut self) { + self.cancellation.cancel(); + if let Some(worker) = self.worker.take() { + let _ = worker.await; + } + } +} + +impl Drop for LandmarkRoamingHandle { + fn drop(&mut self) { + self.cancellation.cancel(); + if let Some(worker) = self.worker.take() { + worker.abort(); + } + } +} + +impl fmt::Debug for LandmarkService { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("LandmarkService") + .field("limits", &self.limits) + .field("persistent", &self.path.is_some()) + .finish_non_exhaustive() + } +} + +impl LandmarkService { + pub fn new( + grid: Arc, + random: Arc, + authorized: BTreeSet, + limits: LandmarkLimits, + path: Option, + ) -> Result { + if !limits.valid() || authorized.iter().any(|id| *id == UUID::zero()) { + return Err(LandmarkError::UnsafeLimits); + } + let persisted = path + .as_deref() + .map(load_catalog) + .transpose()? + .unwrap_or_default(); + if persisted.entries.len() > limits.max_catalog_entries { + return Err(LandmarkError::CorruptCatalog); + } + Ok(Self { + grid, + random, + authorized, + limits, + path, + state: Mutex::new(LandmarkState { + entries: persisted + .entries + .into_iter() + .map(|entry| (entry.stable_id.clone(), entry)) + .collect(), + schedules: persisted + .schedules + .into_iter() + .map(|schedule| (schedule.schedule_id.clone(), schedule)) + .collect(), + handled_offers: BTreeSet::new(), + last_teleport: None, + }), + observations: Mutex::new(VecDeque::new()), + teleport: Arc::new(tokio::sync::Semaphore::new(1)), + }) + } + + pub fn entries(&self) -> Vec { + lock(&self.state).entries.values().cloned().collect() + } + + pub fn schedules(&self) -> Vec { + lock(&self.state).schedules.values().cloned().collect() + } + + pub fn observations(&self) -> Vec { + lock(&self.observations).iter().cloned().collect() + } + + #[must_use] + pub fn is_authorized(&self, avatar_id: UUID) -> bool { + self.authorized.contains(&avatar_id) + } + + pub fn start_roaming_runner( + self: &Arc, + initial_pause: RoamingPause, + behavior: Option, + ) -> LandmarkRoamingHandle { + let (pause, pause_rx) = tokio::sync::watch::channel(initial_pause); + let cancellation = libremetaverse_types::compat::CancellationTokenSource::new(); + let token = cancellation.token(); + let service = Arc::clone(self); + let worker = tokio::spawn(async move { + let mut interval = tokio::time::interval(Duration::from_secs(1)); + interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); + loop { + tokio::select! { + () = token.cancelled() => break, + _ = interval.tick() => {} + } + let pause = *pause_rx.borrow(); + for schedule in service.schedules() { + if token.is_cancellation_requested() { + break; + } + if let Some(behavior) = &behavior { + let _ = behavior.set_roaming(true); + } + let _ = service + .execute_due_roaming( + &schedule.schedule_id, + unix_millis(), + pause, + token.clone(), + ) + .await; + if let Some(behavior) = &behavior { + let _ = behavior.set_roaming(false); + } + } + } + }); + LandmarkRoamingHandle { + pause, + cancellation, + worker: Some(worker), + } + } + + pub async fn ingest_offer( + &self, + offer: LandmarkOffer, + cancellation: CancellationToken, + ) -> Result { + if !self.authorized.contains(&offer.sender_id) { + self.grid + .decline_offer(&offer.offer_id, cancellation) + .await?; + return Err(LandmarkError::Unauthorized); + } + validate_offer_id(&offer.offer_id)?; + { + let mut state = lock(&self.state); + if !state.handled_offers.insert(offer.offer_id.clone()) { + return Ok(OfferDecision::Duplicate); + } + } + let entries = match flatten_offer(&offer, self.limits) { + Ok(entries) => entries, + Err(error) => { + self.grid + .decline_offer(&offer.offer_id, cancellation) + .await?; + return Err(error); + } + }; + if cancellation.is_cancellation_requested() { + return Err(LandmarkError::Cancelled); + } + let catalog_full = { + let state = lock(&self.state); + let additions = entries + .iter() + .filter(|entry| !state.entries.contains_key(&entry.stable_id)) + .count(); + state.entries.len().saturating_add(additions) > self.limits.max_catalog_entries + }; + if catalog_full { + self.grid + .decline_offer(&offer.offer_id, cancellation) + .await?; + return Err(LandmarkError::CatalogFull); + } + self.grid + .accept_offer(&offer.offer_id, cancellation.clone()) + .await?; + self.commit_entries(entries)?; + Ok(OfferDecision::Accepted) + } + + /// Catalogs an offer that the live inventory event already accepted into + /// quarantine. This second phase validates the received server objects; + /// invalid folders remain untouched in quarantine and never enter catalog. + pub fn catalog_received_offer(&self, offer: &LandmarkOffer) -> Result<(), LandmarkError> { + if !self.authorized.contains(&offer.sender_id) { + return Err(LandmarkError::Unauthorized); + } + validate_offer_id(&offer.offer_id)?; + let entries = flatten_offer(offer, self.limits)?; + let catalog_full = { + let state = lock(&self.state); + let additions = entries + .iter() + .filter(|entry| !state.entries.contains_key(&entry.stable_id)) + .count(); + state.entries.len().saturating_add(additions) > self.limits.max_catalog_entries + }; + if catalog_full { + return Err(LandmarkError::CatalogFull); + } + self.commit_entries(entries) + } + + fn commit_entries(&self, entries: Vec) -> Result<(), LandmarkError> { + { + let mut state = lock(&self.state); + for entry in entries { + state + .entries + .entry(entry.stable_id.clone()) + .or_insert(entry); + } + self.persist_locked(&state)?; + } + Ok(()) + } + + pub fn select(&self, selector: &str) -> Result { + let state = lock(&self.state); + if let Some(entry) = state.entries.get(selector) { + return Ok(entry.clone()); + } + let folded = selector.to_lowercase(); + let mut matches = state + .entries + .values() + .filter(|entry| entry.display_name.to_lowercase() == folded); + let first = matches + .next() + .cloned() + .ok_or(LandmarkError::UnknownSelection)?; + if matches.next().is_some() { + Err(LandmarkError::AmbiguousSelection) + } else { + Ok(first) + } + } + + #[allow(clippy::too_many_lines)] + pub async fn teleport( + &self, + requesting_avatar: UUID, + selector: &str, + correlation_id: &str, + trigger: TeleportTrigger, + cancellation: CancellationToken, + ) -> Result { + if !self.authorized.contains(&requesting_avatar) { + return Err(LandmarkError::Unauthorized); + } + validate_offer_id(correlation_id)?; + let entry = self.select(selector)?; + self.observe(LandmarkObservation { + correlation_id: correlation_id.to_owned(), + stable_id: entry.stable_id.clone(), + trigger, + kind: LandmarkObservationKind::Selected, + }); + { + let state = lock(&self.state); + if state + .last_teleport + .is_some_and(|last| last.elapsed() < self.limits.teleport_cooldown) + { + return Err(LandmarkError::Cooldown); + } + } + let _permit = self + .teleport + .clone() + .try_acquire_owned() + .map_err(|_| LandmarkError::Busy)?; + if cancellation.is_cancellation_requested() { + self.observe_teleport(&entry.stable_id, correlation_id, trigger, false); + return Err(LandmarkError::Cancelled); + } + let inventory_id = entry.inventory_uuid()?; + let asset_id = UUID::new_with_string(entry.asset_id.clone()) + .map_err(|_| LandmarkError::CorruptCatalog)?; + if !self + .grid + .verify_landmark( + inventory_id, + asset_id, + entry.permissions_fingerprint, + cancellation.clone(), + ) + .await? + { + self.mark_entry( + &entry.stable_id, + LandmarkValidation::Stale, + LandmarkOutcome::Failed, + )?; + self.observe_teleport(&entry.stable_id, correlation_id, trigger, false); + return Err(LandmarkError::StaleLandmark); + } + if cancellation.is_cancellation_requested() { + self.mark_entry( + &entry.stable_id, + LandmarkValidation::Valid, + LandmarkOutcome::Cancelled, + )?; + self.observe_teleport(&entry.stable_id, correlation_id, trigger, false); + return Err(LandmarkError::Cancelled); + } + let operation = self.grid.teleport_landmark(asset_id, cancellation.clone()); + let succeeded = match tokio::time::timeout(self.limits.teleport_timeout, operation).await { + Ok(Ok(value)) => value, + Ok(Err(error)) => { + self.mark_entry( + &entry.stable_id, + LandmarkValidation::Valid, + LandmarkOutcome::Failed, + )?; + self.observe_teleport(&entry.stable_id, correlation_id, trigger, false); + return Err(error); + } + Err(_) => { + self.mark_entry( + &entry.stable_id, + LandmarkValidation::Valid, + LandmarkOutcome::Failed, + )?; + self.observe_teleport(&entry.stable_id, correlation_id, trigger, false); + return Err(LandmarkError::TeleportFailed); + } + }; + let outcome = if succeeded { + LandmarkOutcome::Teleported + } else { + LandmarkOutcome::Failed + }; + self.mark_entry(&entry.stable_id, LandmarkValidation::Valid, outcome)?; + lock(&self.state).last_teleport = Some(Instant::now()); + if !succeeded { + self.observe_teleport(&entry.stable_id, correlation_id, trigger, false); + return Err(LandmarkError::TeleportFailed); + } + self.observe_teleport(&entry.stable_id, correlation_id, trigger, true); + Ok(TeleportReceipt { + correlation_id: correlation_id.to_owned(), + stable_id: entry.stable_id, + trigger, + succeeded, + }) + } + + pub fn upsert_schedule( + &self, + authorizing_avatar: UUID, + schedule_id: &str, + minimum: Duration, + maximum: Duration, + enabled: bool, + now_unix_millis: u64, + ) -> Result { + if !self.authorized.contains(&authorizing_avatar) + || minimum < self.limits.minimum_roam_interval + || maximum > self.limits.maximum_roam_interval + || maximum < minimum + { + return Err(LandmarkError::InvalidSchedule); + } + validate_offer_id(schedule_id)?; + let delay = self + .random + .interval_seconds(minimum.as_secs(), maximum.as_secs()); + let schedule = RoamingSchedule { + schedule_id: schedule_id.to_owned(), + authorizing_avatar: authorizing_avatar.to_string(), + minimum_interval_seconds: minimum.as_secs(), + maximum_interval_seconds: maximum.as_secs(), + enabled, + next_run_unix_millis: now_unix_millis.saturating_add(delay.saturating_mul(1_000)), + last_stable_id: None, + expires_unix_millis: now_unix_millis.saturating_add(86_400_000), + remaining_runs: 86_400_u64 + .checked_div(minimum.as_secs()) + .unwrap_or(0) + .clamp(1, 1_024), + }; + let mut state = lock(&self.state); + state + .schedules + .insert(schedule.schedule_id.clone(), schedule.clone()); + self.persist_locked(&state)?; + Ok(schedule) + } + + pub fn disable_schedule(&self, schedule_id: &str) -> Result<(), LandmarkError> { + let mut state = lock(&self.state); + let schedule = state + .schedules + .get_mut(schedule_id) + .ok_or(LandmarkError::InvalidSchedule)?; + schedule.enabled = false; + self.persist_locked(&state) + } + + pub fn due_roam_selection( + &self, + schedule_id: &str, + now_unix_millis: u64, + pause: RoamingPause, + ) -> Result, LandmarkError> { + if pause.active() { + return Err(LandmarkError::Paused); + } + let mut state = lock(&self.state); + let schedule = state + .schedules + .get(schedule_id) + .cloned() + .ok_or(LandmarkError::InvalidSchedule)?; + if !schedule.enabled || now_unix_millis < schedule.next_run_unix_millis { + return Ok(None); + } + if now_unix_millis >= schedule.expires_unix_millis || schedule.remaining_runs == 0 { + if let Some(expired) = state.schedules.get_mut(schedule_id) { + expired.enabled = false; + } + self.persist_locked(&state)?; + return Ok(None); + } + let mut candidates = state.entries.keys().cloned().collect::>(); + if candidates.len() > 1 + && let Some(last) = &schedule.last_stable_id + { + candidates.retain(|candidate| candidate != last); + } + if candidates.is_empty() { + return Ok(None); + } + let choice = candidates[self.random.index(candidates.len())].clone(); + let delay = self.random.interval_seconds( + schedule.minimum_interval_seconds, + schedule.maximum_interval_seconds, + ); + let schedule = state + .schedules + .get_mut(schedule_id) + .ok_or(LandmarkError::InvalidSchedule)?; + schedule.last_stable_id = Some(choice.clone()); + schedule.remaining_runs = schedule.remaining_runs.saturating_sub(1); + // Reschedule from now. Missed runs are skipped, never caught up. + schedule.next_run_unix_millis = now_unix_millis.saturating_add(delay.saturating_mul(1_000)); + let principal = UUID::new_with_string(schedule.authorizing_avatar.clone()) + .map_err(|_| LandmarkError::CorruptCatalog)?; + let correlation = format!("roam-{schedule_id}-{now_unix_millis}"); + self.persist_locked(&state)?; + Ok(Some((principal, choice, correlation))) + } + + pub async fn execute_due_roaming( + &self, + schedule_id: &str, + now_unix_millis: u64, + pause: RoamingPause, + cancellation: CancellationToken, + ) -> Result, LandmarkError> { + let Some((principal, selection, correlation)) = + self.due_roam_selection(schedule_id, now_unix_millis, pause)? + else { + return Ok(None); + }; + self.teleport( + principal, + &selection, + &correlation, + TeleportTrigger::Schedule, + cancellation, + ) + .await + .map(Some) + } + + fn mark_entry( + &self, + stable_id: &str, + validation: LandmarkValidation, + outcome: LandmarkOutcome, + ) -> Result<(), LandmarkError> { + let mut state = lock(&self.state); + let entry = state + .entries + .get_mut(stable_id) + .ok_or(LandmarkError::UnknownSelection)?; + entry.validation = validation; + entry.last_outcome = outcome; + self.persist_locked(&state) + } + + fn observe_teleport( + &self, + stable_id: &str, + correlation_id: &str, + trigger: TeleportTrigger, + succeeded: bool, + ) { + self.observe(LandmarkObservation { + correlation_id: correlation_id.to_owned(), + stable_id: stable_id.to_owned(), + trigger, + kind: if succeeded { + LandmarkObservationKind::TeleportSucceeded + } else { + LandmarkObservationKind::TeleportFailed + }, + }); + } + + fn observe(&self, observation: LandmarkObservation) { + let mut observations = lock(&self.observations); + if observations.len() == 1_024 { + observations.pop_front(); + } + observations.push_back(observation); + } + + fn persist_locked(&self, state: &LandmarkState) -> Result<(), LandmarkError> { + let Some(path) = &self.path else { + return Ok(()); + }; + persist_catalog( + path, + &PersistedLandmarks { + version: 1, + entries: state.entries.values().cloned().collect(), + schedules: state.schedules.values().cloned().collect(), + }, + ) + } +} + +fn flatten_offer( + offer: &LandmarkOffer, + limits: LandmarkLimits, +) -> Result, LandmarkError> { + let mut queue = VecDeque::from([(&offer.root, 0_usize)]); + let mut visited = BTreeSet::new(); + let mut assets = BTreeSet::new(); + let mut entries = Vec::new(); + while let Some((node, depth)) = queue.pop_front() { + if depth > limits.max_folder_depth || visited.len() >= limits.max_folder_nodes { + return Err(LandmarkError::TraversalLimit); + } + if node.inventory_id == UUID::zero() + || node.owner_id == UUID::zero() + || !visited.insert(node.inventory_id) + { + return Err(LandmarkError::InvalidOffer); + } + if node.name.is_empty() + || node.name.len() > MAX_LANDMARK_NAME_BYTES + || node.name.chars().any(char::is_control) + { + return Err(LandmarkError::InvalidOffer); + } + match node.kind { + OfferedInventoryKind::Landmark { + asset_id, + permissions_fingerprint, + } => { + if asset_id == UUID::zero() || !assets.insert(asset_id) { + return Err(LandmarkError::InvalidOffer); + } + if assets.len() > limits.max_asset_fetches { + return Err(LandmarkError::TraversalLimit); + } + entries.push(LandmarkEntry { + stable_id: format!("lm-{}", node.inventory_id), + source_sender: offer.sender_id.to_string(), + inventory_id: node.inventory_id.to_string(), + asset_id: asset_id.to_string(), + permissions_fingerprint, + display_name: node.name.clone(), + received_unix_millis: offer.received_unix_millis, + validation: LandmarkValidation::Valid, + last_outcome: LandmarkOutcome::NeverUsed, + }); + } + OfferedInventoryKind::Folder => { + if node.children.is_empty() { + return Err(LandmarkError::InvalidOffer); + } + queue.extend(node.children.iter().map(|child| (child, depth + 1))); + } + OfferedInventoryKind::Link { .. } | OfferedInventoryKind::Other(_) => { + return Err(LandmarkError::UnsupportedInventory); + } + } + } + if entries.is_empty() { + Err(LandmarkError::UnsupportedInventory) + } else { + Ok(entries) + } +} + +fn validate_offer_id(value: &str) -> Result<(), LandmarkError> { + if value.is_empty() + || value.len() > 128 + || !value + .bytes() + .all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'-' | b'_' | b'.')) + { + Err(LandmarkError::InvalidOffer) + } else { + Ok(()) + } +} + +fn lock(mutex: &Mutex) -> std::sync::MutexGuard<'_, T> { + mutex + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) +} + +fn load_catalog(path: &Path) -> Result { + if !path.exists() { + return Ok(PersistedLandmarks::default()); + } + let bytes = std::fs::read(path).map_err(|_| LandmarkError::Persistence)?; + if bytes.len() > 4 * 1024 * 1024 { + return Err(LandmarkError::CorruptCatalog); + } + let persisted: PersistedLandmarks = + serde_json::from_slice(&bytes).map_err(|_| LandmarkError::CorruptCatalog)?; + if persisted.version != 1 { + return Err(LandmarkError::CorruptCatalog); + } + Ok(persisted) +} + +fn persist_catalog(path: &Path, value: &PersistedLandmarks) -> Result<(), LandmarkError> { + let parent = path.parent().ok_or(LandmarkError::Persistence)?; + std::fs::create_dir_all(parent).map_err(|_| LandmarkError::Persistence)?; + let bytes = serde_json::to_vec(value).map_err(|_| LandmarkError::Persistence)?; + let temporary = path.with_extension("tmp"); + std::fs::write(&temporary, bytes).map_err(|_| LandmarkError::Persistence)?; + #[cfg(unix)] + { + use std::os::unix::fs::PermissionsExt; + std::fs::set_permissions(&temporary, std::fs::Permissions::from_mode(0o600)) + .map_err(|_| LandmarkError::Persistence)?; + } + std::fs::rename(temporary, path).map_err(|_| LandmarkError::Persistence) +} + +#[derive(Deserialize)] +#[serde(deny_unknown_fields)] +struct SelectorArguments { + selector: String, +} + +#[derive(Deserialize)] +#[serde(deny_unknown_fields)] +struct ScheduleArguments { + schedule_id: String, + minimum_interval_seconds: u64, + maximum_interval_seconds: u64, + enabled: bool, +} + +/// Policy-authorized tool adapter. It can only use the principal sealed into +/// [`AuthorizedAction`]; tool JSON never carries an avatar or destination UUID. +#[derive(Debug)] +pub struct LandmarkToolBackend { + service: Arc, + behavior: Option, +} + +impl LandmarkToolBackend { + #[must_use] + pub const fn new(service: Arc) -> Self { + Self { + service, + behavior: None, + } + } + + #[must_use] + pub fn with_behavior(mut self, behavior: crate::behavior::BehaviorIngress) -> Self { + self.behavior = Some(behavior); + self + } +} + +impl AuthorizedToolBackend for LandmarkToolBackend { + #[allow(clippy::too_many_lines)] + fn apply( + &self, + action: AuthorizedAction, + cancellation: CancellationToken, + ) -> BackendFuture<'_, Result> { + Box::pin(async move { + let call_id = action.call().call_id.clone(); + let principal = action.authenticated_avatar_id(); + let result: Result = async { + match action.call().name.as_str() { + LANDMARK_LIST_TOOL => { + let empty: BTreeMap = + serde_json::from_str(action.call().arguments_json.as_str()) + .map_err(|_| LandmarkError::InvalidOffer)?; + if empty.is_empty() { + let entries = self + .service + .entries() + .into_iter() + .map(|entry| { + json!({ + "stable_id": entry.stable_id, + "name": crate::perception::sanitize_untrusted(&entry.display_name), + "validation": entry.validation, + "last_outcome": entry.last_outcome, + }) + }) + .collect::>(); + serde_json::to_string(&json!({ + "trust": "untrusted_inventory_data", + "landmarks": entries + })) + .map_err(|_| LandmarkError::CorruptCatalog) + } else { + Err(LandmarkError::InvalidOffer) + } + } + LANDMARK_STATUS_TOOL => { + let empty: BTreeMap = + serde_json::from_str(action.call().arguments_json.as_str()) + .map_err(|_| LandmarkError::InvalidOffer)?; + if empty.is_empty() { + serde_json::to_string(&json!({"schedules": self.service.schedules()})) + .map_err(|_| LandmarkError::CorruptCatalog) + } else { + Err(LandmarkError::InvalidOffer) + } + } + LANDMARK_TELEPORT_TOOL => { + let principal = principal.ok_or(LandmarkError::Unauthorized)?; + let arguments: SelectorArguments = + serde_json::from_str(action.call().arguments_json.as_str()) + .map_err(|_| LandmarkError::InvalidOffer)?; + if let Some(behavior) = &self.behavior { + behavior + .set_roaming(true) + .map_err(|_| LandmarkError::Busy)?; + } + let result = self + .service + .teleport( + principal, + &arguments.selector, + call_id.as_str(), + TeleportTrigger::Command, + cancellation, + ) + .await; + if let Some(behavior) = &self.behavior { + let _ = behavior.set_roaming(false); + } + result.and_then(|receipt| { + serde_json::to_string(&json!({ + "stable_id": receipt.stable_id, + "succeeded": receipt.succeeded, + })) + .map_err(|_| LandmarkError::CorruptCatalog) + }) + } + LANDMARK_SCHEDULE_TOOL => { + let principal = principal.ok_or(LandmarkError::Unauthorized)?; + let arguments: ScheduleArguments = + serde_json::from_str(action.call().arguments_json.as_str()) + .map_err(|_| LandmarkError::InvalidSchedule)?; + self.service + .upsert_schedule( + principal, + &arguments.schedule_id, + Duration::from_secs(arguments.minimum_interval_seconds), + Duration::from_secs(arguments.maximum_interval_seconds), + arguments.enabled, + unix_millis(), + ) + .and_then(|schedule| { + serde_json::to_string(&json!({ + "schedule_id": schedule.schedule_id, + "enabled": schedule.enabled, + "next_run_unix_millis": schedule.next_run_unix_millis, + })) + .map_err(|_| LandmarkError::Persistence) + }) + } + _ => Err(LandmarkError::InvalidOffer), + } + } + .await; + Ok(match result { + Ok(value) => ToolCallOutcome::Completed { + call_id, + result: BoundedText::::new("landmarks.result", value).map_err( + |_| BackendError::Operation { + operation: "bounded landmark result", + }, + )?, + }, + Err(error) => ToolCallOutcome::Rejected { + call_id, + reason: BoundedText::::new( + "landmarks.rejection", + error.to_string(), + ) + .map_err(|_| BackendError::Operation { + operation: "bounded landmark rejection", + })?, + }, + }) + }) + } +} + +#[allow(clippy::too_many_lines)] +pub fn landmark_policy_tools(limits: LandmarkLimits) -> Result, PolicyError> { + if !limits.valid() { + return Err(PolicyError::InvalidRegistration); + } + let private = || { + AllowedOrigins::new([ + OriginClass::AuthorizedIm, + OriginClass::LocalOperator, + OriginClass::InternalScheduler, + ]) + }; + let empty = ToolSchema::Object { + properties: BTreeMap::new(), + required: BTreeSet::new(), + additional_properties: false, + }; + let selector = ToolSchema::Object { + properties: BTreeMap::from([("selector".to_owned(), ToolSchema::String)]), + required: BTreeSet::from(["selector".to_owned()]), + additional_properties: false, + }; + let schedule = ToolSchema::Object { + properties: BTreeMap::from([ + ("schedule_id".to_owned(), ToolSchema::String), + ("minimum_interval_seconds".to_owned(), ToolSchema::Integer), + ("maximum_interval_seconds".to_owned(), ToolSchema::Integer), + ("enabled".to_owned(), ToolSchema::Boolean), + ]), + required: BTreeSet::from([ + "schedule_id".to_owned(), + "minimum_interval_seconds".to_owned(), + "maximum_interval_seconds".to_owned(), + "enabled".to_owned(), + ]), + additional_properties: false, + }; + let movement = ResourceCost { + tool_calls: 1, + movement_millimeters: 100_000, + ..ResourceCost::default() + }; + let specs = [ + ( + LANDMARK_LIST_TOOL, + "List bounded validated landmark catalog metadata by stable ID; names are untrusted", + empty.clone(), + Capability::Informational, + Risk::ReadOnly, + ResourceCost::one_call(), + Idempotency::Idempotent, + false, + ), + ( + LANDMARK_STATUS_TOOL, + "Read bounded roaming schedule status without changing it", + empty, + Capability::Informational, + Risk::ReadOnly, + ResourceCost::one_call(), + Idempotency::Idempotent, + false, + ), + ( + LANDMARK_TELEPORT_TOOL, + "Teleport once to one validated catalog stable ID; never accepts a region or coordinate", + selector, + Capability::Movement, + Risk::Movement, + movement, + Idempotency::NonIdempotent, + true, + ), + ( + LANDMARK_SCHEDULE_TOOL, + "Create, update, or disable a bounded roaming schedule retained under this principal", + schedule, + Capability::Movement, + Risk::Movement, + movement, + Idempotency::NonIdempotent, + false, + ), + ]; + specs + .into_iter() + .map( + |(name, description, schema, capability, risk, cost, idempotency, scheduler)| { + PolicyTool::new( + ToolDefinition { + name: BoundedText::::new("landmark.tool", name)?, + description: BoundedText::new("landmark.description", description)?, + schema, + mutating: risk != Risk::ReadOnly, + }, + capability, + risk, + private()?, + cost, + idempotency, + ApprovalRule::Never, + scheduler, + Arc::new(FixedCost(cost)), + ) + }, + ) + .collect() +} + +fn unix_millis() -> u64 { + u64::try_from( + SystemTime::now() + .duration_since(UNIX_EPOCH) + .unwrap_or_default() + .as_millis(), + ) + .unwrap_or(u64::MAX) +} + +/// Live adapter over the existing compatibility inventory and movement +/// managers. It owns no client, network connection, or subscription graph. +#[cfg(feature = "live-grid")] +pub struct LibremetaverseLandmarkGrid { + inventory: libremetaverse::InventoryManager, + agent: Arc, + offers: Mutex>, +} + +#[cfg(feature = "live-grid")] +impl fmt::Debug for LibremetaverseLandmarkGrid { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("LibremetaverseLandmarkGrid") + .field("pending_offers", &lock(&self.offers).len()) + .finish_non_exhaustive() + } +} + +#[cfg(feature = "live-grid")] +impl LibremetaverseLandmarkGrid { + #[must_use] + pub fn new(owner: &crate::backend::LibremetaverseClientOwner) -> Self { + Self { + inventory: owner.client().inventory(), + agent: owner.agent(), + offers: Mutex::new(BTreeMap::new()), + } + } + + /// Registers the shared offer decision object delivered by the original + /// inventory event. The correlation ID contains no item name or message. + pub fn register_offer( + &self, + correlation_id: String, + offer: libremetaverse::InventoryObjectOfferedEventArgs, + ) -> Result<(), LandmarkError> { + validate_offer_id(&correlation_id)?; + let mut offers = lock(&self.offers); + if offers.len() >= 1_024 || offers.insert(correlation_id, offer).is_some() { + return Err(LandmarkError::DuplicateOffer); + } + Ok(()) + } +} + +#[cfg(feature = "live-grid")] +impl LandmarkGrid for LibremetaverseLandmarkGrid { + fn accept_offer( + &self, + offer_id: &str, + cancellation: CancellationToken, + ) -> LandmarkFuture<'_, ()> { + let offer_id = offer_id.to_owned(); + Box::pin(async move { + if cancellation.is_cancellation_requested() { + return Err(LandmarkError::Cancelled); + } + let mut offer = lock(&self.offers) + .remove(&offer_id) + .ok_or(LandmarkError::InvalidOffer)?; + offer.set_accept(true); + Ok(()) + }) + } + + fn decline_offer( + &self, + offer_id: &str, + cancellation: CancellationToken, + ) -> LandmarkFuture<'_, ()> { + let offer_id = offer_id.to_owned(); + Box::pin(async move { + if cancellation.is_cancellation_requested() { + return Err(LandmarkError::Cancelled); + } + let mut offer = lock(&self.offers) + .remove(&offer_id) + .ok_or(LandmarkError::InvalidOffer)?; + offer.set_accept(false); + Ok(()) + }) + } + + fn verify_landmark( + &self, + inventory_id: UUID, + asset_id: UUID, + permissions_fingerprint: u64, + cancellation: CancellationToken, + ) -> LandmarkFuture<'_, bool> { + Box::pin(async move { + let item = self + .inventory + .fetch_item(inventory_id, self.agent.agent_id(), Some(cancellation)) + .await + .map_err(|_| LandmarkError::StaleLandmark)?; + Ok(item.is_some_and(|item| { + item.asset_type() == AssetType::Landmark + && item.asset_uuid() == asset_id + && permission_fingerprint(item.permissions()) == permissions_fingerprint + })) + }) + } + + fn teleport_landmark( + &self, + asset_id: UUID, + cancellation: CancellationToken, + ) -> LandmarkFuture<'_, bool> { + Box::pin(async move { + self.agent + .teleport_with_uuid_cancellation_token(asset_id, Some(cancellation)) + .await + .map_err(|_| LandmarkError::TeleportFailed) + }) + } +} + +#[cfg(feature = "live-grid")] +#[must_use] +pub fn permission_fingerprint(permissions: libremetaverse::Permissions) -> u64 { + // The compatibility type's stable hash is defined as the xor of all five + // permission masks. Re-checking the entire item type/asset plus this value + // detects permission changes without persisting raw inventory metadata. + u64::from(u32::from_ne_bytes( + permissions.get_hash_code().to_ne_bytes(), + )) +} + +#[cfg(feature = "live-grid")] +#[derive(Clone)] +struct PendingLiveOffer { + offer_id: String, + sender_id: UUID, + display_name: String, + received_unix_millis: u64, + folder: bool, +} + +/// Owns the original inventory event guards and catalogs accepted landmark +/// items only after the server returns authoritative type, asset, owner, and +/// permission state. Unsupported offers remain declined by default. +#[cfg(feature = "live-grid")] +pub struct LibremetaverseLandmarkIntake { + _offer: libremetaverse_types::compat::Subscription, + _received: libremetaverse_types::compat::Subscription, + _folder: libremetaverse_types::compat::Subscription, + pending: Arc>>, + cancellation: libremetaverse_types::compat::CancellationTokenSource, +} + +#[cfg(feature = "live-grid")] +impl fmt::Debug for LibremetaverseLandmarkIntake { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("LibremetaverseLandmarkIntake") + .field("pending", &lock(&self.pending).len()) + .finish_non_exhaustive() + } +} + +#[cfg(feature = "live-grid")] +impl LibremetaverseLandmarkIntake { + #[allow(clippy::too_many_lines)] + pub fn start( + service: &Arc, + owner: &crate::backend::LibremetaverseClientOwner, + ) -> Result { + let inventory = owner.client().inventory(); + let quarantine = Arc::new(Mutex::new(None::)); + let pending = Arc::new(Mutex::new(BTreeMap::::new())); + let cancellation = libremetaverse_types::compat::CancellationTokenSource::new(); + let offer_pending = Arc::clone(&pending); + let offer_service = Arc::clone(service); + let offer_inventory = inventory.clone(); + let offer_quarantine = Arc::clone(&quarantine); + let offer_folder_service = Arc::clone(service); + let offer_folder_inventory = inventory.clone(); + let offer_folder_agent = owner.agent(); + let offer_folder_cancellation = cancellation.token(); + let offer_folder_pending = Arc::clone(&pending); + let offer = inventory.subscribe_inventory_object_offered(Arc::new(move |mut event| { + let message = event.offer(); + let supported = matches!(event.asset_type(), AssetType::Landmark | AssetType::Folder) + && !event.from_task() + && event.object_id() != UUID::zero(); + if !supported || !offer_service.is_authorized(message.from_agent_id) { + event.set_accept(false); + return; + } + let quarantine = { + let mut cached = lock(&offer_quarantine); + if cached.is_none() { + *cached = offer_inventory + .find_folder_for_type_with_asset_type(AssetType::Landmark) + .and_then(|parent| { + offer_inventory.create_folder_with_uuid_string( + parent, + "MetaCrate Received Landmarks".into(), + ) + }) + .ok(); + } + *cached + }; + let Some(quarantine) = quarantine else { + event.set_accept(false); + return; + }; + let pending_offer = PendingLiveOffer { + offer_id: format!("inventory-{}", event.object_id()), + sender_id: message.from_agent_id, + display_name: message.message, + received_unix_millis: unix_millis(), + folder: event.asset_type() == AssetType::Folder, + }; + let mut values = lock(&offer_pending); + let stale_before = unix_millis().saturating_sub(5 * 60 * 1_000); + values.retain(|_, value| value.received_unix_millis >= stale_before); + if values.len() >= 1_024 + || values + .insert(event.object_id(), pending_offer.clone()) + .is_some() + { + event.set_accept(false); + return; + } + event.set_folder_id(quarantine); + event.set_accept(true); + if pending_offer.folder { + let service = Arc::clone(&offer_folder_service); + let inventory = offer_folder_inventory.clone(); + let cancellation = offer_folder_cancellation.clone(); + let pending = Arc::clone(&offer_folder_pending); + let agent_id = offer_folder_agent.agent_id(); + let folder_id = event.object_id(); + tokio::spawn(async move { + tokio::task::yield_now().await; + let Some(pending_offer) = lock(&pending).remove(&folder_id) else { + return; + }; + let _ = resolve_and_catalog_live_folder( + service, + inventory, + agent_id, + folder_id, + pending_offer, + cancellation, + ) + .await; + }); + } + })); + let received_pending = Arc::clone(&pending); + let received_service = Arc::clone(service); + let received = inventory.subscribe_item_received(Arc::new(move |event| { + let item = event.item(); + let Some(pending_offer) = lock(&received_pending).remove(&item.base.uuid()) else { + return; + }; + if pending_offer.folder + || item.asset_type() != AssetType::Landmark + || item.is_link().unwrap_or(true) + { + return; + } + let _ = received_service.catalog_received_offer(&LandmarkOffer { + offer_id: pending_offer.offer_id, + sender_id: pending_offer.sender_id, + root: OfferedInventoryNode { + inventory_id: item.base.uuid(), + owner_id: item.base.owner_id(), + name: if item.base.name().is_empty() { + pending_offer.display_name + } else { + item.base.name() + }, + kind: OfferedInventoryKind::Landmark { + asset_id: item.asset_uuid(), + permissions_fingerprint: permission_fingerprint(item.permissions()), + }, + children: Vec::new(), + }, + received_unix_millis: pending_offer.received_unix_millis, + }); + })); + let folder_pending = Arc::clone(&pending); + let folder_service = Arc::clone(service); + let folder_inventory = inventory; + let folder_subscription_owner = folder_inventory.clone(); + let folder_agent = owner.agent(); + let folder_cancellation = cancellation.token(); + let folder = folder_subscription_owner.subscribe_folder_updated(Arc::new(move |event| { + if !event.success() { + return; + } + let Some(pending_offer) = lock(&folder_pending).remove(&event.folder_id()) else { + return; + }; + if !pending_offer.folder { + return; + } + let service = Arc::clone(&folder_service); + let inventory = folder_inventory.clone(); + let cancellation = folder_cancellation.clone(); + let agent_id = folder_agent.agent_id(); + tokio::spawn(async move { + let _ = resolve_and_catalog_live_folder( + service, + inventory, + agent_id, + event.folder_id(), + pending_offer, + cancellation, + ) + .await; + }); + })); + Ok(Self { + _offer: offer, + _received: received, + _folder: folder, + pending, + cancellation, + }) + } +} + +#[cfg(feature = "live-grid")] +impl Drop for LibremetaverseLandmarkIntake { + fn drop(&mut self) { + self.cancellation.cancel(); + } +} + +#[cfg(feature = "live-grid")] +async fn resolve_and_catalog_live_folder( + service: Arc, + inventory: libremetaverse::InventoryManager, + owner_id: UUID, + root_id: UUID, + pending: PendingLiveOffer, + cancellation: CancellationToken, +) -> Result<(), LandmarkError> { + let limits = service.limits; + let mut queue = VecDeque::from([(root_id, 0_usize)]); + let mut visited = BTreeSet::new(); + let mut children = Vec::new(); + while let Some((folder_id, depth)) = queue.pop_front() { + if cancellation.is_cancellation_requested() { + return Err(LandmarkError::Cancelled); + } + if depth > limits.max_folder_depth + || visited.len() >= limits.max_folder_nodes + || !visited.insert(folder_id) + { + return Err(LandmarkError::TraversalLimit); + } + let contents = inventory + .request_folder_contents_with_uuid_uuid_boolean_boolean_inventory_sort_order_cancellation_token( + folder_id, + owner_id, + true, + true, + libremetaverse::InventorySortOrder::BY_NAME, + Some(cancellation.clone()), + ) + .await + .map_err(|_| LandmarkError::InvalidOffer)?; + for value in contents { + let item = inventory + .fetch_item(value.uuid(), owner_id, Some(cancellation.clone())) + .await + .map_err(|_| LandmarkError::InvalidOffer)?; + let Some(item) = item else { + queue.push_back((value.uuid(), depth + 1)); + continue; + }; + if item.is_link().unwrap_or(true) || item.asset_type() != AssetType::Landmark { + return Err(LandmarkError::UnsupportedInventory); + } + children.push(OfferedInventoryNode { + inventory_id: item.base.uuid(), + owner_id: item.base.owner_id(), + name: item.base.name(), + kind: OfferedInventoryKind::Landmark { + asset_id: item.asset_uuid(), + permissions_fingerprint: permission_fingerprint(item.permissions()), + }, + children: Vec::new(), + }); + if children.len() > limits.max_asset_fetches { + return Err(LandmarkError::TraversalLimit); + } + } + } + service.catalog_received_offer(&LandmarkOffer { + offer_id: pending.offer_id, + sender_id: pending.sender_id, + root: OfferedInventoryNode { + inventory_id: root_id, + owner_id, + name: pending.display_name, + kind: OfferedInventoryKind::Folder, + children, + }, + received_unix_millis: pending.received_unix_millis, + }) +} diff --git a/crates/metacrate-grid-agent/src/lib.rs b/crates/metacrate-grid-agent/src/lib.rs index 011ee6d..53b5377 100644 --- a/crates/metacrate-grid-agent/src/lib.rs +++ b/crates/metacrate-grid-agent/src/lib.rs @@ -11,6 +11,7 @@ pub mod control_plane; pub mod control_runtime; pub mod conversation; pub mod interaction; +pub mod landmarks; pub mod llm; pub mod observability; pub mod perception; @@ -33,6 +34,8 @@ mod conversation_tests; #[cfg(test)] mod interaction_tests; #[cfg(test)] +mod landmark_tests; +#[cfg(test)] mod observability_tests; #[cfg(test)] mod perception_tests; @@ -90,6 +93,18 @@ pub use interaction::{ PolicyLlmResponder, ResponderFuture, ResponsePacer, ResponseRequest, SuppressionReason, VisibleResponse, split_utf8, }; +pub use landmarks::{ + LANDMARK_LIST_TOOL, LANDMARK_SCHEDULE_TOOL, LANDMARK_STATUS_TOOL, LANDMARK_TELEPORT_TOOL, + LandmarkEntry, LandmarkError, LandmarkFuture, LandmarkGrid, LandmarkLimits, + LandmarkObservation, LandmarkObservationKind, LandmarkOffer, LandmarkOutcome, + LandmarkRoamingHandle, LandmarkService, LandmarkToolBackend, LandmarkValidation, OfferDecision, + OfferedInventoryKind, OfferedInventoryNode, RoamingPause, RoamingRandom, RoamingSchedule, + SystemRoamingRandom, TeleportReceipt, TeleportTrigger, landmark_policy_tools, +}; +#[cfg(feature = "live-grid")] +pub use landmarks::{ + LibremetaverseLandmarkGrid, LibremetaverseLandmarkIntake, permission_fingerprint, +}; pub use llm::{ Completion, CompletionMessage, ContentPart, ImageDetail, LlmClient, LlmError, LlmTransportLimits, ToolDefinition, ToolSchema, Usage, diff --git a/crates/metacrate-grid-agent/src/main.rs b/crates/metacrate-grid-agent/src/main.rs index ccc9599..3f9bbcd 100644 --- a/crates/metacrate-grid-agent/src/main.rs +++ b/crates/metacrate-grid-agent/src/main.rs @@ -346,6 +346,7 @@ async fn run_live( println!("grid agent session supervisor started; press Ctrl-C to stop"); let mut signal_error = None; let mut control_failure = None; + let mut active_interactions = 0_usize; loop { tokio::select! { signal = tokio::signal::ctrl_c() => { @@ -359,6 +360,8 @@ async fn run_live( record_session_observation(&live.observability, &event); if let SessionObservation::Transition { status, reason, retry_in } = event { control_target.update_session(status); + live.landmark_roaming + .update_pause(|pause| pause.degraded = !status.agent_ready); control_plane.publish(ControlEventKind::StateChanged { component: "session".into(), state: status.state.as_str().into(), @@ -374,6 +377,18 @@ async fn run_live( } event = live.interaction.next_observation() => { let Some(event) = event else { break; }; + match &event { + metacrate_grid_agent::InteractionObservation::InferenceStarted { .. } => { + active_interactions = active_interactions.saturating_add(1); + } + metacrate_grid_agent::InteractionObservation::InferenceFinished { .. } => { + active_interactions = active_interactions.saturating_sub(1); + } + _ => {} + } + live.landmark_roaming.update_pause(|pause| { + pause.conversation = active_interactions != 0; + }); record_interaction_observation(&live.observability, &event); println!("grid interaction event={event:?}"); } @@ -398,6 +413,8 @@ async fn run_live( let Some(command) = command else { break; }; match command { RuntimeControlCommand::Pause => { + live.landmark_roaming + .update_pause(|pause| pause.operator = true); if let Err(error) = handle.control(SessionControl::Pause).await { control_failure = Some(error.to_string()); break; @@ -408,6 +425,8 @@ async fn run_live( control_failure = Some(error.to_string()); break; } + live.landmark_roaming + .update_pause(|pause| pause.operator = false); } RuntimeControlCommand::ForceReconnect => { if let Err(error) = handle.control(SessionControl::ForceReconnect).await { @@ -435,6 +454,7 @@ async fn run_live( .result_code("requested")?; let _ = live.observability.record(shutdown_event); control_target.mark_stopping(); + live.landmark_roaming.shutdown().await; let session_result = handle.shutdown().await; let interaction_result = live.interaction.shutdown().await; let behavior_result = live.behavior.shutdown().await; @@ -468,6 +488,8 @@ struct LiveInteractions { policy: Arc, audit: Arc, observability: Arc, + _landmark_intake: metacrate_grid_agent::LibremetaverseLandmarkIntake, + landmark_roaming: metacrate_grid_agent::LandmarkRoamingHandle, } #[cfg(feature = "live-grid")] @@ -478,11 +500,13 @@ fn start_live_interactions( ) -> Result> { use metacrate_grid_agent::{ AuthorizedBackendRouter, AuthorizedToolBackend, BehaviorBackend, BehaviorController, - ConversationStore, InteractionCoordinator, LibremetaverseScriptInventory, LlmClient, - LlmTransportLimits, MemoryPolicyAudit, Observability, ObservabilityLimits, - PerceptionBackend, PolicyAuditSink, PolicyGateway, PolicyLimits, PolicyLlmResponder, - ScriptDeliveryBackend, ScriptDeliverySettings, ToolLoopLimits, UnifiedPolicyAudit, - behavior_policy_tools, perception_policy_tools, script_delivery_policy_tool, + ConversationStore, InteractionCoordinator, LandmarkLimits, LandmarkService, + LandmarkToolBackend, LibremetaverseLandmarkGrid, LibremetaverseLandmarkIntake, + LibremetaverseScriptInventory, LlmClient, LlmTransportLimits, MemoryPolicyAudit, + Observability, ObservabilityLimits, PerceptionBackend, PolicyAuditSink, PolicyGateway, + PolicyLimits, PolicyLlmResponder, ScriptDeliveryBackend, ScriptDeliverySettings, + SystemRoamingRandom, ToolLoopLimits, UnifiedPolicyAudit, behavior_policy_tools, + landmark_policy_tools, perception_policy_tools, script_delivery_policy_tool, }; let transport_limits = LlmTransportLimits { @@ -534,6 +558,26 @@ fn start_live_interactions( tools.push(script_delivery_policy_tool( ScriptDeliverySettings::default(), )?); + let landmark_service = Arc::new(LandmarkService::new( + Arc::new(LibremetaverseLandmarkGrid::new(owner)), + Arc::new(SystemRoamingRandom::default()), + config.authorized_avatar_uuids.clone(), + LandmarkLimits::default(), + Some(config.storage_path.join("landmarks.json")), + )?); + let landmark_backend: Arc = Arc::new( + LandmarkToolBackend::new(Arc::clone(&landmark_service)) + .with_behavior(behavior_ingress.clone()), + ); + let landmark_intake = LibremetaverseLandmarkIntake::start(&landmark_service, owner)?; + let landmark_roaming = landmark_service.start_roaming_runner( + metacrate_grid_agent::RoamingPause { + degraded: true, + ..metacrate_grid_agent::RoamingPause::default() + }, + Some(behavior_ingress.clone()), + ); + tools.extend(landmark_policy_tools(LandmarkLimits::default())?); let routes = tools .iter() .map(|tool| tool.definition.name.as_str().to_owned()) @@ -541,6 +585,8 @@ fn start_live_interactions( let backend: Arc = if name == metacrate_grid_agent::SCRIPT_DELIVERY_TOOL { Arc::clone(&script_backend) + } else if name.starts_with("landmark_") { + Arc::clone(&landmark_backend) } else if name.starts_with("behavior_") { Arc::new(BehaviorBackend::new(behavior_ingress.clone())) } else { @@ -601,6 +647,8 @@ fn start_live_interactions( policy: gateway, audit, observability, + _landmark_intake: landmark_intake, + landmark_roaming, }) } diff --git a/crates/metacrate-grid-agent/tests/dependency_policy.rs b/crates/metacrate-grid-agent/tests/dependency_policy.rs index f06ecf6..7365406 100644 --- a/crates/metacrate-grid-agent/tests/dependency_policy.rs +++ b/crates/metacrate-grid-agent/tests/dependency_policy.rs @@ -48,10 +48,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(20); + let mut files = Vec::with_capacity(32); collect_rust_files(&source, &mut files); assert!( - files.len() <= 30, + files.len() <= 32, "source-file count needs a reviewed bound update" ); for path in files { diff --git a/docs/grid-agent-landmarks.md b/docs/grid-agent-landmarks.md new file mode 100644 index 0000000..d2590e5 --- /dev/null +++ b/docs/grid-agent-landmarks.md @@ -0,0 +1,40 @@ +# Landmark intake, teleport, and roaming + +Landmark authority is private. Public chat and non-allow-listed IM can neither +accept offers nor see teleport/schedule tools. The policy layer seals the +authenticated avatar into every mutation; tool schemas contain a catalog +selector or bounded interval, never an avatar, region, coordinate, inventory +folder, or L$ field. + +The live adapter subscribes through the original LibreMetaverse compatibility +events but lives entirely in `metacrate-grid-agent`. Authorized landmark and +folder offers are accepted into `MetaCrate Received Landmarks`, which is a +quarantine and recovery folder. Direct items are cataloged only after the +server returns authoritative item metadata. Folder descendants are fetched +with bounded breadth/depth and asset-fetch counts; a cycle, link, duplicate +asset, bad type, empty folder, stale response, or limit breach leaves the +quarantined inventory untouched and out of the catalog. Task offers and all +unauthorized or arbitrary inventory offers are declined. + +The persisted catalog contains sender UUID, inventory UUID, asset UUID, +permission fingerprint, display name, receipt time, validation state, and last +outcome. Names are untrusted. Persistence is bounded, versioned, atomically +replaced, and mode `0600` on Unix. Teleport selection prefers stable IDs and +requires clarification for duplicate names. Immediately before the native +landmark teleport, the item is fetched again and its type, asset UUID, and +permissions must match. One semaphore, a timeout, cooldown, and lifecycle +cancellation prevent overlapping or stale teleports. + +Roaming schedules retain the authorizing avatar UUID and safe minimum/maximum +intervals. Randomness is injectable. A due run chooses without immediate +repetition when possible and schedules its next deadline from the current time, +so downtime never creates catch-up bursts. Conversation, operator, degraded, +build, and viewport-capture pause reasons all suppress selection. Disabling a +schedule is persistent and cancellation stops in-flight native teleport work. + +Focused verification: + +```console +cargo test --locked -p metacrate-grid-agent --lib landmark_tests +cargo check --locked -p metacrate-grid-agent --all-targets --features live-grid +```