@@ -15,3 +15,4 @@ serde_json = "1"
|
||||
tokio = { version = "1", features = ["rt", "rt-multi-thread", "macros", "fs", "signal", "sync"] }
|
||||
tracing = "0.1"
|
||||
tracing-subscriber = { version = "0.3", features = ["env-filter"] }
|
||||
db = { path = "../libs/db" }
|
||||
|
||||
+39
-16
@@ -1,3 +1,4 @@
|
||||
use db::{Db, DbSection};
|
||||
use kobject_uevent::KobjectUevent;
|
||||
use rule::{engine::Engine, parser::When};
|
||||
use rustc_hash::FxHashMap;
|
||||
@@ -7,52 +8,57 @@ use std::{
|
||||
path::{Path, PathBuf},
|
||||
sync::{Arc, LazyLock, RwLock},
|
||||
};
|
||||
use tokio::sync::Mutex;
|
||||
use tokio::sync::{Mutex, broadcast};
|
||||
|
||||
static DEVICES: LazyLock<RwLock<FxHashMap<String, ArcDevice>>> =
|
||||
LazyLock::new(|| Default::default());
|
||||
static DB: LazyLock<Arc<Db>> = LazyLock::new(|| Db::open(db::PERSIST_PATH));
|
||||
static NOTIFY: LazyLock<broadcast::Sender<String>> = LazyLock::new(|| broadcast::channel(128).0);
|
||||
|
||||
pub type ArcDevice = Arc<Device>;
|
||||
|
||||
#[derive(Debug)]
|
||||
pub struct Device {
|
||||
lock: Mutex<()>,
|
||||
properties: DbSection,
|
||||
engine: Engine,
|
||||
}
|
||||
impl Device {
|
||||
pub fn new() -> Self {
|
||||
let engine = Engine::new();
|
||||
engine.set_property("DEVROOT".into(), dev_root().to_string_lossy().into());
|
||||
pub fn new(name: &str) -> Self {
|
||||
let engine = Engine::with_properties(DB.section(name));
|
||||
engine
|
||||
.properties
|
||||
.insert("DEVROOT", dev_root().to_string_lossy().into());
|
||||
|
||||
Self {
|
||||
lock: Mutex::new(()),
|
||||
properties: DB.section(name),
|
||||
engine,
|
||||
}
|
||||
}
|
||||
|
||||
fn apply_uevent(&self, ev: &KobjectUevent) {
|
||||
for (k, v) in ev.inner().iter() {
|
||||
self.engine.set_property(k.into(), v.into());
|
||||
self.engine.properties.insert(k, v.into());
|
||||
}
|
||||
self.update_derived_properties();
|
||||
}
|
||||
|
||||
fn update_derived_properties(&self) {
|
||||
if let Some(devname) = self.engine.get_property("DEVNAME") {
|
||||
self.engine.set_property(
|
||||
"DEVNODE".into(),
|
||||
format!("{}/{devname}", dev_root().display()),
|
||||
);
|
||||
if let Some(devname) = self.engine.properties.get("DEVNAME") {
|
||||
self.engine
|
||||
.properties
|
||||
.insert("DEVNODE", format!("{}/{devname}", dev_root().display()));
|
||||
}
|
||||
}
|
||||
|
||||
async fn update_devnode_creds(&self) {
|
||||
let Some(devnode) = self.engine.get_property("DEVNODE") else {
|
||||
let Some(devnode) = self.engine.properties.get("DEVNODE") else {
|
||||
return;
|
||||
};
|
||||
let uid = crate::util::uid_by_string(self.engine.get_property("OWNER").as_deref());
|
||||
let gid = crate::util::gid_by_string(self.engine.get_property("GROUP").as_deref());
|
||||
let mode = crate::util::parse_mode(self.engine.get_property("MODE").as_deref());
|
||||
let uid = crate::util::uid_by_string(self.engine.properties.get("OWNER").as_deref());
|
||||
let gid = crate::util::gid_by_string(self.engine.properties.get("GROUP").as_deref());
|
||||
let mode = crate::util::parse_mode(self.engine.properties.get("MODE").as_deref());
|
||||
|
||||
if let Err(err) = crate::util::set_ownership(&devnode, uid, gid).await {
|
||||
tracing::warn!("{devnode}: failed to set ownership: {err}");
|
||||
@@ -67,7 +73,7 @@ impl Device {
|
||||
}
|
||||
|
||||
async fn update_symlinks(&self) {
|
||||
let Some(devnode) = self.engine.get_property("DEVNODE") else {
|
||||
let Some(devnode) = self.engine.properties.get("DEVNODE") else {
|
||||
return;
|
||||
};
|
||||
for i in self.engine.get_list("SYMLINKS") {
|
||||
@@ -104,6 +110,9 @@ pub async fn handle_uevent(ev: &KobjectUevent) {
|
||||
}
|
||||
None => (),
|
||||
}
|
||||
if let Some(devpath) = ev.devpath() {
|
||||
_ = NOTIFY.send(devpath.into());
|
||||
}
|
||||
}
|
||||
|
||||
async fn add(ev: &KobjectUevent) {
|
||||
@@ -112,7 +121,7 @@ async fn add(ev: &KobjectUevent) {
|
||||
let Some(devpath) = ev.devpath() else {
|
||||
return;
|
||||
};
|
||||
let device = Device::new();
|
||||
let device = Device::new(devpath);
|
||||
device.apply_uevent(&ev);
|
||||
device.engine.exec(&crate::rules::get()).await;
|
||||
|
||||
@@ -140,6 +149,7 @@ async fn remove(ev: &KobjectUevent) {
|
||||
|
||||
device.remove_symlinks().await;
|
||||
|
||||
device.properties.clear();
|
||||
DEVICES.write().unwrap().remove(devpath);
|
||||
}
|
||||
|
||||
@@ -221,6 +231,19 @@ pub async fn wait_for_idle() {
|
||||
}
|
||||
}
|
||||
|
||||
pub async fn wait_for_event() -> Option<String> {
|
||||
let mut rx = NOTIFY.subscribe();
|
||||
rx.recv().await.ok()
|
||||
}
|
||||
|
||||
pub async fn get_device(name: &str) -> DbSection {
|
||||
DB.section(name)
|
||||
}
|
||||
|
||||
pub fn sync_db() {
|
||||
DB.save();
|
||||
}
|
||||
|
||||
fn dev_root() -> PathBuf {
|
||||
std::env::var("LXDEVICED_DEV_ROOT")
|
||||
.as_deref()
|
||||
|
||||
@@ -3,6 +3,7 @@ use ipc::{
|
||||
server::{Connection, Listener},
|
||||
};
|
||||
use rule::parser::When;
|
||||
use rustc_hash::FxHashMap;
|
||||
use serde::Serialize;
|
||||
use std::pin::Pin;
|
||||
|
||||
@@ -35,6 +36,8 @@ async fn handle_connection(mut connection: Connection) -> std::io::Result<()> {
|
||||
Request::ReloadRules => handler(reload_rules()),
|
||||
Request::ListDevices(when) => handler(list_devices(when)),
|
||||
Request::WaitForIdle => handler(wait_for_idle()),
|
||||
Request::WaitForEvent => handler(wait_for_event()),
|
||||
Request::GetDevice(name) => handler(get_device(name)),
|
||||
};
|
||||
let reply = handler.await;
|
||||
connection.send(reply).await?;
|
||||
@@ -63,6 +66,16 @@ async fn wait_for_idle() -> Result<(), Error> {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn wait_for_event() -> Result<String, Error> {
|
||||
crate::device::wait_for_event()
|
||||
.await
|
||||
.ok_or(Error::PermissionDenied)
|
||||
}
|
||||
|
||||
async fn get_device(name: String) -> Result<FxHashMap<String, String>, Error> {
|
||||
Ok(crate::device::get_device(&name).await.dump())
|
||||
}
|
||||
|
||||
fn handler<R: Serialize>(
|
||||
fut: impl Future<Output = Result<R, Error>> + Send + 'static,
|
||||
) -> Pin<Box<dyn Future<Output = Result<serde_json::Value, Error>> + Send>> {
|
||||
|
||||
@@ -15,6 +15,7 @@ pub fn start_monitor() -> std::io::Result<()> {
|
||||
crate::device::update_devnode(&ev).await;
|
||||
tokio::spawn(async move {
|
||||
crate::device::handle_uevent(&ev).await;
|
||||
crate::device::sync_db();
|
||||
});
|
||||
}
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user