Skip to main content

homecore_hap/
bridge.rs

1//! `HapBridge` — owns the set of HOMECORE entities exposed as HAP accessories.
2//!
3//! The bridge owns mappings and their event stream. The feature-gated network
4//! lifecycle is started separately with `start_server`.
5
6use std::collections::HashMap;
7use std::sync::{Arc, RwLock};
8
9use homecore::entity::EntityId;
10use tokio::sync::broadcast;
11
12use crate::accessory::{HapAccessoryType, HapCharacteristic, HapCharacteristicValue};
13use crate::error::HapError;
14use crate::mapping::{AccessoryMapping, EntityToAccessoryMapper};
15use crate::mdns::{HapServiceRecord, MdnsAdvertiser, NullAdvertiser};
16
17/// One registered HAP accessory — an entity + its last-known mapping.
18#[derive(Debug, Clone)]
19pub struct ExposedAccessory {
20    pub entity_id: EntityId,
21    pub accessory_type: HapAccessoryType,
22    pub mapping: AccessoryMapping,
23}
24
25/// A characteristic snapshot emitted after a registered entity changes.
26#[derive(Debug, Clone)]
27pub struct CharacteristicEvent {
28    pub entity_id: EntityId,
29    pub accessory_type: HapAccessoryType,
30    pub characteristics: Vec<(HapCharacteristic, HapCharacteristicValue)>,
31}
32
33struct BridgeInner {
34    accessories: HashMap<EntityId, ExposedAccessory>,
35}
36
37/// HOMECORE-to-HAP accessory bridge state.
38///
39/// Call [`HapBridge::add_accessory`] to register entities and
40/// [`HapBridge::running_accessories`] to read back what is currently
41/// registered. Use `start_server` for the bounded TCP lifecycle.
42#[derive(Clone)]
43pub struct HapBridge {
44    inner: Arc<RwLock<BridgeInner>>,
45    advertiser: Arc<dyn MdnsAdvertiser>,
46    events: broadcast::Sender<CharacteristicEvent>,
47    pub service_record: HapServiceRecord,
48}
49
50impl HapBridge {
51    /// Create a bridge with the given service record and a `NullAdvertiser`.
52    pub fn new(service_record: HapServiceRecord) -> Self {
53        Self::with_advertiser(service_record, Arc::new(NullAdvertiser))
54    }
55
56    /// Create a bridge with a custom `MdnsAdvertiser`.
57    pub fn with_advertiser(
58        service_record: HapServiceRecord,
59        advertiser: Arc<dyn MdnsAdvertiser>,
60    ) -> Self {
61        let (events, _) = broadcast::channel(128);
62        Self {
63            inner: Arc::new(RwLock::new(BridgeInner {
64                accessories: HashMap::new(),
65            })),
66            advertiser,
67            events,
68            service_record,
69        }
70    }
71
72    /// Register an entity as a HAP accessory.
73    ///
74    /// The entity's current mapping is computed from `state`; call
75    /// `update_accessory` on each `StateChanged` event to keep it fresh.
76    ///
77    /// Returns `HapError::AlreadyRegistered` if the entity is already
78    /// registered. Call `remove_accessory` first to replace it.
79    pub fn add_accessory(
80        &self,
81        entity_id: &EntityId,
82        state: &homecore::entity::State,
83    ) -> Result<(), HapError> {
84        let mapping = EntityToAccessoryMapper::map(entity_id, state)?;
85        let accessory_type = mapping.accessory_type;
86        let exposed = ExposedAccessory {
87            entity_id: entity_id.clone(),
88            accessory_type,
89            mapping,
90        };
91        let mut inner = self.inner.write().unwrap();
92        if inner.accessories.contains_key(entity_id) {
93            return Err(HapError::AlreadyRegistered(entity_id.as_str().to_owned()));
94        }
95        inner.accessories.insert(entity_id.clone(), exposed);
96        tracing::debug!(entity = %entity_id, ?accessory_type, "HAP accessory registered");
97        Ok(())
98    }
99
100    /// Remove a registered accessory.
101    ///
102    /// Returns `HapError::EntityNotFound` if the entity was not registered.
103    pub fn remove_accessory(&self, entity_id: &EntityId) -> Result<(), HapError> {
104        let mut inner = self.inner.write().unwrap();
105        if inner.accessories.remove(entity_id).is_none() {
106            return Err(HapError::EntityNotFound(entity_id.as_str().to_owned()));
107        }
108        tracing::debug!(entity = %entity_id, "HAP accessory removed");
109        Ok(())
110    }
111
112    /// Refresh a registered accessory and notify event subscribers.
113    pub fn update_accessory(
114        &self,
115        entity_id: &EntityId,
116        state: &homecore::entity::State,
117    ) -> Result<(), HapError> {
118        let mapping = EntityToAccessoryMapper::map(entity_id, state)?;
119        let accessory_type = mapping.accessory_type;
120        {
121            let mut inner = self.inner.write().unwrap();
122            let accessory = inner
123                .accessories
124                .get_mut(entity_id)
125                .ok_or_else(|| HapError::EntityNotFound(entity_id.as_str().to_owned()))?;
126            accessory.accessory_type = accessory_type;
127            accessory.mapping = mapping.clone();
128        }
129        let _ = self.events.send(CharacteristicEvent {
130            entity_id: entity_id.clone(),
131            accessory_type,
132            characteristics: mapping.characteristics,
133        });
134        Ok(())
135    }
136
137    /// Subscribe to bounded characteristic updates. Lagging receivers receive
138    /// Tokio's explicit `Lagged` error and must resynchronize from a snapshot.
139    pub fn subscribe_events(&self) -> broadcast::Receiver<CharacteristicEvent> {
140        self.events.subscribe()
141    }
142
143    /// Snapshot all currently registered accessories.
144    pub fn running_accessories(&self) -> Vec<ExposedAccessory> {
145        self.inner
146            .read()
147            .unwrap()
148            .accessories
149            .values()
150            .cloned()
151            .collect()
152    }
153
154    /// Number of registered accessories.
155    pub fn len(&self) -> usize {
156        self.inner.read().unwrap().accessories.len()
157    }
158
159    pub fn is_empty(&self) -> bool {
160        self.len() == 0
161    }
162
163    /// Start advertisement only.
164    ///
165    /// This legacy lifecycle does not bind a TCP listener. New integrations
166    /// should call `start_server`, which advertises only after binding.
167    pub async fn start(&self) -> Result<(), HapError> {
168        self.advertiser.advertise(&self.service_record).await?;
169        tracing::info!(
170            instance = %self.service_record.instance_name,
171            port = self.service_record.port,
172            "HAP advertisement started without a TCP server"
173        );
174        Ok(())
175    }
176
177    /// Graceful shutdown — retracts mDNS advertisement.
178    pub async fn stop(&self) -> Result<(), HapError> {
179        self.advertiser
180            .retract(&self.service_record.instance_name)
181            .await?;
182        Ok(())
183    }
184}
185
186#[cfg(test)]
187mod tests {
188    use super::*;
189    use homecore::entity::{EntityId, State};
190    use homecore::event::Context;
191
192    fn make_bridge() -> HapBridge {
193        HapBridge::new(HapServiceRecord::bridge(
194            "RuView Sense",
195            51826,
196            "AA:BB:CC:DD:EE:FF",
197        ))
198    }
199
200    fn light_state(name: &str, on: bool, brightness: u8) -> (EntityId, State) {
201        let eid = EntityId::parse(format!("light.{name}")).unwrap();
202        let attrs = serde_json::json!({"brightness": brightness});
203        let s = State::new(
204            eid.clone(),
205            if on { "on" } else { "off" },
206            attrs,
207            Context::default(),
208        );
209        (eid, s)
210    }
211
212    #[test]
213    fn add_remove_roundtrip() {
214        let bridge = make_bridge();
215        let (eid, s) = light_state("kitchen", true, 200);
216
217        assert!(bridge.is_empty());
218        bridge.add_accessory(&eid, &s).unwrap();
219        assert_eq!(bridge.len(), 1);
220
221        let acc = bridge.running_accessories();
222        assert_eq!(acc.len(), 1);
223        assert_eq!(acc[0].entity_id, eid);
224        assert_eq!(acc[0].accessory_type, HapAccessoryType::Lightbulb);
225
226        bridge.remove_accessory(&eid).unwrap();
227        assert!(bridge.is_empty());
228    }
229
230    #[test]
231    fn add_duplicate_returns_error() {
232        let bridge = make_bridge();
233        let (eid, s) = light_state("kitchen", true, 200);
234        bridge.add_accessory(&eid, &s).unwrap();
235        let err = bridge.add_accessory(&eid, &s).unwrap_err();
236        assert!(matches!(err, HapError::AlreadyRegistered(_)));
237    }
238
239    #[test]
240    fn remove_nonexistent_returns_error() {
241        let bridge = make_bridge();
242        let eid = EntityId::parse("light.ghost").unwrap();
243        let err = bridge.remove_accessory(&eid).unwrap_err();
244        assert!(matches!(err, HapError::EntityNotFound(_)));
245    }
246
247    #[tokio::test]
248    async fn update_emits_characteristic_event() {
249        let bridge = make_bridge();
250        let (eid, initial) = light_state("kitchen", false, 10);
251        bridge.add_accessory(&eid, &initial).unwrap();
252        let mut events = bridge.subscribe_events();
253        let (_, updated) = light_state("kitchen", true, 200);
254        bridge.update_accessory(&eid, &updated).unwrap();
255
256        let event = events.recv().await.unwrap();
257        assert_eq!(event.entity_id, eid);
258        assert!(event
259            .characteristics
260            .contains(&(HapCharacteristic::On, HapCharacteristicValue::Bool(true))));
261    }
262
263    #[tokio::test]
264    async fn start_stop_with_null_advertiser() {
265        let bridge = make_bridge();
266        bridge.start().await.unwrap();
267        bridge.stop().await.unwrap();
268    }
269}