Oregami
Repositories/oxedyne/fe2o3

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

1use crate::lib_tui::{
2 draw::text::TextLines,
3};
4
5use oxedyne_fe2o3_core::{
6 prelude::*,
7 channels::Simplex,
8};
9
10use 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)]
26pub 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
37impl<
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)]
74pub 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
82pub 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
90impl<
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}