diff options
| author | Luke Hoersten <[email protected]> | 2026-07-28 18:09:57 -0500 |
|---|---|---|
| committer | Luke Hoersten <[email protected]> | 2026-07-28 18:09:57 -0500 |
| commit | cf5eabb9647e6d53ec29a3f7de599a38dbfd5490 (patch) | |
| tree | 89fcd6f00437c168f042760cc605277a957efc16 /src/controller | |
| parent | 4129dcaeceab060d9f1cac599e9e29cab0f44fd7 (diff) | |
Merge the two device state files and reduce mutation
The registry (device names) and service-state (sync results) were two
JSON files keyed the same way, split only for a TypeScript-compatibility
that is moot now that the fabric identity cannot migrate between
implementations. Merge them into one devices.json with a DeviceRecord per
node, which also removes the parallel-map join in status and collapses the
per-command load/store pairs. identity.json stays separate as a write-once
secret.
Also, from a reduction/immutability review:
- Fix with_timeout passing a literal "{what}" instead of the error context
- Drop needless mut: CompactDuration and parse_wall_clock become immutable
expressions; run_status filters with a predicate instead of a mutator;
sync_one extracts write_zone_and_dst; InspectOutcome gains Default plus a
failed() constructor
- Share a join_ids helper and a devices_path helper; add IdentityReport
From; minor combinator tidy-ups
35 unit tests, clippy clean.
Diffstat (limited to 'src/controller')
| -rw-r--r-- | src/controller/mod.rs | 17 | ||||
| -rw-r--r-- | src/controller/ops.rs | 218 |
2 files changed, 117 insertions, 118 deletions
diff --git a/src/controller/mod.rs b/src/controller/mod.rs index c85efdd..739615c 100644 --- a/src/controller/mod.rs +++ b/src/controller/mod.rs @@ -100,7 +100,8 @@ pub enum IdentityStatus { /// Inspects the identity pair on disk. A fabric blob is any k_* file in the /// matter/ KV directory. pub fn identity_status(storage: &Path) -> IdentityStatus { - let identity_exists = identity_path(storage).exists(); + let identity_file = identity_path(storage); + let identity_exists = identity_file.exists(); let fabric_exists = std::fs::read_dir(matter_kv_path(storage)) .map(|entries| { entries.flatten().any(|e| { @@ -119,7 +120,7 @@ pub fn identity_status(storage: &Path) -> IdentityStatus { (false, true) => { IdentityStatus::Inconsistent("fabric storage exists but identity.json is missing") } - (true, true) => match std::fs::read_to_string(identity_path(storage)) + (true, true) => match std::fs::read_to_string(&identity_file) .ok() .and_then(|raw| serde_json::from_str::<Identity>(&raw).ok()) { @@ -349,8 +350,8 @@ fn ensure_fabric<C: Crypto>( let identity_file = identity_path(&config.storage_path); let existing = matter.with_state(|state| state.fabrics.iter().map(|f| f.fab_idx()).next()); - let registry: crate::state::NodeRegistry = - state::load_or_default(&state::registry_path(&config.storage_path)); + let registry: crate::state::Devices = + state::load_or_default(&state::devices_path(&config.storage_path)); if let Some(fab_idx) = existing { match std::fs::read_to_string(&identity_file) { @@ -602,13 +603,7 @@ fn prepare_storage_dir(path: &Path) -> anyhow::Result<()> { /// so is safe to claim and tighten. #[cfg(unix)] fn is_our_storage_dir(path: &Path) -> bool { - const MARKERS: [&str; 5] = [ - "identity.json", - "matter", - ".lock", - "nodes.json", - "service-state.json", - ]; + const MARKERS: [&str; 4] = ["identity.json", "matter", ".lock", "devices.json"]; if MARKERS.iter().any(|m| path.join(m).exists()) { return true; } diff --git a/src/controller/ops.rs b/src/controller/ops.rs index 7693650..7db1431 100644 --- a/src/controller/ops.rs +++ b/src/controller/ops.rs @@ -17,7 +17,11 @@ use rs_matter::transport::exchange::Exchange; use jiff::Timestamp; use crate::pairing::Onboarding; -use crate::state::{NodeInfo, NodeRegistry}; + +/// The device-registry file path for this run's storage. +fn devices_path<C: Crypto>(ctx: &Ctx<'_, C>) -> std::path::PathBuf { + state::devices_path(&ctx.config.storage_path) +} const BROWSE_TIMEOUT_MS: u32 = 30_000; /// Per-phase bound for commissioning (PASE handshake + invokes, CASE + complete). @@ -38,7 +42,7 @@ async fn with_timeout<T>( embassy_time::Duration::from_secs(secs) )); match embassy_futures::select::select(&mut fut, &mut timer).await { - embassy_futures::select::Either::First(result) => result.ctx("{what}"), + embassy_futures::select::Either::First(result) => result.ctx(what), embassy_futures::select::Either::Second(()) => bail!("{what} timed out after {secs}s"), } } @@ -124,30 +128,18 @@ impl Op for CommissionOp { ) .await?; - // Cache the device identity for the local registry; best-effort. + // Record the device: cached names plus the initial connection. let (vendor_name, product_name) = read_device_names(ctx, device_node_id).await; - - let registry_file = state::registry_path(&ctx.config.storage_path); - let mut registry: NodeRegistry = state::load_or_default(®istry_file); - registry.nodes.insert( - device_node_id.to_string(), - NodeInfo { - vendor_name: vendor_name.clone(), - product_name: product_name.clone(), - }, - ); - state::store(®istry_file, ®istry)?; + state::update_device(&devices_path(ctx), device_node_id, |device| { + device.vendor_name = vendor_name.clone(); + device.product_name = product_name.clone(); + device.last_successful_connection = Some(Timestamp::now()); + })?; let mut identity = ctx.identity.clone(); identity.next_device_node_id += 1; state::store(&identity_path(&ctx.config.storage_path), &identity)?; - state::update_node_state( - &state::service_state_path(&ctx.config.storage_path), - device_node_id, - |node| node.last_successful_connection = Some(Timestamp::now()), - )?; - Ok(CommissionOutcome { node_id: device_node_id.into(), fabric_id: ctx.identity.fabric_id.into(), @@ -314,23 +306,23 @@ impl Op for SyncOp { type Out = Vec<SyncOutcome>; async fn run<C: Crypto>(self, ctx: &Ctx<'_, C>) -> anyhow::Result<Vec<SyncOutcome>> { - let state_path = state::service_state_path(&ctx.config.storage_path); + let state_path = devices_path(ctx); let mut outcomes = Vec::with_capacity(self.targets.len()); for node_id in self.targets { - state::update_node_state(&state_path, node_id, |n| { - n.last_attempted_sync = Some(Timestamp::now()); + state::update_device(&state_path, node_id, |d| { + d.last_attempted_sync = Some(Timestamp::now()); })?; let outcome = match sync_one(ctx, node_id, self.manual_time).await { Ok(outcome) => outcome, Err(error) => SyncOutcome::failed(node_id, format!("{error:#}")), }; - state::update_node_state(&state_path, node_id, |n| { + state::update_device(&state_path, node_id, |d| { if outcome.success { - n.last_successful_sync = Some(Timestamp::now()); - n.last_successful_connection = Some(Timestamp::now()); - n.last_error = None; + d.last_successful_sync = Some(Timestamp::now()); + d.last_successful_connection = Some(Timestamp::now()); + d.last_error = None; } else { - n.last_error = outcome.error.clone(); + d.last_error = outcome.error.clone(); } })?; if let Some(error) = &outcome.error { @@ -398,67 +390,11 @@ async fn sync_one<C: Crypto>( .ctx("SetUTCTime rejected")?; log::info!("Node {node_id}: SetUTCTime {utc_write}"); - let mut time_zone_written = None; - let mut dst_entries_written = 0; - if has_time_zone { - let zone_entries = build_time_zone_list(&tz, &ctx.config.timezone, now_ts); - let response = connect(ctx, node_id) - .await? - .time_synchronization() - .set_time_zone(ROOT_ENDPOINT_ID, |b| { - let mut list = b.time_zone()?; - for entry in &zone_entries { - list = list - .push()? - .offset(entry.offset_seconds)? - .valid_at(entry.valid_at.0)? - .name(Some(&entry.name))? - .end()?; - } - list.end()?.end() - }) - .await - .ctx("SetTimeZone rejected")?; - let dst_required = response - .response() - .map(|r| r.dst_offset_required().unwrap_or(true)) - .unwrap_or(true); - response.complete().await.ctx("SetTimeZone completion")?; - time_zone_written = Some(ctx.config.timezone.clone()); - log::info!("Node {node_id}: SetTimeZone {}", ctx.config.timezone); - - if dst_required { - let max_entries = connect(ctx, node_id) - .await? - .time_synchronization() - .dst_offset_list_max_size_read(ROOT_ENDPOINT_ID) - .await - .ctx("read dstOffsetListMaxSize")?; - let dst_entries = build_dst_offset_list(&tz, usize::from(max_entries), now_ts); - connect(ctx, node_id) - .await? - .time_synchronization() - .set_dst_offset(ROOT_ENDPOINT_ID, |b| { - let mut list = b.dst_offset()?; - for entry in &dst_entries { - list = list - .push()? - .offset(entry.offset_seconds)? - .valid_starting(entry.valid_starting.0)? - .valid_until(match entry.valid_until { - Some(until) => rs_matter::tlv::Nullable::some(until.0), - None => rs_matter::tlv::Nullable::none(), - })? - .end()?; - } - list.end()?.end() - }) - .await - .ctx("SetDSTOffset rejected")?; - dst_entries_written = dst_entries.len(); - log::info!("Node {node_id}: SetDSTOffset with {dst_entries_written} entries"); - } - } + let (time_zone_written, dst_entries_written) = if has_time_zone { + write_zone_and_dst(ctx, node_id, &tz, now_ts).await? + } else { + (None, 0) + }; // Verify by reading the clock back. let after = connect(ctx, node_id) @@ -495,6 +431,78 @@ async fn sync_one<C: Crypto>( }) } +/// Writes SetTimeZone and (when the device still needs it) SetDSTOffset, +/// returning the zone name written and the number of DST entries. Split out +/// of `sync_one` so the two writes read as one cohesive step. +async fn write_zone_and_dst<C: Crypto>( + ctx: &Ctx<'_, C>, + node_id: u64, + tz: &jiff::tz::TimeZone, + now_ts: Timestamp, +) -> anyhow::Result<(Option<String>, usize)> { + let zone_entries = build_time_zone_list(tz, &ctx.config.timezone, now_ts); + let response = connect(ctx, node_id) + .await? + .time_synchronization() + .set_time_zone(ROOT_ENDPOINT_ID, |b| { + let mut list = b.time_zone()?; + for entry in &zone_entries { + list = list + .push()? + .offset(entry.offset_seconds)? + .valid_at(entry.valid_at.0)? + .name(Some(&entry.name))? + .end()?; + } + list.end()?.end() + }) + .await + .ctx("SetTimeZone rejected")?; + let dst_required = response + .response() + .map(|r| r.dst_offset_required().unwrap_or(true)) + .unwrap_or(true); + response.complete().await.ctx("SetTimeZone completion")?; + log::info!("Node {node_id}: SetTimeZone {}", ctx.config.timezone); + + if !dst_required { + return Ok((Some(ctx.config.timezone.clone()), 0)); + } + + let max_entries = connect(ctx, node_id) + .await? + .time_synchronization() + .dst_offset_list_max_size_read(ROOT_ENDPOINT_ID) + .await + .ctx("read dstOffsetListMaxSize")?; + let dst_entries = build_dst_offset_list(tz, usize::from(max_entries), now_ts); + connect(ctx, node_id) + .await? + .time_synchronization() + .set_dst_offset(ROOT_ENDPOINT_ID, |b| { + let mut list = b.dst_offset()?; + for entry in &dst_entries { + list = list + .push()? + .offset(entry.offset_seconds)? + .valid_starting(entry.valid_starting.0)? + .valid_until(match entry.valid_until { + Some(until) => rs_matter::tlv::Nullable::some(until.0), + None => rs_matter::tlv::Nullable::none(), + })? + .end()?; + } + list.end()?.end() + }) + .await + .ctx("SetDSTOffset rejected")?; + log::info!( + "Node {node_id}: SetDSTOffset with {} entries", + dst_entries.len() + ); + Ok((Some(ctx.config.timezone.clone()), dst_entries.len())) +} + /// Pushes the configured fabric label to the device when it differs from the /// stored one. Best-effort: a cosmetic label must not fail a clock sync. /// Labels are unique per device, so a conflict with another admin's label is @@ -577,14 +585,7 @@ impl Op for DecommissionOp { let result = decommission_one(ctx, node_id).await; match result { Ok(()) => { - let registry_file = state::registry_path(&ctx.config.storage_path); - let mut registry: NodeRegistry = state::load_or_default(®istry_file); - registry.nodes.remove(&node_id.to_string()); - state::store(®istry_file, ®istry)?; - state::remove_node_state( - &state::service_state_path(&ctx.config.storage_path), - node_id, - )?; + state::remove_device(&devices_path(ctx), node_id)?; outcomes.push(DecommissionOutcome { node_id: node_id.into(), success: true, @@ -633,7 +634,7 @@ pub struct InspectOp { pub targets: Vec<u64>, } -#[derive(Debug, serde::Serialize)] +#[derive(Debug, Default, serde::Serialize)] #[serde(rename_all = "camelCase")] pub struct InspectOutcome { pub node_id: crate::output::Id64, @@ -645,6 +646,17 @@ pub struct InspectOutcome { pub our_fabric_index: Option<u8>, } +impl InspectOutcome { + /// A node whose live inspection could not complete. + fn failed(node_id: u64, error: String) -> Self { + Self { + node_id: node_id.into(), + error: Some(error), + ..Default::default() + } + } +} + #[derive(Debug, serde::Serialize)] #[serde(rename_all = "camelCase")] pub struct TimeSyncCaps { @@ -665,15 +677,7 @@ impl Op for InspectOp { Ok(outcome) => outcomes.push(outcome), Err(error) => { log::error!("Node {node_id}: inspect failed: {error:#}"); - outcomes.push(InspectOutcome { - node_id: node_id.into(), - error: Some(format!("{error:#}")), - vendor_name: None, - product_name: None, - time_sync: None, - our_fabric_label: None, - our_fabric_index: None, - }); + outcomes.push(InspectOutcome::failed(node_id, format!("{error:#}"))); } } } |
