fix(find): plugin never successfully interrupting

This commit is contained in:
Michael Aaron Murphy 2022-03-28 00:56:05 +02:00 • committed by Michael Murphy
parent 95e96c8958
commit f7823707fc

View file

@ -11,6 +11,7 @@ use std::rc::Rc;
use tokio::io::AsyncBufReadExt; use tokio::io::AsyncBufReadExt;
use tokio::process::{Child, ChildStdout, Command}; use tokio::process::{Child, ChildStdout, Command};
#[derive(Debug)]
enum Event { enum Event {
Activate(u32), Activate(u32),
Search(String), Search(String),
@ -63,7 +64,6 @@ pub async fn main() {
let tx = interrupt_tx.clone(); let tx = interrupt_tx.clone();
async move { async move {
if active.get() { if active.get() {
tracing::debug!("sending interrupt");
let _ = tx.try_send(()); let _ = tx.try_send(());
} }
} }
@ -145,8 +145,6 @@ impl SearchContext {
/// Submits the query to `fdfind` and actively monitors the search results while handling interrupts. /// Submits the query to `fdfind` and actively monitors the search results while handling interrupts.
async fn search(&mut self, search: String) { async fn search(&mut self, search: String) {
self.search_results.clear(); self.search_results.clear();
tracing::debug!("searching for {}", search);
let (mut child, mut stdout) = match query(&search).await { let (mut child, mut stdout) = match query(&search).await {
Ok((child, stdout)) => (child, tokio::io::BufReader::new(stdout).lines()), Ok((child, stdout)) => (child, tokio::io::BufReader::new(stdout).lines()),
Err(why) => { Err(why) => {
@ -170,34 +168,45 @@ impl SearchContext {
} }
}; };
let mut id = 0; let timeout = async {
let mut append; tokio::time::sleep(std::time::Duration::from_secs(3)).await;
};
'stream: loop { let listener = async {
let interrupt = async { let mut id = 0;
let _ = self.interrupt_rx.recv_async().await; let mut append;
Ok(None)
};
match crate::or(interrupt, stdout.next_line()).await { 'stream: loop {
Ok(Some(line)) => append = line, let interrupt = async {
Ok(None) => break 'stream, let _ = self.interrupt_rx.recv_async().await;
Err(why) => { Ok(None)
tracing::error!("error on stdout line read: {}", why); };
match crate::or(interrupt, stdout.next_line()).await {
Ok(Some(line)) => append = line,
Ok(None) => break 'stream,
Err(why) => {
tracing::error!("error on stdout line read: {}", why);
break 'stream;
}
}
self.append(id, append).await;
id += 1;
if id == 10 {
break 'stream; break 'stream;
} }
} }
};
self.append(id, append).await; futures::pin_mut!(timeout);
futures::pin_mut!(listener);
id += 1; let _ = futures::future::select(timeout, listener).await;
if id == 10 { let _ = child.kill().await;
break 'stream;
}
}
let _ = child.kill();
let _ = child.wait().await; let _ = child.wait().await;
} }
} }