a11y: Add magnifier toggle

This commit is contained in:
Victoria Brekenfeld 2025-02-20 16:03:09 +01:00 committed by Michael Murphy
parent 257784ebc8
commit 2b88f35991
21 changed files with 461 additions and 191 deletions

View file

@ -0,0 +1,139 @@
// Copyright 2023 System76 <info@system76.com>
// SPDX-License-Identifier: GPL-3.0-only
use cosmic::iced::futures::FutureExt;
use cosmic::{
iced::{
self,
futures::{self, select, SinkExt, StreamExt},
Subscription,
},
iced_futures::stream,
};
use cosmic_dbus_a11y::*;
use std::fmt::Debug;
use tokio::sync::mpsc::{UnboundedReceiver, UnboundedSender};
use zbus::Connection;
#[derive(Debug, Clone)]
pub enum DBusUpdate {
Error(String),
Status(bool),
Init(bool, UnboundedSender<DBusRequest>),
}
pub enum DBusRequest {
Status(bool),
}
#[derive(Debug)]
pub enum State {
Ready,
Waiting(Connection, u8, bool, UnboundedReceiver<DBusRequest>),
Finished,
}
pub fn subscription() -> iced::Subscription<DBusUpdate> {
struct MyId;
Subscription::run_with_id(
std::any::TypeId::of::<MyId>(),
stream::channel(50, move |mut output| async move {
let mut state = State::Ready;
loop {
state = start_listening(state, &mut output).await;
}
}),
)
}
async fn start_listening(
state: State,
output: &mut futures::channel::mpsc::Sender<DBusUpdate>,
) -> State {
match state {
State::Ready => {
let conn = match Connection::session().await.map_err(|e| e.to_string()) {
Ok(conn) => conn,
Err(e) => {
_ = output.send(DBusUpdate::Error(e)).await;
return State::Finished;
}
};
let (tx, rx) = tokio::sync::mpsc::unbounded_channel();
let mut enabled = false;
if let Ok(proxy) = StatusProxy::new(&conn).await {
if let Ok(status) = proxy.screen_reader_enabled().await {
enabled = status;
}
}
_ = output.send(DBusUpdate::Init(enabled, tx)).await;
State::Waiting(conn, 20, enabled, rx)
}
State::Waiting(conn, mut retry, mut enabled, mut rx) => {
let Ok(proxy) = StatusProxy::new(&conn).await else {
if retry == 0 {
tracing::error!("Accessibility Status is unavailable.");
return State::Finished;
} else {
_ = tokio::time::sleep(tokio::time::Duration::from_secs(
2_u64.pow(retry as u32),
))
.await;
retry -= 1;
return State::Waiting(conn, retry, enabled, rx);
}
};
retry = 20;
let mut watch_changes = proxy.receive_screen_reader_enabled_changed().await;
if let Ok(status) = proxy.screen_reader_enabled().await {
if enabled != status {
_ = output.send(DBusUpdate::Status(enabled));
}
enabled = status;
}
loop {
if let Ok(status) = proxy.screen_reader_enabled().await {
if enabled != status {
_ = output.send(DBusUpdate::Status(enabled));
}
enabled = status;
}
let mut next_change = Box::pin(watch_changes.next()).fuse();
let mut next_request = Box::pin(rx.recv()).fuse();
select! {
v = next_request => {
match v {
Some(DBusRequest::Status(is_enabled)) => {
// Set status
enabled = is_enabled;
_ = proxy.set_is_enabled(is_enabled).await;
_ = proxy.set_screen_reader_enabled(is_enabled).await;
}
None => return State::Finished,
}
}
v = next_change => {
match v {
Some(f) => {
if let Ok(enabled) = f.get().await {
_ = output.send(DBusUpdate::Status(enabled));
}
}
None => break,
};
}
}
}
State::Waiting(conn, retry, enabled, rx)
}
State::Finished => iced::futures::future::pending().await,
}
}

View file

@ -1,139 +1,2 @@
// Copyright 2023 System76 <info@system76.com>
// SPDX-License-Identifier: GPL-3.0-only
use cosmic::iced::futures::FutureExt;
use cosmic::{
iced::{
self,
futures::{self, select, SinkExt, StreamExt},
Subscription,
},
iced_futures::stream,
};
use cosmic_dbus_a11y::*;
use std::fmt::Debug;
use tokio::sync::mpsc::{UnboundedReceiver, UnboundedSender};
use zbus::Connection;
#[derive(Debug, Clone)]
pub enum Update {
Error(String),
Status(bool),
Init(bool, UnboundedSender<A11yRequest>),
}
pub enum A11yRequest {
Status(bool),
}
#[derive(Debug)]
pub enum State {
Ready,
Waiting(Connection, u8, bool, UnboundedReceiver<A11yRequest>),
Finished,
}
pub fn subscription() -> iced::Subscription<Update> {
struct MyId;
Subscription::run_with_id(
std::any::TypeId::of::<MyId>(),
stream::channel(50, move |mut output| async move {
let mut state = State::Ready;
loop {
state = start_listening(state, &mut output).await;
}
}),
)
}
async fn start_listening(
state: State,
output: &mut futures::channel::mpsc::Sender<Update>,
) -> State {
match state {
State::Ready => {
let conn = match Connection::session().await.map_err(|e| e.to_string()) {
Ok(conn) => conn,
Err(e) => {
_ = output.send(Update::Error(e)).await;
return State::Finished;
}
};
let (tx, rx) = tokio::sync::mpsc::unbounded_channel();
let mut enabled = false;
if let Ok(proxy) = StatusProxy::new(&conn).await {
if let Ok(status) = proxy.screen_reader_enabled().await {
enabled = status;
}
}
_ = output.send(Update::Init(enabled, tx)).await;
State::Waiting(conn, 20, enabled, rx)
}
State::Waiting(conn, mut retry, mut enabled, mut rx) => {
let Ok(proxy) = StatusProxy::new(&conn).await else {
if retry == 0 {
tracing::error!("Accessibility Status is unavailable.");
return State::Finished;
} else {
_ = tokio::time::sleep(tokio::time::Duration::from_secs(
2_u64.pow(retry as u32),
))
.await;
retry -= 1;
return State::Waiting(conn, retry, enabled, rx);
}
};
retry = 20;
let mut watch_changes = proxy.receive_screen_reader_enabled_changed().await;
if let Ok(status) = proxy.screen_reader_enabled().await {
if enabled != status {
_ = output.send(Update::Status(enabled));
}
enabled = status;
}
loop {
if let Ok(status) = proxy.screen_reader_enabled().await {
if enabled != status {
_ = output.send(Update::Status(enabled));
}
enabled = status;
}
let mut next_change = Box::pin(watch_changes.next()).fuse();
let mut next_request = Box::pin(rx.recv()).fuse();
select! {
v = next_request => {
match v {
Some(A11yRequest::Status(is_enabled)) => {
// Set status
enabled = is_enabled;
_ = proxy.set_is_enabled(is_enabled).await;
_ = proxy.set_screen_reader_enabled(is_enabled).await;
}
None => return State::Finished,
}
}
v = next_change => {
match v {
Some(f) => {
if let Ok(enabled) = f.get().await {
_ = output.send(Update::Status(enabled));
}
}
None => break,
};
}
}
}
State::Waiting(conn, retry, enabled, rx)
}
State::Finished => iced::futures::future::pending().await,
}
}
pub mod dbus;
pub mod wayland;

View file

@ -0,0 +1,96 @@
// Copyright 2023 System76 <info@system76.com>
// SPDX-License-Identifier: GPL-3.0-only
use anyhow;
use cctk::sctk::reexports::calloop::channel::SyncSender;
use cosmic::iced::{
self,
futures::{self, channel::mpsc, SinkExt, StreamExt},
stream, Subscription,
};
use once_cell::sync::Lazy;
use tokio::sync::Mutex;
mod thread;
pub static WAYLAND_RX: Lazy<Mutex<Option<mpsc::Receiver<AccessibilityEvent>>>> =
Lazy::new(|| Mutex::new(None));
#[derive(Debug, Clone)]
pub enum WaylandUpdate {
State(AccessibilityEvent),
Started(SyncSender<AccessibilityRequest>),
Errored,
}
#[derive(Debug, Clone, Copy)]
pub enum AccessibilityEvent {
Magnifier(bool),
}
#[derive(Debug, Clone, Copy)]
pub enum AccessibilityRequest {
Magnifier(bool),
}
pub fn a11y_subscription() -> iced::Subscription<WaylandUpdate> {
Subscription::run_with_id(
std::any::TypeId::of::<WaylandUpdate>(),
stream::channel(50, move |mut output| async move {
let mut state = State::Waiting;
loop {
state = start_listening(state, &mut output).await;
}
}),
)
}
async fn start_listening(
state: State,
output: &mut futures::channel::mpsc::Sender<WaylandUpdate>,
) -> State {
match state {
State::Waiting => {
let mut guard = WAYLAND_RX.lock().await;
let rx = {
if guard.is_none() {
if let Ok(WaylandWatcher { rx, tx }) = WaylandWatcher::new() {
*guard = Some(rx);
_ = output.send(WaylandUpdate::Started(tx)).await;
} else {
_ = output.send(WaylandUpdate::Errored).await;
return State::Error;
}
}
guard.as_mut().unwrap()
};
if let Some(w) = rx.next().await {
_ = output.send(WaylandUpdate::State(w)).await;
State::Waiting
} else {
_ = output.send(WaylandUpdate::Errored).await;
State::Error
}
}
State::Error => cosmic::iced::futures::future::pending().await,
}
}
pub enum State {
Waiting,
Error,
}
pub struct WaylandWatcher {
rx: mpsc::Receiver<AccessibilityEvent>,
tx: SyncSender<AccessibilityRequest>,
}
impl WaylandWatcher {
pub fn new() -> anyhow::Result<Self> {
let (tx, rx) = mpsc::channel(20);
let tx = thread::spawn_a11y(tx)?;
Ok(Self { tx, rx })
}
}

View file

@ -0,0 +1,123 @@
// Copyright 2025 System76 <info@system76.com>
// SPDX-License-Identifier: GPL-3.0-only
use calloop::channel::*;
use cctk::{
sctk::{
self,
reexports::{
calloop::{self, channel},
calloop_wayland_source::WaylandSource,
},
registry::RegistryState,
},
wayland_client::{self, globals::GlobalListContents, protocol::wl_registry, Dispatch, Proxy},
};
use cosmic::iced::futures::{self, SinkExt};
use cosmic_protocols::a11y::v1::client::cosmic_a11y_manager_v1;
use futures::{channel::mpsc, executor::block_on};
use wayland_client::{globals::registry_queue_init, Connection};
use super::{AccessibilityEvent, AccessibilityRequest};
pub fn spawn_a11y(
tx: mpsc::Sender<AccessibilityEvent>,
) -> anyhow::Result<SyncSender<AccessibilityRequest>> {
let (a11y_tx, a11y_rx) = calloop::channel::sync_channel(100);
let conn = Connection::connect_to_env()?;
std::thread::spawn(move || {
struct State {
loop_signal: calloop::LoopSignal,
tx: mpsc::Sender<AccessibilityEvent>,
global: cosmic_a11y_manager_v1::CosmicA11yManagerV1,
magnifier: bool,
}
impl Dispatch<cosmic_a11y_manager_v1::CosmicA11yManagerV1, ()> for State {
fn event(
state: &mut Self,
_proxy: &cosmic_a11y_manager_v1::CosmicA11yManagerV1,
event: <cosmic_a11y_manager_v1::CosmicA11yManagerV1 as Proxy>::Event,
_data: &(),
_conn: &Connection,
_qhandle: &sctk::reexports::client::QueueHandle<Self>,
) {
match event {
cosmic_a11y_manager_v1::Event::Magnifier { active } => {
let magnifier = active
.into_result()
.unwrap_or(cosmic_a11y_manager_v1::ActiveState::Disabled)
== cosmic_a11y_manager_v1::ActiveState::Enabled;
if magnifier != state.magnifier {
if block_on(state.tx.send(AccessibilityEvent::Magnifier(magnifier)))
.is_err()
{
state.loop_signal.stop();
state.loop_signal.wakeup();
};
state.magnifier = magnifier;
}
}
_ => unreachable!(),
}
}
}
impl Dispatch<wl_registry::WlRegistry, GlobalListContents> for State {
fn event(
_state: &mut Self,
_proxy: &wl_registry::WlRegistry,
_event: <wl_registry::WlRegistry as Proxy>::Event,
_data: &GlobalListContents,
_conn: &Connection,
_qhandle: &sctk::reexports::client::QueueHandle<Self>,
) {
// We don't care about any dynamic globals
}
}
let mut event_loop = calloop::EventLoop::<State>::try_new().unwrap();
let loop_handle = event_loop.handle();
let (globals, event_queue) = registry_queue_init(&conn).unwrap();
let qhandle = event_queue.handle();
WaylandSource::new(conn, event_queue)
.insert(loop_handle.clone())
.unwrap();
let registry_state = RegistryState::new(&globals);
let global = registry_state
.bind_one::<cosmic_a11y_manager_v1::CosmicA11yManagerV1, _, _>(&qhandle, 1..=1, ())
.unwrap();
loop_handle
.insert_source(a11y_rx, |request, _, state| match request {
channel::Event::Msg(AccessibilityRequest::Magnifier(val)) => {
state.global.set_magnifier(if val {
cosmic_a11y_manager_v1::ActiveState::Enabled
} else {
cosmic_a11y_manager_v1::ActiveState::Disabled
});
}
channel::Event::Closed => {
state.loop_signal.stop();
state.loop_signal.wakeup();
}
})
.unwrap();
let mut state = State {
loop_signal: event_loop.get_signal(),
tx,
global,
magnifier: false,
};
event_loop.run(None, &mut state, |_| {}).unwrap();
});
Ok(a11y_tx)
}