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/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-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/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/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); + } +} 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); +} 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/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) -} 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()); +}