use std::net::IpAddr;
use std::sync::Arc;
use std::thread::{self, JoinHandle};
use parking_lot::RwLock;
use sonos_event_manager::SonosEventManager;
use sonos_stream::events::{EnrichedEvent, EventData};
use sonos_api::ServiceScope;
use crate::decoder::{decode_event, decode_topology_event, PropertyChange, TopologyChanges};
use crate::iter::EventFanout;
use crate::model::SpeakerId;
use crate::property::{GroupMembership, Property, Scope};
use crate::state::{
is_pair_watched, ChangeEvent, ChangeSource, StateStore, WatchCounts, WriteStamp,
};
pub(crate) fn spawn_state_event_worker(
event_manager: Arc<SonosEventManager>,
store: Arc<RwLock<StateStore>>,
watched: Arc<RwLock<WatchCounts>>,
fanout: Arc<EventFanout>,
ip_to_speaker: Arc<RwLock<std::collections::HashMap<IpAddr, SpeakerId>>>,
) -> JoinHandle<()> {
thread::spawn(move || {
tracing::info!("State event worker started, waiting for events...");
run_event_loop(
event_manager.iter(),
&store,
&watched,
&fanout,
&ip_to_speaker,
);
tracing::info!("State event worker stopped");
})
}
const PANIC_ESCALATION_INTERVAL: u64 = 10;
fn run_event_loop<I>(
events: I,
store: &Arc<RwLock<StateStore>>,
watched: &Arc<RwLock<WatchCounts>>,
fanout: &EventFanout,
ip_to_speaker: &Arc<RwLock<std::collections::HashMap<IpAddr, SpeakerId>>>,
) where
I: Iterator<Item = EnrichedEvent>,
{
let mut panic_count: u64 = 0;
for event in events {
let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
handle_event(&event, store, watched, fanout, ip_to_speaker);
}));
if result.is_err() {
panic_count += 1;
tracing::error!(
"Panic while processing event from {} for service {:?}; \
skipping this event and continuing (panic #{} for this worker)",
event.speaker_ip,
event.service,
panic_count
);
if panic_count % PANIC_ESCALATION_INTERVAL == 0 {
tracing::error!(
"State event worker has now panicked {} times — state updates \
are being dropped and this is a bug that needs fixing",
panic_count
);
}
}
}
}
fn handle_event(
event: &EnrichedEvent,
store: &Arc<RwLock<StateStore>>,
watched: &Arc<RwLock<WatchCounts>>,
fanout: &EventFanout,
ip_to_speaker: &Arc<RwLock<std::collections::HashMap<IpAddr, SpeakerId>>>,
) {
tracing::debug!(
"Received event from {} for service {:?}",
event.speaker_ip,
event.service
);
#[cfg(test)]
if event.speaker_ip == tests::PANIC_TRIGGER_IP {
panic!("injected test panic while decoding event");
}
let stamp = WriteStamp::observed_at(ChangeSource::Event, event.observed_at);
if let EventData::ZoneGroupTopology(ref zgt_event) = event.event_data {
tracing::debug!("Processing ZoneGroupTopology event");
let topology_changes = decode_topology_event(zgt_event);
apply_topology_changes(
store,
watched,
fanout,
ip_to_speaker,
topology_changes,
stamp,
);
return;
}
let speaker_id = {
let ip_map = ip_to_speaker.read();
tracing::debug!(
"ip_to_speaker map has {} entries: {:?}",
ip_map.len(),
ip_map.keys().collect::<Vec<_>>()
);
match ip_map.get(&event.speaker_ip) {
Some(id) => id.clone(),
None => {
tracing::warn!(
"Received event from unknown speaker IP: {} (not in ip_to_speaker map)",
event.speaker_ip
);
return;
}
}
};
tracing::debug!(
"Mapped IP {} to speaker_id {}",
event.speaker_ip,
speaker_id.as_str()
);
if event.service.scope() == ServiceScope::PerCoordinator {
let coordinator_lookup = {
let s = store.read();
s.speaker_to_group
.get(&speaker_id)
.and_then(|gid| s.groups.get(gid))
.map(|group| group.coordinator_id == speaker_id)
};
let is_coordinator = coordinator_lookup.unwrap_or_else(|| {
tracing::debug!(
"No group data for {} while handling PerCoordinator {:?} event; \
treating it as its own coordinator",
speaker_id.as_str(),
event.service
);
true
});
if !is_coordinator {
tracing::debug!(
"Skipping PerCoordinator {:?} event from non-coordinator {}",
event.service,
speaker_id.as_str()
);
return;
}
}
let decoded = decode_event(event, speaker_id.clone());
tracing::debug!(
"Decoded {} property changes from event",
decoded.changes.len()
);
for change in &decoded.changes {
tracing::debug!("Applying change: {:?}", change);
apply_property_change(store, watched, fanout, &speaker_id, change, stamp);
}
if event.service.scope() == ServiceScope::PerCoordinator {
let members = {
let s = store.read();
resolve_group_members(&s, &speaker_id)
};
if !members.is_empty() {
notify_group_members(watched, fanout, &members, &decoded.changes, stamp);
}
}
}
fn apply_topology_changes(
store: &Arc<RwLock<StateStore>>,
watched: &Arc<RwLock<WatchCounts>>,
fanout: &EventFanout,
ip_to_speaker: &Arc<RwLock<std::collections::HashMap<IpAddr, SpeakerId>>>,
changes: TopologyChanges,
stamp: WriteStamp,
) {
tracing::debug!(
"Applying topology changes: {} groups, {} memberships",
changes.groups.len(),
changes.memberships.len()
);
if changes.groups.is_empty() {
tracing::warn!(
"Ignoring ZoneGroupTopology event with no zone groups \
({} memberships, {} boot_seqs, {} IPs, {} satellites): treating it as a \
partial event rather than clearing cached group state",
changes.memberships.len(),
changes.boot_seqs.len(),
changes.speaker_ips.len(),
changes.satellite_ids.len()
);
return;
}
let (membership_changes, ip_updates) = {
let mut store = store.write();
store.clear_groups();
for group in changes.groups {
tracing::debug!(
"Adding group {} with {} members",
group.id.as_str(),
group.member_ids.len()
);
store.add_group(group);
}
let mut changed_memberships = Vec::new();
for (speaker_id, membership) in changes.memberships {
let outcome = store.set(&speaker_id, membership.clone(), stamp);
changed_memberships.push((speaker_id, outcome, membership));
}
for (speaker_id, boot_seq) in changes.boot_seqs {
if let Some(speaker) = store.speakers.get_mut(&speaker_id) {
speaker.boot_seq = boot_seq;
}
}
let mut changed_ips = Vec::new();
for (speaker_id, new_ip) in &changes.speaker_ips {
if let Some(old_ip) = store.update_speaker_ip_address(speaker_id, *new_ip) {
tracing::info!(
"Speaker {} IP changed: {} -> {}",
speaker_id.as_str(),
old_ip,
new_ip
);
changed_ips.push((old_ip, *new_ip, speaker_id.clone()));
}
}
store.satellite_ids = changes.satellite_ids.into_iter().collect();
(changed_memberships, changed_ips)
};
if !ip_updates.is_empty() {
let mut map = ip_to_speaker.write();
for (old_ip, new_ip, speaker_id) in ip_updates {
map.remove(&old_ip);
map.insert(new_ip, speaker_id);
}
}
let watched_set = watched.read();
for (speaker_id, outcome, membership) in membership_changes {
if outcome.changed() && is_pair_watched(&watched_set, &speaker_id, GroupMembership::KEY) {
tracing::debug!(
"GroupMembership changed for {}, emitting event",
speaker_id.as_str()
);
fanout.send(ChangeEvent::new(
speaker_id,
PropertyChange::GroupMembership(membership),
stamp,
));
}
}
}
fn resolve_group_members(store: &StateStore, speaker_id: &SpeakerId) -> Vec<SpeakerId> {
store
.speaker_to_group
.get(speaker_id)
.and_then(|gid| store.groups.get(gid))
.filter(|group| group.coordinator_id == *speaker_id && group.member_ids.len() > 1)
.map(|group| {
group
.member_ids
.iter()
.filter(|id| *id != speaker_id)
.cloned()
.collect()
})
.unwrap_or_default()
}
fn notify_group_members(
watched: &Arc<RwLock<WatchCounts>>,
fanout: &EventFanout,
members: &[SpeakerId],
changes: &[PropertyChange],
stamp: WriteStamp,
) {
let watched_set = watched.read();
for member_id in members {
for change in changes {
if change.scope() == Scope::Speaker {
let key = change.key();
if is_pair_watched(&watched_set, member_id, key) {
tracing::debug!(
"Notifying member {} of coordinator change for {}",
member_id.as_str(),
key
);
fanout.send(ChangeEvent::new(member_id.clone(), change.clone(), stamp));
}
}
}
}
}
fn apply_property_change(
store: &Arc<RwLock<StateStore>>,
watched: &Arc<RwLock<WatchCounts>>,
fanout: &EventFanout,
speaker_id: &SpeakerId,
change: &PropertyChange,
stamp: WriteStamp,
) {
let key = change.key();
let outcome = {
let mut store = store.write();
change.apply(&mut store, speaker_id, stamp)
};
if outcome.changed() {
let is_watched = is_pair_watched(&watched.read(), speaker_id, key);
if is_watched {
tracing::debug!(
"Property {} changed for {}, emitting event",
key,
speaker_id.as_str()
);
fanout.send(ChangeEvent::new(speaker_id.clone(), change.clone(), stamp));
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::iter::ChangeIterator;
use crate::model::GroupId;
use crate::property::{GroupInfo, Property, Volume};
use crate::state::retain_direct_watch;
use sonos_api::Service;
pub(super) const PANIC_TRIGGER_IP: IpAddr =
IpAddr::V4(std::net::Ipv4Addr::new(203, 0, 113, 255));
fn test_stamp() -> WriteStamp {
WriteStamp::now(ChangeSource::Event)
}
#[test]
fn test_apply_property_change_volume() {
let store = Arc::new(RwLock::new(StateStore::new()));
let watched = Arc::new(RwLock::new(WatchCounts::new()));
let fanout = Arc::new(EventFanout::new());
let iter = ChangeIterator::new(&fanout);
let speaker_id = SpeakerId::new("test-speaker");
{
let mut s = store.write();
s.add_speaker(crate::model::SpeakerInfo {
id: speaker_id.clone(),
name: "Test".to_string(),
room_name: "Test".to_string(),
ip_address: "192.168.1.100".parse().unwrap(),
port: 1400,
model_name: "Test".to_string(),
software_version: "1.0".to_string(),
boot_seq: 0,
satellites: vec![],
});
}
apply_property_change(
&store,
&watched,
&fanout,
&speaker_id,
&PropertyChange::Volume(Volume(50)),
test_stamp(),
);
assert!(iter.try_recv().is_none());
let stored: Option<Volume> = store.read().get(&speaker_id);
assert_eq!(stored, Some(Volume(50)));
}
#[test]
fn test_apply_property_change_with_watch() {
let store = Arc::new(RwLock::new(StateStore::new()));
let watched = Arc::new(RwLock::new(WatchCounts::new()));
let fanout = Arc::new(EventFanout::new());
let iter = ChangeIterator::new(&fanout);
let speaker_id = SpeakerId::new("test-speaker");
{
let mut s = store.write();
s.add_speaker(crate::model::SpeakerInfo {
id: speaker_id.clone(),
name: "Test".to_string(),
room_name: "Test".to_string(),
ip_address: "192.168.1.100".parse().unwrap(),
port: 1400,
model_name: "Test".to_string(),
software_version: "1.0".to_string(),
boot_seq: 0,
satellites: vec![],
});
}
retain_direct_watch(&watched, &speaker_id, Volume::KEY);
apply_property_change(
&store,
&watched,
&fanout,
&speaker_id,
&PropertyChange::Volume(Volume(75)),
test_stamp(),
);
let event = iter.try_recv().unwrap();
assert_eq!(event.speaker_id, speaker_id);
assert_eq!(event.property_key(), Volume::KEY);
assert_eq!(event.service(), Service::RenderingControl);
}
fn make_speaker_info(id: &str, name: &str, ip: &str) -> crate::model::SpeakerInfo {
crate::model::SpeakerInfo {
id: SpeakerId::new(id),
name: name.to_string(),
room_name: name.to_string(),
ip_address: ip.parse().unwrap(),
port: 1400,
model_name: "Test".to_string(),
software_version: "1.0".to_string(),
boot_seq: 0,
satellites: vec![],
}
}
#[test]
fn test_apply_property_change_group_volume() {
let store = Arc::new(RwLock::new(StateStore::new()));
let watched = Arc::new(RwLock::new(WatchCounts::new()));
let fanout = Arc::new(EventFanout::new());
let speaker_id = SpeakerId::new("RINCON_111");
let group_id = GroupId::new("RINCON_111:1");
{
let mut s = store.write();
s.add_speaker(make_speaker_info(
"RINCON_111",
"Living Room",
"192.168.1.101",
));
s.add_group(GroupInfo::new(
group_id.clone(),
speaker_id.clone(),
vec![speaker_id.clone()],
));
}
apply_property_change(
&store,
&watched,
&fanout,
&speaker_id,
&PropertyChange::GroupVolume(crate::property::GroupVolume(75)),
test_stamp(),
);
let s = store.read();
let stored: Option<crate::property::GroupVolume> = s.get_group(&group_id);
assert_eq!(stored, Some(crate::property::GroupVolume(75)));
}
#[test]
fn test_apply_property_change_group_volume_no_group() {
let store = Arc::new(RwLock::new(StateStore::new()));
let watched = Arc::new(RwLock::new(WatchCounts::new()));
let fanout = Arc::new(EventFanout::new());
let speaker_id = SpeakerId::new("RINCON_111");
{
let mut s = store.write();
s.add_speaker(make_speaker_info(
"RINCON_111",
"Living Room",
"192.168.1.101",
));
}
apply_property_change(
&store,
&watched,
&fanout,
&speaker_id,
&PropertyChange::GroupVolume(crate::property::GroupVolume(50)),
test_stamp(),
);
let s = store.read();
assert!(s.group_props.is_empty());
}
#[test]
fn test_apply_topology_changes_updates_groups() {
let store = Arc::new(RwLock::new(StateStore::new()));
let watched = Arc::new(RwLock::new(WatchCounts::new()));
let fanout = Arc::new(EventFanout::new());
{
let mut s = store.write();
s.add_speaker(make_speaker_info(
"RINCON_111",
"Living Room",
"192.168.1.101",
));
s.add_speaker(make_speaker_info("RINCON_222", "Kitchen", "192.168.1.102"));
}
let group_id = GroupId::new("RINCON_111:1");
let speaker1 = SpeakerId::new("RINCON_111");
let speaker2 = SpeakerId::new("RINCON_222");
let changes = TopologyChanges {
groups: vec![GroupInfo::new(
group_id.clone(),
speaker1.clone(),
vec![speaker1.clone(), speaker2.clone()],
)],
memberships: vec![
(
speaker1.clone(),
GroupMembership::new(group_id.clone(), true),
),
(
speaker2.clone(),
GroupMembership::new(group_id.clone(), false),
),
],
boot_seqs: vec![],
speaker_ips: vec![],
satellite_ids: vec![],
};
let ip_to_speaker = Arc::new(RwLock::new(std::collections::HashMap::new()));
apply_topology_changes(
&store,
&watched,
&fanout,
&ip_to_speaker,
changes,
test_stamp(),
);
let s = store.read();
assert_eq!(s.groups.len(), 1);
let group = s.groups.get(&group_id).unwrap();
assert_eq!(group.coordinator_id, speaker1);
assert_eq!(group.member_ids.len(), 2);
assert!(group.member_ids.contains(&speaker1));
assert!(group.member_ids.contains(&speaker2));
}
#[test]
fn test_apply_topology_changes_updates_group_membership() {
let store = Arc::new(RwLock::new(StateStore::new()));
let watched = Arc::new(RwLock::new(WatchCounts::new()));
let fanout = Arc::new(EventFanout::new());
{
let mut s = store.write();
s.add_speaker(make_speaker_info(
"RINCON_111",
"Living Room",
"192.168.1.101",
));
s.add_speaker(make_speaker_info("RINCON_222", "Kitchen", "192.168.1.102"));
}
let group_id = GroupId::new("RINCON_111:1");
let speaker1 = SpeakerId::new("RINCON_111");
let speaker2 = SpeakerId::new("RINCON_222");
let changes = TopologyChanges {
groups: vec![GroupInfo::new(
group_id.clone(),
speaker1.clone(),
vec![speaker1.clone(), speaker2.clone()],
)],
memberships: vec![
(
speaker1.clone(),
GroupMembership::new(group_id.clone(), true),
),
(
speaker2.clone(),
GroupMembership::new(group_id.clone(), false),
),
],
boot_seqs: vec![],
speaker_ips: vec![],
satellite_ids: vec![],
};
let ip_to_speaker = Arc::new(RwLock::new(std::collections::HashMap::new()));
apply_topology_changes(
&store,
&watched,
&fanout,
&ip_to_speaker,
changes,
test_stamp(),
);
let s = store.read();
let membership1: Option<GroupMembership> = s.get(&speaker1);
assert!(membership1.is_some());
let m1 = membership1.unwrap();
assert_eq!(m1.group_id, group_id);
assert!(m1.is_coordinator);
let membership2: Option<GroupMembership> = s.get(&speaker2);
assert!(membership2.is_some());
let m2 = membership2.unwrap();
assert_eq!(m2.group_id, group_id);
assert!(!m2.is_coordinator);
}
#[test]
fn test_apply_topology_changes_emits_events_for_watched_properties() {
let store = Arc::new(RwLock::new(StateStore::new()));
let watched = Arc::new(RwLock::new(WatchCounts::new()));
let fanout = Arc::new(EventFanout::new());
let iter = ChangeIterator::new(&fanout);
let speaker1 = SpeakerId::new("RINCON_111");
let speaker2 = SpeakerId::new("RINCON_222");
{
let mut s = store.write();
s.add_speaker(make_speaker_info(
"RINCON_111",
"Living Room",
"192.168.1.101",
));
s.add_speaker(make_speaker_info("RINCON_222", "Kitchen", "192.168.1.102"));
}
retain_direct_watch(&watched, &speaker1, GroupMembership::KEY);
let group_id = GroupId::new("RINCON_111:1");
let changes = TopologyChanges {
groups: vec![GroupInfo::new(
group_id.clone(),
speaker1.clone(),
vec![speaker1.clone(), speaker2.clone()],
)],
memberships: vec![
(
speaker1.clone(),
GroupMembership::new(group_id.clone(), true),
),
(
speaker2.clone(),
GroupMembership::new(group_id.clone(), false),
),
],
boot_seqs: vec![],
speaker_ips: vec![],
satellite_ids: vec![],
};
let ip_to_speaker = Arc::new(RwLock::new(std::collections::HashMap::new()));
apply_topology_changes(
&store,
&watched,
&fanout,
&ip_to_speaker,
changes,
test_stamp(),
);
let event = iter.try_recv().unwrap();
assert_eq!(event.speaker_id, speaker1);
assert_eq!(event.property_key(), GroupMembership::KEY);
assert_eq!(event.service(), Service::ZoneGroupTopology);
assert!(iter.try_recv().is_none());
}
#[test]
fn test_apply_topology_changes_clears_old_groups() {
let store = Arc::new(RwLock::new(StateStore::new()));
let watched = Arc::new(RwLock::new(WatchCounts::new()));
let fanout = Arc::new(EventFanout::new());
let speaker1 = SpeakerId::new("RINCON_111");
let speaker2 = SpeakerId::new("RINCON_222");
{
let mut s = store.write();
s.add_speaker(make_speaker_info(
"RINCON_111",
"Living Room",
"192.168.1.101",
));
s.add_speaker(make_speaker_info("RINCON_222", "Kitchen", "192.168.1.102"));
let old_group_id = GroupId::new("OLD_GROUP:1");
s.add_group(GroupInfo::new(
old_group_id.clone(),
speaker1.clone(),
vec![speaker1.clone()],
));
}
{
let s = store.read();
assert_eq!(s.groups.len(), 1);
assert!(s.groups.contains_key(&GroupId::new("OLD_GROUP:1")));
}
let new_group_id = GroupId::new("NEW_GROUP:1");
let changes = TopologyChanges {
groups: vec![GroupInfo::new(
new_group_id.clone(),
speaker2.clone(),
vec![speaker1.clone(), speaker2.clone()],
)],
memberships: vec![
(
speaker1.clone(),
GroupMembership::new(new_group_id.clone(), false),
),
(
speaker2.clone(),
GroupMembership::new(new_group_id.clone(), true),
),
],
boot_seqs: vec![],
speaker_ips: vec![],
satellite_ids: vec![],
};
let ip_to_speaker = Arc::new(RwLock::new(std::collections::HashMap::new()));
apply_topology_changes(
&store,
&watched,
&fanout,
&ip_to_speaker,
changes,
test_stamp(),
);
let s = store.read();
assert_eq!(s.groups.len(), 1);
assert!(!s.groups.contains_key(&GroupId::new("OLD_GROUP:1")));
assert!(s.groups.contains_key(&new_group_id));
}
#[test]
fn test_apply_topology_changes_updates_speaker_to_group_mapping() {
let store = Arc::new(RwLock::new(StateStore::new()));
let watched = Arc::new(RwLock::new(WatchCounts::new()));
let fanout = Arc::new(EventFanout::new());
let speaker1 = SpeakerId::new("RINCON_111");
let speaker2 = SpeakerId::new("RINCON_222");
{
let mut s = store.write();
s.add_speaker(make_speaker_info(
"RINCON_111",
"Living Room",
"192.168.1.101",
));
s.add_speaker(make_speaker_info("RINCON_222", "Kitchen", "192.168.1.102"));
}
let group_id = GroupId::new("RINCON_111:1");
let changes = TopologyChanges {
groups: vec![GroupInfo::new(
group_id.clone(),
speaker1.clone(),
vec![speaker1.clone(), speaker2.clone()],
)],
memberships: vec![
(
speaker1.clone(),
GroupMembership::new(group_id.clone(), true),
),
(
speaker2.clone(),
GroupMembership::new(group_id.clone(), false),
),
],
boot_seqs: vec![],
speaker_ips: vec![],
satellite_ids: vec![],
};
let ip_to_speaker = Arc::new(RwLock::new(std::collections::HashMap::new()));
apply_topology_changes(
&store,
&watched,
&fanout,
&ip_to_speaker,
changes,
test_stamp(),
);
let s = store.read();
assert_eq!(s.speaker_to_group.get(&speaker1), Some(&group_id));
assert_eq!(s.speaker_to_group.get(&speaker2), Some(&group_id));
}
#[test]
fn test_partial_topology_event_does_not_clear_groups() {
let store = Arc::new(RwLock::new(StateStore::new()));
let watched = Arc::new(RwLock::new(WatchCounts::new()));
let fanout = Arc::new(EventFanout::new());
let iter = ChangeIterator::new(&fanout);
let speaker1 = SpeakerId::new("RINCON_111");
let group_id = GroupId::new("RINCON_111:1");
{
let mut s = store.write();
s.add_speaker(make_speaker_info(
"RINCON_111",
"Living Room",
"192.168.1.101",
));
s.add_group(GroupInfo::new(
group_id.clone(),
speaker1.clone(),
vec![speaker1.clone()],
));
s.set_group(&group_id, crate::property::GroupVolume(42), test_stamp());
}
retain_direct_watch(&watched, &speaker1, GroupMembership::KEY);
let partial = TopologyChanges {
groups: vec![],
memberships: vec![],
boot_seqs: vec![],
speaker_ips: vec![],
satellite_ids: vec![],
};
let ip_to_speaker = Arc::new(RwLock::new(std::collections::HashMap::new()));
apply_topology_changes(
&store,
&watched,
&fanout,
&ip_to_speaker,
partial,
test_stamp(),
);
let s = store.read();
assert_eq!(s.groups.len(), 1);
assert!(s.groups.contains_key(&group_id));
assert_eq!(s.speaker_to_group.get(&speaker1), Some(&group_id));
assert_eq!(
s.get_group::<crate::property::GroupVolume>(&group_id),
Some(crate::property::GroupVolume(42))
);
assert!(iter.try_recv().is_none());
}
fn volume_event(ip: IpAddr, volume: u8, observed_at: std::time::Instant) -> EnrichedEvent {
use sonos_stream::events::RenderingControlState;
use sonos_stream::{EventSource, RegistrationId};
EnrichedEvent::observed_at(
RegistrationId::new(1),
ip,
Service::RenderingControl,
EventSource::UPnPNotification {
subscription_id: "uuid:test".to_string(),
},
EventData::RenderingControl(RenderingControlState {
master_volume: Some(volume.to_string()),
master_mute: None,
bass: None,
treble: None,
loudness: None,
lf_volume: None,
rf_volume: None,
lf_mute: None,
rf_mute: None,
balance: None,
other_channels: std::collections::HashMap::new(),
}),
observed_at,
)
}
#[allow(clippy::type_complexity)]
fn one_speaker_worker_fixture() -> (
Arc<RwLock<StateStore>>,
Arc<RwLock<WatchCounts>>,
Arc<EventFanout>,
ChangeIterator,
Arc<RwLock<std::collections::HashMap<IpAddr, SpeakerId>>>,
SpeakerId,
IpAddr,
) {
let store = Arc::new(RwLock::new(StateStore::new()));
let watched = Arc::new(RwLock::new(WatchCounts::new()));
let speaker_id = SpeakerId::new("RINCON_111");
let speaker_ip: IpAddr = "192.0.2.11".parse().unwrap();
store
.write()
.add_speaker(make_speaker_info("RINCON_111", "Living Room", "192.0.2.11"));
let ip_to_speaker = Arc::new(RwLock::new(std::collections::HashMap::from([(
speaker_ip,
speaker_id.clone(),
)])));
let fanout = Arc::new(EventFanout::new());
let iter = ChangeIterator::new(&fanout);
(
store,
watched,
fanout,
iter,
ip_to_speaker,
speaker_id,
speaker_ip,
)
}
fn tick() {
thread::sleep(std::time::Duration::from_millis(20));
}
#[test]
fn test_late_notify_does_not_clobber_fresher_local_action() {
let (store, watched, fanout, _iter, ip_to_speaker, speaker_id, speaker_ip) =
one_speaker_worker_fixture();
let event = volume_event(speaker_ip, 10, std::time::Instant::now());
tick();
store.write().set(
&speaker_id,
Volume(40),
WriteStamp::now(ChangeSource::LocalAction),
);
tick();
handle_event(&event, &store, &watched, &fanout, &ip_to_speaker);
assert_eq!(
store.read().get::<Volume>(&speaker_id),
Some(Volume(40)),
"a NOTIFY observed before a local write must not overwrite it — \
this is the volume-snaps-back symptom"
);
}
#[test]
fn test_current_fetch_not_rejected_by_older_notify() {
let (store, watched, fanout, _iter, ip_to_speaker, speaker_id, speaker_ip) =
one_speaker_worker_fixture();
let event = volume_event(speaker_ip, 10, std::time::Instant::now());
tick();
let fetch_observed = std::time::Instant::now();
tick();
handle_event(&event, &store, &watched, &fanout, &ip_to_speaker);
assert_eq!(
store.read().get::<Volume>(&speaker_id),
Some(Volume(10)),
"precondition: the event should have been applied"
);
let outcome = store.write().set(
&speaker_id,
Volume(40),
WriteStamp::observed_at(ChangeSource::Fetch, fetch_observed),
);
assert_eq!(
outcome,
crate::state::WriteOutcome::Changed,
"a fetch that observed the device after the event must be accepted, \
not rejected as stale"
);
assert_eq!(store.read().get::<Volume>(&speaker_id), Some(Volume(40)));
}
#[test]
fn test_genuinely_newer_event_still_wins() {
let (store, watched, fanout, iter, ip_to_speaker, speaker_id, speaker_ip) =
one_speaker_worker_fixture();
retain_direct_watch(&watched, &speaker_id, Volume::KEY);
store.write().set(
&speaker_id,
Volume(40),
WriteStamp::now(ChangeSource::LocalAction),
);
tick();
let event_observed = std::time::Instant::now();
handle_event(
&volume_event(speaker_ip, 10, event_observed),
&store,
&watched,
&fanout,
&ip_to_speaker,
);
assert_eq!(
store.read().get::<Volume>(&speaker_id),
Some(Volume(10)),
"an event observed after the stored write must still be applied"
);
let notified = iter.try_recv().expect("newer event must still notify");
assert_eq!(notified.property_key(), Volume::KEY);
assert_eq!(
notified.timestamp, event_observed,
"the emitted event must carry the observation instant, not the apply instant"
);
}
#[test]
fn test_worker_survives_decoder_panic() {
use sonos_stream::events::RenderingControlState;
use sonos_stream::{EnrichedEvent, EventSource, RegistrationId};
let store = Arc::new(RwLock::new(StateStore::new()));
let watched = Arc::new(RwLock::new(WatchCounts::new()));
let speaker_id = SpeakerId::new("RINCON_111");
let speaker_ip: IpAddr = "192.168.1.101".parse().unwrap();
{
let mut s = store.write();
s.add_speaker(make_speaker_info(
"RINCON_111",
"Living Room",
"192.168.1.101",
));
}
retain_direct_watch(&watched, &speaker_id, Volume::KEY);
let fanout = Arc::new(EventFanout::new());
let iter = ChangeIterator::new(&fanout);
let ip_to_speaker = Arc::new(RwLock::new(std::collections::HashMap::from([(
speaker_ip,
speaker_id.clone(),
)])));
let make_event = |ip: IpAddr, volume: &str| {
EnrichedEvent::new(
RegistrationId::new(1),
ip,
Service::RenderingControl,
EventSource::UPnPNotification {
subscription_id: "uuid:test".to_string(),
},
EventData::RenderingControl(RenderingControlState {
master_volume: Some(volume.to_string()),
master_mute: None,
bass: None,
treble: None,
loudness: None,
lf_volume: None,
rf_volume: None,
lf_mute: None,
rf_mute: None,
balance: None,
other_channels: std::collections::HashMap::new(),
}),
)
};
let events = vec![
make_event(PANIC_TRIGGER_IP, "10"),
make_event(speaker_ip, "37"),
];
run_event_loop(
events.into_iter(),
&store,
&watched,
&fanout,
&ip_to_speaker,
);
let event = iter.try_recv().unwrap();
assert_eq!(event.speaker_id, speaker_id);
assert_eq!(event.property_key(), Volume::KEY);
assert_eq!(store.read().get::<Volume>(&speaker_id), Some(Volume(37)));
}
#[test]
fn test_apply_topology_changes_no_event_when_membership_unchanged() {
let store = Arc::new(RwLock::new(StateStore::new()));
let watched = Arc::new(RwLock::new(WatchCounts::new()));
let fanout = Arc::new(EventFanout::new());
let iter = ChangeIterator::new(&fanout);
let speaker1 = SpeakerId::new("RINCON_111");
let group_id = GroupId::new("RINCON_111:1");
{
let mut s = store.write();
s.add_speaker(make_speaker_info(
"RINCON_111",
"Living Room",
"192.168.1.101",
));
s.set(
&speaker1,
GroupMembership::new(group_id.clone(), true),
test_stamp(),
);
}
retain_direct_watch(&watched, &speaker1, GroupMembership::KEY);
let changes = TopologyChanges {
groups: vec![GroupInfo::new(
group_id.clone(),
speaker1.clone(),
vec![speaker1.clone()],
)],
memberships: vec![(
speaker1.clone(),
GroupMembership::new(group_id.clone(), true),
)],
boot_seqs: vec![],
speaker_ips: vec![],
satellite_ids: vec![],
};
let ip_to_speaker = Arc::new(RwLock::new(std::collections::HashMap::new()));
apply_topology_changes(
&store,
&watched,
&fanout,
&ip_to_speaker,
changes,
test_stamp(),
);
assert!(iter.try_recv().is_none());
}
#[test]
fn test_per_coordinator_notifies_members_without_data_copy() {
use crate::property::PlaybackState;
let store = Arc::new(RwLock::new(StateStore::new()));
let watched = Arc::new(RwLock::new(WatchCounts::new()));
let fanout = Arc::new(EventFanout::new());
let iter = ChangeIterator::new(&fanout);
let coordinator = SpeakerId::new("RINCON_COORD");
let member = SpeakerId::new("RINCON_MEMBER");
let group_id = GroupId::new("RINCON_COORD:1");
{
let mut s = store.write();
s.add_speaker(make_speaker_info(
"RINCON_COORD",
"Bedroom",
"192.168.1.101",
));
s.add_speaker(make_speaker_info(
"RINCON_MEMBER",
"Kitchen",
"192.168.1.102",
));
s.add_group(GroupInfo::new(
group_id.clone(),
coordinator.clone(),
vec![coordinator.clone(), member.clone()],
));
}
retain_direct_watch(&watched, &coordinator, PlaybackState::KEY);
retain_direct_watch(&watched, &member, PlaybackState::KEY);
let changes = vec![PropertyChange::PlaybackState(PlaybackState::Playing)];
for change in &changes {
apply_property_change(
&store,
&watched,
&fanout,
&coordinator,
change,
test_stamp(),
);
}
let members = {
let s = store.read();
resolve_group_members(&s, &coordinator)
};
notify_group_members(&watched, &fanout, &members, &changes, test_stamp());
let event1 = iter.try_recv().unwrap();
assert_eq!(event1.speaker_id, coordinator);
assert_eq!(event1.property_key(), PlaybackState::KEY);
let event2 = iter.try_recv().unwrap();
assert_eq!(event2.speaker_id, member);
assert_eq!(event2.property_key(), PlaybackState::KEY);
assert!(
matches!(
event2.change,
PropertyChange::PlaybackState(PlaybackState::Playing)
),
"member notification must carry the coordinator's value, got {:?}",
event2.change
);
assert!(iter.try_recv().is_none());
let s = store.read();
let coord_state: Option<PlaybackState> = s.get(&coordinator);
assert_eq!(coord_state, Some(PlaybackState::Playing));
let member_state: Option<PlaybackState> = s.get(&member);
assert_eq!(member_state, None);
let resolved_state: Option<PlaybackState> = s.get_resolved(&member);
assert_eq!(resolved_state, Some(PlaybackState::Playing));
}
#[test]
fn test_per_coordinator_no_notification_for_standalone() {
use crate::property::PlaybackState;
let store = Arc::new(RwLock::new(StateStore::new()));
let watched = Arc::new(RwLock::new(WatchCounts::new()));
let fanout = Arc::new(EventFanout::new());
let iter = ChangeIterator::new(&fanout);
let speaker = SpeakerId::new("RINCON_STANDALONE");
let group_id = GroupId::new("RINCON_STANDALONE:1");
{
let mut s = store.write();
s.add_speaker(make_speaker_info(
"RINCON_STANDALONE",
"Bedroom",
"192.168.1.101",
));
s.add_group(GroupInfo::new(
group_id.clone(),
speaker.clone(),
vec![speaker.clone()],
));
}
retain_direct_watch(&watched, &speaker, PlaybackState::KEY);
let changes = vec![PropertyChange::PlaybackState(PlaybackState::Playing)];
for change in &changes {
apply_property_change(&store, &watched, &fanout, &speaker, change, test_stamp());
}
let members = {
let s = store.read();
resolve_group_members(&s, &speaker)
};
assert!(members.is_empty());
let event = iter.try_recv().unwrap();
assert_eq!(event.speaker_id, speaker);
assert!(iter.try_recv().is_none());
}
#[test]
fn test_per_speaker_service_not_notified() {
let store = Arc::new(RwLock::new(StateStore::new()));
let watched = Arc::new(RwLock::new(WatchCounts::new()));
let fanout = Arc::new(EventFanout::new());
let iter = ChangeIterator::new(&fanout);
let coordinator = SpeakerId::new("RINCON_COORD");
let member = SpeakerId::new("RINCON_MEMBER");
let group_id = GroupId::new("RINCON_COORD:1");
{
let mut s = store.write();
s.add_speaker(make_speaker_info(
"RINCON_COORD",
"Bedroom",
"192.168.1.101",
));
s.add_speaker(make_speaker_info(
"RINCON_MEMBER",
"Kitchen",
"192.168.1.102",
));
s.add_group(GroupInfo::new(
group_id.clone(),
coordinator.clone(),
vec![coordinator.clone(), member.clone()],
));
}
retain_direct_watch(&watched, &coordinator, Volume::KEY);
retain_direct_watch(&watched, &member, Volume::KEY);
apply_property_change(
&store,
&watched,
&fanout,
&coordinator,
&PropertyChange::Volume(Volume(80)),
test_stamp(),
);
let event = iter.try_recv().unwrap();
assert_eq!(event.speaker_id, coordinator);
assert_eq!(event.property_key(), Volume::KEY);
assert!(iter.try_recv().is_none());
let s = store.read();
let coord_vol: Option<Volume> = s.get(&coordinator);
let member_vol: Option<Volume> = s.get(&member);
assert_eq!(coord_vol, Some(Volume(80)));
assert_eq!(member_vol, None);
}
#[test]
fn test_resolve_group_members_empty_for_non_coordinator() {
let mut store = StateStore::new();
let coordinator = SpeakerId::new("RINCON_COORD");
let member = SpeakerId::new("RINCON_MEMBER");
let group_id = GroupId::new("RINCON_COORD:1");
store.add_speaker(make_speaker_info(
"RINCON_COORD",
"Bedroom",
"192.168.1.101",
));
store.add_speaker(make_speaker_info(
"RINCON_MEMBER",
"Kitchen",
"192.168.1.102",
));
store.add_group(GroupInfo::new(
group_id,
coordinator,
vec![SpeakerId::new("RINCON_COORD"), member.clone()],
));
let members = resolve_group_members(&store, &member);
assert!(members.is_empty());
}
#[test]
fn test_notify_group_members_only_notifies_watched() {
use crate::property::PlaybackState;
let watched = Arc::new(RwLock::new(WatchCounts::new()));
let fanout = Arc::new(EventFanout::new());
let iter = ChangeIterator::new(&fanout);
let member_watched = SpeakerId::new("RINCON_WATCHED");
let member_unwatched = SpeakerId::new("RINCON_UNWATCHED");
retain_direct_watch(&watched, &member_watched, PlaybackState::KEY);
let changes = vec![PropertyChange::PlaybackState(PlaybackState::Playing)];
let members = vec![member_watched.clone(), member_unwatched.clone()];
notify_group_members(&watched, &fanout, &members, &changes, test_stamp());
let event = iter.try_recv().unwrap();
assert_eq!(event.speaker_id, member_watched);
assert_eq!(event.property_key(), PlaybackState::KEY);
assert!(iter.try_recv().is_none());
}
}