use bevy::ecs::change_detection::DetectChanges;
use bevy::ecs::change_detection::Ref;
use bevy::ecs::observer::On;
use bevy::ecs::system::Commands;
use bevy::ecs::system::Query;
use bevy::ecs::system::Res;
use bevy::ecs::system::ResMut;
use bevy::prelude::Added;
use bevy::prelude::Entity;
use bevy::prelude::Resource;
use crate::Bindings;
use crate::Claim;
use crate::ClaimChanged;
use crate::ConfiguredDeviceConnectionChanged;
use crate::DeviceArrived;
use crate::DeviceDeparted;
use crate::DeviceKey;
use crate::DeviceResolution;
use crate::DeviceStateLookup;
use crate::Devices;
use crate::DiscoveryControl;
use crate::DiscoveryFinished;
use crate::DiscoveryProgressChanged;
use crate::IdentityChanged;
use crate::IdentityVerdict;
use crate::Presence;
use crate::PresenceChanged;
use crate::ReapplyConfiguration;
use crate::RecoveryPolicy;
use crate::RecoveryPolicyChanged;
use crate::RetireRole;
use crate::RoleAvailable;
use crate::RoleAwaiting;
use crate::RoleKey;
use crate::SchemeName;
use crate::StartupDiscoveryChanged;
use crate::UnregisteredSchemeReported;
use crate::WaitingWork;
use crate::devices::DepartureAnnouncements;
use crate::discovery::DiscoveryTransition;
use crate::discovery::DiscoveryTransitionJournal;
use crate::registration::Reporters;
#[derive(Debug, Default, Resource)]
pub(crate) struct AnnouncedEdges {
schemes: Vec<SchemeName>,
awaiting: Vec<RoleKey>,
available: Vec<RoleKey>,
}
pub(crate) fn announce_device_changes(
mut commands: Commands,
arrived: Query<(Entity, &DeviceKey), Added<DeviceKey>>,
presences: Query<(Entity, Ref<'_, Presence>)>,
claims: Query<(Entity, Ref<'_, Claim>)>,
verdicts: Query<(Entity, Ref<'_, IdentityVerdict>)>,
) {
for (device, key) in &arrived {
commands.trigger(DeviceArrived {
device,
key: key.clone(),
});
}
for (device, presence) in &presences {
if presence.is_changed() {
commands.trigger(PresenceChanged {
device,
presence: *presence,
});
}
}
for (device, claim) in &claims {
if claim.is_changed() {
commands.trigger(ClaimChanged {
device,
claim: (*claim).clone(),
});
}
}
for (device, verdict) in &verdicts {
if verdict.is_changed() {
commands.trigger(IdentityChanged {
device,
verdict: (*verdict).clone(),
});
}
}
}
pub(crate) fn announce_binding_changes(
mut commands: Commands,
recoveries: Query<(Entity, &RoleKey, Ref<'_, RecoveryPolicy>)>,
) {
for (binding, role, recovery) in &recoveries {
if recovery.is_changed() {
commands.trigger(RecoveryPolicyChanged {
binding,
role: role.clone(),
recovery: *recovery,
});
}
}
}
pub(crate) fn announce_reconciled_facts(
mut commands: Commands,
mut departure_announcements: ResMut<DepartureAnnouncements>,
devices: Res<Devices>,
mut announced_edges: ResMut<AnnouncedEdges>,
) {
for departed_device in departure_announcements.departed.drain(..) {
commands.trigger(DeviceDeparted {
key: departed_device.key,
departure: departed_device.departure,
});
}
for connection_change in departure_announcements.connections.drain(..) {
commands.trigger(ConfiguredDeviceConnectionChanged {
key: connection_change.key,
connection: connection_change.connection,
});
}
for scheme in devices.unregistered_schemes() {
if !announced_edges.schemes.contains(scheme) {
announced_edges.schemes.push(scheme.clone());
commands.trigger(UnregisteredSchemeReported {
scheme: scheme.clone(),
});
}
}
}
fn endpoint_has_live_device(devices: &Devices, device: &DeviceKey) -> bool {
let DeviceResolution::Resolved(device_id) = devices.resolve(device) else {
return false;
};
matches!(
devices.state(device_id),
DeviceStateLookup::Retained(reconciled_device_state)
if reconciled_device_state.presence == Presence::Present
)
}
pub(crate) fn announce_role_availability(
mut commands: Commands,
bindings: Res<Bindings>,
devices: Res<Devices>,
mut announced_edges: ResMut<AnnouncedEdges>,
) {
let mut awaiting = Vec::new();
let mut available = Vec::new();
for role in bindings.registered_roles() {
let Ok(binding) = bindings.binding(role) else {
continue;
};
if endpoint_has_live_device(&devices, &binding.endpoint.device) {
available.push(role.clone());
} else {
awaiting.push(role.clone());
}
}
for role in &awaiting {
if !announced_edges.awaiting.contains(role) {
commands.trigger(RoleAwaiting { role: role.clone() });
}
}
for role in &available {
if !announced_edges.available.contains(role) {
commands.trigger(RoleAvailable { role: role.clone() });
}
}
announced_edges.awaiting = awaiting;
announced_edges.available = available;
}
pub(crate) fn on_role_awaiting(
role_awaiting: On<RoleAwaiting>,
bindings: Res<Bindings>,
reporters: Res<Reporters>,
mut discovery_control: ResMut<DiscoveryControl>,
) {
let Ok(binding) = bindings.binding(&role_awaiting.role) else {
return;
};
let Some(covering_reporter) = reporters
.registered_reporters()
.find(|registered_reporter| {
registered_reporter
.coverage
.establishes_absence_for(&binding.endpoint.device)
})
.map(|registered_reporter| registered_reporter.reporter)
else {
return;
};
drop(discovery_control.enable(covering_reporter));
}
pub(crate) fn on_retire_role(retire_role: On<RetireRole>, mut bindings: ResMut<Bindings>) {
drop(bindings.retire(&retire_role.role));
}
pub(crate) fn on_reapply_configuration(
reapply_configuration: On<ReapplyConfiguration>,
roles: Query<&RoleKey>,
mut bindings: ResMut<Bindings>,
) {
let Ok(role) = roles.get(reapply_configuration.binding).cloned() else {
return;
};
let Ok(binding) = bindings.binding(&role) else {
return;
};
if binding.recovery != RecoveryPolicy::ReapplyOnRequest
|| bindings.waiting_work(&role) != WaitingWork::ApplicationRequestOwed
{
return;
}
bindings.request_reapply(&role);
}
pub(crate) fn announce_discovery_transitions(
mut commands: Commands,
mut discovery_transition_journal: ResMut<DiscoveryTransitionJournal>,
) {
for discovery_transition in discovery_transition_journal.drain() {
match discovery_transition {
DiscoveryTransition::Progressed {
batch,
reporter,
progress,
completed,
total,
running,
queued,
} => {
commands.trigger(DiscoveryProgressChanged {
batch,
reporter,
progress,
completed,
total,
running,
queued,
});
},
DiscoveryTransition::Finished {
batch,
reporter,
outcome,
} => {
commands.trigger(DiscoveryFinished {
batch,
reporter,
outcome,
});
},
DiscoveryTransition::StartupChanged { startup } => {
commands.trigger(StartupDiscoveryChanged { state: startup });
},
}
}
}
#[cfg(test)]
mod tests {
use std::error::Error;
use std::sync::Arc;
use std::sync::atomic::AtomicUsize;
use std::sync::atomic::Ordering;
use std::time::Duration;
use bevy::app::App;
use crate::AuthoritativeReporterCoverage;
use crate::Binding;
use crate::Bindings;
use crate::CoveredDeviceIdentitySpace;
use crate::DeviceEndpoint;
use crate::DeviceIdSource;
use crate::DeviceKey;
use crate::DeviceKind;
use crate::DeviceReporter;
use crate::DeviceScan;
use crate::DiscoveryCadence;
use crate::DiscoveryControl;
use crate::DiscoveryWork;
use crate::EndpointId;
use crate::LastKnownGoodConfiguration;
use crate::MainThreadDiscoveryJob;
use crate::OnAbort;
use crate::OnSessionLoss;
use crate::RecoveryPolicy;
use crate::ReporterActivation;
use crate::ReporterCoverage;
use crate::ReporterId;
use crate::ReporterRegistration;
use crate::RequestedConfiguration;
use crate::RetryOn;
use crate::RiggingAppExt;
use crate::RiggingPlugin;
use crate::RoleKey;
use crate::RoleState;
use crate::binding::ApplyDeadline;
use crate::registration::DriverId;
use crate::scheme::AuthoredId;
const NEVER_DUE_UNAIDED: Duration = Duration::from_hours(1);
struct CountingReporter {
scans: Arc<AtomicUsize>,
}
impl DeviceReporter for CountingReporter {
fn discover(&mut self) -> DiscoveryWork {
self.scans.fetch_add(1, Ordering::Relaxed);
DiscoveryWork::Immediate(MainThreadDiscoveryJob::new(|_| {
DeviceScan::Complete(Vec::new())
}))
}
}
fn app_with_disabled_reporter_covering(
kind: DeviceKind,
) -> (App, ReporterId, Arc<AtomicUsize>) {
let mut app = App::new();
app.add_plugins(RiggingPlugin);
let scans = Arc::new(AtomicUsize::new(0));
let reporter = app.add_device_reporter(
CountingReporter {
scans: Arc::clone(&scans),
},
ReporterRegistration::optional(
DiscoveryCadence::Periodic {
interval: NEVER_DUE_UNAIDED,
},
ReporterActivation::Disabled,
ReporterCoverage::EstablishesAbsence(AuthoritativeReporterCoverage::one(
CoveredDeviceIdentitySpace::AllKeysOfKind { kind },
)),
),
);
(app, reporter, scans)
}
fn unresolved_camera_binding(role: &str) -> Result<Binding, Box<dyn Error>> {
Ok(Binding {
role: RoleKey::new(role)?,
endpoint: DeviceEndpoint {
device: DeviceKey {
kind: DeviceKind::Camera,
id: DeviceIdSource::Authored {
value: AuthoredId::new(role)?,
},
},
id: EndpointId::Whole,
},
driver: DriverId(0),
recovery: RecoveryPolicy::Forget,
retry: RetryOn::NewRevision,
on_abort: OnAbort::default(),
on_loss: OnSessionLoss::default(),
state: RoleState::default(),
requested: RequestedConfiguration::new(()),
last_known_good: LastKnownGoodConfiguration::default(),
apply_deadline: ApplyDeadline::ProcessDefault,
})
}
#[test]
fn a_role_awaiting_edge_enables_and_runs_the_covering_disabled_reporter()
-> Result<(), Box<dyn Error>> {
let (mut app, reporter, scans) = app_with_disabled_reporter_covering(DeviceKind::Camera);
app.world_mut()
.resource_mut::<Bindings>()
.register(unresolved_camera_binding("camera-tool")?)?;
app.update();
assert_eq!(
app.world()
.resource::<DiscoveryControl>()
.activation(reporter),
ReporterActivation::Enabled,
"the RoleAwaiting edge for an unresolved covered claim must enable the reporter"
);
app.update();
assert_eq!(
scans.load(Ordering::Relaxed),
1,
"the enable must carry one requested run, collected on the next frame"
);
Ok(())
}
#[test]
fn a_claim_no_reporter_covers_arms_nothing() -> Result<(), Box<dyn Error>> {
let (mut app, reporter, scans) = app_with_disabled_reporter_covering(DeviceKind::Display);
app.world_mut()
.resource_mut::<Bindings>()
.register(unresolved_camera_binding("camera-tool")?)?;
app.update();
app.update();
assert_eq!(
app.world()
.resource::<DiscoveryControl>()
.activation(reporter),
ReporterActivation::Disabled,
"a display-only reporter proves nothing about a camera key and must stay disabled"
);
assert_eq!(scans.load(Ordering::Relaxed), 0);
Ok(())
}
}