oxedyne/fe2o3/fe2o3_tui/src/lib_tui/proc.rs
2.9 KiB, 5 runs
created by r1870400018:1197, which is this file's identity for as long as the history lasts, whatever it is later renamed to
download · who wrote it · its history
| 1 | use crate::lib_tui::{ |
| 2 | draw::text::TextLines, |
| 3 | }; |
| 4 | |
| 5 | use oxedyne_fe2o3_core::{ |
| 6 | prelude::*, |
| 7 | channels::Simplex, |
| 8 | }; |
| 9 | |
| 10 | use std::{ |
| 11 | fmt::Debug, |
| 12 | marker::PhantomData, |
| 13 | io::{ |
| 14 | BufRead, |
| 15 | Write, |
| 16 | }, |
| 17 | sync::{ |
| 18 | Arc, |
| 19 | RwLock, |
| 20 | }, |
| 21 | thread, |
| 22 | }; |
| 23 | |
| 24 | |
| 25 | #[derive(Clone, Debug)] |
| 26 | pub struct Process< |
| 27 | R: BufRead + Debug + Send + Sync, |
| 28 | W: Write + Debug + Send + Sync, |
| 29 | > { |
| 30 | label: String, |
| 31 | stream_in: Arc<RwLock<W>>, |
| 32 | stream_out: Arc<RwLock<R>>, |
| 33 | output: Arc<RwLock<TextLines>>, |
| 34 | _phantom: PhantomData<(R, W)>, |
| 35 | } |
| 36 | |
| 37 | impl< |
| 38 | R: BufRead + Debug + Send + Sync + 'static, |
| 39 | W: Write + Debug + Send + Sync + 'static, |
| 40 | > |
| 41 | Process<R, W> |
| 42 | { |
| 43 | fn new( |
| 44 | label: String, |
| 45 | stream_in: W, |
| 46 | stream_out: R, |
| 47 | ) |
| 48 | -> Self |
| 49 | { |
| 50 | Self { |
| 51 | label, |
| 52 | stream_in: Arc::new(RwLock::new(stream_in)), |
| 53 | stream_out: Arc::new(RwLock::new(stream_out)), |
| 54 | output: Arc::new(RwLock::new(TextLines::default())), |
| 55 | _phantom: PhantomData, |
| 56 | } |
| 57 | } |
| 58 | |
| 59 | fn write(&self, byts: &[u8]) -> Outcome<()> { |
| 60 | let mut stream = lock_write!(self.stream_in); |
| 61 | res!(stream.write_all(byts)); |
| 62 | res!(stream.flush()); |
| 63 | Ok(()) |
| 64 | } |
| 65 | |
| 66 | fn read(&self, buf: &mut String) -> Outcome<usize> { |
| 67 | let mut stream = lock_write!(self.stream_out); |
| 68 | let byts = res!(stream.read_line(buf)); |
| 69 | Ok(byts) |
| 70 | } |
| 71 | } |
| 72 | |
| 73 | #[derive(Clone, Debug)] |
| 74 | pub enum Msg< |
| 75 | R: BufRead + Debug + Send + Sync + 'static, |
| 76 | W: Write + Debug + Send + Sync + 'static, |
| 77 | > { |
| 78 | AddProcess(Process<R, W>), |
| 79 | Input(usize, String), |
| 80 | } |
| 81 | |
| 82 | pub struct ProcessManager< |
| 83 | R: BufRead + Debug + Send + Sync + 'static, |
| 84 | W: Write + Debug + Send + Sync + 'static, |
| 85 | > { |
| 86 | procs: Vec<Process<R, W>>, |
| 87 | chan_in: Simplex<Msg<R, W>>, |
| 88 | } |
| 89 | |
| 90 | impl< |
| 91 | R: BufRead + Debug + Send + Sync + 'static, |
| 92 | W: Write + Debug + Send + Sync + 'static, |
| 93 | > |
| 94 | ProcessManager<R, W> |
| 95 | { |
| 96 | fn new(chan_in: Simplex<Msg<R, W>>) -> Self { |
| 97 | Self { |
| 98 | procs: Vec::new(), |
| 99 | chan_in, |
| 100 | } |
| 101 | } |
| 102 | |
| 103 | fn add_process(&mut self, proc: Process<R, W>) -> Outcome<()> { |
| 104 | self.chan_in.send(Msg::AddProcess(proc)) |
| 105 | } |
| 106 | |
| 107 | fn write(&self, ind: usize, input: String) -> Outcome<()> { |
| 108 | self.chan_in.send(Msg::Input(ind, input)) |
| 109 | } |
| 110 | |
| 111 | //fn listen(&mut self) { |
| 112 | // let procs = self.procs.clone(); |
| 113 | // thread::spawn(move || { |
| 114 | // loop { |
| 115 | // for proc in &procs { |
| 116 | // let mut buf = String::new(); |
| 117 | // let byts = res!(proc.read_output(&mut buf)); |
| 118 | // if bytes_read > 0 { |
| 119 | // let mut output = lock_write!(proc.output); |
| 120 | // output.add_text(buf.trim()); |
| 121 | // } |
| 122 | // buf.clear(); |
| 123 | // } |
| 124 | // } |
| 125 | // }); |
| 126 | //} |
| 127 | } |