From 6f8eec91354f0f3ee7dca227a8c99cb59fe9d5b2 Mon Sep 17 00:00:00 2001 From: Claude Date: Fri, 11 Sep 2026 17:10:11 +0000 Subject: [PATCH] feat(ocpp-cp): auto-trip a v201 variable monitor on a threshold-crossing SetVariables write (M7) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Closes #573. Completes the OCPP 2.0.1 variable-monitoring loop that #570 (the trip_variable_monitor injection seam) left half-done: a monitor now trips autonomously when a CSMS SetVariables write moves a monitored variable's Actual value across an installed threshold, or by >= a configured delta. Because the trip is raised from inside an inbound-CALL handler, it is drained through the RemoteCommand queue (new V201MonitorTrip variant) rather than emitted inline, so the SetVariablesResponse CALLRESULT is flushed before the outbound NotifyEvent (no receive-loop re-entrancy) — the same discipline NotifyReport / NotifyMonitoringReport already follow. - V201DeviceModel::monitors_tripped_by_write selects only the write-driven monitors a numeric crossing hits (Upper/LowerThreshold cross either direction; Delta on a single-write jump; Periodic never). Hysteresis falls out of crossing detection; non-numeric / NaN / inf values cross nothing but are still reported verbatim. - The SetVariables handler captures the prior Actual value and enqueues one V201MonitorTrip per crossed variable; rejected / unknown / non-Actual / no-op-rewrite writes trip nothing. - The #545 NotifyEvent emitter core is extracted into ChargePoint::emit_monitor_notify_event, shared by the inline seam and the new consumer path — no duplicated trigger/correlation/seqNo logic. Tests: pure device-model selector + predicate + parse tests; handler tests via v201_drain_commands covering all threshold/delta/hysteresis/rejected/ periodic/no-op/non-numeric arms; and an end-to-end mock-CSMS test proving the NotifyEvent is emitted after the response, off the command queue. Co-Authored-By: Claude Opus 4.8 (1M context) Claude-Session: https://claude.ai/code/session_01BSdgXz4N5iiSh1UiXPi7za --- crates/ocpp-cp/src/lib.rs | 635 ++++++++++++++++++++++-- crates/ocpp-cp/src/v201_device_model.rs | 311 +++++++++++- 2 files changed, 893 insertions(+), 53 deletions(-) diff --git a/crates/ocpp-cp/src/lib.rs b/crates/ocpp-cp/src/lib.rs index 2827892..a092f8a 100644 --- a/crates/ocpp-cp/src/lib.rs +++ b/crates/ocpp-cp/src/lib.rs @@ -84,7 +84,7 @@ use v201_certificate_store::V201CertificateStore; use v201_charging_profiles::V201TxProfileStore; use v201_cost::V201CostStore; use v201_customer_information::V201CustomerInformationStore; -use v201_device_model::V201DeviceModel; +use v201_device_model::{TrippedMonitor, V201DeviceModel}; use v201_display_message::V201DisplayMessageStore; use v201_firmware_update::V201FirmwareUpdateStore; use v201_log_upload::V201LogUploadStore; @@ -165,9 +165,9 @@ use ocpp_types::v201::{ NotifyEVChargingNeedsStatusEnumType, OCSPRequestDataType, OperationalStatusEnumType, PublishFirmwareStatusEnumType, RegistrationStatusEnumType, ReportDataType, RequestStartStopStatusEnumType, ReservationUpdateStatusEnumType, ReserveNowStatusEnumType, - ResetStatusEnumType, SetNetworkProfileStatusEnumType, SetVariableResultType, StatusInfoType, - TriggerMessageStatusEnumType, UnlockStatusEnumType, UpdateFirmwareStatusEnumType, - UploadLogStatusEnumType, + ResetStatusEnumType, SetNetworkProfileStatusEnumType, SetVariableResultType, + SetVariableStatusEnumType, StatusInfoType, TriggerMessageStatusEnumType, UnlockStatusEnumType, + UpdateFirmwareStatusEnumType, UploadLogStatusEnumType, }; use serde::{Deserialize, Serialize}; use std::collections::HashMap; @@ -699,6 +699,28 @@ enum RemoteCommand { request_id: i32, monitor_data: Vec, }, + /// Auto-originate a 2.0.1 `NotifyEvent` for the variable monitor(s) an inbound + /// `SetVariables` write tripped by crossing their threshold / delta (Part 2, + /// monitoring; Issue #573). The autonomous, inbound-CALL-triggered twin of the + /// [`trip_variable_monitor`](ChargePoint::trip_variable_monitor) injection + /// seam: because the trip is raised from *inside* the `SetVariables` handler, + /// it must be drained here off the command-consumer task — so the + /// `SetVariablesResponse` CALLRESULT is flushed before the outbound + /// `NotifyEvent` CALL, with no receive-loop re-entrancy (the same discipline as + /// [`V201NotifyReport`](Self::V201NotifyReport) / + /// [`V201NotifyMonitoringReport`](Self::V201NotifyMonitoringReport)), unlike the + /// CP-initiated seam which emits inline. `monitors` is the already-selected + /// crossed subset (computed under the device-model write lock on the CALL path + /// via [`V201DeviceModel::monitors_tripped_by_write`](crate::v201_device_model::V201DeviceModel::monitors_tripped_by_write), + /// so the send touches no shared state) and `actual_value` the new value + /// threaded verbatim onto each event. Emitted through the shared + /// [`emit_monitor_notify_event`](ChargePoint::emit_monitor_notify_event) core, + /// so trigger derivation, `CustomMonitor` tagging, correlation, and the + /// monotonic `seqNo` / `eventId` streams are not duplicated. + V201MonitorTrip { + monitors: Vec, + actual_value: String, + }, /// Stream the installed charging profiles a CSMS asked for with an `Accepted` /// OCPP 2.0.1 `GetChargingProfiles` (Part 2). The synchronous /// `GetChargingProfiles.conf` only reports `Accepted` / `NoProfiles`; the @@ -1916,40 +1938,121 @@ impl ChargePoint { // input. The write is applied on the CALL path (a device-model store // update is cheap and must be visible to the CALLRESULT and to any // subsequent `GetVariables`), serialized by the model's write lock. + // + // Autonomous monitor trip (Issue #573): an accepted write that actually + // changes a monitored variable's `Actual` value and crosses an installed + // threshold / delta monitor auto-originates a `NotifyEvent`. The crossed + // monitors are selected under the same write lock (so the read is + // consistent with the write), but the `NotifyEvent` is emitted off the + // command-consumer task via `RemoteCommand::V201MonitorTrip` — never + // inline — so this `SetVariablesResponse` is flushed first (no + // receive-loop re-entrancy). Rejected / unknown / non-`Actual` / + // value-unchanged writes trip nothing. // Ports `ocpp.v201.call.SetVariables`. if matches!(protocol_version, OcppVersion::V201) { let device_model = v201_device_model.clone(); + let command_sender = command_sender.clone(); d.on(move |req: V201SetVariablesRequest| { let device_model = device_model.clone(); + let command_sender = command_sender.clone(); async move { - let mut model = device_model.write().await; + // Monitor trips to enqueue once the write lock is dropped and + // the CALLRESULT is on its way: the crossed monitors and the + // new value each `NotifyEvent` reports verbatim. + let mut trips: Vec<(Vec, String)> = Vec::new(); + // One result per requested entry, in request order. Each // result echoes the CSMS's original (un-normalized) // `component` / `variable` / `attributeType`; an omitted // `attributeType` resolves to `Actual` for the write but is // echoed back as `None`, mirroring the read seam. - let set_variable_result: Vec = req - .set_variable_data - .iter() - .map(|data| { - let attribute = - data.attribute_type.unwrap_or(AttributeEnumType::Actual); - let attribute_status = model.set( - &data.component, - &data.variable, - attribute, - &data.attribute_value, + let set_variable_result: Vec = { + let mut model = device_model.write().await; + req.set_variable_data + .iter() + .map(|data| { + let attribute = + data.attribute_type.unwrap_or(AttributeEnumType::Actual); + // Snapshot the prior `Actual` value *before* the + // write so a threshold / delta crossing can be + // detected; only the `Actual` attribute drives + // monitors, so other attributes need no snapshot. + let previous = if attribute == AttributeEnumType::Actual { + model + .get( + &data.component, + &data.variable, + AttributeEnumType::Actual, + ) + .1 + } else { + None + }; + let attribute_status = model.set( + &data.component, + &data.variable, + attribute, + &data.attribute_value, + ); + // A write that took effect on `Actual` and + // actually changed the value may cross a monitor. + // `Rejected` / `Unknown*` / `NotSupported` writes + // never mutate, and a no-op rewrite crosses + // nothing, so neither trips. + let took_effect = matches!( + attribute_status, + SetVariableStatusEnumType::Accepted + | SetVariableStatusEnumType::RebootRequired + ); + if took_effect && attribute == AttributeEnumType::Actual { + if let Some(prev) = previous.as_deref() { + if prev != data.attribute_value { + let crossed = model.monitors_tripped_by_write( + &data.component.name, + &data.variable.name, + prev, + &data.attribute_value, + ); + if !crossed.is_empty() { + trips.push((crossed, data.attribute_value.clone())); + } + } + } + } + SetVariableResultType { + attribute_status, + component: data.component.clone(), + variable: data.variable.clone(), + attribute_type: data.attribute_type, + attribute_status_info: None, + custom_data: None, + } + }) + .collect() + // Write lock dropped here, before the trips are enqueued. + }; + + // Drain the collected trips onto the command-consumer task so + // each `NotifyEvent` is emitted *after* this + // `SetVariablesResponse`, never inline. A gone consumer (CP + // shutting down) drops the trip best-effort; the write is + // already applied and the CALLRESULT is honest. + for (monitors, actual_value) in trips { + if command_sender + .send(RemoteCommand::V201MonitorTrip { + monitors, + actual_value, + }) + .is_err() + { + warn!( + "v201 SetVariables: consumer gone, cannot emit NotifyEvent for a \ + tripped variable monitor" ); - SetVariableResultType { - attribute_status, - component: data.component.clone(), - variable: data.variable.clone(), - attribute_type: data.attribute_type, - attribute_status_info: None, - custom_data: None, - } - }) - .collect(); + break; + } + } + Ok(V201SetVariablesResponse { set_variable_result, custom_data: None, @@ -5406,6 +5509,24 @@ impl ChargePoint { cp.send_v201_notify_monitoring_report(request_id, monitor_data) .await; } + RemoteCommand::V201MonitorTrip { + monitors, + actual_value, + } => { + // Emit the auto-trip `NotifyEvent` off the CALL path + // (Issue #573). Best-effort: the triggering + // `SetVariablesResponse` was already returned, so a + // transport / CALLERROR failure here is logged, not + // surfaced. + if let Err(e) = + cp.emit_monitor_notify_event(&monitors, &actual_value).await + { + warn!( + "v201 auto-trip: failed to emit NotifyEvent for a \ + SetVariables threshold crossing: {e}" + ); + } + } RemoteCommand::V201ReportChargingProfiles { request_id, profiles, @@ -6846,8 +6967,54 @@ impl ChargePoint { return Ok(MonitorTripOutcome::NoMonitor); } + let outcome = self + .emit_monitor_notify_event(&matched, actual_value) + .await?; + + info!( + component = %component, + variable = %variable, + // Presence/count only — `actual_value` is opaque and never logged. + outcome = ?outcome, + "originated NotifyEvent for a variable-monitor trip; CSMS acknowledged" + ); + + Ok(outcome) + } + + /// Build and emit a single OCPP 2.0.1 `NotifyEvent` CALL carrying one + /// [`EventDataType`](ocpp_types::v201::EventDataType) per monitor in + /// `monitors`, stamped with the next monotonic `seqNo` and a per-event + /// monotonic `eventId`, all reporting the opaque `actual_value` verbatim. + /// + /// The shared emitter core behind both variable-monitor trip paths — the + /// CP-initiated injection seam + /// [`trip_variable_monitor`](Self::trip_variable_monitor) and the autonomous + /// `SetVariables`-crossing path drained via + /// [`RemoteCommand::V201MonitorTrip`] — so trigger derivation + /// ([`v201_command::v201_event_trigger_for_monitor`]), `CustomMonitor` + /// tagging, `variableMonitoringId` correlation, and the monotonic `seqNo` / + /// `eventId` streams live in exactly one place (Issue #573; no duplicated trip + /// logic). + /// + /// An empty `monitors` slice emits nothing and claims **no** `seqNo` + /// (returns [`MonitorTripOutcome::NoMonitor`]), so a caller that finds no + /// match never burns a sequence position and leaves no phantom gap in the + /// CSMS's view of the stream. The `actual_value` is opaque, caller-supplied + /// text threaded to the wire verbatim; an over-long value surfaces as an + /// [`OcppError`] from [`call`](Self::call)'s outbound validation, never a + /// panic or a silent truncation. + async fn emit_monitor_notify_event( + &self, + monitors: &[TrippedMonitor], + actual_value: &str, + ) -> OcppResult { + if monitors.is_empty() { + return Ok(MonitorTripOutcome::NoMonitor); + } + let generated_at = v201_now(); - let event_data: Vec<_> = matched + let event_data: Vec<_> = monitors .iter() .map(|m| { v201_command::v201_monitor_event_data( @@ -6869,15 +7036,6 @@ impl ChargePoint { // the outbound request and the ack. Discard it on success. let _ack = self.call(request).await?; - info!( - component = %component, - variable = %variable, - events, - seq_no, - // Presence/count only — `actual_value` is opaque and never logged. - "originated NotifyEvent for a variable-monitor trip; CSMS acknowledged" - ); - Ok(MonitorTripOutcome::Emitted { events, seq_no }) } @@ -13051,6 +13209,411 @@ mod tests { )); } + // --- SetVariables auto-trips a variable monitor (Issue #573) --- + // + // A CSMS `SetVariables` write that moves a monitored variable's Actual value + // across an installed threshold / delta monitor auto-originates a + // `NotifyEvent`, drained off the command-consumer task (never inline) so the + // `SetVariablesResponse` is flushed first. `OCPPCommCtrlr` / `HeartbeatInterval` + // seeds to "300" and is writable, so it is the fixture for these. + + /// Build a V201 `SetVariables` CALL writing one component-variable's `Actual` + /// attribute (attributeType omitted → `Actual`). + fn make_v201_set_variables(component: &str, variable: &str, value: &str) -> CallMessage { + make_call(V201SetVariablesRequest { + set_variable_data: vec![ocpp_types::v201::SetVariableDataType { + attribute_value: value.to_string(), + component: ocpp_types::v201::ComponentType { + name: component.to_string(), + instance: None, + evse: None, + custom_data: None, + }, + variable: ocpp_types::v201::VariableType { + name: variable.to_string(), + instance: None, + custom_data: None, + }, + attribute_type: None, + custom_data: None, + }], + custom_data: None, + }) + } + + /// Install one CSMS monitor of `kind` with threshold/delta `value` on + /// `OCPPCommCtrlr` / `variable`, returning its station-assigned id. The valued + /// twin of [`install_one_monitor`] for the auto-trip tests, which need a + /// specific threshold. + async fn install_valued_monitor( + cp: &ChargePoint, + variable: &str, + kind: ocpp_types::v201::MonitorEnumType, + value: f64, + ) -> i32 { + use ocpp_types::v201::SetMonitoringStatusEnumType; + let resp = cp + .handle_message(Message::Call(make_v201_set_variable_monitoring(vec![ + v201_monitor_data("OCPPCommCtrlr", variable, kind, value, 3), + ]))) + .await + .unwrap(); + match resp.unwrap() { + Message::CallResult(r) => { + let body: ocpp_messages::v201::SetVariableMonitoringResponse = + r.payload_as().unwrap(); + assert_eq!( + body.set_monitoring_result[0].status, + SetMonitoringStatusEnumType::Accepted + ); + body.set_monitoring_result[0].id.expect("accepted → id") + } + other => panic!("expected CallResult, got: {other:?}"), + } + } + + /// Project the drained commands into just the auto-trips as + /// `(sorted monitor ids, reported actualValue)`, dropping any other queued + /// side effect — the concise shape the assertions below compare against. + fn monitor_trips(commands: &[RemoteCommand]) -> Vec<(Vec, String)> { + commands + .iter() + .filter_map(|c| match c { + RemoteCommand::V201MonitorTrip { + monitors, + actual_value, + } => Some(( + monitors.iter().map(|m| m.id).collect::>(), + actual_value.clone(), + )), + _ => None, + }) + .collect() + } + + #[tokio::test] + async fn v201_set_variables_upper_threshold_fires_on_each_crossing_with_hysteresis() { + use ocpp_types::v201::MonitorEnumType; + let cp = ChargePoint::new(ChargePointConfig::for_version(OcppVersion::V201)).unwrap(); + // Alerting band = value > 900. Seed is "300" (below the band). + let id = install_valued_monitor( + &cp, + "HeartbeatInterval", + MonitorEnumType::UpperThreshold, + 900.0, + ) + .await; + + for value in ["1000", "1100", "800", "700"] { + cp.handle_message(Message::Call(make_v201_set_variables( + "OCPPCommCtrlr", + "HeartbeatInterval", + value, + ))) + .await + .unwrap(); + } + + // "300"→"1000" enters the band (trip); "1000"→"1100" stays inside + // (hysteresis, no trip); "1100"→"800" leaves the band (trip); + // "800"→"700" stays below (no trip). Two crossings → two trips. + let commands = v201_drain_commands(&cp).await; + assert_eq!( + monitor_trips(&commands), + vec![ + (vec![id], "1000".to_string()), + (vec![id], "800".to_string()) + ] + ); + } + + #[tokio::test] + async fn v201_set_variables_lower_threshold_is_symmetric() { + use ocpp_types::v201::MonitorEnumType; + let cp = ChargePoint::new(ChargePointConfig::for_version(OcppVersion::V201)).unwrap(); + // Alerting band = value < 200. Seed "300" is above the band. + let id = install_valued_monitor( + &cp, + "HeartbeatInterval", + MonitorEnumType::LowerThreshold, + 200.0, + ) + .await; + + for value in ["100", "150", "300"] { + cp.handle_message(Message::Call(make_v201_set_variables( + "OCPPCommCtrlr", + "HeartbeatInterval", + value, + ))) + .await + .unwrap(); + } + + // "300"→"100" enters (trip); "100"→"150" stays below 200 (no trip); + // "150"→"300" leaves (trip). + let commands = v201_drain_commands(&cp).await; + assert_eq!( + monitor_trips(&commands), + vec![(vec![id], "100".to_string()), (vec![id], "300".to_string())] + ); + } + + #[tokio::test] + async fn v201_set_variables_delta_fires_on_jump_not_subdelta() { + use ocpp_types::v201::MonitorEnumType; + let cp = ChargePoint::new(ChargePointConfig::for_version(OcppVersion::V201)).unwrap(); + let id = + install_valued_monitor(&cp, "HeartbeatInterval", MonitorEnumType::Delta, 50.0).await; + + for value in ["360", "370", "500"] { + cp.handle_message(Message::Call(make_v201_set_variables( + "OCPPCommCtrlr", + "HeartbeatInterval", + value, + ))) + .await + .unwrap(); + } + + // "300"→"360" jumps 60 ≥ 50 (trip); "360"→"370" jumps 10 < 50 (no trip); + // "370"→"500" jumps 130 (trip). A delta is measured per single write. + let commands = v201_drain_commands(&cp).await; + assert_eq!( + monitor_trips(&commands), + vec![(vec![id], "360".to_string()), (vec![id], "500".to_string())] + ); + } + + #[tokio::test] + async fn v201_set_variables_multiple_monitors_ride_one_id_sorted_trip() { + use ocpp_types::v201::MonitorEnumType; + let cp = ChargePoint::new(ChargePointConfig::for_version(OcppVersion::V201)).unwrap(); + // Two monitors on the same variable: an upper threshold and a delta, both + // crossed by the same "300"→"1000" write. + let upper = install_valued_monitor( + &cp, + "HeartbeatInterval", + MonitorEnumType::UpperThreshold, + 900.0, + ) + .await; + let delta = + install_valued_monitor(&cp, "HeartbeatInterval", MonitorEnumType::Delta, 50.0).await; + + cp.handle_message(Message::Call(make_v201_set_variables( + "OCPPCommCtrlr", + "HeartbeatInterval", + "1000", + ))) + .await + .unwrap(); + + // Both cross on the one write → a single trip carrying both monitors, + // id-sorted (matching `monitors_for_variable`), so the emitter builds one + // NotifyEvent with two correlated events. + let mut ids = vec![upper, delta]; + ids.sort_unstable(); + let commands = v201_drain_commands(&cp).await; + assert_eq!(monitor_trips(&commands), vec![(ids, "1000".to_string())]); + } + + #[tokio::test] + async fn v201_set_variables_rejected_write_never_trips() { + use ocpp_types::v201::{MonitorEnumType, SetMonitoringStatusEnumType}; + let cp = ChargePoint::new(ChargePointConfig::for_version(OcppVersion::V201)).unwrap(); + // SecurityCtrlr/MaxCertificateChainSize is read-only (seed "3"); install a + // monitor on it that *would* cross, then attempt a write. The write is + // Rejected → no mutation → no trip. + let install = cp + .handle_message(Message::Call(make_v201_set_variable_monitoring(vec![ + v201_monitor_data( + "SecurityCtrlr", + "MaxCertificateChainSize", + MonitorEnumType::UpperThreshold, + 2.0, + 3, + ), + ]))) + .await + .unwrap(); + match install.unwrap() { + Message::CallResult(r) => { + let body: ocpp_messages::v201::SetVariableMonitoringResponse = + r.payload_as().unwrap(); + assert_eq!( + body.set_monitoring_result[0].status, + SetMonitoringStatusEnumType::Accepted, + "a monitor can be installed on a read-only variable", + ); + } + other => panic!("expected CallResult, got: {other:?}"), + } + let resp = cp + .handle_message(Message::Call(make_call(V201SetVariablesRequest { + set_variable_data: vec![ocpp_types::v201::SetVariableDataType { + attribute_value: "10".to_string(), + component: ocpp_types::v201::ComponentType { + name: "SecurityCtrlr".to_string(), + instance: None, + evse: None, + custom_data: None, + }, + variable: ocpp_types::v201::VariableType { + name: "MaxCertificateChainSize".to_string(), + instance: None, + custom_data: None, + }, + attribute_type: None, + custom_data: None, + }], + custom_data: None, + }))) + .await + .unwrap(); + match resp.unwrap() { + Message::CallResult(r) => { + let body: ocpp_messages::v201::SetVariablesResponse = r.payload_as().unwrap(); + assert_eq!( + body.set_variable_result[0].attribute_status, + SetVariableStatusEnumType::Rejected + ); + } + other => panic!("expected CallResult, got: {other:?}"), + } + assert!( + monitor_trips(&v201_drain_commands(&cp).await).is_empty(), + "a Rejected write must not trip" + ); + } + + #[tokio::test] + async fn v201_set_variables_non_monitored_write_never_trips() { + // A write that changes a variable no monitor watches trips nothing. + let cp = ChargePoint::new(ChargePointConfig::for_version(OcppVersion::V201)).unwrap(); + cp.handle_message(Message::Call(make_v201_set_variables( + "OCPPCommCtrlr", + "HeartbeatInterval", + "1000", + ))) + .await + .unwrap(); + assert!(monitor_trips(&v201_drain_commands(&cp).await).is_empty()); + } + + #[tokio::test] + async fn v201_set_variables_periodic_monitor_never_trips() { + use ocpp_types::v201::MonitorEnumType; + // Periodic monitors are time-driven, not write-driven — a crossing write + // never trips them. + let cp = ChargePoint::new(ChargePointConfig::for_version(OcppVersion::V201)).unwrap(); + install_valued_monitor(&cp, "HeartbeatInterval", MonitorEnumType::Periodic, 60.0).await; + cp.handle_message(Message::Call(make_v201_set_variables( + "OCPPCommCtrlr", + "HeartbeatInterval", + "1000", + ))) + .await + .unwrap(); + assert!(monitor_trips(&v201_drain_commands(&cp).await).is_empty()); + } + + #[tokio::test] + async fn v201_set_variables_no_op_rewrite_does_not_trip() { + use ocpp_types::v201::MonitorEnumType; + // Rewriting the seed value ("300") is Accepted but changes nothing, so it + // crosses nothing — no phantom trip. + let cp = ChargePoint::new(ChargePointConfig::for_version(OcppVersion::V201)).unwrap(); + install_valued_monitor( + &cp, + "HeartbeatInterval", + MonitorEnumType::UpperThreshold, + 900.0, + ) + .await; + cp.handle_message(Message::Call(make_v201_set_variables( + "OCPPCommCtrlr", + "HeartbeatInterval", + "300", + ))) + .await + .unwrap(); + assert!(monitor_trips(&v201_drain_commands(&cp).await).is_empty()); + } + + #[tokio::test] + async fn v201_set_variables_non_numeric_value_trips_nothing_without_panic() { + use ocpp_types::v201::MonitorEnumType; + // Trust boundary: a monitored variable can be written a non-numeric value. + // It is stored verbatim and cannot cross a numeric threshold, so it trips + // nothing — and never panics on the parse. + let cp = ChargePoint::new(ChargePointConfig::for_version(OcppVersion::V201)).unwrap(); + install_valued_monitor(&cp, "HeartbeatInterval", MonitorEnumType::Delta, 50.0).await; + cp.handle_message(Message::Call(make_v201_set_variables( + "OCPPCommCtrlr", + "HeartbeatInterval", + "not-a-number", + ))) + .await + .unwrap(); + assert!(monitor_trips(&v201_drain_commands(&cp).await).is_empty()); + } + + #[tokio::test] + async fn v201_set_variables_crossing_emits_notify_event_after_response() { + use ocpp_types::v201::MonitorEnumType; + // End-to-end over a real socket: the crossing NotifyEvent is emitted off + // the command-consumer task *after* the SetVariablesResponse, correlated + // to the installed monitor. + let (addr, mut rx) = spawn_mock_csms_capturing(notify_event_routes()).await; + let cp = ChargePoint::new(ChargePointConfig { + central_system_url: format!("ws://{addr}"), + ..ChargePointConfig::for_version(OcppVersion::V201) + }) + .unwrap(); + cp.connect().await.unwrap(); + let id = install_valued_monitor( + &cp, + "HeartbeatInterval", + MonitorEnumType::UpperThreshold, + 900.0, + ) + .await; + + // The write is answered synchronously (Accepted) … + let resp = cp + .handle_message(Message::Call(make_v201_set_variables( + "OCPPCommCtrlr", + "HeartbeatInterval", + "1000", + ))) + .await + .unwrap(); + match resp.unwrap() { + Message::CallResult(r) => { + let body: ocpp_messages::v201::SetVariablesResponse = r.payload_as().unwrap(); + assert_eq!( + body.set_variable_result[0].attribute_status, + SetVariableStatusEnumType::Accepted + ); + } + other => panic!("expected CallResult, got: {other:?}"), + } + + // … and the NotifyEvent follows on the wire, correlated to the monitor. + let payload = recv_notify_event(&mut rx).await; + assert_eq!(payload["seqNo"], 0); + let events = payload["eventData"].as_array().unwrap(); + assert_eq!(events.len(), 1); + let e = &events[0]; + assert_eq!(e["trigger"], "Alerting"); + assert_eq!(e["eventNotificationType"], "CustomMonitor"); + assert_eq!(e["actualValue"], "1000"); + assert_eq!(e["variableMonitoringId"], id); + assert_eq!(e["component"]["name"], "OCPPCommCtrlr"); + assert_eq!(e["variable"]["name"], "HeartbeatInterval"); + } + #[tokio::test] async fn request_notify_ev_charging_schedule_surfaces_accepted_status() { let mut routes = std::collections::HashMap::new(); diff --git a/crates/ocpp-cp/src/v201_device_model.rs b/crates/ocpp-cp/src/v201_device_model.rs index e02796f..b496363 100644 --- a/crates/ocpp-cp/src/v201_device_model.rs +++ b/crates/ocpp-cp/src/v201_device_model.rs @@ -136,6 +136,57 @@ struct MonitorEntry { monitor: VariableMonitoringType, } +impl MonitorEntry { + /// Project this entry into the [`TrippedMonitor`] an emitter needs — the + /// station-assigned id, the monitor kind, and the display-form + /// component / variable. Shared by the lookup paths so both project a + /// matched monitor identically. + fn to_tripped(&self) -> TrippedMonitor { + TrippedMonitor { + id: self.monitor.id, + kind: self.monitor.kind, + component: self.component.clone(), + variable: self.variable.clone(), + } + } +} + +/// Parse an opaque, stored variable value as a finite `f64` for numeric monitor +/// evaluation, or `None` when it is not a finite number. +/// +/// Variable values are opaque strings on the wire (a monitored variable's value +/// is not guaranteed numeric). Threshold and delta monitors are numeric, so a +/// value that does not parse as a finite `f64` — non-numeric text, or an +/// explicit `NaN` / `inf` — participates in no numeric crossing. Leading and +/// trailing ASCII whitespace is tolerated to match lenient CSMS input. +fn parse_finite_f64(value: &str) -> Option { + value.trim().parse::().ok().filter(|v| v.is_finite()) +} + +/// Whether a single write moving a numeric value `prev → next` should trip a +/// monitor of `kind` configured with magnitude `value`. +/// +/// The match is exhaustive without a wildcard so a newly-added +/// [`MonitorEnumType`] is a compile error to classify here rather than a silent +/// mis-trip. +fn monitor_write_trips(kind: MonitorEnumType, value: f64, prev: f64, next: f64) -> bool { + match kind { + // Alerting bands. A crossing is the two samples sitting on opposite + // sides of the threshold. Strict `>` (upper) / `<` (lower) treat a value + // exactly on the threshold as *outside* the alerting band in both + // directions, so a write onto the threshold from inside the band fires + // (leaving) while one from outside does not (no phantom entry). + MonitorEnumType::UpperThreshold => (prev > value) != (next > value), + MonitorEnumType::LowerThreshold => (prev < value) != (next < value), + // A delta monitor fires when this single write jumps by at least the + // configured magnitude. `value` is a magnitude; compare the absolute + // change against its absolute value defensively. + MonitorEnumType::Delta => (next - prev).abs() >= value.abs(), + // Time-driven kinds are never write-driven. + MonitorEnumType::Periodic | MonitorEnumType::PeriodicClockAligned => false, + } +} + /// One installed variable monitor that a /// [`trip`](V201DeviceModel::monitors_for_variable) matched — the minimal /// projection of a `MonitorEntry` a `NotifyEvent` emitter needs to build a @@ -769,8 +820,79 @@ impl V201DeviceModel { /// An empty result means no monitor watches that identity — the caller emits /// nothing (a no-op trip), never a panic. pub fn monitors_for_variable(&self, component: &str, variable: &str) -> Vec { - // Build the lookup key through the same normalization as install / - // snapshot, from a name-only component / variable (no instance, no EVSE). + let key = Self::monitor_lookup_key(component, variable); + let mut matches: Vec = self + .monitors + .values() + .filter(|entry| entry.key == key) + .map(MonitorEntry::to_tripped) + .collect(); + matches.sort_by_key(|m| m.id); + matches + } + + /// The subset of installed monitors on the (`component`, `variable`) identity + /// that a `SetVariables` write moving the variable's `Actual` value from + /// `previous` to `new` should **autonomously trip** (Issue #573) — the + /// write-driven counterpart to the injection-seam lookup + /// [`monitors_for_variable`](Self::monitors_for_variable), sorted by id so the + /// emitted event order is deterministic. + /// + /// Only the write-driven monitor kinds are ever selected: + /// - [`UpperThreshold`](MonitorEnumType::UpperThreshold) — trips when the + /// value **crosses** the threshold in either direction (enters the alerting + /// band `value > threshold`, or leaves it). + /// - [`LowerThreshold`](MonitorEnumType::LowerThreshold) — symmetric + /// (`value < threshold`). + /// - [`Delta`](MonitorEnumType::Delta) — trips when the single write jumps by + /// at least the configured delta (`|new − previous| ≥ value`). + /// + /// [`Periodic`](MonitorEnumType::Periodic) / + /// [`PeriodicClockAligned`](MonitorEnumType::PeriodicClockAligned) monitors are + /// time-driven, never write-driven, so they are never returned here. + /// + /// **Hysteresis** falls out of crossing detection: a write that leaves the + /// value on the same side of a threshold (or moves it by less than the delta) + /// crosses nothing and returns no monitor, so repeated same-band writes raise + /// no duplicate event. **Numeric evaluation only:** thresholds and deltas are + /// numeric ([`value: f64`](VariableMonitoringType::value)), so a `previous` / + /// `new` that does not parse as a finite `f64` can neither cross a threshold + /// nor exceed a delta and trips nothing — the opaque string is still reported + /// verbatim by the emitter, but a non-numeric (or `NaN`/infinite) value never + /// fabricates a numeric crossing. + pub fn monitors_tripped_by_write( + &self, + component: &str, + variable: &str, + previous: &str, + new: &str, + ) -> Vec { + // A threshold can only be crossed between two finite numeric samples; an + // unparseable value trips nothing (the value is still reported verbatim). + let (Some(prev), Some(next)) = (parse_finite_f64(previous), parse_finite_f64(new)) else { + return Vec::new(); + }; + let key = Self::monitor_lookup_key(component, variable); + let mut matches: Vec = self + .monitors + .values() + .filter(|entry| entry.key == key) + .filter(|entry| { + monitor_write_trips(entry.monitor.kind, entry.monitor.value, prev, next) + }) + .map(MonitorEntry::to_tripped) + .collect(); + matches.sort_by_key(|m| m.id); + matches + } + + /// The normalized [`VariableKey`] a name-only (`component`, `variable`) + /// identity resolves to — the same normalization the install / snapshot paths + /// use, so lookups are case-insensitive and match the station-wide, + /// un-instanced identities the simulator seeds. Shared by + /// [`monitors_for_variable`](Self::monitors_for_variable) and + /// [`monitors_tripped_by_write`](Self::monitors_tripped_by_write). + fn monitor_lookup_key(component: &str, variable: &str) -> VariableKey { let component_ty = ComponentType { name: component.to_string(), instance: None, @@ -782,21 +904,7 @@ impl V201DeviceModel { instance: None, custom_data: None, }; - let key = VariableKey::from_request(&component_ty, &variable_ty); - - let mut matches: Vec = self - .monitors - .values() - .filter(|entry| entry.key == key) - .map(|entry| TrippedMonitor { - id: entry.monitor.id, - kind: entry.monitor.kind, - component: entry.component.clone(), - variable: entry.variable.clone(), - }) - .collect(); - matches.sort_by_key(|m| m.id); - matches + VariableKey::from_request(&component_ty, &variable_ty) } /// Whether a monitor of the given [`MonitorEnumType`] falls under a requested @@ -2271,4 +2379,173 @@ mod tests { .monitors_for_variable("OCPPCommCtrlr", "NoSuchVariable") .is_empty()); } + + // --- monitors_tripped_by_write: autonomous write-driven trips (#573) ---- + + #[test] + fn monitors_tripped_by_write_selects_only_crossed_write_driven_monitors() { + let mut model = V201DeviceModel::with_standard_profile(); + // An upper threshold (900), a delta (50), and a periodic (60) on the same + // variable. Only the write-driven kinds a crossing hits are selected. + let results = model.install_monitors(&[ + monitor_data( + "OCPPCommCtrlr", + "HeartbeatInterval", + MonitorEnumType::UpperThreshold, + 900.0, + 3, + ), + monitor_data( + "OCPPCommCtrlr", + "HeartbeatInterval", + MonitorEnumType::Delta, + 50.0, + 3, + ), + monitor_data( + "OCPPCommCtrlr", + "HeartbeatInterval", + MonitorEnumType::Periodic, + 60.0, + 3, + ), + ]); + let ids: Vec = results.iter().map(|r| r.id.unwrap()).collect(); + let (upper, delta) = (ids[0], ids[1]); + + // "300"→"1000": crosses the upper threshold *and* jumps ≥ 50 → both, + // id-sorted; the periodic monitor is never selected. + let tripped = + model.monitors_tripped_by_write("OCPPCommCtrlr", "HeartbeatInterval", "300", "1000"); + assert_eq!(tripped.iter().map(|m| m.id).collect::>(), { + let mut v = vec![upper, delta]; + v.sort_unstable(); + v + }); + + // "300"→"320": below the threshold and a sub-delta jump → nothing. + assert!(model + .monitors_tripped_by_write("OCPPCommCtrlr", "HeartbeatInterval", "300", "320") + .is_empty()); + + // Non-numeric new value can cross nothing (still reported verbatim by the + // emitter, but no numeric trip). + assert!(model + .monitors_tripped_by_write("OCPPCommCtrlr", "HeartbeatInterval", "300", "oops") + .is_empty()); + + // A variable no monitor watches trips nothing. + assert!(model + .monitors_tripped_by_write("OCPPCommCtrlr", "NoSuchVariable", "1", "100000") + .is_empty()); + } + + #[test] + fn monitor_write_trips_predicate_covers_bands_deltas_and_time_driven() { + // Upper threshold 10: a crossing in either direction fires; staying on one + // side does not; sitting exactly on the threshold is treated as outside + // the alerting band (so entering from the value itself does not fire, and + // leaving onto it does). + assert!(monitor_write_trips( + MonitorEnumType::UpperThreshold, + 10.0, + 5.0, + 15.0 + )); + assert!(monitor_write_trips( + MonitorEnumType::UpperThreshold, + 10.0, + 15.0, + 5.0 + )); + assert!(!monitor_write_trips( + MonitorEnumType::UpperThreshold, + 10.0, + 11.0, + 12.0 + )); + assert!(!monitor_write_trips( + MonitorEnumType::UpperThreshold, + 10.0, + 5.0, + 10.0 + )); + assert!(monitor_write_trips( + MonitorEnumType::UpperThreshold, + 10.0, + 15.0, + 10.0 + )); + + // Lower threshold 10: symmetric (alerting band = value < 10). + assert!(monitor_write_trips( + MonitorEnumType::LowerThreshold, + 10.0, + 15.0, + 5.0 + )); + assert!(monitor_write_trips( + MonitorEnumType::LowerThreshold, + 10.0, + 5.0, + 15.0 + )); + assert!(!monitor_write_trips( + MonitorEnumType::LowerThreshold, + 10.0, + 5.0, + 8.0 + )); + + // Delta 50: a jump ≥ 50 (either sign) fires; a smaller jump does not. + assert!(monitor_write_trips( + MonitorEnumType::Delta, + 50.0, + 300.0, + 360.0 + )); + assert!(monitor_write_trips( + MonitorEnumType::Delta, + 50.0, + 360.0, + 300.0 + )); + assert!(monitor_write_trips( + MonitorEnumType::Delta, + 50.0, + 300.0, + 350.0 + )); + assert!(!monitor_write_trips( + MonitorEnumType::Delta, + 50.0, + 300.0, + 349.0 + )); + + // Time-driven kinds are never write-driven. + assert!(!monitor_write_trips( + MonitorEnumType::Periodic, + 0.0, + 1.0, + 1_000.0 + )); + assert!(!monitor_write_trips( + MonitorEnumType::PeriodicClockAligned, + 0.0, + 1.0, + 1_000.0 + )); + } + + #[test] + fn parse_finite_f64_rejects_non_finite_and_non_numeric() { + assert_eq!(parse_finite_f64("300"), Some(300.0)); + assert_eq!(parse_finite_f64(" 42.5 "), Some(42.5)); + assert_eq!(parse_finite_f64("-7"), Some(-7.0)); + assert_eq!(parse_finite_f64(""), None); + assert_eq!(parse_finite_f64("not-a-number"), None); + assert_eq!(parse_finite_f64("NaN"), None); + assert_eq!(parse_finite_f64("inf"), None); + } }