调整任务停止逻辑

This commit is contained in:
lbl8603
2024-06-25 22:48:06 +08:00
parent 6798c652f9
commit 3773c09f57
4 changed files with 42 additions and 53 deletions
+3 -5
View File
@@ -18,22 +18,20 @@ pub(crate) use windows::*;
/// 仅仅是停止tun,不停止vnt /// 仅仅是停止tun,不停止vnt
#[derive(Clone, Default)] #[derive(Clone, Default)]
pub struct DeviceStop { pub struct DeviceStop {
f: Arc<Mutex<Option<Box<dyn FnOnce() -> bool + Send>>>>, f: Arc<Mutex<Option<Box<dyn FnOnce() + Send>>>>,
stopped: Arc<AtomicCell<bool>>, stopped: Arc<AtomicCell<bool>>,
} }
impl DeviceStop { impl DeviceStop {
pub fn set_stop_fn<F>(&self, f: F) pub fn set_stop_fn<F>(&self, f: F)
where where
F: FnOnce() -> bool + Send + 'static, F: FnOnce() + Send + 'static,
{ {
self.f.lock().replace(Box::new(f)); self.f.lock().replace(Box::new(f));
} }
pub fn stop(&self) -> bool { pub fn stop(&self) {
if let Some(f) = self.f.lock().take() { if let Some(f) = self.f.lock().take() {
f() f()
} else {
false
} }
} }
pub fn stopped(&self) { pub fn stopped(&self) {
+10 -23
View File
@@ -35,36 +35,23 @@ pub(crate) fn start_simple(
compressor: Compressor, compressor: Compressor,
device_stop: DeviceStop, device_stop: DeviceStop,
) -> anyhow::Result<()> { ) -> anyhow::Result<()> {
let stop_all = Arc::new(AtomicCell::new(true));
let poll = Poll::new()?; let poll = Poll::new()?;
let waker = Arc::new(Waker::new(poll.registry(), STOP)?); let waker = Arc::new(Waker::new(poll.registry(), STOP)?);
let _waker = waker.clone(); let _waker = waker.clone();
let device_cell = Arc::new(AtomicCell::new(Some(waker)));
let worker = { let worker = {
let device_cell = device_cell.clone();
stop_manager.add_listener("tun_device".into(), move || { stop_manager.add_listener("tun_device".into(), move || {
if let Some(waker) = device_cell.take() { if let Err(e) = waker.wake() {
if let Err(e) = waker.wake() { log::warn!("{:?}", e);
log::warn!("{:?}", e);
}
} }
})? })?
}; };
{ let worker_cell = Arc::new(AtomicCell::new(Some(worker)));
let stop_all = stop_all.clone(); let _worker_cell = worker_cell.clone();
device_stop.set_stop_fn(move || { device_stop.set_stop_fn(move || {
if let Some(waker) = device_cell.take() { if let Some(worker) = _worker_cell.take() {
stop_all.store(false); worker.stop_self()
if let Err(e) = waker.wake() { }
log::warn!("{:?}", e); });
return false;
}
true
} else {
false
}
});
}
if let Err(e) = start_simple0( if let Err(e) = start_simple0(
poll, poll,
context, context,
@@ -82,7 +69,7 @@ pub(crate) fn start_simple(
log::error!("{:?}", e); log::error!("{:?}", e);
}; };
device_stop.stopped(); device_stop.stopped();
if stop_all.load() { if let Some(worker) = worker_cell.take() {
worker.stop_all(); worker.stop_all();
} }
drop(_waker); drop(_waker);
+9 -18
View File
@@ -28,30 +28,21 @@ pub(crate) fn start_simple(
compressor: Compressor, compressor: Compressor,
device_stop: DeviceStop, device_stop: DeviceStop,
) -> anyhow::Result<()> { ) -> anyhow::Result<()> {
let device_cell = Arc::new(AtomicCell::new(Some(device.clone())));
let stop_all = Arc::new(AtomicCell::new(true));
let worker = { let worker = {
let device_cell = device_cell.clone(); let device = device.clone();
stop_manager.add_listener("tun_device".into(), move || { stop_manager.add_listener("tun_device".into(), move || {
if let Some(device) = device_cell.take() { if let Err(e) = device.shutdown() {
if let Err(e) = device.shutdown() { log::warn!("{:?}", e);
log::warn!("{:?}", e);
}
} }
})? })?
}; };
let worker_cell = Arc::new(AtomicCell::new(Some(worker)));
{ {
let stop_all = stop_all.clone(); let worker_cell = worker_cell.clone();
device_stop.set_stop_fn(move || { device_stop.set_stop_fn(move || {
if let Some(device) = device_cell.take() { if let Some(worker) = worker_cell.take() {
stop_all.store(false); worker.stop_self()
if let Err(e) = device.shutdown() {
log::warn!("{:?}", e);
return false;
}
true
} else {
false
} }
}); });
} }
@@ -71,7 +62,7 @@ pub(crate) fn start_simple(
log::error!("{:?}", e); log::error!("{:?}", e);
} }
device_stop.stopped(); device_stop.stopped();
if stop_all.load() { if let Some(worker) = worker_cell.take() {
worker.stop_all(); worker.stop_all();
} }
Ok(()) Ok(())
+20 -7
View File
@@ -28,7 +28,7 @@ impl StopManager {
self.inner.add_listener(name, f) self.inner.add_listener(name, f)
} }
pub fn stop(&self) { pub fn stop(&self) {
self.inner.stop(""); self.inner.stop();
} }
pub fn wait(&self) { pub fn wait(&self) {
self.inner.wait(); self.inner.wait();
@@ -81,14 +81,11 @@ impl StopManagerInner {
guard.1.push((name.clone(), Box::new(f))); guard.1.push((name.clone(), Box::new(f)));
Ok(Worker::new(name, self.clone())) Ok(Worker::new(name, self.clone()))
} }
fn stop(&self, skip_name: &str) { fn stop(&self) {
self.state.store(true, Ordering::Release); self.state.store(true, Ordering::Release);
let mut guard = self.listeners.lock(); let mut guard = self.listeners.lock();
guard.0 = true; guard.0 = true;
for (name, listener) in guard.1.drain(..) { for (_name, listener) in guard.1.drain(..) {
if &name == skip_name {
continue;
}
listener(); listener();
} }
} }
@@ -136,6 +133,19 @@ impl Worker {
} }
fn release0(&self) { fn release0(&self) {
let inner = &self.inner; let inner = &self.inner;
let worker_name = &self.name;
{
let mut mutex_guard = inner.listeners.lock();
if let Some(pos) = mutex_guard
.1
.iter()
.position(|(name, _)| name == worker_name)
{
let (_, listener) = mutex_guard.1.remove(pos);
listener();
}
}
let count = inner.worker_num.fetch_sub(1, Ordering::AcqRel); let count = inner.worker_num.fetch_sub(1, Ordering::AcqRel);
if count == 1 { if count == 1 {
for x in inner.park_threads.lock().drain(..) { for x in inner.park_threads.lock().drain(..) {
@@ -145,7 +155,10 @@ impl Worker {
} }
} }
pub fn stop_all(self) { pub fn stop_all(self) {
self.inner.stop(&self.name) self.inner.stop()
}
pub fn stop_self(self) {
drop(self)
} }
} }