use std::collections::HashMap;
use std::pin::Pin;
use std::task::{Context, Poll};
use std::time::Duration;
use futures_util::{Stream, StreamExt, stream::SelectAll};
use zbus::zvariant::OwnedFd;
use crate::error::{CpdbError, Result};
use crate::events::{DiscoveryEvent, PrinterSnapshot};
use crate::media::MediaCollection;
use crate::options::OptionsCollection;
use crate::proxy::PrintBackendProxy;
use crate::types::PrinterState;
pub const BACKEND_INACTIVITY_TIMEOUT: Duration = Duration::from_secs(30);
const RETRY_INTERVAL_MS: Duration = Duration::from_millis(200);
type RawPrinterTuple = (String, String, String, String, String, bool, String, String);
macro_rules! retry_dbus {
($proxy:expr, $method:ident($($arg:expr),*)) => {{
let __proxy = &($proxy);
let mut __result = __proxy.$method($($arg),*).await;
if let Err(zbus::Error::MethodError(ref __n, _, _)) = __result {
if __n.as_str() == "org.freedesktop.DBus.Error.UnknownMethod" {
tokio::time::sleep(RETRY_INTERVAL_MS).await;
__result = __proxy.$method($($arg),*).await;
}
}
__result
}};
}
#[derive(Clone)]
pub struct CpdbClient {
backends: Vec<BackendHandle>,
}
#[derive(Clone)]
struct BackendHandle {
service_name: String,
proxy: PrintBackendProxy<'static>,
}
impl CpdbClient {
pub async fn new() -> Result<Self> {
let connection = zbus::Connection::session().await.map_err(CpdbError::from)?;
let dbus = zbus::fdo::DBusProxy::new(&connection)
.await
.map_err(CpdbError::from)?;
let names = dbus
.list_activatable_names()
.await
.map_err(CpdbError::from)?;
let backend_names: Vec<String> = names
.iter()
.filter(|n| n.starts_with("org.openprinting.Backend."))
.map(|n| n.to_string())
.collect();
let mut backends = Vec::new();
for name in &backend_names {
let bus_name = match zbus::names::BusName::try_from(name.clone()) {
Ok(n) => n,
Err(_) => continue,
};
match PrintBackendProxy::builder(&connection)
.destination(bus_name)?
.path("/")?
.build()
.await
{
Ok(proxy) => {
backends.push(BackendHandle {
service_name: name.clone(),
proxy,
});
}
Err(e) => {
log::warn!("cpdb-rs: skipping backend {}: {}", name, e);
}
}
}
Ok(Self { backends })
}
pub fn backend_count(&self) -> usize {
self.backends.len()
}
pub async fn get_all_printers(&self) -> Result<Vec<PrinterSnapshot>> {
let mut printers = Vec::new();
for bh in &self.backends {
let result = retry_dbus!(bh.proxy, get_all_printers());
let (_count, raw_printers) = match result {
Ok(v) => v,
Err(e) => {
log::error!(
"cpdb-rs: error fetching printers from {}: {}",
bh.service_name,
e
);
continue;
}
};
for val in raw_printers {
let tuple: std::result::Result<RawPrinterTuple, zbus::zvariant::Error> =
val.0.try_into();
let raw = match tuple {
Ok(r) => r,
Err(e) => {
log::error!("cpdb-rs: error unpacking printer variant: {}", e);
continue;
}
};
printers.push(PrinterSnapshot {
id: raw.0,
name: raw.1,
info: raw.2,
location: raw.3,
make_model: raw.4,
accepting_jobs: raw.5,
state: PrinterState::from(raw.6),
backend: raw.7,
});
}
}
Ok(printers)
}
pub async fn get_filtered_printers(&self) -> Result<Vec<PrinterSnapshot>> {
let mut printers = Vec::new();
for bh in &self.backends {
let result = retry_dbus!(bh.proxy, get_filtered_printer_list());
let (_count, raw_printers) = match result {
Ok(v) => v,
Err(e) => {
log::error!(
"cpdb-rs: error fetching filtered printers from {}: {}",
bh.service_name,
e
);
continue;
}
};
for val in raw_printers {
let tuple: std::result::Result<RawPrinterTuple, zbus::zvariant::Error> =
val.0.try_into();
let raw = match tuple {
Ok(r) => r,
Err(e) => {
log::error!("cpdb-rs: error unpacking filtered printer variant: {}", e);
continue;
}
};
printers.push(PrinterSnapshot {
id: raw.0,
name: raw.1,
info: raw.2,
location: raw.3,
make_model: raw.4,
accepting_jobs: raw.5,
state: PrinterState::from(raw.6),
backend: raw.7,
});
}
}
Ok(printers)
}
pub async fn discovery_stream(&self) -> Result<DiscoveryStream> {
let mut all: SelectAll<futures_util::stream::BoxStream<'static, DiscoveryEvent>> =
SelectAll::new();
for bh in &self.backends {
let added = bh
.proxy
.receive_printer_added()
.await
.map_err(CpdbError::from)?;
all.push(
added
.filter_map(|sig| async move {
match sig.args() {
Ok(a) => Some(DiscoveryEvent::PrinterAdded(PrinterSnapshot {
id: a.printer_id.to_string(),
name: a.printer_name.to_string(),
info: a.printer_info.to_string(),
location: a.printer_location.to_string(),
make_model: a.printer_make_and_model.to_string(),
accepting_jobs: a.printer_is_accepting_jobs,
state: PrinterState::from(a.printer_state),
backend: a.backend_name.to_string(),
})),
Err(e) => {
log::debug!("cpdb-rs: failed to parse PrinterAdded signal: {}", e);
None
}
}
})
.boxed(),
);
let removed = bh
.proxy
.receive_printer_removed()
.await
.map_err(CpdbError::from)?;
all.push(
removed
.filter_map(|sig| async move {
match sig.args() {
Ok(a) => Some(DiscoveryEvent::PrinterRemoved {
id: a.printer_id.to_string(),
backend: a.backend_name.to_string(),
}),
Err(e) => {
log::debug!(
"cpdb-rs: failed to parse PrinterRemoved signal: {}",
e
);
None
}
}
})
.boxed(),
);
let changed = bh
.proxy
.receive_printer_state_changed()
.await
.map_err(CpdbError::from)?;
all.push(
changed
.filter_map(|sig| async move {
match sig.args() {
Ok(a) => Some(DiscoveryEvent::PrinterStateChanged {
id: a.printer_id.to_string(),
backend: a.backend_name.to_string(),
state: PrinterState::from(a.printer_state),
accepting_jobs: a.printer_is_accepting_jobs,
}),
Err(e) => {
log::debug!(
"cpdb-rs: failed to parse PrinterStateChanged signal: {}",
e
);
None
}
}
})
.boxed(),
);
}
for bh in &self.backends {
if let Err(e) = bh.proxy.do_listing(true).await {
log::warn!(
"cpdb-rs: do_listing failed for {}, discovery stream may be empty: {}",
bh.service_name,
e
);
}
}
let client = self.clone();
let keep_alive_task = tokio::spawn(async move {
let ping_interval = Duration::from_secs(BACKEND_INACTIVITY_TIMEOUT.as_secs() / 3);
loop {
tokio::time::sleep(ping_interval).await;
client.keep_alive_all().await;
}
});
Ok(DiscoveryStream {
inner: all,
keep_alive_task,
})
}
pub async fn keep_alive_all(&self) {
for bh in &self.backends {
let _ = bh.proxy.keep_alive().await;
}
}
pub async fn get_printer_details(
&self,
printer_id: &str,
backend: &str,
) -> Result<(OptionsCollection, MediaCollection)> {
let proxy = self.proxy_for(backend)?;
let (_n_opts, raw_opts, _n_media, raw_media) =
retry_dbus!(proxy, get_all_options(printer_id)).map_err(CpdbError::from)?;
Ok((
OptionsCollection::from_dbus(raw_opts),
MediaCollection::from_dbus(raw_media),
))
}
pub async fn get_translations(
&self,
printer_id: &str,
backend: &str,
locale: &str,
) -> Result<HashMap<String, String>> {
let proxy = self.proxy_for(backend)?;
retry_dbus!(proxy, get_all_translations(printer_id, locale)).map_err(CpdbError::from)
}
pub async fn get_default_printer(&self, backend: &str) -> Result<String> {
let proxy = self.proxy_for(backend)?;
retry_dbus!(proxy, get_default_printer()).map_err(CpdbError::from)
}
pub async fn print_fd(
&self,
printer_id: &str,
backend: &str,
settings: &[(&str, &str)],
title: &str,
) -> Result<(String, OwnedFd)> {
let proxy = self.proxy_for(backend)?;
let (job_id, fd) = retry_dbus!(
proxy,
print_fd(printer_id, settings.len() as i32, settings, title)
)
.map_err(CpdbError::from)?;
Ok((job_id, fd))
}
pub async fn print_socket(
&self,
printer_id: &str,
backend: &str,
settings: &[(&str, &str)],
title: &str,
) -> Result<(String, String)> {
let proxy = self.proxy_for(backend)?;
let (job_id, socket_path) = retry_dbus!(
proxy,
print_socket(printer_id, settings.len() as i32, settings, title)
)
.map_err(CpdbError::from)?;
Ok((job_id, socket_path))
}
pub async fn show_remote_printers(&self, visible: bool) {
for b in &self.backends {
let _ = b.proxy.show_remote_printers(visible).await;
}
}
pub async fn show_temporary_printers(&self, visible: bool) {
for b in &self.backends {
let _ = b.proxy.show_temporary_printers(visible).await;
}
}
fn proxy_for(&self, backend: &str) -> Result<&PrintBackendProxy<'static>> {
self.backends
.iter()
.find(|b| {
if b.service_name == backend {
return true;
}
if let Some(idx) = b.service_name.rfind('.') {
&b.service_name[idx + 1..] == backend
} else {
false
}
})
.map(|b| &b.proxy)
.ok_or_else(|| {
CpdbError::BackendError(format!("No backend found matching '{}'", backend))
})
}
}
pub struct DiscoveryStream {
inner: SelectAll<futures_util::stream::BoxStream<'static, DiscoveryEvent>>,
keep_alive_task: tokio::task::JoinHandle<()>,
}
impl Stream for DiscoveryStream {
type Item = DiscoveryEvent;
fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
Pin::new(&mut self.inner).poll_next(cx)
}
}
impl Drop for DiscoveryStream {
fn drop(&mut self) {
self.keep_alive_task.abort();
}
}