From cb823b59d5835deecdd3dd6b956452c58e826994 Mon Sep 17 00:00:00 2001 From: "hh.(SII)" Date: Sat, 26 Sep 2026 02:16:24 +0800 Subject: [PATCH 1/3] feat(simulation): add local virtual lab runtime --- crates/lab-cli/src/commands/serve.rs | 16 +- crates/osdl-core/src/adapter/mod.rs | 1 + crates/osdl-core/src/adapter/simulation.rs | 63 ++++ crates/osdl-core/src/config.rs | 228 ++++++++++++ crates/osdl-core/src/engine.rs | 107 ++++++ crates/osdl-core/src/lib.rs | 3 +- crates/osdl-core/src/transport/mod.rs | 1 + crates/osdl-core/src/transport/simulation.rs | 373 +++++++++++++++++++ crates/osdl-core/tests/integration.rs | 61 +++ docs/recipes/configs/simulation.yaml | 71 ++++ docs/recipes/simulation-mode.md | 46 +++ 11 files changed, 968 insertions(+), 2 deletions(-) create mode 100644 crates/osdl-core/src/adapter/simulation.rs create mode 100644 crates/osdl-core/src/transport/simulation.rs create mode 100644 docs/recipes/configs/simulation.yaml create mode 100644 docs/recipes/simulation-mode.md diff --git a/crates/lab-cli/src/commands/serve.rs b/crates/lab-cli/src/commands/serve.rs index 693da77..73d4fa5 100644 --- a/crates/lab-cli/src/commands/serve.rs +++ b/crates/lab-cli/src/commands/serve.rs @@ -12,8 +12,11 @@ use std::path::{Path, PathBuf}; use anyhow::{anyhow, Context}; use clap::Args; use osdl_core::adapter::onvif::OnvifAdapter; +use osdl_core::adapter::simulation::SimulationAdapter; use osdl_core::adapter::unilabos::UniLabOsAdapter; -use osdl_core::config::{AdapterConfig, EspNowDongleConfig, MqttConfig, OsdlConfig}; +use osdl_core::config::{ + AdapterConfig, EspNowDongleConfig, MqttConfig, OsdlConfig, SimulationConfig, +}; use osdl_core::driver::registry::DriverRegistry; use osdl_core::path_expand; use osdl_core::{EmbeddedBroker, EventStore, MdnsAdvertiser, OsdlEngine}; @@ -81,6 +84,12 @@ pub struct ServeArgs { /// non-loopback address; optional on loopback. #[arg(long, env = "OSDL_AUTH_TOKEN", hide_env_values = true)] pub auth_token: Option, + + /// Start a deterministic local simulation world with virtual devices. + /// The world uses the same Lab Action Model and gRPC surface as physical + /// devices, so local development does not require hardware. + #[arg(long, env = "OSDL_SIMULATION")] + pub simulation: bool, } /// Synchronous entrypoint called from `main`. Handles `--detach` *before* @@ -372,6 +381,7 @@ pub async fn run(args: ServeArgs) -> anyhow::Result<()> { let adapters: Vec> = vec![ Box::new(UniLabOsAdapter::new(DriverRegistry::with_builtins())), Box::new(OnvifAdapter::new()), + Box::new(SimulationAdapter::new()), ]; let mut engine = OsdlEngine::new(config, adapters).with_store(store); let handle = engine.handle(); @@ -488,6 +498,10 @@ fn build_config(args: &ServeArgs) -> anyhow::Result { } } + if args.simulation { + cfg.simulation = Some(SimulationConfig::default()); + } + Ok(cfg) } diff --git a/crates/osdl-core/src/adapter/mod.rs b/crates/osdl-core/src/adapter/mod.rs index d2b27f2..87fe28b 100644 --- a/crates/osdl-core/src/adapter/mod.rs +++ b/crates/osdl-core/src/adapter/mod.rs @@ -1,4 +1,5 @@ pub mod onvif; +pub mod simulation; pub mod unilabos; use crate::protocol::*; diff --git a/crates/osdl-core/src/adapter/simulation.rs b/crates/osdl-core/src/adapter/simulation.rs new file mode 100644 index 0000000..67b533d --- /dev/null +++ b/crates/osdl-core/src/adapter/simulation.rs @@ -0,0 +1,63 @@ +//! Protocol adapter for the built-in local simulation backend. +//! +//! Simulation deliberately uses the existing byte-oriented transport path: +//! the adapter serializes a small JSON action envelope and decodes a JSON +//! telemetry envelope. This keeps Lab Action Model calls identical for real +//! and virtual devices while leaving the physics/runtime implementation in +//! the transport/backend layer. + +use crate::adapter::{DeviceMatch, ProtocolAdapter}; +use crate::protocol::DeviceCommand; +use serde_json::Value; +use std::collections::HashMap; + +pub const PLATFORM: &str = "simulation"; + +#[derive(Debug, Default)] +pub struct SimulationAdapter; + +impl SimulationAdapter { + pub fn new() -> Self { + Self + } +} + +impl ProtocolAdapter for SimulationAdapter { + fn platform(&self) -> &str { + PLATFORM + } + + fn load_registry(&mut self, _path: &str) -> Result<(), String> { + // Simulation devices are declared by the simulation world, not by a + // physical registry directory. + Ok(()) + } + + fn match_hardware(&self, _hardware_id: &str) -> Option { + None + } + + fn encode_command(&self, _device_type: &str, cmd: &DeviceCommand) -> Result, String> { + serde_json::to_vec(cmd).map_err(|e| format!("simulation: encode command: {e}")) + } + + fn decode_response(&self, _device_type: &str, bytes: &[u8]) -> Option> { + let envelope: Value = serde_json::from_slice(bytes).ok()?; + let properties = envelope.get("properties")?.as_object()?; + let mut decoded = properties + .iter() + .map(|(key, value)| (key.clone(), value.clone())) + .collect::>(); + decoded.insert("simulation".into(), Value::Bool(true)); + if let Some(engine) = envelope.get("engine") { + decoded.insert("simulation_engine".into(), engine.clone()); + } + if let Some(world_id) = envelope.get("world_id") { + decoded.insert("simulation_world".into(), world_id.clone()); + } + if let Some(action) = envelope.get("last_action") { + decoded.insert("last_action".into(), action.clone()); + } + Some(decoded) + } +} diff --git a/crates/osdl-core/src/config.rs b/crates/osdl-core/src/config.rs index f3fde9c..7c518b3 100644 --- a/crates/osdl-core/src/config.rs +++ b/crates/osdl-core/src/config.rs @@ -1,4 +1,5 @@ use crate::media::{mediamtx::MediaGatewayConfig, MediaSourceConfig}; +use crate::protocol::ActionSchema; use serde::{Deserialize, Serialize}; use std::collections::HashMap; use std::path::PathBuf; @@ -51,6 +52,233 @@ pub struct OsdlConfig { /// will fall back to the system temp directory. #[serde(default, skip_serializing_if = "Option::is_none")] pub data_dir: Option, + /// Optional local simulation world. Simulation devices use the same + /// Device/Transport/ProtocolAdapter path as physical hardware, so an + /// Agent and the UI can develop against them without a connected lab. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub simulation: Option, +} + +/// Configuration for a local simulation world. +/// +/// `engine` is deliberately a string at this boundary. It is the runtime +/// capability negotiated by a future backend adapter (for example `rapier`, +/// `mujoco`, `isaac-sim`, or `gazebo`). The built-in `kinematic` backend is +/// deterministic and available in every OpenSDL build; unsupported engines +/// fail with an actionable error instead of silently falling back to it. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct SimulationConfig { + /// Stable world identifier used in simulated transport/device ids. + #[serde(default = "default_simulation_world_id")] + pub world_id: String, + /// Physics/runtime backend name. `kinematic` is the built-in backend. + #[serde(default = "default_simulation_engine")] + pub engine: String, + /// Fixed update frequency for telemetry and deterministic stepping. + #[serde(default = "default_simulation_tick_hz")] + pub tick_hz: u32, + /// Seed reserved for deterministic physics backends. + #[serde(default)] + pub seed: u64, + /// Virtual devices exposed through the normal OpenSDL device contract. + #[serde(default = "default_simulation_devices")] + pub devices: Vec, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct SimulationDeviceConfig { + /// Local id inside the simulation world, e.g. `heater-1`. + pub id: String, + pub device_type: String, + #[serde(default)] + pub role: Option, + #[serde(default)] + pub description: String, + #[serde(default)] + pub actions: Vec, + #[serde(default)] + pub properties: HashMap, + /// Initial scene position in metres. The UI uses this to place entities. + #[serde(default)] + pub position: [f64; 3], +} + +fn default_simulation_world_id() -> String { + "lab-sim".into() +} + +fn default_simulation_engine() -> String { + "kinematic".into() +} + +fn default_simulation_tick_hz() -> u32 { + 10 +} + +fn default_simulation_devices() -> Vec { + vec![ + SimulationDeviceConfig { + id: "heater-1".into(), + device_type: "simulation.heater".into(), + role: Some("heater".into()), + description: "Virtual temperature-controlled hotplate".into(), + actions: vec![ + action_schema( + "set_temperature", + "Set the target temperature", + serde_json::json!({"type":"object","properties":{"temperature":{"type":"number","unit":"°C"}},"required":["temperature"]}), + ), + action_schema( + "start", + "Start heating", + serde_json::json!({"type":"object","properties":{}}), + ), + action_schema( + "stop", + "Stop heating", + serde_json::json!({"type":"object","properties":{}}), + ), + ], + properties: HashMap::from([ + ("temperature".into(), serde_json::json!(22.0)), + ("target_temperature".into(), serde_json::json!(22.0)), + ("running".into(), serde_json::json!(false)), + ]), + position: [-1.6, 0.0, 0.0], + }, + SimulationDeviceConfig { + id: "stirrer-1".into(), + device_type: "simulation.stirrer".into(), + role: Some("stirrer".into()), + description: "Virtual magnetic stirrer".into(), + actions: vec![ + action_schema( + "set_speed", + "Set stirring speed", + serde_json::json!({"type":"object","properties":{"speed":{"type":"number","unit":"rpm"}},"required":["speed"]}), + ), + action_schema( + "start", + "Start stirring", + serde_json::json!({"type":"object","properties":{}}), + ), + action_schema( + "stop", + "Stop stirring", + serde_json::json!({"type":"object","properties":{}}), + ), + ], + properties: HashMap::from([ + ("speed".into(), serde_json::json!(0.0)), + ("running".into(), serde_json::json!(false)), + ]), + position: [0.0, 0.0, 0.0], + }, + SimulationDeviceConfig { + id: "valve-1".into(), + device_type: "simulation.valve".into(), + role: Some("valve".into()), + description: "Virtual fluid control valve".into(), + actions: vec![ + action_schema( + "open", + "Open the valve", + serde_json::json!({"type":"object","properties":{}}), + ), + action_schema( + "close", + "Close the valve", + serde_json::json!({"type":"object","properties":{}}), + ), + ], + properties: HashMap::from([("state".into(), serde_json::json!("closed"))]), + position: [1.6, 0.0, 0.0], + }, + SimulationDeviceConfig { + id: "probe-1".into(), + device_type: "simulation.sensor".into(), + role: Some("temperature_sensor".into()), + description: "Virtual temperature probe".into(), + actions: vec![action_schema( + "read", + "Read the current measurement", + serde_json::json!({"type":"object","properties":{}}), + )], + properties: HashMap::from([ + ("temperature".into(), serde_json::json!(22.0)), + ("unit".into(), serde_json::json!("°C")), + ]), + position: [0.0, 0.0, 1.8], + }, + ] +} + +fn action_schema(name: &str, description: &str, params: serde_json::Value) -> ActionSchema { + ActionSchema { + name: name.into(), + description: description.into(), + params, + } +} + +impl Default for SimulationConfig { + fn default() -> Self { + Self { + world_id: default_simulation_world_id(), + engine: default_simulation_engine(), + tick_hz: default_simulation_tick_hz(), + seed: 0, + devices: default_simulation_devices(), + } + } +} + +impl SimulationConfig { + pub fn validate(&self) -> Result<(), String> { + if self.world_id.trim().is_empty() { + return Err("simulation world_id must not be empty".into()); + } + if self.tick_hz == 0 || self.tick_hz > 240 { + return Err("simulation tick_hz must be between 1 and 240".into()); + } + if self.engine.trim().is_empty() { + return Err("simulation engine must not be empty".into()); + } + if self.devices.is_empty() { + return Err("simulation must define at least one device".into()); + } + let mut ids = std::collections::HashSet::new(); + for device in &self.devices { + if device.id.trim().is_empty() || !ids.insert(&device.id) { + return Err(format!( + "simulation device id '{}' is empty or duplicated", + device.id + )); + } + } + Ok(()) + } +} + +#[cfg(test)] +mod simulation_tests { + use super::SimulationConfig; + + #[test] + fn default_simulation_is_valid_and_deterministic() { + let config = SimulationConfig::default(); + config.validate().expect("default simulation config"); + assert_eq!(config.engine, "kinematic"); + assert_eq!(config.devices.len(), 4); + } + + #[test] + fn empty_engine_is_rejected() { + let mut config = SimulationConfig::default(); + config.engine.clear(); + let error = config.validate().expect_err("an empty engine is invalid"); + assert!(error.contains("engine")); + } } /// One physical bus (e.g., RS-485) reached through a single transport, diff --git a/crates/osdl-core/src/engine.rs b/crates/osdl-core/src/engine.rs index bac97f5..2853fbe 100644 --- a/crates/osdl-core/src/engine.rs +++ b/crates/osdl-core/src/engine.rs @@ -13,6 +13,10 @@ use crate::transport::espnow_dongle::{ }; use crate::transport::mqtt_serial::MqttSerialTransport; use crate::transport::onvif::{transport_id_for as onvif_transport_id, OnvifTransport}; +use crate::transport::simulation::{ + device_id_for as simulation_device_id, transport_id_for as simulation_transport_id, + SimulationTransport, +}; use crate::transport::{Transport, TransportRx}; use rumqttc::AsyncClient; @@ -475,6 +479,16 @@ impl OsdlEngine { device_count: 0, }); + if let Err(error) = self.start_simulation().await { + log::error!("Failed to start simulation world: {error}"); + self.stop_simulation().await; + let _ = self + .handle + .status_tx + .send(OsdlStatus::Error { message: error }); + return; + } + // Start configured ESP-NOW dongles (USB-CDC). Each dongle owns a // serial read loop and emits REG events for registration-driven // device discovery. @@ -497,6 +511,7 @@ impl OsdlEngine { // that case. if *stop_rx.borrow() { log::info!("OSDL engine stop already requested before run() entered the loop"); + self.stop_simulation().await; let _ = self.handle.status_tx.send(OsdlStatus::Disconnected); return; } @@ -655,6 +670,8 @@ impl OsdlEngine { #[cfg(feature = "espnow")] self.stop_espnow_dongles().await; + self.stop_simulation().await; + if let Some(p) = media_proc.take() { p.shutdown().await; // Drop the snapshot too so a fresh `run()` doesn't observe @@ -665,6 +682,96 @@ impl OsdlEngine { let _ = self.handle.status_tx.send(OsdlStatus::Disconnected); } + /// Register the configured local simulation world before entering the + /// event loop. Virtual devices use the same transport and adapter path as + /// physical devices, which keeps the Agent/UI contract identical. + async fn start_simulation(&self) -> Result<(), String> { + let Some(config) = self.handle.config.simulation.clone() else { + return Ok(()); + }; + config.validate()?; + + for device_config in &config.devices { + let transport_id = simulation_transport_id(&config.world_id, &device_config.id); + let transport = Arc::new(SimulationTransport::new( + config.world_id.clone(), + config.engine.clone(), + device_config, + config.tick_hz, + self.handle.transport_rx_tx.clone(), + )?); + transport.start().await?; + self.handle + .register_transport(transport_id.clone(), transport) + .await; + + let device_id = simulation_device_id(&config.world_id, &device_config.id); + let mut properties = device_config.properties.clone(); + properties.insert("simulation".into(), serde_json::Value::Bool(true)); + properties.insert( + "simulation_engine".into(), + serde_json::Value::String(config.engine.clone()), + ); + properties.insert( + "simulation_world".into(), + serde_json::Value::String(config.world_id.clone()), + ); + properties.insert("position".into(), serde_json::json!(device_config.position)); + + let description = if device_config.description.is_empty() { + format!("Virtual {}", device_config.device_type) + } else { + device_config.description.clone() + }; + self.handle + .register_device(Device { + id: device_id, + transport_id, + device_type: device_config.device_type.clone(), + adapter: crate::adapter::simulation::PLATFORM.into(), + description, + online: true, + properties, + actions: device_config.actions.clone(), + role: device_config.role.clone(), + }) + .await?; + } + + log::info!( + "Simulation world '{}' started with '{}' backend ({} devices)", + config.world_id, + config.engine, + config.devices.len() + ); + Ok(()) + } + + async fn stop_simulation(&self) { + let Some(config) = self.handle.config.simulation.as_ref() else { + return; + }; + let transport_list = { + let transports = self.handle.transports.read().await; + config + .devices + .iter() + .filter_map(|device| { + let id = simulation_transport_id(&config.world_id, &device.id); + transports + .get(&id) + .cloned() + .map(|transport| (id, transport)) + }) + .collect::>() + }; + for (id, transport) in transport_list { + if let Err(error) = transport.stop().await { + log::warn!("Failed to stop simulation transport {id}: {error}"); + } + } + } + /// Spawn mediamtx if any media sources are configured. Emits /// `MediaSourceOnline` for each source so the host learns the URLs. /// Returns the process handle so the caller can shut it down. diff --git a/crates/osdl-core/src/lib.rs b/crates/osdl-core/src/lib.rs index c9d10f4..b61f248 100644 --- a/crates/osdl-core/src/lib.rs +++ b/crates/osdl-core/src/lib.rs @@ -14,7 +14,7 @@ pub mod store; pub mod transport; pub use broker::EmbeddedBroker; -pub use config::OsdlConfig; +pub use config::{OsdlConfig, SimulationConfig, SimulationDeviceConfig}; pub use engine::{EngineHandle, OsdlEngine, OsdlStatus}; pub use event::OsdlEvent; pub use mdns::MdnsAdvertiser; @@ -23,3 +23,4 @@ pub use protocol::{ CommandResult, CommandStatus, Device, DeviceCommand, DeviceStatus, Node, NodeRegistration, }; pub use store::EventStore; +pub use transport::simulation::{SimulationBackend, SimulationState, SimulationTransport}; diff --git a/crates/osdl-core/src/transport/mod.rs b/crates/osdl-core/src/transport/mod.rs index 0c67cc3..2d106b7 100644 --- a/crates/osdl-core/src/transport/mod.rs +++ b/crates/osdl-core/src/transport/mod.rs @@ -19,6 +19,7 @@ pub mod direct_serial; pub mod espnow_dongle; pub mod mqtt_serial; pub mod onvif; +pub mod simulation; pub mod tcp; use async_trait::async_trait; diff --git a/crates/osdl-core/src/transport/simulation.rs b/crates/osdl-core/src/transport/simulation.rs new file mode 100644 index 0000000..266f070 --- /dev/null +++ b/crates/osdl-core/src/transport/simulation.rs @@ -0,0 +1,373 @@ +//! Deterministic local simulation transport. +//! +//! The built-in backend is intentionally small: it provides stable device +//! state and telemetry for developing Lab Action Model workflows without a +//! connected instrument. A future Rapier, Bullet, MuJoCo, Isaac, or Gazebo +//! integration can implement the same contract behind a new backend while +//! retaining the engine's device/action/event surface. + +use super::{Transport, TransportRx}; +use crate::config::SimulationDeviceConfig; +use async_trait::async_trait; +use serde_json::{json, Value}; +use std::collections::HashMap; +use std::sync::{ + atomic::{AtomicBool, Ordering}, + Arc, +}; +use std::time::{SystemTime, UNIX_EPOCH}; +use tokio::sync::{mpsc, Mutex}; +use tokio::task::JoinHandle; + +#[derive(Clone)] +pub struct SimulationTransport { + inner: Arc, +} + +struct Inner { + transport_id: String, + world_id: String, + device_id: String, + role: Option, + position: [f64; 3], + tick_hz: u32, + state: Mutex, + backend: Arc, + rx_tx: mpsc::UnboundedSender, + tick_task: Mutex>>, + connected: AtomicBool, +} + +#[derive(Debug, Clone)] +pub struct SimulationState { + pub properties: HashMap, + pub step: u64, + pub last_action: Option, +} + +/// Runtime seam for physics integrations. The transport owns identity, +/// telemetry, and the OpenSDL event path; a backend owns world evolution and +/// action semantics. A future Rapier/Bullet/MuJoCo/Isaac/Gazebo adapter can +/// implement this seam without changing the device contract. +#[async_trait] +pub trait SimulationBackend: Send + Sync { + fn engine_id(&self) -> &str; + async fn apply_action(&self, state: &mut SimulationState, action: &str, params: &Value); + async fn advance(&self, state: &mut SimulationState, dt: f64); +} + +struct KinematicBackend; + +#[async_trait] +impl SimulationBackend for KinematicBackend { + fn engine_id(&self) -> &str { + "kinematic" + } + + async fn apply_action(&self, state: &mut SimulationState, action: &str, params: &Value) { + state.last_action = Some(action.to_string()); + match action { + "start" | "resume" => { + state.properties.insert("running".into(), Value::Bool(true)); + } + "stop" | "pause" => { + state + .properties + .insert("running".into(), Value::Bool(false)); + } + "open" => { + state + .properties + .insert("state".into(), Value::String("open".into())); + } + "close" => { + state + .properties + .insert("state".into(), Value::String("closed".into())); + } + "set_temperature" | "set_target_temperature" => { + if let Some(value) = number_param(params, &["temperature", "target", "value"]) { + state + .properties + .insert("target_temperature".into(), json!(value)); + } + } + "set_speed" | "set_stir_speed" => { + if let Some(value) = number_param(params, &["speed", "rpm", "value"]) { + state + .properties + .insert("speed".into(), json!(value.max(0.0))); + } + } + "set_position" | "move" => { + if let Some(value) = number_param(params, &["position", "value"]) { + state + .properties + .insert("position_value".into(), json!(value)); + } + } + _ => { + state + .properties + .insert("last_action_params".into(), params.clone()); + } + } + } + + async fn advance(&self, state: &mut SimulationState, dt: f64) { + state.step = state.step.saturating_add(1); + if let (Some(Value::Number(current)), Some(Value::Number(target))) = ( + state.properties.get("temperature"), + state.properties.get("target_temperature"), + ) { + let current = current.as_f64().unwrap_or(22.0); + let target = target.as_f64().unwrap_or(current); + let running = state + .properties + .get("running") + .and_then(Value::as_bool) + .unwrap_or(true); + let next = if running { + current + (target - current) * (1.0 - (-0.8 * dt).exp()) + } else { + current + }; + state.properties.insert("temperature".into(), json!(next)); + } + if let Some(speed) = state.properties.get("speed").and_then(Value::as_f64) { + let running = state + .properties + .get("running") + .and_then(Value::as_bool) + .unwrap_or(false); + state.properties.insert( + "angular_velocity".into(), + json!(if running { speed } else { 0.0 }), + ); + } + } +} + +fn backend_for(engine: &str) -> Result, String> { + match engine { + "kinematic" => Ok(Arc::new(KinematicBackend)), + other => Err(format!( + "simulation engine '{other}' is not available in this OpenSDL build; use 'kinematic' or install a runtime adapter" + )), + } +} + +impl SimulationTransport { + pub fn new( + world_id: String, + engine: String, + device: &SimulationDeviceConfig, + tick_hz: u32, + rx_tx: mpsc::UnboundedSender, + ) -> Result { + let backend = backend_for(&engine)?; + Self::with_backend(world_id, device, tick_hz, rx_tx, backend) + } + + /// Construct a transport from an externally provided physics backend. + /// Desktop and headless hosts can register a Rapier, Bullet, MuJoCo, + /// Isaac, or Gazebo implementation here without changing OpenSDL's + /// device/action/event contract. + pub fn with_backend( + world_id: String, + device: &SimulationDeviceConfig, + tick_hz: u32, + rx_tx: mpsc::UnboundedSender, + backend: Arc, + ) -> Result { + if tick_hz == 0 || tick_hz > 240 { + return Err("simulation tick_hz must be between 1 and 240".into()); + } + let transport_id = transport_id_for(&world_id, &device.id); + let engine = backend.engine_id().to_string(); + let mut properties = device.properties.clone(); + properties.insert("simulation".into(), Value::Bool(true)); + properties.insert("simulation_engine".into(), Value::String(engine.clone())); + properties.insert("simulation_world".into(), Value::String(world_id.clone())); + properties.insert("position".into(), json!(device.position)); + + Ok(Self { + inner: Arc::new(Inner { + transport_id, + world_id: world_id.clone(), + device_id: device_id_for(&world_id, &device.id), + role: device.role.clone(), + position: device.position, + tick_hz, + state: Mutex::new(SimulationState { + properties, + step: 0, + last_action: None, + }), + backend, + rx_tx, + tick_task: Mutex::new(None), + connected: AtomicBool::new(false), + }), + }) + } + + async fn emit_status(&self) { + let state = self.inner.state.lock().await.clone(); + let mut properties = state.properties; + properties.insert("step".into(), json!(state.step)); + let envelope = json!({ + "device_id": self.inner.device_id, + "world_id": self.inner.world_id, + "engine": self.inner.backend.engine_id(), + "role": self.inner.role, + "position": self.inner.position, + "timestamp": now_millis(), + "last_action": state.last_action, + "properties": properties, + }); + let _ = self.inner.rx_tx.send(TransportRx { + transport_id: self.inner.transport_id.clone(), + data: serde_json::to_vec(&envelope).unwrap_or_default(), + }); + } + + async fn apply_action(&self, action: &str, params: &Value) { + let mut state = self.inner.state.lock().await; + self.inner + .backend + .apply_action(&mut state, action, params) + .await; + } + + async fn advance(&self) { + let dt = 1.0 / self.inner.tick_hz as f64; + let mut state = self.inner.state.lock().await; + self.inner.backend.advance(&mut state, dt).await; + } +} + +#[async_trait] +impl Transport for SimulationTransport { + fn transport_type(&self) -> &str { + "simulation" + } + + fn description(&self) -> String { + format!( + "Simulation {} ({})", + self.inner.device_id, + self.inner.backend.engine_id() + ) + } + + async fn send(&self, bytes: &[u8]) -> Result<(), String> { + if !self.is_connected() { + return Err("simulation transport is not running".into()); + } + let command: Value = + serde_json::from_slice(bytes).map_err(|e| format!("simulation: parse command: {e}"))?; + let action = command + .get("action") + .and_then(Value::as_str) + .ok_or("simulation: command missing action")?; + let params = command.get("params").cloned().unwrap_or_else(|| json!({})); + self.apply_action(action, ¶ms).await; + self.emit_status().await; + Ok(()) + } + + fn is_connected(&self) -> bool { + self.inner.connected.load(Ordering::Acquire) + } + + async fn start(&self) -> Result<(), String> { + if self.inner.connected.swap(true, Ordering::AcqRel) { + return Ok(()); + } + self.emit_status().await; + let this = self.clone(); + let interval = std::time::Duration::from_secs_f64(1.0 / this.inner.tick_hz as f64); + let task = tokio::spawn(async move { + let mut ticker = tokio::time::interval(interval); + ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); + loop { + ticker.tick().await; + if !this.is_connected() { + break; + } + this.advance().await; + this.emit_status().await; + } + }); + *self.inner.tick_task.lock().await = Some(task); + Ok(()) + } + + async fn stop(&self) -> Result<(), String> { + self.inner.connected.store(false, Ordering::Release); + if let Some(task) = self.inner.tick_task.lock().await.take() { + task.abort(); + } + Ok(()) + } +} + +pub fn transport_id_for(world_id: &str, device_id: &str) -> String { + format!("simulation:{world_id}:{device_id}") +} + +pub fn device_id_for(world_id: &str, device_id: &str) -> String { + format!("sim:{world_id}:{device_id}") +} + +fn number_param(params: &Value, names: &[&str]) -> Option { + names + .iter() + .find_map(|name| params.get(*name).and_then(Value::as_f64)) +} + +fn now_millis() -> i64 { + SystemTime::now() + .duration_since(UNIX_EPOCH) + .map(|duration| duration.as_millis() as i64) + .unwrap_or_default() +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::config::SimulationConfig; + use crate::protocol::DeviceCommand; + + #[tokio::test] + async fn emits_telemetry_and_applies_actions() { + let config = SimulationConfig::default(); + let device = config.devices.first().expect("default heater"); + let (tx, mut rx) = mpsc::unbounded_channel(); + let transport = + SimulationTransport::new(config.world_id, config.engine, device, config.tick_hz, tx) + .expect("kinematic backend"); + transport.start().await.expect("start"); + let initial = rx.recv().await.expect("initial telemetry"); + assert_eq!( + initial.transport_id, + transport_id_for("lab-sim", "heater-1") + ); + + let command = DeviceCommand { + command_id: "command-1".into(), + device_id: device_id_for("lab-sim", &device.id), + action: "set_temperature".into(), + params: json!({"temperature": 80.0}), + }; + transport + .send(&serde_json::to_vec(&command).expect("encode")) + .await + .expect("send"); + let update = rx.recv().await.expect("action telemetry"); + let payload: Value = serde_json::from_slice(&update.data).expect("json"); + assert_eq!(payload["properties"]["target_temperature"], json!(80.0)); + transport.stop().await.expect("stop"); + } +} diff --git a/crates/osdl-core/tests/integration.rs b/crates/osdl-core/tests/integration.rs index a03d69a..9cc6281 100644 --- a/crates/osdl-core/tests/integration.rs +++ b/crates/osdl-core/tests/integration.rs @@ -112,6 +112,67 @@ fn test_engine_creation() { let _rx = engine.subscribe_events(); } +#[tokio::test] +async fn test_simulation_world_uses_the_device_contract() { + use osdl_core::adapter::simulation::SimulationAdapter; + use osdl_core::config::{OsdlConfig, SimulationConfig}; + use osdl_core::event::OsdlEvent; + use osdl_core::{EventStore, OsdlEngine}; + + let config = OsdlConfig { + mqtt: None, + simulation: Some(SimulationConfig::default()), + ..Default::default() + }; + let adapters: Vec> = vec![Box::new(SimulationAdapter::new())]; + let mut engine = OsdlEngine::new(config, adapters).with_store(EventStore::in_memory().unwrap()); + let handle = engine.handle(); + let mut events = handle.subscribe_events(); + let task = tokio::spawn(async move { engine.run().await }); + + let device = tokio::time::timeout(std::time::Duration::from_secs(1), async { + loop { + if let Some(device) = handle.get_device("sim:lab-sim:heater-1").await { + break device; + } + tokio::task::yield_now().await; + } + }) + .await + .expect("simulation device registered"); + assert_eq!(device.adapter, "simulation"); + assert_eq!(device.role.as_deref(), Some("heater")); + + let result = handle + .send_command(osdl_core::protocol::DeviceCommand { + command_id: "simulation-command".into(), + device_id: device.id.clone(), + action: "set_temperature".into(), + params: serde_json::json!({"temperature": 80.0}), + }) + .await + .expect("simulation command"); + assert_eq!(result.status, osdl_core::protocol::CommandStatus::Pending); + + let status = tokio::time::timeout(std::time::Duration::from_secs(1), async { + loop { + if let Ok(OsdlEvent::DeviceStatus(status)) = events.recv().await { + if status.device_id == device.id + && status.properties.get("target_temperature") == Some(&serde_json::json!(80.0)) + { + break status; + } + } + } + }) + .await + .expect("simulation status"); + assert_eq!(status.properties["simulation"], serde_json::json!(true)); + + handle.request_stop(); + task.await.expect("engine task"); +} + #[test] fn test_event_store_logging() { use osdl_core::event::OsdlEvent; diff --git a/docs/recipes/configs/simulation.yaml b/docs/recipes/configs/simulation.yaml new file mode 100644 index 0000000..7542743 --- /dev/null +++ b/docs/recipes/configs/simulation.yaml @@ -0,0 +1,71 @@ +# Deterministic local development world. No MQTT broker or hardware is needed. +mqtt: null +simulation: + world_id: lab-sim + engine: kinematic + tick_hz: 10 + seed: 42 + devices: + - id: heater-1 + device_type: simulation.heater + role: heater + description: Virtual temperature-controlled hotplate + position: [-1.6, 0.0, 0.0] + actions: + - name: set_temperature + description: Set the target temperature + params: + type: object + properties: + temperature: + type: number + unit: °C + required: [temperature] + - name: start + description: Start heating + params: {type: object, properties: {}} + - name: stop + description: Stop heating + params: {type: object, properties: {}} + properties: + temperature: 22.0 + target_temperature: 22.0 + running: false + - id: stirrer-1 + device_type: simulation.stirrer + role: stirrer + description: Virtual magnetic stirrer + position: [0.0, 0.0, 0.0] + actions: + - name: set_speed + description: Set stirring speed + params: + type: object + properties: + speed: + type: number + unit: rpm + required: [speed] + - name: start + description: Start stirring + params: {type: object, properties: {}} + - name: stop + description: Stop stirring + params: {type: object, properties: {}} + properties: + speed: 0.0 + running: false + - id: valve-1 + device_type: simulation.valve + role: valve + description: Virtual fluid control valve + position: [1.6, 0.0, 0.0] + actions: + - name: open + description: Open the valve + params: {type: object, properties: {}} + - name: close + description: Close the valve + params: {type: object, properties: {}} + properties: + state: closed diff --git a/docs/recipes/simulation-mode.md b/docs/recipes/simulation-mode.md new file mode 100644 index 0000000..b683e79 --- /dev/null +++ b/docs/recipes/simulation-mode.md @@ -0,0 +1,46 @@ +# Local simulation mode + +OpenSDL can boot a deterministic virtual lab when no physical devices are +connected. The virtual devices use the normal `Device`, `DeviceCommand`, +`DeviceStatus`, event, and gRPC contracts, so a Runner, Agent, or UI can be +developed against them without a hardware-specific code path. + +## Quick start + +```bash +cargo run --bin lab -- serve --simulation +``` + +The default world exposes a virtual heater, magnetic stirrer, valve, and +temperature probe. Inspect them from another shell: + +```bash +cargo run --bin lab -- device list +cargo run --bin lab -- send sim:lab-sim:heater-1 set_temperature \ + --params '{"temperature":80}' +cargo run --bin lab -- send sim:lab-sim:heater-1 start +``` + +`--simulation` is equivalent to `OSDL_SIMULATION=true`. It is a local +development mode and does not discover or control physical hardware. + +## A custom world + +Use [`configs/simulation.yaml`](configs/simulation.yaml) as a starting point: + +```bash +cargo run --bin lab -- serve --config docs/recipes/configs/simulation.yaml +``` + +`engine: kinematic` is the deterministic backend shipped with OpenSDL. The +configuration already reserves the runtime boundary for `rapier`, `bullet`, +`mujoco`, `isaac-sim`, and `gazebo`; those engines must be provided by a +validated local backend adapter before they can be selected. OpenSDL fails +with an explicit error for an unavailable engine instead of silently changing +the experiment's semantics. + +The React + Three workbench treats the simulation snapshot as an observer and +action surface. Physics remains authoritative in the local OpenSDL runtime; +the browser does not become a second physics engine. This makes an eventual +desktop Runner integration and a headless lab host share the same world +state. From 7c4d872cc4decc3102cd5b6e76ef3d50dca2dd6e Mon Sep 17 00:00:00 2001 From: "hh.(SII)" Date: Sat, 26 Sep 2026 03:34:45 +0800 Subject: [PATCH 2/3] feat(simulation): expose backend and hub asset seams --- crates/osdl-core/src/config.rs | 36 +++++++++++++++++ crates/osdl-core/src/engine.rs | 38 ++++++++++++++++-- crates/osdl-core/src/lib.rs | 7 +++- crates/osdl-core/src/transport/simulation.rs | 42 +++++++++++++++++++- 4 files changed, 116 insertions(+), 7 deletions(-) diff --git a/crates/osdl-core/src/config.rs b/crates/osdl-core/src/config.rs index 7c518b3..f6e1f0a 100644 --- a/crates/osdl-core/src/config.rs +++ b/crates/osdl-core/src/config.rs @@ -98,11 +98,25 @@ pub struct SimulationDeviceConfig { pub actions: Vec, #[serde(default)] pub properties: HashMap, + /// Optional immutable Hub asset identity used by visual workbenches. + #[serde(default)] + pub asset_ref: Option, /// Initial scene position in metres. The UI uses this to place entities. #[serde(default)] pub position: [f64; 3], } +/// Pinned Hub identity for a device model. The source bytes remain in the Hub; +/// simulation telemetry carries this reference so clients can resolve a +/// verified preview without copying model data through OpenSDL. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct SimulationAssetRef { + pub namespace: String, + pub name: String, + #[serde(default)] + pub version: Option, +} + fn default_simulation_world_id() -> String { "lab-sim".into() } @@ -144,6 +158,7 @@ fn default_simulation_devices() -> Vec { ("target_temperature".into(), serde_json::json!(22.0)), ("running".into(), serde_json::json!(false)), ]), + asset_ref: None, position: [-1.6, 0.0, 0.0], }, SimulationDeviceConfig { @@ -172,6 +187,7 @@ fn default_simulation_devices() -> Vec { ("speed".into(), serde_json::json!(0.0)), ("running".into(), serde_json::json!(false)), ]), + asset_ref: None, position: [0.0, 0.0, 0.0], }, SimulationDeviceConfig { @@ -192,6 +208,7 @@ fn default_simulation_devices() -> Vec { ), ], properties: HashMap::from([("state".into(), serde_json::json!("closed"))]), + asset_ref: None, position: [1.6, 0.0, 0.0], }, SimulationDeviceConfig { @@ -208,6 +225,7 @@ fn default_simulation_devices() -> Vec { ("temperature".into(), serde_json::json!(22.0)), ("unit".into(), serde_json::json!("°C")), ]), + asset_ref: None, position: [0.0, 0.0, 1.8], }, ] @@ -255,6 +273,24 @@ impl SimulationConfig { device.id )); } + if let Some(asset) = &device.asset_ref { + if asset.namespace.trim().is_empty() || asset.name.trim().is_empty() { + return Err(format!( + "simulation device '{}' has an incomplete asset_ref", + device.id + )); + } + if asset + .version + .as_deref() + .is_some_and(|version| version.trim().is_empty()) + { + return Err(format!( + "simulation device '{}' has an empty asset_ref version", + device.id + )); + } + } } Ok(()) } diff --git a/crates/osdl-core/src/engine.rs b/crates/osdl-core/src/engine.rs index 2853fbe..fe6278b 100644 --- a/crates/osdl-core/src/engine.rs +++ b/crates/osdl-core/src/engine.rs @@ -15,7 +15,7 @@ use crate::transport::mqtt_serial::MqttSerialTransport; use crate::transport::onvif::{transport_id_for as onvif_transport_id, OnvifTransport}; use crate::transport::simulation::{ device_id_for as simulation_device_id, transport_id_for as simulation_transport_id, - SimulationTransport, + BuiltinSimulationBackendFactory, SimulationBackendFactory, SimulationTransport, }; use crate::transport::{Transport, TransportRx}; @@ -66,6 +66,8 @@ pub struct EngineHandle { /// Sender for transports to push received bytes back to the engine. transport_rx_tx: mpsc::UnboundedSender, + /// Factory for the physics runtime used by configured simulation worlds. + simulation_backend_factory: Arc, /// External command injection — drop a `DeviceCommand` in here and the /// engine's main loop dispatches it via `send_command`. cmd_inject_tx: mpsc::UnboundedSender, @@ -346,6 +348,7 @@ impl OsdlEngine { devices: Arc::new(RwLock::new(HashMap::new())), transports: Arc::new(RwLock::new(HashMap::new())), transport_rx_tx, + simulation_backend_factory: Arc::new(BuiltinSimulationBackendFactory), cmd_inject_tx, events_tx, media_sources: Arc::new(RwLock::new(HashMap::new())), @@ -374,6 +377,17 @@ impl OsdlEngine { self } + /// Inject a physics backend factory while retaining OpenSDL's device and + /// event contract. Hosts can use this to register native runtimes such as + /// Rapier, Bullet, MuJoCo, Isaac Sim, or Gazebo. + pub fn with_simulation_backend_factory( + mut self, + factory: Arc, + ) -> Self { + self.handle.simulation_backend_factory = factory; + self + } + /// Cheap-to-clone handle exposing inspection / command / event APIs. /// Most callers (gRPC server, orchestrator, CLI) talk to the engine /// through this handle rather than `&OsdlEngine` directly. @@ -693,12 +707,18 @@ impl OsdlEngine { for device_config in &config.devices { let transport_id = simulation_transport_id(&config.world_id, &device_config.id); - let transport = Arc::new(SimulationTransport::new( + let backend = self.handle.simulation_backend_factory.create( + &config.engine, + &config.world_id, + device_config, + )?; + let backend_engine = backend.engine_id().to_string(); + let transport = Arc::new(SimulationTransport::with_backend( config.world_id.clone(), - config.engine.clone(), device_config, config.tick_hz, self.handle.transport_rx_tx.clone(), + backend, )?); transport.start().await?; self.handle @@ -710,12 +730,22 @@ impl OsdlEngine { properties.insert("simulation".into(), serde_json::Value::Bool(true)); properties.insert( "simulation_engine".into(), - serde_json::Value::String(config.engine.clone()), + serde_json::Value::String(backend_engine), ); properties.insert( "simulation_world".into(), serde_json::Value::String(config.world_id.clone()), ); + if let Some(asset_ref) = &device_config.asset_ref { + properties.insert( + "simulation_asset".into(), + serde_json::json!({ + "namespace": asset_ref.namespace, + "name": asset_ref.name, + "version": asset_ref.version, + }), + ); + } properties.insert("position".into(), serde_json::json!(device_config.position)); let description = if device_config.description.is_empty() { diff --git a/crates/osdl-core/src/lib.rs b/crates/osdl-core/src/lib.rs index b61f248..fe759a8 100644 --- a/crates/osdl-core/src/lib.rs +++ b/crates/osdl-core/src/lib.rs @@ -14,7 +14,7 @@ pub mod store; pub mod transport; pub use broker::EmbeddedBroker; -pub use config::{OsdlConfig, SimulationConfig, SimulationDeviceConfig}; +pub use config::{OsdlConfig, SimulationAssetRef, SimulationConfig, SimulationDeviceConfig}; pub use engine::{EngineHandle, OsdlEngine, OsdlStatus}; pub use event::OsdlEvent; pub use mdns::MdnsAdvertiser; @@ -23,4 +23,7 @@ pub use protocol::{ CommandResult, CommandStatus, Device, DeviceCommand, DeviceStatus, Node, NodeRegistration, }; pub use store::EventStore; -pub use transport::simulation::{SimulationBackend, SimulationState, SimulationTransport}; +pub use transport::simulation::{ + BuiltinSimulationBackendFactory, SimulationBackend, SimulationBackendFactory, SimulationState, + SimulationTransport, +}; diff --git a/crates/osdl-core/src/transport/simulation.rs b/crates/osdl-core/src/transport/simulation.rs index 266f070..ed0d0f1 100644 --- a/crates/osdl-core/src/transport/simulation.rs +++ b/crates/osdl-core/src/transport/simulation.rs @@ -56,6 +56,21 @@ pub trait SimulationBackend: Send + Sync { async fn advance(&self, state: &mut SimulationState, dt: f64); } +/// Creates a physics backend for one configured simulation world. +/// +/// Hosts that embed OpenSDL can inject a factory for a native runtime while +/// keeping device identity, commands, telemetry, and events on the OpenSDL +/// contract. The default implementation only exposes the deterministic +/// kinematic backend shipped in this crate. +pub trait SimulationBackendFactory: Send + Sync { + fn create( + &self, + engine: &str, + world_id: &str, + device: &SimulationDeviceConfig, + ) -> Result, String>; +} + struct KinematicBackend; #[async_trait] @@ -157,6 +172,20 @@ fn backend_for(engine: &str) -> Result, String> { } } +#[derive(Debug, Default)] +pub struct BuiltinSimulationBackendFactory; + +impl SimulationBackendFactory for BuiltinSimulationBackendFactory { + fn create( + &self, + engine: &str, + _world_id: &str, + _device: &SimulationDeviceConfig, + ) -> Result, String> { + backend_for(engine) + } +} + impl SimulationTransport { pub fn new( world_id: String, @@ -165,7 +194,8 @@ impl SimulationTransport { tick_hz: u32, rx_tx: mpsc::UnboundedSender, ) -> Result { - let backend = backend_for(&engine)?; + let backend = + BuiltinSimulationBackendFactory::default().create(&engine, &world_id, device)?; Self::with_backend(world_id, device, tick_hz, rx_tx, backend) } @@ -190,6 +220,16 @@ impl SimulationTransport { properties.insert("simulation_engine".into(), Value::String(engine.clone())); properties.insert("simulation_world".into(), Value::String(world_id.clone())); properties.insert("position".into(), json!(device.position)); + if let Some(asset_ref) = &device.asset_ref { + properties.insert( + "simulation_asset".into(), + json!({ + "namespace": asset_ref.namespace, + "name": asset_ref.name, + "version": asset_ref.version, + }), + ); + } Ok(Self { inner: Arc::new(Inner { From ba120e9af317547258a1a23d03a6cc60fe6976b3 Mon Sep 17 00:00:00 2001 From: "hh.(SII)" Date: Sat, 26 Sep 2026 04:09:29 +0800 Subject: [PATCH 3/3] docs(simulation): document Hub model bindings --- docs/recipes/configs/simulation.yaml | 6 ++++++ docs/recipes/simulation-mode.md | 18 ++++++++++++++++++ 2 files changed, 24 insertions(+) diff --git a/docs/recipes/configs/simulation.yaml b/docs/recipes/configs/simulation.yaml index 7542743..435d804 100644 --- a/docs/recipes/configs/simulation.yaml +++ b/docs/recipes/configs/simulation.yaml @@ -31,6 +31,12 @@ simulation: temperature: 22.0 target_temperature: 22.0 running: false + # Optional pinned Hub model. The preview is resolved by the workbench; + # OpenSDL only carries this immutable identity in telemetry. + # asset_ref: + # namespace: scienceol + # name: ur5e + # version: 1.0.0 - id: stirrer-1 device_type: simulation.stirrer role: stirrer diff --git a/docs/recipes/simulation-mode.md b/docs/recipes/simulation-mode.md index b683e79..19ef357 100644 --- a/docs/recipes/simulation-mode.md +++ b/docs/recipes/simulation-mode.md @@ -44,3 +44,21 @@ action surface. Physics remains authoritative in the local OpenSDL runtime; the browser does not become a second physics engine. This makes an eventual desktop Runner integration and a headless lab host share the same world state. + +## Hub device models + +Each virtual device may pin a published Hub model without copying model bytes +into the OpenSDL state: + +```yaml +asset_ref: + namespace: scienceol + name: ur5e + version: 1.0.0 +``` + +The reference is included in simulation telemetry as `simulation_asset`. The +Web workbench resolves the immutable version through the Hub catalog, loads +its verified GLB preview in Three, and falls back to a primitive when the +preview is unavailable. If `asset_ref` is omitted, the workbench uses the +asset's published `bindings.deviceTypes` declaration to choose a model.