Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 10 additions & 1 deletion crates/osdl-core/src/engine.rs
Original file line number Diff line number Diff line change
Expand Up @@ -753,6 +753,15 @@ impl OsdlEngine {
} else {
device_config.description.clone()
};
let mut actions = device_config.actions.clone();
if actions.iter().all(|action| action.name != "reset") {
actions.push(ActionSchema {
name: "reset".into(),
description: "Restore this virtual device to its configured initial state"
.into(),
params: serde_json::json!({"type":"object","properties":{}}),
});
}
self.handle
.register_device(Device {
id: device_id,
Expand All @@ -762,7 +771,7 @@ impl OsdlEngine {
description,
online: true,
properties,
actions: device_config.actions.clone(),
actions,
role: device_config.role.clone(),
})
.await?;
Expand Down
55 changes: 55 additions & 0 deletions crates/osdl-core/src/transport/simulation.rs
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,9 @@ struct Inner {

#[derive(Debug, Clone)]
pub struct SimulationState {
/// Properties declared by the simulation device configuration. Backends
/// use this snapshot to make reset deterministic across repeated runs.
pub initial_properties: HashMap<String, Value>,
pub properties: HashMap<String, Value>,
pub step: u64,
pub last_action: Option<String>,
Expand All @@ -54,6 +57,14 @@ 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);

/// Restore the device's configured state. Physics integrations can
/// override this to reset their world state alongside the OpenSDL view.
async fn reset(&self, state: &mut SimulationState) {
state.properties = state.initial_properties.clone();
state.step = 0;
state.last_action = Some("reset".into());
Comment on lines +63 to +66

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

问题 (bug_risk): 重置状态的修改与重置遥测的发送不是原子的:reset 释放状态互斥锁后,ticker 可能会在 emit_status 执行前推进状态并加入遥测消息,因此重置响应报告的可能是非零步数或演化后的属性,而不是重置快照。新测试同样会消费下一条遥测消息,因此可能读到这次介入的 tick,并间歇性失败。

触发条件: 模拟 ticker 在动作应用与状态发送之间触发时。

建议修复: 在单一的命令/遥测同步机制下,确保重置及其对应的状态发送按顺序执行;同时让测试等待 last_action == "reset" 的遥测消息,而不是消费任意一条排队中的消息。

Original comment in English

issue (bug_risk): Reset state mutation and reset telemetry emission are not atomic: after reset releases the state mutex, the ticker can advance the state and enqueue telemetry before emit_status runs, so the reset response reports a nonzero step or evolved properties instead of the reset snapshot. The new test also consumes whichever telemetry message is next, so it can read that intervening tick and fail intermittently.

Triggers: When the simulation ticker fires between action application and status emission.

Suggested fix: Keep reset and its corresponding status emission ordered under a single command/telemetry synchronization mechanism, and make the test wait for telemetry with last_action == "reset" rather than consuming an arbitrary queued message.

}
}

/// Creates a physics backend for one configured simulation world.
Expand Down Expand Up @@ -82,6 +93,7 @@ impl SimulationBackend for KinematicBackend {
async fn apply_action(&self, state: &mut SimulationState, action: &str, params: &Value) {
state.last_action = Some(action.to_string());
match action {
"reset" => self.reset(state).await,
"start" | "resume" => {
state.properties.insert("running".into(), Value::Bool(true));
}
Expand Down Expand Up @@ -257,6 +269,7 @@ impl SimulationTransport {
position: device.position,
tick_hz,
state: Mutex::new(SimulationState {
initial_properties: properties.clone(),
properties,
step: 0,
last_action: None,
Expand Down Expand Up @@ -496,4 +509,46 @@ mod tests {
assert!(payload["properties"].get("simulation_asset").is_none());
transport.stop().await.expect("stop");
}

#[tokio::test]
async fn reset_restores_configured_properties_and_step() {
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 _ = rx.recv().await.expect("initial telemetry");

let set_temperature = DeviceCommand {
command_id: "reset-temperature-command".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(&set_temperature).expect("encode"))
.await
.expect("set temperature");
let _ = rx.recv().await.expect("temperature telemetry");

let reset = DeviceCommand {
command_id: "reset-command".into(),
device_id: device_id_for("lab-sim", &device.id),
action: "reset".into(),
params: json!({}),
};
transport
.send(&serde_json::to_vec(&reset).expect("encode"))
.await
.expect("reset");
let telemetry = rx.recv().await.expect("reset telemetry");
let payload: Value = serde_json::from_slice(&telemetry.data).expect("json");
assert_eq!(payload["properties"]["temperature"], json!(22.0));
assert_eq!(payload["properties"]["target_temperature"], json!(22.0));
assert_eq!(payload["properties"]["step"], json!(0));
assert_eq!(payload["last_action"], json!("reset"));
transport.stop().await.expect("stop");
}
}