diff --git a/mux/src/lib.rs b/mux/src/lib.rs index 566d54e55..dbaf123a3 100644 --- a/mux/src/lib.rs +++ b/mux/src/lib.rs @@ -97,7 +97,7 @@ pub enum MuxNotification { }, } -static SUB_ID: AtomicUsize = AtomicUsize::new(0); +static LAST_SUBSCRIBER_ID: AtomicUsize = AtomicUsize::new(0); pub struct Mux { tabs: RwLock>>, @@ -137,6 +137,8 @@ fn send_actions_to_mux(pane: &Weak, dead: &Arc, actions: V histogram!("send_actions_to_mux.rate").record(1.); } +/// This is the parsing loop for the given pane. +/// It reads all data sent to `rx` (from pane PTY) and handles all terminal events for this pane. fn parse_buffered_data(pane: Weak, dead: &Arc, mut rx: FileDescriptor) { let mut buf = vec![0; configuration().mux_output_parser_buffer_size]; let mut parser = termwiz::escape::parser::Parser::new(); @@ -163,21 +165,23 @@ fn parse_buffered_data(pane: Weak, dead: &Arc, mut rx: Fil Action::CSI(CSI::Mode(Mode::SetDecPrivateMode(DecPrivateMode::Code( DecPrivateModeCode::SynchronizedOutput, )))) => { + // Synchronized output frame started: + // => We hold off ~all actions that applies changes to the terminal. hold = true; - // Flush prior actions - if !actions.is_empty() { - send_actions_to_mux(&pane, &dead, std::mem::take(&mut actions)); - action_size = 0; - } + // => We also flush prior actions + flush = true; } Action::CSI(CSI::Mode(Mode::ResetDecPrivateMode( DecPrivateMode::Code(DecPrivateModeCode::SynchronizedOutput), ))) => { + // Synchronized output frame ended: + // => We flush out all pending actions to the terminal. hold = false; flush = true; } Action::CSI(CSI::Device(dev)) if matches!(**dev, Device::SoftReset) => { + // Soft reset requested hold = false; flush = true; } @@ -307,6 +311,7 @@ fn read_from_pane_pty( } }; + // Spawn parser thread for this pane std::thread::spawn({ let dead = Arc::clone(&dead); move || parse_buffered_data(pane, &dead, rx) @@ -316,6 +321,8 @@ fn read_from_pane_pty( tx.write_all(banner.as_bytes()).ok(); } + // Loop until the pane or the main mux thread is dead. + // Read data from the pane pty and send it to the parser thread via tx/rx. while !dead.load(Ordering::Relaxed) { match reader.read(&mut buf) { Ok(size) if size == 0 => { @@ -329,9 +336,10 @@ fn read_from_pane_pty( Ok(size) => { histogram!("read_from_pane_pty.bytes.rate").record(size as f64); log::trace!("read_pty pane {pane_id} read {size} bytes"); + // Send received data to this pane parser thread. if let Err(err) = tx.write_all(&buf[..size]) { error!( - "read_pty failed to write to parser: pane {} {:?}", + "read_pty failed to write to parser for pane {}: {:?}", pane_id, err ); break; @@ -693,7 +701,7 @@ impl Mux { where F: Fn(MuxNotification) -> bool + 'static + Send + Sync, { - let sub_id = SUB_ID.fetch_add(1, Ordering::Relaxed); + let sub_id = LAST_SUBSCRIBER_ID.fetch_add(1, Ordering::Relaxed); self.subscribers .write() .insert(sub_id, Box::new(subscriber)); @@ -1005,6 +1013,7 @@ impl Mux { Ok(()) } + /// Returns the ID of the window containing the given tab ID, if any. pub fn window_containing_tab(&self, tab_id: TabId) -> Option { for w in self.windows.read().values() { for t in w.iter() {