use crate::config::SecretString; use crate::control_plane::*; use libremetaverse_types::compat::CancellationToken; use serde_json::json; use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; use std::sync::{Arc, Mutex}; use std::time::Duration; use tokio::io::{AsyncReadExt, AsyncWriteExt}; use tokio::net::TcpStream; #[derive(Default)] struct FakeTarget { mutations: Mutex>, shutdown: AtomicBool, active: AtomicUsize, maximum_active: AtomicUsize, } impl FakeTarget { fn mutation_count(&self) -> usize { self.mutations.lock().expect("mutations").len() } } impl ControlTarget for FakeTarget { #[allow(clippy::too_many_lines)] // Exhaustive fake covers every response family. fn execute( &self, context: ControlContext, request: ControlRequest, cancellation: CancellationToken, ) -> ControlFuture<'_> { Box::pin(async move { let active = self.active.fetch_add(1, Ordering::AcqRel) + 1; self.maximum_active.fetch_max(active, Ordering::AcqRel); let result = match request { ControlRequest::Health => Ok(ControlPayload::Health(HealthView { service_state: "running".into(), ready: true, uptime_seconds: 42, protocol_version: CONTROL_PROTOCOL_VERSION, })), ControlRequest::Runtime => Ok(ControlPayload::Runtime(RuntimeView { grid_state: "online".into(), generation: 7, transport_connected: true, agent_ready: true, region_id: Some("00000000-0000-4000-8000-000000000001".into()), region_handle: Some(9_007_199_254_740_992), region_name: Some("Safe Region".into()), position: Some([128.0, 128.0, 24.0]), behavior_mode: "available".into(), control_queue_used: 1, control_queue_capacity: 64, budget_tool_calls_used: 3, budget_movement_millimeters_used: 2_000, active_build_transaction: None, build_progress: None, build_orphan_ids: Vec::new(), active_visual_capture: None, visual_progress: None, visual_image_sha256: None, })), ControlRequest::ListSessions { page } => Ok(ControlPayload::Sessions(page_values( &page, (0..7) .map(|index| SessionMetadataView { session_id: format!("session-{index}"), avatar_id: format!("00000000-0000-4000-8000-{index:012x}"), channel: "direct_im".into(), created_unix_millis: 10, last_active_unix_millis: 20, turns: 2, bytes: 32, }) .collect(), ))), ControlRequest::ListScheduledJobs { page } => { Ok(ControlPayload::ScheduledJobs(page_values( &page, vec![ScheduledJobView { job_id: "roam-1".into(), kind: "landmark_roam".into(), enabled: true, next_run_unix_millis: Some(50), }], ))) } ControlRequest::ListPendingApprovals { page } => { Ok(ControlPayload::PendingApprovals(page_values( &page, vec![PendingApprovalView { approval_id: 1, tool: "build".into(), principal: "operator:local".into(), expires_unix_seconds: 99, movement_millimeters: 0, inventory_operations: 0, build_prims: 1, }], ))) } ControlRequest::ListAuditEvents { page } => { Ok(ControlPayload::AuditEvents(page_values( &page, vec![AuditEventView { sequence: 1, unix_millis: 1, principal: "operator:local".into(), operation: "pause_autonomy".into(), outcome: "completed".into(), authorization_id: None, }], ))) } ControlRequest::InjectOperatorMessage { message } if message == "wait" => { cancellation.cancelled().await; Err(ControlError { code: ControlErrorCode::Cancelled, message: "operation cancelled", retryable: false, }) } ControlRequest::GracefulShutdown => { self.shutdown.store(true, Ordering::Release); self.mutations .lock() .expect("mutations") .push("graceful_shutdown".into()); Ok(ControlPayload::Completed) } other => { assert_eq!(context.role, ControlRole::Operator); self.mutations .lock() .expect("mutations") .push(format!("{other:?}")); Ok(ControlPayload::Completed) } }; self.active.fetch_sub(1, Ordering::AcqRel); result }) } } fn page_values(page: &PageRequest, values: Vec) -> Page { let start = usize::try_from(page.cursor.unwrap_or(0)) .expect("cursor") .min(values.len()); let end = start .saturating_add(usize::from(page.limit)) .min(values.len()); let next_cursor = (end < values.len()).then(|| u64::try_from(end).expect("cursor")); Page { items: values.into_iter().skip(start).take(end - start).collect(), next_cursor, } } fn token(value: &str) -> SecretString { SecretString::new("control.test_token", value).expect("token") } fn limits() -> ControlLimits { ControlLimits { request_timeout: Duration::from_secs(5), idle_timeout: Duration::from_secs(5), write_timeout: Duration::from_secs(1), ..ControlLimits::default() } } fn plane(target: Arc, configured: ControlLimits) -> Arc { ControlPlane::new( target, token("operator-secret"), Some(token("observer-secret")), configured, ) .expect("plane") } fn request(id: &str, request: ControlRequest) -> ControlRequestEnvelope { ControlRequestEnvelope { version: CONTROL_PROTOCOL_VERSION, request_id: id.into(), request, } } enum Client { InProcess(InProcessControlClient), Tcp(TcpControlClient), } impl Client { async fn request(&self, envelope: ControlRequestEnvelope) -> ControlResponseEnvelope { match self { Self::InProcess(client) => client.request(envelope).await, Self::Tcp(client) => client.request(envelope).await.expect("TCP response"), } } } enum Transport { InProcess, Tcp, } struct Fixture { client: Client, plane: Arc, target: Arc, server: Option, } async fn fixture(transport: Transport, role: ControlRole, configured: ControlLimits) -> Fixture { let target = Arc::new(FakeTarget::default()); let plane = plane(target.clone(), configured); let value = if role == ControlRole::Operator { "operator-secret" } else { "observer-secret" }; match transport { Transport::InProcess => Fixture { client: Client::InProcess(plane.connect(value).expect("connect")), plane, target, server: None, }, Transport::Tcp => { let server = TcpControlServer::bind( plane.clone(), TcpControlConfig { listen: "127.0.0.1:0".parse().expect("address"), limits: configured, }, ) .await .expect("bind"); let client = TcpControlClient::connect(server.local_addr(), value, configured) .await .expect("connect"); Fixture { client: Client::Tcp(client), plane, target, server: Some(server), } } } } async fn close(fixture: Fixture) { drop(fixture.client); if let Some(server) = fixture.server { server.shutdown().await.expect("shutdown"); } } async fn conformance(transport: Transport) { let fixture = fixture(transport, ControlRole::Operator, limits()).await; let health = fixture .client .request(request("health-1", ControlRequest::Health)) .await; assert_eq!(health.version, CONTROL_PROTOCOL_VERSION); assert!(matches!( health.result, Ok(ControlPayload::Health(HealthView { ready: true, .. })) )); let page = fixture .client .request(request( "sessions-1", ControlRequest::ListSessions { page: PageRequest { cursor: Some(2), limit: 3, }, }, )) .await; let Ok(ControlPayload::Sessions(page)) = page.result else { panic!("session page") }; assert_eq!(page.items.len(), 3); assert_eq!(page.items[0].session_id, "session-2"); assert_eq!(page.next_cursor, Some(5)); assert!( fixture .client .request(request("pause-1", ControlRequest::PauseAutonomy)) .await .result .is_ok() ); let replay = fixture .client .request(request("pause-1", ControlRequest::ResumeAutonomy)) .await; assert_eq!( replay.result.expect_err("replay").code, ControlErrorCode::Replay ); assert_eq!(fixture.target.mutation_count(), 1); let shutdown = fixture .client .request(request("shutdown-1", ControlRequest::GracefulShutdown)) .await; assert!(shutdown.result.is_ok()); assert!(fixture.target.shutdown.load(Ordering::Acquire)); close(fixture).await; } #[tokio::test] async fn same_conformance_suite_runs_in_process() { conformance(Transport::InProcess).await; } #[tokio::test] async fn same_conformance_suite_runs_over_loopback_tcp() { conformance(Transport::Tcp).await; } #[tokio::test] async fn observer_is_read_only_and_tokens_are_redacted() { let fixture = fixture(Transport::InProcess, ControlRole::Observer, limits()).await; assert!(matches!( fixture .client .request(request("read", ControlRequest::Runtime)) .await .result, Ok(ControlPayload::Runtime(_)) )); let denied = fixture .client .request(request("mutate", ControlRequest::PauseAutonomy)) .await; assert_eq!( denied.result.expect_err("denied").code, ControlErrorCode::PermissionDenied ); assert_eq!(fixture.target.mutation_count(), 0); assert!(!format!("{:?}", fixture.plane).contains("operator-secret")); let Err(error) = fixture.plane.connect("wrong") else { panic!("bad token was accepted") }; assert_eq!(error.code, ControlErrorCode::AuthenticationFailed); close(fixture).await; } #[tokio::test] async fn version_bounds_and_malformed_requests_are_typed() { let fixture = fixture(Transport::InProcess, ControlRole::Operator, limits()).await; let mut wrong = request("version", ControlRequest::Health); wrong.version = 99; assert_eq!( fixture .client .request(wrong) .await .result .expect_err("version") .code, ControlErrorCode::VersionMismatch ); let invalid = fixture .client .request(request( "page", ControlRequest::ListSessions { page: PageRequest { cursor: None, limit: 101, }, }, )) .await; assert_eq!( invalid.result.expect_err("page").code, ControlErrorCode::InvalidRequest ); let secret_payload = fixture .client .request(request( "audit", ControlRequest::ListAuditEvents { page: PageRequest::default(), }, )) .await; let rendered = serde_json::to_string(&secret_payload).expect("JSON"); for secret in ["operator-secret", "observer-secret", "Bearer ", "CAPS/"] { assert!(!rendered.contains(secret)); } close(fixture).await; } #[tokio::test] async fn cancellation_is_per_request_and_concurrent_operators_stay_responsive() { for transport in [Transport::InProcess, Transport::Tcp] { let fixture = fixture(transport, ControlRole::Operator, limits()).await; let slow = match &fixture.client { Client::InProcess(value) => Client::InProcess(value.clone()), Client::Tcp(value) => Client::Tcp(value.clone()), }; let task = tokio::spawn(async move { slow.request(request( "slow", ControlRequest::InjectOperatorMessage { message: "wait".into(), }, )) .await }); tokio::time::timeout(Duration::from_secs(1), async { while fixture.target.active.load(Ordering::Acquire) == 0 { tokio::task::yield_now().await; } }) .await .expect("slow request became active"); let concurrent_operator = match &fixture.server { Some(server) => Client::Tcp( TcpControlClient::connect(server.local_addr(), "operator-secret", limits()) .await .expect("second TCP operator"), ), None => Client::InProcess( fixture .plane .connect("operator-secret") .expect("second in-process operator"), ), }; let health = concurrent_operator .request(request("health-during", ControlRequest::Health)) .await; assert!(health.result.is_ok()); let cancelled = fixture .client .request(request( "cancel", ControlRequest::CancelRequest { target_request_id: "slow".into(), }, )) .await; assert!(matches!( cancelled.result, Ok(ControlPayload::Cancelled { .. }) )); let slow = task.await.expect("slow task"); assert_eq!( slow.result.expect_err("cancelled").code, ControlErrorCode::Cancelled ); assert!(fixture.target.maximum_active.load(Ordering::Acquire) >= 2); drop(concurrent_operator); close(fixture).await; } } #[tokio::test] async fn event_stream_resubscribes_with_explicit_gap_and_slow_consumers_are_cut_off() { let mut configured = limits(); configured.event_queue = 2; configured.event_history = 3; let fixture = fixture(Transport::InProcess, ControlRole::Operator, configured).await; let mut in_subscription = None; if let Client::InProcess(client) = &fixture.client { let (_, subscription) = client .subscribe(request( "subscribe", ControlRequest::SubscribeEvents { after_sequence: Some(0), }, )) .await .expect("subscribe"); in_subscription = Some(subscription); } for index in 0..10 { fixture.plane.publish(ControlEventKind::StateChanged { component: "grid".into(), state: format!("state-{index}"), }); } assert_eq!(fixture.plane.retained_event_count(), 3); let first = in_subscription .as_mut() .expect("subscription") .recv() .await .expect("buffered"); assert_eq!(first.sequence, 1); drop(in_subscription); let client = fixture.plane.connect("operator-secret").expect("reconnect"); let (_, mut resumed) = client .subscribe(request( "resubscribe", ControlRequest::SubscribeEvents { after_sequence: Some(1), }, )) .await .expect("resubscribe"); let gap = resumed.recv().await.expect("gap"); assert!(matches!( gap.event, ControlEventKind::Gap { first_available: 9, last_missed: 8 } )); close(fixture).await; } #[tokio::test] async fn zero_client_event_flood_remains_bounded_and_headless() { let target = Arc::new(FakeTarget::default()); let mut configured = limits(); configured.event_history = 4; let plane = plane(target, configured); assert_eq!(plane.active_connections(), 0); for index in 0..10_000 { plane.publish(ControlEventKind::StateChanged { component: "load".into(), state: index.to_string(), }); } assert_eq!(plane.retained_event_count(), 4); assert_eq!(plane.active_connections(), 0); } #[tokio::test] async fn subscriptions_have_a_hard_global_bound_and_release_on_drop() { let mut configured = limits(); configured.max_subscriptions = 2; let target = Arc::new(FakeTarget::default()); let plane = plane(target, configured); let client = plane.connect("operator-secret").expect("client"); let (_, first) = client .subscribe(request( "subscribe-1", ControlRequest::SubscribeEvents { after_sequence: None, }, )) .await .expect("first subscription"); let (_, second) = client .subscribe(request( "subscribe-2", ControlRequest::SubscribeEvents { after_sequence: None, }, )) .await .expect("second subscription"); let error = client .subscribe(request( "subscribe-3", ControlRequest::SubscribeEvents { after_sequence: None, }, )) .await .err() .expect("subscription bound"); assert_eq!(error.code, ControlErrorCode::Busy); drop(first); client .subscribe(request( "subscribe-4", ControlRequest::SubscribeEvents { after_sequence: None, }, )) .await .expect("released subscription"); drop(second); } #[tokio::test] #[allow(clippy::too_many_lines)] // One raw-transport matrix shares one listener. async fn tcp_rejects_bad_tokens_oversized_and_partial_frames() { let configured = limits(); let target = Arc::new(FakeTarget::default()); let plane = plane(target, configured); let server = TcpControlServer::bind( plane.clone(), TcpControlConfig { listen: "127.0.0.1:0".parse().expect("address"), limits: configured, }, ) .await .expect("server"); let Err(error) = TcpControlClient::connect(server.local_addr(), "wrong", configured).await else { panic!("bad token was accepted") }; assert_eq!(error.code, ControlErrorCode::AuthenticationFailed); let mut wrong_version = TcpStream::connect(server.local_addr()) .await .expect("version connect"); let hello = serde_json::to_vec(&json!({ "frame":"hello", "body":{"version":99,"token":"operator-secret"} })) .expect("hello JSON"); wrong_version .write_all( &u32::try_from(hello.len()) .expect("frame length") .to_be_bytes(), ) .await .expect("version header"); wrong_version.write_all(&hello).await.expect("version body"); let mut response_header = [0_u8; 4]; wrong_version .read_exact(&mut response_header) .await .expect("version response header"); let mut response = vec![0_u8; usize::try_from(u32::from_be_bytes(response_header)).expect("response length")]; wrong_version .read_exact(&mut response) .await .expect("version response"); let response: serde_json::Value = serde_json::from_slice(&response).expect("version JSON"); assert_eq!(response["body"]["error"]["code"], json!("version_mismatch")); let subscriber = TcpControlClient::connect(server.local_addr(), "operator-secret", configured) .await .expect("subscriber"); let subscription_response = subscriber .request(request( "events", ControlRequest::SubscribeEvents { after_sequence: None, }, )) .await .expect("subscription response"); assert!(matches!( subscription_response.result, Ok(ControlPayload::Subscribed { .. }) )); plane.publish(ControlEventKind::StateChanged { component: "session".into(), state: "online".into(), }); let event = tokio::time::timeout(Duration::from_secs(1), subscriber.next_event()) .await .expect("event deadline") .expect("event"); assert!(matches!(event.event, ControlEventKind::StateChanged { .. })); let last_sequence = event.sequence; drop(subscriber); plane.publish(ControlEventKind::StateChanged { component: "session".into(), state: "reconnecting".into(), }); let reconnected = TcpControlClient::connect(server.local_addr(), "operator-secret", configured) .await .expect("reconnected subscriber"); reconnected .request(request( "events-resumed", ControlRequest::SubscribeEvents { after_sequence: Some(last_sequence), }, )) .await .expect("resubscription response"); let replayed = tokio::time::timeout(Duration::from_secs(1), reconnected.next_event()) .await .expect("replay deadline") .expect("replayed event"); assert_eq!(replayed.sequence, last_sequence.saturating_add(1)); let mut oversized = TcpStream::connect(server.local_addr()) .await .expect("raw connect"); oversized .write_all(&(u32::try_from(configured.max_frame_bytes).expect("bound") + 1).to_be_bytes()) .await .expect("header"); let mut byte = [0_u8; 1]; assert_eq!( tokio::time::timeout(Duration::from_secs(1), oversized.read(&mut byte)) .await .expect("closed") .expect("read"), 0 ); let mut partial = TcpStream::connect(server.local_addr()) .await .expect("raw connect"); partial .write_all(&100_u32.to_be_bytes()) .await .expect("header"); partial.write_all(b"{}").await.expect("partial"); partial.shutdown().await.expect("close"); drop(reconnected); server.shutdown().await.expect("shutdown"); } #[tokio::test] async fn plaintext_server_rejects_non_loopback_before_binding() { let configured = limits(); let target = Arc::new(FakeTarget::default()); let plane = plane(target, configured); let Err(error) = TcpControlServer::bind( plane, TcpControlConfig { listen: "0.0.0.0:7943".parse().expect("address"), limits: configured, }, ) .await else { panic!("plaintext remote listener was accepted") }; assert_eq!(error.code, ControlErrorCode::InvalidRequest); } #[test] fn plain_tcp_is_loopback_only_and_protocol_json_is_stable() { let configured = limits(); assert!( !TcpControlConfig { listen: "0.0.0.0:7943".parse().expect("address"), limits: configured } .listen .ip() .is_loopback() ); let encoded = serde_json::to_value(request("health", ControlRequest::Health)).expect("JSON"); assert_eq!( encoded, json!({"version":1,"request_id":"health","request":{"method":"health"}}) ); }