From 0470c70cb4dbf98421d993250482da3e46b7b456 Mon Sep 17 00:00:00 2001 From: Stefanie Jane Date: Sat, 12 Sep 2026 20:26:24 -0700 Subject: [PATCH 1/4] fix(devices): restore observed names when clearing overrides Keep the latest driver name separate from the effective display name so clearing a user override restores current hardware metadata immediately. Rediscovery and guarded refresh update that observation, while metadata-only refreshes preserve it instead of recycling a customized display name. Cover name reset through the registry and device API, including rediscovery, SMBus remapping, rejected refreshes, and preservation of other settings. --- crates/hypercolor-core/src/device/registry.rs | 26 ++- crates/hypercolor-core/tests/device_tests.rs | 172 ++++++++++++++++++ .../src/discovery/device_helpers.rs | 2 +- .../tests/device_name_reset_tests.rs | 64 +++++++ 4 files changed, 261 insertions(+), 3 deletions(-) create mode 100644 crates/hypercolor-daemon/tests/device_name_reset_tests.rs diff --git a/crates/hypercolor-core/src/device/registry.rs b/crates/hypercolor-core/src/device/registry.rs index d9b2e0db9..7816c177a 100644 --- a/crates/hypercolor-core/src/device/registry.rs +++ b/crates/hypercolor-core/src/device/registry.rs @@ -27,6 +27,9 @@ pub struct TrackedDevice { /// Full device metadata. pub info: DeviceInfo, + /// Latest driver-provided name before the user override is applied. + observed_name: String, + /// Current lifecycle state. pub state: DeviceState, @@ -40,6 +43,16 @@ pub struct TrackedDevice { pub revision: u64, } +impl TrackedDevice { + /// Device metadata without the user-provided name override. + #[must_use] + pub fn observed_info(&self) -> DeviceInfo { + let mut info = self.info.clone(); + info.name.clone_from(&self.observed_name); + info + } +} + // ── DeviceRegistry ─────────────────────────────────────────────────────── /// Thread-safe registry for tracking all known devices. @@ -217,9 +230,11 @@ impl DeviceRegistry { // Keep the canonical registry ID stable across rediscovery. updated_info.id = existing_id; preserve_resolved_device_shape(&mut updated_info, &entry.info); + let observed_name = updated_info.name.clone(); apply_user_settings_to_info(&mut updated_info, &entry.user_settings); - let entry_changed = - entry.info != updated_info || entry.connect_behavior != connect_behavior; + let entry_changed = entry.info != updated_info + || entry.observed_name != observed_name + || entry.connect_behavior != connect_behavior; let fingerprint_changed = inner.id_to_fingerprint.get(&existing_id) != Some(&fingerprint); let metadata_changed = !metadata.is_empty() @@ -236,6 +251,7 @@ impl DeviceRegistry { .get_mut(&existing_id) .expect("existing device was resolved above"); entry.info = updated_info; + entry.observed_name = observed_name; entry.connect_behavior = connect_behavior; bump_device_revision(entry); } @@ -280,6 +296,7 @@ impl DeviceRegistry { let mut updated_info = info; updated_info.id = existing_id; preserve_resolved_device_shape(&mut updated_info, &entry.info); + entry.observed_name.clone_from(&updated_info.name); apply_user_settings_to_info(&mut updated_info, &entry.user_settings); debug!( device_id = %existing_id, @@ -319,6 +336,7 @@ impl DeviceRegistry { let name = tracked_info.name.clone(); let tracked = TrackedDevice { + observed_name: name.clone(), info: tracked_info, state: DeviceState::Known, connect_behavior, @@ -424,6 +442,7 @@ impl DeviceRegistry { let mut updated_info = info; updated_info.id = *id; + entry.observed_name.clone_from(&updated_info.name); apply_user_settings_to_info(&mut updated_info, &entry.user_settings); entry.info = updated_info; bump_device_revision(entry); @@ -456,6 +475,7 @@ impl DeviceRegistry { return None; } let mut info = discovered.info; + entry.observed_name.clone_from(&info.name); apply_user_settings_to_info(&mut info, &entry.user_settings); entry.info = info; entry.connect_behavior = discovered.connect_behavior; @@ -521,6 +541,7 @@ impl DeviceRegistry { let entry = inner.devices.get_mut(id)?; let mut updated_info = entry.info.clone(); + updated_info.name.clone_from(&entry.observed_name); apply_user_settings_to_info(&mut updated_info, &settings); if entry.user_settings == settings && entry.info == updated_info { return Some(entry.clone()); @@ -735,6 +756,7 @@ impl DeviceRegistry { .expect("device presence was checked above"); if let Some(settings) = inherited_settings { entry.user_settings = settings; + entry.info.name.clone_from(&entry.observed_name); apply_user_settings_to_info(&mut entry.info, &entry.user_settings); } bump_device_revision(entry); diff --git a/crates/hypercolor-core/tests/device_tests.rs b/crates/hypercolor-core/tests/device_tests.rs index faa215766..8348be43a 100644 --- a/crates/hypercolor-core/tests/device_tests.rs +++ b/crates/hypercolor-core/tests/device_tests.rs @@ -2055,3 +2055,175 @@ async fn orchestrator_reappeared_device_keeps_stable_id_when_scanner_emits_new_i .expect("stable registry entry should remain"); assert_eq!(tracked.info.id, existing_id); } + +#[tokio::test] +async fn registry_name_clear_restores_latest_observation_and_preserves_settings() { + let registry = DeviceRegistry::new(); + let fingerprint = DeviceFingerprint::from_persisted("bridge:name-reset".to_owned()); + let id = registry + .add_with_fingerprint(mock_device_info("Original"), fingerprint.clone()) + .await; + registry + .update_user_settings( + &id, + Some("Custom".to_owned()), + Some(false), + Some(0.37), + None, + ) + .await; + let before = registry.get(&id).await.expect("registered"); + assert_eq!( + registry + .add_with_fingerprint(mock_device_info("Latest"), fingerprint) + .await, + id + ); + let observed = registry.get(&id).await.expect("rediscovered"); + assert_eq!(observed.info.name, "Custom"); + assert!(observed.revision > before.revision); + let settings = DeviceUserSettings { + name: None, + ..observed.user_settings + }; + let cleared = registry + .replace_user_settings(&id, settings.clone()) + .await + .expect("cleared"); + assert_eq!(cleared.info.name, "Latest"); + assert!(!cleared.user_settings.enabled); + assert_eq!(cleared.user_settings.brightness, 0.37); + assert_eq!(registry.list().await[0].info.name, "Latest"); + let generation = registry.generation(); + let replay = registry + .replace_user_settings(&id, settings) + .await + .expect("replayed"); + assert_eq!(replay.revision, cleared.revision); + assert_eq!(registry.generation(), generation); +} + +#[tokio::test] +async fn registry_metadata_refresh_keeps_latest_name_behind_override() { + let registry = DeviceRegistry::new(); + let fingerprint = DeviceFingerprint::from_persisted("bridge:name-refresh".to_owned()); + let id = registry + .add_with_fingerprint(mock_device_info("Original"), fingerprint.clone()) + .await; + registry + .update_user_settings(&id, Some("Custom".to_owned()), None, None, None) + .await; + registry + .update_info(&id, mock_device_info("Connected name")) + .await + .expect("metadata update"); + let settings = registry.get(&id).await.expect("registered").user_settings; + assert_eq!( + registry + .replace_user_settings( + &id, + DeviceUserSettings { + name: None, + ..settings + } + ) + .await + .expect("clear") + .info + .name, + "Connected name" + ); + registry + .update_user_settings(&id, Some("Custom again".to_owned()), None, None, None) + .await; + let mut observation = DiscoveredDevice { + info: DeviceInfo { + id, + ..mock_device_info("Refreshed name") + }, + fingerprint, + connect_behavior: DiscoveryConnectBehavior::AutoConnect, + metadata: HashMap::new(), + claim: None, + }; + registry + .refresh_discovered(&id, observation.clone()) + .await + .expect("refresh"); + let tracked = registry.get(&id).await.expect("refreshed device"); + assert_eq!(tracked.info.name, "Custom again"); + assert_eq!(tracked.observed_info().name, "Refreshed name"); + let mut metadata_only = observation.clone(); + metadata_only.info = tracked.observed_info(); + metadata_only + .metadata + .insert("status".to_owned(), "ready".to_owned()); + registry + .refresh_discovered(&id, metadata_only) + .await + .expect("metadata-only refresh"); + observation.info.name = "Rejected name".to_owned(); + observation.info.id = DeviceId::new(); + assert!( + registry + .refresh_discovered(&id, observation) + .await + .is_none() + ); + let settings = registry.get(&id).await.expect("registered").user_settings; + assert_eq!( + registry + .replace_user_settings( + &id, + DeviceUserSettings { + name: None, + ..settings + } + ) + .await + .expect("clear refreshed name") + .info + .name, + "Refreshed name" + ); +} + +#[tokio::test] +async fn registry_name_clear_after_smbus_remap_uses_new_observed_name() { + let registry = DeviceRegistry::new(); + let id = registry + .add_with_fingerprint_and_metadata( + asus_dram_device_info(0x71), + DeviceFingerprint::from_persisted("smbus:/dev/i2c-9:71".to_owned()), + asus_dram_metadata(0x71), + ) + .await; + registry.set_state(&id, DeviceState::Connected).await; + registry + .update_user_settings(&id, Some("My RAM".to_owned()), None, None, None) + .await; + let remapped = registry + .add_with_fingerprint_and_metadata( + asus_dram_device_info(0x73), + DeviceFingerprint::from_persisted("smbus:/dev/i2c-9:73".to_owned()), + asus_dram_metadata(0x73), + ) + .await; + assert_eq!(remapped, id); + let settings = registry.get(&id).await.expect("remapped").user_settings; + assert_eq!( + registry + .replace_user_settings( + &id, + DeviceUserSettings { + name: None, + ..settings + } + ) + .await + .expect("clear") + .info + .name, + "ASUS Aura DRAM (SMBus 0x73)" + ); +} diff --git a/crates/hypercolor-daemon/src/discovery/device_helpers.rs b/crates/hypercolor-daemon/src/discovery/device_helpers.rs index 70c36a701..ed63c2f54 100644 --- a/crates/hypercolor-daemon/src/discovery/device_helpers.rs +++ b/crates/hypercolor-daemon/src/discovery/device_helpers.rs @@ -75,7 +75,7 @@ pub(super) async fn refresh_connected_device_info( .fingerprint_for_id(&device_id) .await .context("connected device has no registered fingerprint")?; - let mut info = maybe_info.unwrap_or(tracked.info); + let mut info = maybe_info.unwrap_or_else(|| tracked.observed_info()); info.id = device_id; runtime .device_registry diff --git a/crates/hypercolor-daemon/tests/device_name_reset_tests.rs b/crates/hypercolor-daemon/tests/device_name_reset_tests.rs new file mode 100644 index 000000000..01e2906a9 --- /dev/null +++ b/crates/hypercolor-daemon/tests/device_name_reset_tests.rs @@ -0,0 +1,64 @@ +use std::sync::Arc; + +use axum::body::to_bytes; +use axum::extract::{Path, State}; +use hypercolor_daemon::api::devices::get_device; +use hypercolor_daemon::app_state::AppState; +use hypercolor_types::device::{ + ConnectionType, DeviceCapabilities, DeviceFamily, DeviceId, DeviceInfo, DeviceOrigin, + DeviceUserSettings, +}; + +#[tokio::test] +async fn device_api_reports_observed_name_after_override_is_cleared() { + let directory = tempfile::tempdir().expect("isolated state directory"); + let state = Arc::new(AppState::new_with_data_dir(directory.path().to_path_buf())); + let id = state + .device_registry + .add(DeviceInfo { + id: DeviceId::new(), + name: "Hardware name".to_owned(), + vendor: "Fixture".to_owned(), + family: DeviceFamily::new_static("fixture", "Fixture"), + model: None, + connection_type: ConnectionType::Network, + origin: DeviceOrigin::native("fixture", "fixture", ConnectionType::Network), + segments: Vec::new(), + firmware_version: None, + capabilities: DeviceCapabilities::default(), + }) + .await; + state + .device_registry + .update_user_settings(&id, Some("Custom name".to_owned()), None, None, None) + .await + .expect("rename"); + for (clear, expected) in [(false, "Custom name"), (true, "Hardware name")] { + if clear { + let settings = state + .device_registry + .get(&id) + .await + .expect("tracked") + .user_settings; + state + .device_registry + .replace_user_settings( + &id, + DeviceUserSettings { + name: None, + ..settings + }, + ) + .await + .expect("clear override"); + } + let response = get_device(State(Arc::clone(&state)), Path(id.to_string())).await; + assert_eq!(response.status(), http::StatusCode::OK); + let bytes = to_bytes(response.into_body(), 1024 * 1024) + .await + .expect("response body"); + let body: serde_json::Value = serde_json::from_slice(&bytes).expect("response JSON"); + assert_eq!(body["data"]["name"], expected); + } +} From 449fee8cee5f39c731a0f3905f145be51cd756ac Mon Sep 17 00:00:00 2001 From: Stefanie Jane Date: Sat, 12 Sep 2026 21:52:58 -0700 Subject: [PATCH 2/4] feat(daemon): guard exact runtime zone color mutations Apply color through the existing scene candidate and revision fence, preserving zone metadata, membership, brightness and output power. Replace only the layer stack and reject stale or display-owned targets. Return actual persistence evidence from both stores. Default-scene layers require their exact Written runtime payload; superseded or later overwritten snapshots cannot certify the original operation. Co-Authored-By: Nova (GPT-6) --- crates/hypercolor-core/src/scene/mod.rs | 27 + crates/hypercolor-daemon/src/domain/commit.rs | 6 +- .../hypercolor-daemon/src/domain/context.rs | 91 ++- crates/hypercolor-daemon/src/domain/mod.rs | 1 + .../src/domain/runtime_zone.rs | 95 +++ crates/hypercolor-daemon/src/domain/scene.rs | 43 ++ .../tests/runtime_zone_tests.rs | 599 ++++++++++++++++++ 7 files changed, 837 insertions(+), 25 deletions(-) create mode 100644 crates/hypercolor-daemon/src/domain/runtime_zone.rs create mode 100644 crates/hypercolor-daemon/tests/runtime_zone_tests.rs diff --git a/crates/hypercolor-core/src/scene/mod.rs b/crates/hypercolor-core/src/scene/mod.rs index e040656c6..72391d88f 100644 --- a/crates/hypercolor-core/src/scene/mod.rs +++ b/crates/hypercolor-core/src/scene/mod.rs @@ -1179,6 +1179,33 @@ impl SceneManager { }) } + /// Replace only a zone's authored layer stack, preserving its other fields. + /// + /// # Errors + /// + /// Refuses missing targets, duplicate identities or invalid layer values + /// before changing the stack. + pub fn replace_zone_layer_stack( + &mut self, + scene_id: SceneId, + zone_id: ZoneId, + layers: Vec, + ) -> Result<(&Zone, u64), LayerMutationError> { + let mut ids = HashSet::new(); + for layer in &layers { + if !ids.insert(layer.id) { + return Err(LayerMutationError::DuplicateLayer { layer_id: layer.id }); + } + layer + .validate() + .map_err(|errors| LayerMutationError::InvalidLayer { errors })?; + } + self.mutate_zone_layers(scene_id, zone_id, None, |zone| { + zone.layers = layers; + Ok(()) + }) + } + pub fn remove_zone_layer( &mut self, scene_id: SceneId, diff --git a/crates/hypercolor-daemon/src/domain/commit.rs b/crates/hypercolor-daemon/src/domain/commit.rs index 9fd230db0..ede0fff39 100644 --- a/crates/hypercolor-daemon/src/domain/commit.rs +++ b/crates/hypercolor-daemon/src/domain/commit.rs @@ -42,9 +42,9 @@ pub enum CommitDurability { /// The payload replaced the destination and its durability barrier /// completed. Written, - /// A newer admitted generation won the destination first. The newer - /// payload is authoritative and already contains this commit's - /// changes, so this commit's own bytes will never be written. + /// A newer admitted generation won the destination first. This commit's + /// own bytes will never be written; the newer payload may have overwritten + /// its fields and is not proof that this commit's changes became durable. Superseded, /// The write did not prove durable on this attempt. The payload /// stays the destination's newest admitted intent and the retry diff --git a/crates/hypercolor-daemon/src/domain/context.rs b/crates/hypercolor-daemon/src/domain/context.rs index 892a9e605..9a830ac06 100644 --- a/crates/hypercolor-daemon/src/domain/context.rs +++ b/crates/hypercolor-daemon/src/domain/context.rs @@ -260,6 +260,25 @@ pub(crate) enum RuntimeSessionPersistenceError { RetryArmed(runtime_state::RuntimeSessionError), } +/// Evidence from the scene-store and runtime-session save boundary. +#[derive(Debug)] +pub enum RuntimeSessionSaveOutcome { + /// No runtime snapshot could be reserved; neither save was attempted. + BeforeAdmission { + /// Reservation error, before any runtime payload admission. + error: runtime_state::RuntimeSessionError, + }, + /// A reservation was obtained and both persistence stages were attempted. + Attempted { + /// Scene-store save outcome; a failed prerequisite is never hidden. + scene_store: anyhow::Result>, + /// Exact projection supplied to the runtime writer, not a later read. + projection: RuntimeSessionSnapshot, + /// Runtime pointer write result. Superseded does not prove these bytes. + snapshot: Result, + }, +} + struct RuntimeSessionSave { pending: runtime_state::RuntimeSnapshotSave, snapshot: RuntimeSessionSnapshot, @@ -439,8 +458,33 @@ impl RuntimeSessionService { .await } - /// Persist the current scene store before the runtime-session pointer. + /// Persist the scene store and runtime pointer, logging failures for callers + /// that do not consume persistence evidence. pub async fn save(&self) { + match self.save_with_outcome().await { + RuntimeSessionSaveOutcome::BeforeAdmission { error } => { + tracing::warn!(path = %self.path.display(), %error, + "Failed to reserve runtime session snapshot"); + } + RuntimeSessionSaveOutcome::Attempted { + scene_store, + snapshot, + .. + } => { + if let Err(error) = scene_store { + tracing::warn!(%error, "Failed to persist scene store before runtime snapshot save"); + } + if let Err(error) = snapshot { + tracing::warn!(path = %self.path.display(), %error, + "Failed to persist runtime session snapshot"); + } + } + } + } + + /// Persist the current scene store before the runtime-session pointer and + /// return each actual outcome without inferring durability from newer state. + pub async fn save_with_outcome(&self) -> RuntimeSessionSaveOutcome { let driver_host = self.driver_host.upgrade(); let refresh_inventory = async move { if let Some(driver_host) = driver_host { @@ -453,25 +497,14 @@ impl RuntimeSessionService { .await { Ok(save) => save, - Err(error) => { - tracing::warn!( - path = %self.path.display(), - %error, - "Failed to reserve runtime session snapshot" - ); - return; - } + Err(error) => return RuntimeSessionSaveOutcome::BeforeAdmission { error }, }; - if let Err(error) = self.save_scene_store_snapshot().await { - tracing::warn!(%error, "Failed to persist scene store before runtime snapshot save"); - } - - if let Err(error) = save.commit() { - tracing::warn!( - path = %self.path.display(), - %error, - "Failed to persist runtime session snapshot" - ); + let scene_store = self.save_scene_store_snapshot().await; + let (projection, snapshot) = save.commit_observed(); + RuntimeSessionSaveOutcome::Attempted { + scene_store, + projection, + snapshot, } } @@ -492,8 +525,8 @@ impl RuntimeSessionService { .await } - async fn save_scene_store_snapshot(&self) -> anyhow::Result<()> { - self.projection.scenes.save_snapshot().await + async fn save_scene_store_snapshot(&self) -> anyhow::Result> { + self.projection.scenes.persist_snapshot().await } #[cfg(all(test, feature = "persistence-test-hooks"))] @@ -506,6 +539,15 @@ impl RuntimeSessionService { impl RuntimeSessionSave { fn commit(self) -> Result { + self.commit_observed().1 + } + + fn commit_observed( + self, + ) -> ( + RuntimeSessionSnapshot, + Result, + ) { let Self { pending, snapshot, @@ -513,7 +555,7 @@ impl RuntimeSessionSave { } = self; let outcome = runtime_state::save_reserved(pending, &snapshot); drop(publication_guard); - outcome + (snapshot, outcome) } } @@ -721,6 +763,11 @@ impl SceneContext { self.runtime_session.save().await; } + /// Return persistence evidence for the current runtime-session projection. + pub async fn save_runtime_session_with_outcome(&self) -> RuntimeSessionSaveOutcome { + self.runtime_session.save_with_outcome().await + } + /// Remove explicitly forgotten controller outputs from every saved scene and layout. pub(crate) async fn forget_layout_targets( &self, diff --git a/crates/hypercolor-daemon/src/domain/mod.rs b/crates/hypercolor-daemon/src/domain/mod.rs index ece764f60..e423876b3 100644 --- a/crates/hypercolor-daemon/src/domain/mod.rs +++ b/crates/hypercolor-daemon/src/domain/mod.rs @@ -45,6 +45,7 @@ mod macos_screen_parity; pub mod openrgb_diagnostics; pub mod openrgb_setup; pub mod output; +pub mod runtime_zone; pub mod scene; pub mod scene_tree; pub mod spatial; diff --git a/crates/hypercolor-daemon/src/domain/runtime_zone.rs b/crates/hypercolor-daemon/src/domain/runtime_zone.rs new file mode 100644 index 000000000..5be2c0a9b --- /dev/null +++ b/crates/hypercolor-daemon/src/domain/runtime_zone.rs @@ -0,0 +1,95 @@ +//! Exact, revision-fenced lighting mutations for native domain consumers. + +use hypercolor_color::Rgb; +use hypercolor_types::scene::{SceneId, SceneKind, Zone, ZoneId}; + +use super::DomainError; +use super::commit::{CommitDurability, SceneCommit}; +use super::context::{RuntimeSessionSaveOutcome, SceneContext}; +use crate::persistence::AtomicWriteOutcome; + +/// A zone observed together with its active scene and scene commit revision. +#[derive(Debug, Clone, Copy)] +pub struct RuntimeZoneTarget { + /// Immutable active scene identity. + pub scene_id: SceneId, + /// Immutable zone identity within the scene. + pub zone_id: ZoneId, + /// Revision captured from the same owned scene snapshot. + pub revision: u64, +} + +/// Original scene admission and the separate runtime projection save attempt. +/// +/// Named scenes persist through the scene store; default-scene layers persist +/// through the runtime projection. Neither Superseded nor a newer unrelated +/// snapshot proves that this operation's layers were written or remain current. +#[derive(Debug)] +pub struct RuntimeZoneColorOutcome { + /// Exact identity and observed revision used for admission. + pub target: RuntimeZoneTarget, + /// Original scene kind, which determines whether the scene store owns it. + pub scene_kind: SceneKind, + /// The admitted zone, not a subsequent read of a potentially newer tree. + pub zone: Zone, + /// Original scene commit and its durability evidence. + pub commit: SceneCommit, + /// Actual runtime projection persistence attempt after admission. + pub runtime_session: RuntimeSessionSaveOutcome, +} + +impl RuntimeZoneColorOutcome { + /// Whether an actual Written payload contains this operation's layers. + /// + /// A named scene uses its original scene-store commit. A default scene + /// requires the exact layer identities and content in a Written runtime + /// projection. This proves persistence, not current output or device apply. + #[must_use] + pub fn has_written_layers(&self) -> bool { + if !self.target.scene_id.is_default() { + return self.scene_kind == SceneKind::Named + && self.commit.durability() == CommitDurability::Written; + } + matches!(&self.runtime_session, + RuntimeSessionSaveOutcome::Attempted { + projection, + snapshot: Ok(AtomicWriteOutcome::Written), + .. + } if projection.default_scene_zones.iter().any(|zone| + zone.id == self.zone.id && zone.layers == self.zone.layers) + ) + } +} + +/// Replace only the selected live lighting zone's layers with an opaque color. +/// +/// RGB is encoded sRGB, explicitly linearized for the engine's ColorFill source. +/// The caller owns authorization; this boundary owns scene/zone/revision and +/// role validation. No scene activation, output wake or brightness write occurs. +/// +/// # Errors +/// +/// Returns target validation or commit-CAS errors before admission. Persistence +/// failures after scene admission remain visible in the returned outcomes. +pub async fn set_color( + ctx: &SceneContext, + target: RuntimeZoneTarget, + color: Rgb, +) -> Result { + let mut mutation = ctx.begin_mutation().await; + let zone = mutation.set_runtime_zone_color(target, color)?; + let scene_kind = mutation + .scenes() + .get(&target.scene_id) + .ok_or_else(|| DomainError::not_found(super::ResourceKind::Scene, target.scene_id))? + .kind; + let commit = ctx.commit(mutation).await?; + let runtime_session = ctx.save_runtime_session_with_outcome().await; + Ok(RuntimeZoneColorOutcome { + target, + scene_kind, + zone, + commit, + runtime_session, + }) +} diff --git a/crates/hypercolor-daemon/src/domain/scene.rs b/crates/hypercolor-daemon/src/domain/scene.rs index 5995ddcaa..e6d5a6402 100644 --- a/crates/hypercolor-daemon/src/domain/scene.rs +++ b/crates/hypercolor-daemon/src/domain/scene.rs @@ -1391,6 +1391,49 @@ impl SceneMutation { Ok(zone) } + /// Replace the layers of one exact live lighting zone in this candidate. + /// + /// # Errors + /// + /// Refuses stale revisions, a different active scene, snapshot scenes, + /// missing zones and display-owned zones before changing the candidate. + pub fn set_runtime_zone_color( + &mut self, + target: super::runtime_zone::RuntimeZoneTarget, + color: hypercolor_color::Rgb, + ) -> Result { + super::scene_tree::check_scene_revision(self, Some(target.revision))?; + let scene_id = self.active_scene_for_runtime_mutation("setting a zone color")?; + if scene_id != target.scene_id { + return Err(DomainError::conflict( + "the requested scene is no longer active", + )); + } + super::scene_tree::ensure_live_zone_mutable(self, target.zone_id)?; + let linear = color.to_linear(); + let layer = SceneLayer { + id: SceneLayerId::new(), + name: None, + source: hypercolor_types::layer::LayerSource::ColorFill { + rgba: [linear.r, linear.g, linear.b, linear.a], + }, + blend: hypercolor_types::layer::BlendMode::Replace, + opacity: 1.0, + transform: hypercolor_types::layer::LayerTransform::default(), + adjust: hypercolor_types::layer::LayerAdjust::default(), + bindings: Vec::new(), + enabled: true, + }; + let (zone, _) = self + .candidate + .replace_zone_layer_stack(scene_id, target.zone_id, vec![layer]) + .map_err(|error| DomainError::Internal(anyhow::anyhow!("{error:?}")))?; + let zone = zone.clone(); + self.persists_scene_content = true; + self.record_layer_change(scene_id, &zone, LayerStackChangeKind::Updated); + Ok(zone) + } + /// Drop one layer out of a zone's stack. pub fn remove_layer( &mut self, diff --git a/crates/hypercolor-daemon/tests/runtime_zone_tests.rs b/crates/hypercolor-daemon/tests/runtime_zone_tests.rs new file mode 100644 index 000000000..89f8fdcb3 --- /dev/null +++ b/crates/hypercolor-daemon/tests/runtime_zone_tests.rs @@ -0,0 +1,599 @@ +//! Exact native lighting targets and evidence from the owning persistence path. + +use hypercolor_color::{LinearRgba, Rgb}; +use hypercolor_core::scene::default_primary_zone; +use hypercolor_daemon::app_state::AppState; +use hypercolor_daemon::domain::DomainError; +use hypercolor_daemon::domain::commit::CommitDurability; +use hypercolor_daemon::domain::context::RuntimeSessionSaveOutcome; +use hypercolor_daemon::domain::runtime_zone::{RuntimeZoneTarget, set_color}; +use hypercolor_daemon::domain::scene::commit_scene; +#[cfg(feature = "persistence-test-hooks")] +use hypercolor_daemon::persistence::AtomicFileWriter; +use hypercolor_daemon::persistence::AtomicWriteOutcome; +use hypercolor_types::event::{HypercolorEvent, SceneChangeReason}; +use hypercolor_types::layer::{ + BlendMode, LayerAdjust, LayerSource, LayerTransform, SceneLayer, SceneLayerId, +}; +use hypercolor_types::scene::{Scene, SceneId, SceneKind, SceneMutationMode, ZoneId, ZoneRole}; + +fn layer() -> SceneLayer { + SceneLayer { + id: SceneLayerId::new(), + name: Some("old effect state".to_owned()), + source: LayerSource::ColorFill { + rgba: [0.25, 0.5, 0.75, 0.4], + }, + blend: BlendMode::Replace, + opacity: 0.6, + transform: LayerTransform::default(), + adjust: LayerAdjust::default(), + bindings: Vec::new(), + enabled: true, + } +} + +async fn fixture() -> (AppState, tempfile::TempDir, Scene, RuntimeZoneTarget) { + let dir = tempfile::tempdir().expect("fixture directory"); + let state = AppState::new_with_data_dir(dir.path().join("data")); + let mut scene = state + .scene_manager + .snapshot() + .await + .active_scene() + .expect("default scene") + .clone(); + scene.id = SceneId::new(); + scene.kind = SceneKind::Named; + scene.mutation_mode = SceneMutationMode::Live; + let mut primary = + default_primary_zone(state.spatial_engine.snapshot().layout().as_ref().clone()); + "retain primary metadata".clone_into(&mut primary.name); + primary.description = Some("do not replace this zone".to_owned()); + primary.brightness = 0.37; + primary.enabled = false; + primary.layers = vec![layer(), layer()]; + primary + .layout + .zones + .push(hypercolor_types::spatial::Output { + id: "retained-membership".to_owned(), + name: "desk strip".to_owned(), + device_id: "mock:desk".into(), + zone_name: Some("left".to_owned()), + position: hypercolor_types::spatial::NormalizedPosition::new(0.3, 0.7), + size: hypercolor_types::spatial::NormalizedPosition::new(0.6, 0.2), + rotation: 27.0, + scale: 0.8, + display_order: 3, + orientation: None, + topology: hypercolor_types::spatial::LedTopology::Strip { + count: 5, + direction: hypercolor_types::spatial::StripDirection::LeftToRight, + }, + led_positions: Vec::new(), + led_mapping: None, + sampling_mode: None, + edge_behavior: None, + shape: None, + shape_preset: None, + attachment: None, + brightness: Some(0.65), + }); + let mut independent = primary.clone(); + independent.id = ZoneId::new(); + independent.layout.zones.clear(); + independent.role = ZoneRole::Custom; + independent.layers = vec![layer()]; + scene.zones = vec![primary, independent]; + let mut mutation = state.scene_manager.begin_mutation().await; + mutation.create_scene(scene.clone()).expect("seed scene"); + mutation + .activate(scene.id, None, SceneChangeReason::UserActivate) + .expect("activate fixture"); + let commit = commit_scene(&state.domains.scene, mutation) + .await + .expect("seed commit"); + let target = RuntimeZoneTarget { + scene_id: scene.id, + zone_id: scene.zones[0].id, + revision: commit.revision(), + }; + (state, dir, scene, target) +} + +#[tokio::test] +async fn color_changes_only_layers_and_preserves_paused_output_and_primary_metadata() { + let (state, _dir, before, target) = fixture().await; + state + .output_power + .set_global_brightness(&state.event_bus, 0.42) + .await + .expect("brightness"); + state + .output_power + .set_manual_pause(&state.event_bus, true, [1, 2, 3]) + .await; + let power = state.output_power.snapshot(); + let mut events = state.event_bus.subscribe_all(); + let color = Rgb::new(34, 129, 207); + let result = set_color(&state.domains.scene, target, color) + .await + .expect("exact color"); + assert_eq!(result.commit.durability(), CommitDurability::Written); + assert!(result.has_written_layers()); + assert!(matches!( + result.runtime_session, + RuntimeSessionSaveOutcome::Attempted { + scene_store: Ok(Some(AtomicWriteOutcome::Written)), + snapshot: Ok(AtomicWriteOutcome::Written), + .. + } + )); + let mut after = state + .scene_manager + .snapshot() + .await + .get(&target.scene_id) + .expect("scene retained") + .clone(); + assert_eq!(result.zone, after.zones[0]); + assert_eq!(after.zones[0].layers.len(), 1); + let new_layer = &after.zones[0].layers[0]; + assert!( + before.zones[0] + .layers + .iter() + .all(|old| old.id != new_layer.id) + ); + let LayerSource::ColorFill { rgba } = new_layer.source else { + panic!("constant color") + }; + assert_eq!( + LinearRgba::new(rgba[0], rgba[1], rgba[2], rgba[3]).to_encoded(), + color.to_rgba() + ); + assert!(rgba[1] < 0.3, "encoded middle gray must be linearized"); + after.zones[0].layers = before.zones[0].layers.clone(); + after.zones[0].layers_version = before.zones[0].layers_version; + assert_eq!( + after, before, + "all other scene and zone fields are preserved" + ); + assert_eq!(state.output_power.snapshot(), power); + let mut layer_events = 0; + let mut zone_events = 0; + while let Ok(event) = events.try_recv() { + match event.event { + HypercolorEvent::LayerStackChanged { .. } => layer_events += 1, + HypercolorEvent::ZoneChanged { .. } => zone_events += 1, + other => panic!("unexpected side-effect event {other:?}"), + } + } + assert_eq!((layer_events, zone_events), (1, 1)); +} + +#[tokio::test] +async fn wrong_scene_stale_revision_and_missing_zone_leave_the_scene_unchanged() { + let (state, _dir, before, target) = fixture().await; + assert!(matches!( + set_color( + &state.domains.scene, + RuntimeZoneTarget { + scene_id: SceneId::new(), + ..target + }, + Rgb::WHITE + ) + .await, + Err(DomainError::Conflict { .. }) + )); + assert!(matches!( + set_color( + &state.domains.scene, + RuntimeZoneTarget { + revision: target.revision + 1, + ..target + }, + Rgb::WHITE + ) + .await, + Err(DomainError::PreconditionFailed { .. }) + )); + assert!(matches!( + set_color( + &state.domains.scene, + RuntimeZoneTarget { + zone_id: ZoneId::new(), + ..target + }, + Rgb::WHITE + ) + .await, + Err(DomainError::NotFound { .. }) + )); + assert_eq!( + state.scene_manager.snapshot().await.get(&target.scene_id), + Some(&before) + ); + assert_eq!(state.scene_manager.revision(), target.revision); +} + +#[tokio::test] +async fn display_and_snapshot_guards_are_checked_on_the_owned_candidate() { + let (state, _dir, before, target) = fixture().await; + let mut mutation = state.scene_manager.begin_mutation().await; + let mut changed = before.clone(); + changed.mutation_mode = SceneMutationMode::Snapshot; + mutation.update_scene(changed).expect("snapshot fixture"); + let commit = commit_scene(&state.domains.scene, mutation) + .await + .expect("snapshot commit"); + assert!(matches!( + set_color( + &state.domains.scene, + RuntimeZoneTarget { + revision: commit.revision(), + ..target + }, + Rgb::WHITE + ) + .await, + Err(DomainError::Conflict { .. }) + )); + let mut changed = before; + changed.zones[0].role = ZoneRole::Display; + changed.zones[0].display_target = Some(hypercolor_types::scene::DisplayFaceTarget { + device_id: hypercolor_types::device::DeviceId::new(), + blend_mode: BlendMode::Replace, + opacity: 1.0, + }); + let mut mutation = state.scene_manager.begin_mutation().await; + mutation + .update_scene(changed.clone()) + .expect("display fixture"); + let commit = commit_scene(&state.domains.scene, mutation) + .await + .expect("display commit"); + assert!(matches!( + set_color( + &state.domains.scene, + RuntimeZoneTarget { + revision: commit.revision(), + ..target + }, + Rgb::WHITE + ) + .await, + Err(DomainError::Validation { .. }) + )); + assert_eq!( + state.scene_manager.snapshot().await.get(&target.scene_id), + Some(&changed) + ); +} + +#[tokio::test] +async fn a_concurrent_scene_activation_rejects_an_already_prepared_color_candidate() { + let (state, _dir, before, target) = fixture().await; + let mut color = state.scene_manager.begin_mutation().await; + color + .set_runtime_zone_color(target, Rgb::WHITE) + .expect("prepare before activation"); + let mut activation = state.scene_manager.begin_mutation().await; + activation + .activate(SceneId::DEFAULT, None, SceneChangeReason::UserActivate) + .expect("activate default"); + commit_scene(&state.domains.scene, activation) + .await + .expect("competing commit"); + assert!(matches!( + commit_scene(&state.domains.scene, color).await, + Err(DomainError::Conflict { .. }) + )); + assert_eq!( + state.scene_manager.snapshot().await.get(&target.scene_id), + Some(&before) + ); +} + +#[cfg(feature = "persistence-test-hooks")] +#[tokio::test] +async fn scene_reservation_failure_does_not_admit_the_color() { + use hypercolor_daemon::persistence::set_injected_serialization_failures; + let (state, _dir, before, target) = fixture().await; + set_injected_serialization_failures(1); + let result = set_color(&state.domains.scene, target, Rgb::WHITE).await; + set_injected_serialization_failures(0); + assert!(matches!(result, Err(DomainError::Internal(_)))); + assert_eq!(state.scene_manager.revision(), target.revision); + assert_eq!( + state.scene_manager.snapshot().await.get(&target.scene_id), + Some(&before) + ); +} + +#[cfg(feature = "persistence-test-hooks")] +#[tokio::test] +async fn an_admitted_failed_write_remains_retrying_without_a_success_event() { + let (state, _dir, _before, target) = fixture().await; + let writer = AtomicFileWriter::new(&state.data_dir.join("scenes.json")).expect("scene writer"); + writer.set_injected_replace_failures(usize::MAX); + let mut events = state.event_bus.subscribe_all(); + let result = set_color(&state.domains.scene, target, Rgb::WHITE).await; + writer.set_injected_replace_failures(0); + writer.kick(); + let result = result.expect("admitted result is retained"); + assert_eq!(result.commit.durability(), CommitDurability::Retrying); + assert!(result.commit.retry_error().is_some()); + assert!(matches!( + result.runtime_session, + RuntimeSessionSaveOutcome::Attempted { + scene_store: Err(_), + .. + } + )); + assert_eq!(state.scene_manager.revision(), result.commit.revision()); + assert!( + events.try_recv().is_err(), + "no unproven applied announcement" + ); +} + +#[cfg(feature = "persistence-test-hooks")] +#[tokio::test] +async fn runtime_write_failure_does_not_erase_original_scene_commit_evidence() { + let (state, _dir, _before, target) = fixture().await; + let writer = AtomicFileWriter::new(&state.runtime_state_path).expect("runtime writer"); + writer.set_injected_replace_failures(usize::MAX); + let result = set_color(&state.domains.scene, target, Rgb::WHITE).await; + writer.set_injected_replace_failures(0); + writer.kick(); + let result = result.expect("scene admitted"); + assert_eq!(result.commit.durability(), CommitDurability::Written); + assert!(matches!( + result.runtime_session, + RuntimeSessionSaveOutcome::Attempted { + snapshot: Err(_), + .. + } + )); +} + +async fn default_target(state: &AppState, scene: &Scene) -> RuntimeZoneTarget { + let mut default = state + .scene_manager + .snapshot() + .await + .get(&SceneId::DEFAULT) + .expect("default") + .clone(); + default.zones = scene.zones.clone(); + let mut mutation = state.scene_manager.begin_mutation().await; + mutation.update_scene(default).expect("seed default zones"); + mutation + .activate(SceneId::DEFAULT, None, SceneChangeReason::UserActivate) + .expect("default active"); + let commit = commit_scene(&state.domains.scene, mutation) + .await + .expect("default fixture"); + RuntimeZoneTarget { + scene_id: SceneId::DEFAULT, + zone_id: scene.zones[0].id, + revision: commit.revision(), + } +} + +#[tokio::test] +async fn default_scene_color_requires_the_written_runtime_payload() { + let (state, _dir, before, _) = fixture().await; + let target = default_target(&state, &before).await; + let result = set_color(&state.domains.scene, target, Rgb::new(10, 20, 30)) + .await + .expect("default color"); + assert!(result.has_written_layers()); + let disk = hypercolor_daemon::runtime_state::load(&state.runtime_state_path) + .expect("read runtime") + .expect("saved projection"); + assert_eq!( + disk.default_scene_zones + .iter() + .find(|z| z.id == target.zone_id) + .expect("persisted zone") + .layers, + result.zone.layers + ); +} + +#[cfg(feature = "persistence-test-hooks")] +#[tokio::test] +async fn default_scene_does_not_claim_durability_from_a_named_scene_store_write() { + let (state, _dir, before, _) = fixture().await; + let target = default_target(&state, &before).await; + let writer = AtomicFileWriter::new(&state.runtime_state_path).expect("runtime writer"); + writer.set_injected_replace_failures(usize::MAX); + let result = set_color(&state.domains.scene, target, Rgb::WHITE).await; + writer.set_injected_replace_failures(0); + writer.kick(); + let result = result.expect("default color admitted"); + assert_eq!(result.commit.durability(), CommitDurability::Written); + assert!(!result.has_written_layers()); +} + +#[tokio::test] +async fn overwritten_default_layers_are_not_proven_by_a_later_written_projection() { + use hypercolor_daemon::domain::runtime_zone::RuntimeZoneColorOutcome; + let (state, _dir, before, _) = fixture().await; + let target = default_target(&state, &before).await; + let mut first = state.scene_manager.begin_mutation().await; + let zone = first + .set_runtime_zone_color(target, Rgb::WHITE) + .expect("first color"); + let scene_kind = first.scenes().get(&target.scene_id).expect("default").kind; + let commit = commit_scene(&state.domains.scene, first) + .await + .expect("first commit"); + let second = set_color( + &state.domains.scene, + RuntimeZoneTarget { + revision: commit.revision(), + ..target + }, + Rgb::BLACK, + ) + .await + .expect("overwrite before first save"); + assert!(second.has_written_layers()); + let runtime_session = state.domains.runtime_session.save_with_outcome().await; + let first = RuntimeZoneColorOutcome { + target, + scene_kind, + zone, + commit, + runtime_session, + }; + assert!( + !first.has_written_layers(), + "later Written contains different layer identities and content" + ); +} + +#[tokio::test] +async fn nondefault_ephemeral_scene_has_no_named_store_durability_claim() { + let (state, _dir, mut before, target) = fixture().await; + before.kind = SceneKind::Ephemeral; + before.id = SceneId::new(); + let target = RuntimeZoneTarget { + scene_id: before.id, + ..target + }; + let mut mutation = state.scene_manager.begin_mutation().await; + mutation.create_scene(before).expect("ephemeral fixture"); + mutation + .activate(target.scene_id, None, SceneChangeReason::UserActivate) + .expect("ephemeral active"); + let commit = commit_scene(&state.domains.scene, mutation) + .await + .expect("ephemeral commit"); + let result = set_color( + &state.domains.scene, + RuntimeZoneTarget { + revision: commit.revision(), + ..target + }, + Rgb::WHITE, + ) + .await + .expect("volatile color still applies"); + assert_eq!(result.commit.durability(), CommitDurability::Written); + assert!( + !result.has_written_layers(), + "ephemeral non-default content is not in either durable store" + ); +} + +#[tokio::test] +async fn an_exact_but_superseded_runtime_payload_is_not_a_written_receipt() { + use hypercolor_daemon::domain::runtime_zone::RuntimeZoneColorOutcome; + use hypercolor_daemon::runtime_state; + let (state, _dir, before, _) = fixture().await; + let target = default_target(&state, &before).await; + let mut mutation = state.scene_manager.begin_mutation().await; + let zone = mutation + .set_runtime_zone_color(target, Rgb::WHITE) + .expect("candidate"); + let scene_kind = mutation + .scenes() + .get(&target.scene_id) + .expect("default") + .kind; + let commit = commit_scene(&state.domains.scene, mutation) + .await + .expect("scene admission"); + let pending = + runtime_state::reserve_save(&state.runtime_state_path).expect("older reservation"); + let projection = state.domains.runtime_session.snapshot().await; + let newer = state.domains.runtime_session.save_with_outcome().await; + assert!(matches!( + newer, + RuntimeSessionSaveOutcome::Attempted { + snapshot: Ok(AtomicWriteOutcome::Written), + .. + } + )); + let snapshot = runtime_state::save_reserved(pending, &projection); + assert!(matches!(snapshot, Ok(AtomicWriteOutcome::Superseded))); + let result = RuntimeZoneColorOutcome { + target, + scene_kind, + zone, + commit, + runtime_session: RuntimeSessionSaveOutcome::Attempted { + scene_store: Ok(Some(AtomicWriteOutcome::Written)), + projection, + snapshot, + }, + }; + assert!( + !result.has_written_layers(), + "even matching content needs the original writer's Written evidence" + ); +} + +#[tokio::test] +async fn runtime_reservation_failure_is_returned_without_attempting_either_save() { + use hypercolor_daemon::app_state::AppStateBuilder; + let dir = tempfile::tempdir().expect("fixture directory"); + let blocked = dir.path().join("blocked"); + std::fs::write(&blocked, b"not a directory").expect("block runtime destination"); + let state = AppStateBuilder::new(dir.path().join("data")) + .with_runtime_state_path(blocked.join("runtime.json")) + .build(); + assert!(matches!( + state.domains.runtime_session.save_with_outcome().await, + RuntimeSessionSaveOutcome::BeforeAdmission { .. } + )); +} + +#[tokio::test] +async fn core_stack_replacement_validates_the_whole_input_before_mutating() { + use hypercolor_core::scene::LayerMutationError; + let (state, _dir, before, target) = fixture().await; + let mut manager = state.scene_manager.snapshot().await; + let valid = layer(); + let mut invalid = layer(); + invalid.opacity = f32::NAN; + assert!(matches!( + manager.replace_zone_layer_stack( + target.scene_id, + target.zone_id, + vec![valid.clone(), invalid] + ), + Err(LayerMutationError::InvalidLayer { .. }) + )); + assert!(matches!( + manager.replace_zone_layer_stack( + target.scene_id, + target.zone_id, + vec![valid.clone(), valid.clone()] + ), + Err(LayerMutationError::DuplicateLayer { .. }) + )); + assert!(matches!( + manager.replace_zone_layer_stack(SceneId::new(), target.zone_id, vec![valid.clone()]), + Err(LayerMutationError::SceneMissing) + )); + assert!(matches!( + manager.replace_zone_layer_stack(target.scene_id, ZoneId::new(), vec![valid]), + Err(LayerMutationError::ZoneMissing) + )); + assert_eq!(manager.get(&target.scene_id), Some(&before)); + let (zone, _) = manager + .replace_zone_layer_stack(target.scene_id, target.zone_id, Vec::new()) + .expect("explicit empty stack"); + assert!(zone.layers.is_empty()); + let mut restored = manager.get(&target.scene_id).expect("scene").clone(); + restored.zones[0].layers = before.zones[0].layers.clone(); + restored.zones[0].layers_version = before.zones[0].layers_version; + assert_eq!(restored, before); +} From 0c4d5f3c56848406bf460ea292e56317b0dac1e0 Mon Sep 17 00:00:00 2001 From: Stefanie Jane Date: Sat, 12 Sep 2026 22:42:23 -0700 Subject: [PATCH 3/4] fix(ui): make desktop bridge tests compile for wasm Keep the native test notifier out of browser builds so it cannot collide with the real browser event publisher. Move existing installation helpers before the test module to satisfy the all-target lint gate. --- crates/hypercolor-ui/src/tauri_bridge.rs | 36 ++++++++++++------------ 1 file changed, 18 insertions(+), 18 deletions(-) diff --git a/crates/hypercolor-ui/src/tauri_bridge.rs b/crates/hypercolor-ui/src/tauri_bridge.rs index ee280b1e5..60d4805af 100644 --- a/crates/hypercolor-ui/src/tauri_bridge.rs +++ b/crates/hypercolor-ui/src/tauri_bridge.rs @@ -119,7 +119,7 @@ fn notify_verified_daemon_connection_change() { } } -#[cfg(test)] +#[cfg(all(test, not(target_arch = "wasm32")))] fn notify_verified_daemon_connection_change() {} /// Initialize the bundled app's process-memory daemon transport. @@ -874,6 +874,23 @@ fn js_error_string(value: JsValue) -> String { }) } +/// Installation guidance from the desktop host, when the app owns the session. +#[cfg(target_arch = "wasm32")] +pub async fn openrgb_install_hints() +-> Result>, String> { + let Some(invoke) = tauri_invoke() else { + return Ok(None); + }; + let value = invoke_command(&invoke, "openrgb_install_hints", None).await?; + serde_json_from_js_value(value).map(Some) +} + +#[cfg(not(target_arch = "wasm32"))] +pub async fn openrgb_install_hints() +-> Result>, String> { + Ok(None) +} + #[cfg(test)] mod tests { use std::{ @@ -1182,20 +1199,3 @@ mod tests { } } } - -/// Installation guidance from the desktop host, when the app owns the session. -#[cfg(target_arch = "wasm32")] -pub async fn openrgb_install_hints() --> Result>, String> { - let Some(invoke) = tauri_invoke() else { - return Ok(None); - }; - let value = invoke_command(&invoke, "openrgb_install_hints", None).await?; - serde_json_from_js_value(value).map(Some) -} - -#[cfg(not(target_arch = "wasm32"))] -pub async fn openrgb_install_hints() --> Result>, String> { - Ok(None) -} From 1201eeac358fe978561691c2edcbd15a6c50efc6 Mon Sep 17 00:00:00 2001 From: Stefanie Jane Date: Sat, 12 Sep 2026 22:42:32 -0700 Subject: [PATCH 4/4] feat(ui): add an owned incremental browser response reader Pull browser response chunks only on demand and copy no more than the consumer's reserved capacity into Rust. Retain one browser-owned chunk, release the reader at EOF, and abort pending reads on cancellation. Exercise the real browser stream in WASM tests, including large chunks, failed reads, invalid values, empty responses, and reader ownership. Native fetch integration remains a separate transport change. --- crates/hypercolor-ui/Cargo.toml | 5 + .../hypercolor-ui/src/api/browser_response.rs | 158 ++++++++++++++++++ crates/hypercolor-ui/src/api/mod.rs | 2 + .../tests/browser_response_tests.rs | 150 +++++++++++++++++ 4 files changed, 315 insertions(+) create mode 100644 crates/hypercolor-ui/src/api/browser_response.rs create mode 100644 crates/hypercolor-ui/tests/browser_response_tests.rs diff --git a/crates/hypercolor-ui/Cargo.toml b/crates/hypercolor-ui/Cargo.toml index 3667773be..7be17b865 100644 --- a/crates/hypercolor-ui/Cargo.toml +++ b/crates/hypercolor-ui/Cargo.toml @@ -38,12 +38,17 @@ console_error_panic_hook = "0.1" log = "0.4" console_log = "1" web-sys = { version = "0.3", features = [ + "AbortController", + "AbortSignal", "Blob", "BlobPropertyBag", "File", "FileList", "FileReader", "FormData", + "ReadableStream", + "ReadableStreamDefaultReader", + "Response", "HtmlCanvasElement", "HtmlInputElement", "HtmlMediaElement", diff --git a/crates/hypercolor-ui/src/api/browser_response.rs b/crates/hypercolor-ui/src/api/browser_response.rs new file mode 100644 index 000000000..088d6705f --- /dev/null +++ b/crates/hypercolor-ui/src/api/browser_response.rs @@ -0,0 +1,158 @@ +//! Single-reader ownership of an incremental browser response. + +use std::{ + future::Future, + num::NonZeroUsize, + pin::Pin, + task::{Context, Poll}, +}; + +use js_sys::{Reflect, Uint8Array}; +use wasm_bindgen::{JsCast, JsValue, prelude::wasm_bindgen}; +use wasm_bindgen_futures::JsFuture; +use web_sys::{AbortController, ReadableStreamDefaultReader, Response}; + +use super::http_transport::{HttpBodySource, HttpStreamError}; + +#[wasm_bindgen(inline_js = r#" +export function cancelResponseReader(reader) { + // Cancellation is cleanup; an already failed stream can reject its promise. + reader.cancel().catch(() => {}); + reader.releaseLock(); +} +"#)] +extern "C" { + #[wasm_bindgen(js_name = cancelResponseReader)] + fn cancel_response_reader(reader: &ReadableStreamDefaultReader); +} + +/// Keeps at most one browser-provided chunk and one caller-sized Rust copy. +/// Browser network buffering and the size of its chunk remain browser-owned. +/// Content-Length is not an exact source length because fetch may decode content. +pub struct BrowserResponseSource { + reader: Option, + controller: Option, + pending: Option, + buffered: Option<(Uint8Array, u32)>, + complete: bool, +} + +impl BrowserResponseSource { + /// Acquire the response's single reader without reading ahead. + /// + /// # Errors + /// Rejects a response whose body is already locked, and aborts the exchange. + pub fn new(response: &Response, controller: AbortController) -> Result { + let reader = response + .body() + .map(|body| { + ReadableStreamDefaultReader::new(&body).map_err(|error| { + controller.abort(); + browser_error(error) + }) + }) + .transpose()?; + Ok(Self { + complete: reader.is_none(), + controller: reader.as_ref().map(|_| controller), + reader, + pending: None, + buffered: None, + }) + } + + fn poll_read( + &mut self, + context: &mut Context<'_>, + maximum: NonZeroUsize, + ) -> Poll>, HttpStreamError>> { + loop { + if let Some((bytes, offset)) = &mut self.buffered { + let remaining = (bytes.length() - *offset) as usize; + let length = remaining.min(maximum.get()); + let end = *offset + length as u32; + let mut output = vec![0; length]; + bytes.subarray(*offset, end).copy_to(&mut output); + *offset = end; + if end == bytes.length() { + self.buffered = None; + } + return Poll::Ready(Ok(Some(output))); + } + if self.complete { + return Poll::Ready(Ok(None)); + } + if self.pending.is_none() { + self.pending = Some(JsFuture::from( + self.reader + .as_ref() + .expect("active response owns its reader") + .read(), + )); + } + let result = match Pin::new(self.pending.as_mut().expect("read exists")).poll(context) { + Poll::Pending => return Poll::Pending, + Poll::Ready(result) => result, + }; + self.pending = None; + let result = result.map_err(browser_error)?; + let done = Reflect::get(&result, &JsValue::from_str("done")).map_err(browser_error)?; + if done.as_bool() == Some(true) { + self.complete = true; + if let Some(reader) = self.reader.take() { + reader.release_lock(); + } + self.controller = None; + return Poll::Ready(Ok(None)); + } + let value = + Reflect::get(&result, &JsValue::from_str("value")).map_err(browser_error)?; + let bytes = value + .dyn_into::() + .map_err(|_| HttpStreamError::InvalidChunk)?; + if bytes.length() != 0 { + self.buffered = Some((bytes, 0)); + } + } + } +} + +impl HttpBodySource for BrowserResponseSource { + fn exact_length(&self) -> Option { + None + } + + fn poll_chunk( + &mut self, + context: &mut Context<'_>, + maximum: NonZeroUsize, + ) -> Poll>, HttpStreamError>> { + let result = self.poll_read(context, maximum); + if matches!(result, Poll::Ready(Err(_))) { + self.cancel(); + } + result + } + + fn cancel(&mut self) { + self.complete = true; + if let Some(controller) = self.controller.take() { + controller.abort(); + } + if let Some(reader) = self.reader.take() { + cancel_response_reader(&reader); + } + self.pending = None; + self.buffered = None; + } +} + +impl Drop for BrowserResponseSource { + fn drop(&mut self) { + self.cancel(); + } +} + +fn browser_error(error: JsValue) -> HttpStreamError { + HttpStreamError::Transport(format!("browser response read failed: {error:?}")) +} diff --git a/crates/hypercolor-ui/src/api/mod.rs b/crates/hypercolor-ui/src/api/mod.rs index 043d37e5b..3ef578a34 100644 --- a/crates/hypercolor-ui/src/api/mod.rs +++ b/crates/hypercolor-ui/src/api/mod.rs @@ -12,6 +12,8 @@ use crate::app::WsContext; pub mod assets; #[cfg(target_arch = "wasm32")] pub mod browser_body; +#[cfg(target_arch = "wasm32")] +pub mod browser_response; pub mod client; pub mod config; pub mod controls; diff --git a/crates/hypercolor-ui/tests/browser_response_tests.rs b/crates/hypercolor-ui/tests/browser_response_tests.rs new file mode 100644 index 000000000..54a899e14 --- /dev/null +++ b/crates/hypercolor-ui/tests/browser_response_tests.rs @@ -0,0 +1,150 @@ +#![cfg(target_arch = "wasm32")] + +use hypercolor_ui::api::{ + browser_response::BrowserResponseSource, + http_transport::{HttpBody, HttpBodySource, HttpCancellation}, +}; +use std::{ + future::Future, + num::NonZeroUsize, + task::{Context, Waker}, +}; +use wasm_bindgen::prelude::*; +use wasm_bindgen_test::*; +use web_sys::{AbortController, Response}; + +wasm_bindgen_test_configure!(run_in_browser); + +#[wasm_bindgen(inline_js = r#" +export function fixture(mode) { + let pulls=0,cancels=0; + const stream=new ReadableStream({pull(controller){ + pulls++; + if(mode==='pending') return new Promise(()=>{}); + if(mode==='failure') {controller.error(new Error('fixture failure'));return;} + if(mode==='invalid') {controller.enqueue('not bytes');return;} + if(pulls===1) controller.enqueue(new Uint8Array(100*1024).fill(91)); + else controller.close(); + },cancel(){cancels++;}}, {highWaterMark:0}); + const response=new Response(stream, {headers:{'Content-Length':'7','Content-Encoding':'gzip'}}); + response.fixturePulls=()=>pulls; response.fixtureCancels=()=>cancels; + return response; +} +export function pulls(response){return response.fixturePulls();} +export function cancels(response){return response.fixtureCancels();} +export function locked(response){return response.body.locked;} +"#)] +extern "C" { + fn fixture(mode: &str) -> Response; + fn pulls(response: &Response) -> u32; + fn cancels(response: &Response) -> u32; + fn locked(response: &Response) -> bool; +} + +fn body(response: &Response, controller: AbortController) -> HttpBody { + let source = BrowserResponseSource::new(response, controller).expect("unlocked response"); + assert_eq!( + source.exact_length(), + None, + "decoded response is not Content-Length bytes" + ); + HttpBody::new(Box::new(source), HttpCancellation::new()) +} + +#[wasm_bindgen_test] +async fn browser_chunks_are_pulled_on_demand_and_copied_within_requested_capacity() { + let response = fixture("normal"); + let controller = AbortController::new().expect("controller"); + let signal = controller.signal(); + let mut body = body(&response, controller); + assert_eq!(pulls(&response), 0); + assert!(locked(&response)); + let maximum = NonZeroUsize::new(1024).expect("capacity"); + for _ in 0..100 { + let bytes = body + .read_chunk(maximum) + .await + .expect("chunk") + .expect("body bytes"); + assert_eq!(bytes, vec![91; 1024]); + assert_eq!( + pulls(&response), + 1, + "no second browser read while buffered bytes remain" + ); + } + assert!(body.read_chunk(maximum).await.expect("EOF").is_none()); + assert_eq!(pulls(&response), 2); + assert!(!locked(&response)); + drop(body); + assert!(!signal.aborted(), "completed response needs no abort"); +} + +#[wasm_bindgen_test] +async fn pending_read_drop_aborts_and_unlocks_the_browser_body() { + let response = fixture("pending"); + let controller = AbortController::new().expect("controller"); + let signal = controller.signal(); + let mut body = body(&response, controller); + let mut read = Box::pin(body.read_chunk(NonZeroUsize::new(1024).expect("capacity"))); + assert!( + read.as_mut() + .poll(&mut Context::from_waker(Waker::noop())) + .is_pending() + ); + drop(read); + drop(body); + assert!(signal.aborted()); + assert!(!locked(&response)); + assert_eq!(cancels(&response), 1); +} + +#[wasm_bindgen_test] +async fn browser_failure_and_nonbyte_chunks_abort_without_leaking_a_reader() { + for mode in ["failure", "invalid"] { + let response = fixture(mode); + let controller = AbortController::new().expect("controller"); + let signal = controller.signal(); + let mut body = body(&response, controller); + assert!( + body.read_chunk(NonZeroUsize::new(1024).expect("capacity")) + .await + .is_err() + ); + assert!(signal.aborted()); + assert!(!locked(&response)); + } +} + +#[wasm_bindgen_test] +fn acquiring_an_already_locked_response_refuses_and_aborts_the_new_exchange() { + let response = fixture("normal"); + let first = BrowserResponseSource::new(&response, AbortController::new().expect("controller")) + .expect("first owner"); + let second = AbortController::new().expect("controller"); + let signal = second.signal(); + assert!(BrowserResponseSource::new(&response, second).is_err()); + assert!(signal.aborted()); + assert!( + locked(&response), + "failed acquisition cannot steal the first owner" + ); + drop(first); + assert!(!locked(&response)); +} + +#[wasm_bindgen_test] +async fn absent_response_body_finishes_without_aborting_a_completed_exchange() { + let response = Response::new().expect("empty response"); + let controller = AbortController::new().expect("controller"); + let signal = controller.signal(); + let mut body = body(&response, controller); + assert!( + body.read_chunk(NonZeroUsize::new(1024).expect("capacity")) + .await + .expect("empty body") + .is_none() + ); + drop(body); + assert!(!signal.aborted()); +}