diff --git a/src/context/timeout.rs b/src/context/timeout.rs index 43f5bc8bfe..9bf5c99095 100644 --- a/src/context/timeout.rs +++ b/src/context/timeout.rs @@ -76,6 +76,7 @@ pub fn trigger(token: &mut CleanLockToken) { } else { break; }; - event::trigger(timeout.scheme_id, timeout.event_id, EVENT_READ); + drop(registry); + event::trigger(timeout.scheme_id, timeout.event_id, EVENT_READ, token); } } diff --git a/src/event.rs b/src/event.rs index 9342fcd0f8..918fa3e8c2 100644 --- a/src/event.rs +++ b/src/event.rs @@ -78,7 +78,7 @@ impl EventQueue { let flags = sync(RegKey { scheme, number }, token)?; if !flags.is_empty() { - trigger(scheme, number, flags); + trigger(scheme, number, flags, token); } } @@ -215,13 +215,10 @@ fn trigger_inner( } } -pub fn trigger(scheme: SchemeId, number: usize, flags: EventFlags) { - //TODO: propagate this lock token - let mut token = unsafe { CleanLockToken::new() }; - +pub fn trigger(scheme: SchemeId, number: usize, flags: EventFlags, token: &mut CleanLockToken) { // First trigger with the original file let mut todo = Vec::new(); - trigger_inner(scheme, number, flags, &mut todo, &mut token); + trigger_inner(scheme, number, flags, &mut todo, token); // Handle triggers on queues //TODO: can this be done with limited allocations? @@ -233,7 +230,7 @@ pub fn trigger(scheme: SchemeId, number: usize, flags: EventFlags) { queue_id.into(), EventFlags::EVENT_READ, &mut todo, - &mut token, + token, ); done.insert(queue_id); } diff --git a/src/ptrace.rs b/src/ptrace.rs index 4bd414cc56..0a62fb422c 100644 --- a/src/ptrace.rs +++ b/src/ptrace.rs @@ -28,13 +28,13 @@ pub struct SessionData { file_id: usize, } impl SessionData { - fn add_event(&mut self, event: PtraceEvent) { + fn add_event(&mut self, event: PtraceEvent, token: &mut CleanLockToken) { self.events.push_back(event); // Notify nonblocking tracers if self.events.len() == 1 { // If the list of events was previously empty, alert now - proc_trigger_event(self.file_id, EVENT_READ); + proc_trigger_event(self.file_id, EVENT_READ, token); } } @@ -119,12 +119,12 @@ pub fn close_tracee(session: &Session, token: &mut CleanLockToken) { session.tracer.notify(token); let data = session.data.lock(); - proc_trigger_event(data.file_id, EVENT_READ); + proc_trigger_event(data.file_id, EVENT_READ, token); } /// Trigger a notification to the event: scheme -fn proc_trigger_event(file_id: usize, flags: EventFlags) { - event::trigger(GlobalSchemes::Proc.scheme_id(), file_id, flags); +fn proc_trigger_event(file_id: usize, flags: EventFlags, token: &mut CleanLockToken) { + event::trigger(GlobalSchemes::Proc.scheme_id(), file_id, flags, token); } /// Dispatch an event to any tracer tracing `self`. This will cause @@ -140,7 +140,7 @@ pub fn send_event(event: PtraceEvent, token: &mut CleanLockToken) -> Option<()> } // Add event to queue - data.add_event(event); + data.add_event(event, token); // Notify tracer session.tracer.notify(token); @@ -219,7 +219,7 @@ pub fn breakpoint_callback( .reached = true; // Add event to queue - data.add_event(event.unwrap_or(ptrace_event!(match_flags))); + data.add_event(event.unwrap_or(ptrace_event!(match_flags)), token); // Wake up sleeping tracer session.tracer.notify(token); diff --git a/src/scheme/acpi.rs b/src/scheme/acpi.rs index dfee8efd0f..e9539112d2 100644 --- a/src/scheme/acpi.rs +++ b/src/scheme/acpi.rs @@ -3,7 +3,7 @@ use core::{ sync::atomic::{self, AtomicUsize}, }; -use alloc::boxed::Box; +use alloc::{boxed::Box, vec::Vec}; use hashbrown::{hash_map::DefaultHashBuilder, HashMap}; use spin::{Mutex, Once}; @@ -61,13 +61,17 @@ pub fn register_kstop(token: &mut CleanLockToken) -> bool { *KSTOP_FLAG.lock() = true; let mut waiters_awoken = KSTOP_WAITCOND.notify(token); - let handles = HANDLES.read(token.token()); + let fds: Vec = { + HANDLES + .read(token.token()) + .iter() + .filter(|(_, handle)| handle.kind == HandleKind::ShutdownPipe) + .map(|(fd, _)| *fd) + .collect() + }; - for (&fd, _) in handles - .iter() - .filter(|(_, handle)| handle.kind == HandleKind::ShutdownPipe) - { - event::trigger(GlobalSchemes::Acpi.scheme_id(), fd, EVENT_READ); + for fd in fds { + event::trigger(GlobalSchemes::Acpi.scheme_id(), fd, EVENT_READ, token); waiters_awoken += 1; } diff --git a/src/scheme/debug.rs b/src/scheme/debug.rs index 8491b55589..7909a501eb 100644 --- a/src/scheme/debug.rs +++ b/src/scheme/debug.rs @@ -33,8 +33,9 @@ pub fn debug_input(data: u8, token: &mut CleanLockToken) { // Notify readers of input updates pub fn debug_notify(token: &mut CleanLockToken) { - for (id, _handle) in HANDLES.read(token.token()).iter() { - event::trigger(GlobalSchemes::Debug.scheme_id(), *id, EVENT_READ); + let ids: Vec = { HANDLES.read(token.token()).iter().map(|x| *x.0).collect() }; + for id in ids { + event::trigger(GlobalSchemes::Debug.scheme_id(), id, EVENT_READ, token); } } diff --git a/src/scheme/irq.rs b/src/scheme/irq.rs index 630143f36c..6fd8bc5c73 100644 --- a/src/scheme/irq.rs +++ b/src/scheme/irq.rs @@ -61,14 +61,18 @@ const INO_PHANDLE: u64 = 0x8003_0000_0000_0000; /// Add to the input queue pub fn irq_trigger(irq: u8, token: &mut CleanLockToken) { COUNTS.lock()[irq as usize] += 1; + let fds: Vec = { + HANDLES + .read(token.token()) + .iter() + .filter_map(|(fd, handle)| Some((fd, handle.as_irq_handle()?))) + .filter(|&(_, (_, handle_irq))| handle_irq == irq) + .map(|(f, _)| *f) + .collect() + }; - for (fd, _) in HANDLES - .read(token.token()) - .iter() - .filter_map(|(fd, handle)| Some((fd, handle.as_irq_handle()?))) - .filter(|&(_, (_, handle_irq))| handle_irq == irq) - { - event::trigger(GlobalSchemes::Irq.scheme_id(), *fd, EVENT_READ); + for fd in fds { + event::trigger(GlobalSchemes::Irq.scheme_id(), fd, EVENT_READ, token); } } diff --git a/src/scheme/pipe.rs b/src/scheme/pipe.rs index 0f4ef74ce2..e2128f89e3 100644 --- a/src/scheme/pipe.rs +++ b/src/scheme/pipe.rs @@ -123,13 +123,13 @@ impl KernelScheme for PipeScheme { let can_remove = if is_write_not_read { pipe.writer_is_alive.store(false, Ordering::SeqCst); - event::trigger(scheme_id, key, EVENT_READ); + event::trigger(scheme_id, key, EVENT_READ, token); pipe.read_condition.notify(token); !pipe.reader_is_alive.load(Ordering::SeqCst) } else { pipe.reader_is_alive.store(false, Ordering::SeqCst); - event::trigger(scheme_id, key | WRITE_NOT_READ_BIT, EVENT_WRITE); + event::trigger(scheme_id, key | WRITE_NOT_READ_BIT, EVENT_WRITE, token); pipe.write_condition.notify(token); !pipe.writer_is_alive.load(Ordering::SeqCst) @@ -257,6 +257,7 @@ impl KernelScheme for PipeScheme { GlobalSchemes::Pipe.scheme_id(), key | WRITE_NOT_READ_BIT, EVENT_WRITE, + token, ); pipe.write_condition.notify(token); @@ -319,7 +320,7 @@ impl KernelScheme for PipeScheme { } if bytes_written > 0 { - event::trigger(GlobalSchemes::Pipe.scheme_id(), key, EVENT_READ); + event::trigger(GlobalSchemes::Pipe.scheme_id(), key, EVENT_READ, token); pipe.read_condition.notify(token); return Ok(bytes_written); @@ -390,7 +391,7 @@ impl KernelScheme for PipeScheme { let fds_written = vec.len() - before_len; if fds_written > 0 { - event::trigger(GlobalSchemes::Pipe.scheme_id(), key, EVENT_READ); + event::trigger(GlobalSchemes::Pipe.scheme_id(), key, EVENT_READ, token); pipe.read_condition.notify(token); return Ok(fds_written); @@ -454,6 +455,7 @@ impl KernelScheme for PipeScheme { GlobalSchemes::Pipe.scheme_id(), key | WRITE_NOT_READ_BIT, EVENT_WRITE, + token, ); pipe.write_condition.notify(token); diff --git a/src/scheme/serio.rs b/src/scheme/serio.rs index 5aeb1fb146..5fdcd8b47d 100644 --- a/src/scheme/serio.rs +++ b/src/scheme/serio.rs @@ -42,8 +42,16 @@ pub fn serio_input(index: usize, data: u8, token: &mut CleanLockToken) { INPUT[index].send(data, token); - for (id, _handle) in HANDLES.read(token.token()).iter() { - event::trigger(GlobalSchemes::Serio.scheme_id(), *id, EVENT_READ); + let ids: Vec = { + HANDLES + .read(token.token()) + .iter() + .map(|(id, _)| *id) + .collect() + }; + + for id in ids { + event::trigger(GlobalSchemes::Serio.scheme_id(), id, EVENT_READ, token); } } diff --git a/src/scheme/user.rs b/src/scheme/user.rs index 23e0725916..2cf843fc47 100644 --- a/src/scheme/user.rs +++ b/src/scheme/user.rs @@ -169,7 +169,7 @@ impl UserInner { unsafe { self.todo.condition.notify_signal(token) }; // Tell the scheme handler to read - event::trigger(self.root_id, self.scheme_id.get(), EVENT_READ); + event::trigger(self.root_id, self.scheme_id.get(), EVENT_READ, token); //TODO: wait for all todo and done to be processed? Ok(()) @@ -251,7 +251,7 @@ impl UserInner { } self.todo.send(sqe, token); - event::trigger(self.root_id, self.scheme_id.get(), EVENT_READ); + event::trigger(self.root_id, self.scheme_id.get(), EVENT_READ, token); } loop { @@ -345,7 +345,7 @@ impl UserInner { }, token, ); - event::trigger(self.root_id, self.scheme_id.get(), EVENT_READ); + event::trigger(self.root_id, self.scheme_id.get(), EVENT_READ, token); // 1. If cancellation was requested and arrived // before the scheme processed the request, an @@ -758,7 +758,7 @@ impl UserInner { }, token, ); - event::trigger(self.root_id, self.scheme_id.get(), EVENT_READ); + event::trigger(self.root_id, self.scheme_id.get(), EVENT_READ, token); Ok(()) } @@ -889,7 +889,7 @@ impl UserInner { } } ParsedCqe::TriggerFevent { number, flags } => { - event::trigger(self.scheme_id, number, flags) + event::trigger(self.scheme_id, number, flags, token) } } Ok(()) @@ -1574,7 +1574,12 @@ impl KernelScheme for UserScheme { token, ); - event::trigger(self.inner.root_id, self.inner.scheme_id.get(), EVENT_READ); + event::trigger( + self.inner.root_id, + self.inner.scheme_id.get(), + EVENT_READ, + token, + ); Ok(()) } diff --git a/src/syscall/process.rs b/src/syscall/process.rs index 1cd89f54c5..8ea1f7638d 100644 --- a/src/syscall/process.rs +++ b/src/syscall/process.rs @@ -60,6 +60,7 @@ pub fn exit_this_context(excp: Option, token: &mut CleanLock GlobalSchemes::Proc.scheme_id(), owner.get(), EventFlags::EVENT_READ, + token, ); } {