Implement authorized landmark roaming (#130)
Some checks failed
CI / rust-skia (Rust only) (push) Successful in 2m46s
CI / required (push) Failing after 53s

This commit is contained in:
2026-08-18 11:08:50 +02:00
parent a749111657
commit 08696cd2fe
7 changed files with 2413 additions and 7 deletions

View File

@@ -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<metacrate_grid_agent::PolicyGateway>,
audit: Arc<metacrate_grid_agent::MemoryPolicyAudit>,
observability: Arc<metacrate_grid_agent::Observability>,
_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<LiveInteractions, Box<dyn Error>> {
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<dyn AuthorizedToolBackend> = 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<dyn AuthorizedToolBackend> =
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,
})
}