use std::io::{BufReader, Read};
use std::thread;
use iced::futures::SinkExt;
use iced::stream;
use parking_lot::Mutex;
use plushie_renderer_engine::Codec;
use plushie_widget_sdk::protocol::IncomingMessage;
use plushie_widget_sdk::runtime::StdinEvent;
pub(crate) type TransportReader = BufReader<Box<dyn Read + Send>>;
pub(crate) static STDIN_RX: Mutex<Option<tokio::sync::mpsc::Receiver<StdinEvent>>> =
Mutex::new(None);
pub(crate) fn stdin_subscription() -> impl iced::futures::Stream<Item = StdinEvent> {
stream::channel(32, async |mut sender| {
let mut rx = STDIN_RX
.lock()
.take()
.expect("stdin_subscription: no receiver (called more than once?)");
while let Some(event) = rx.recv().await {
if sender.send(event).await.is_err() {
break;
}
}
})
}
pub(crate) fn spawn_stdin_reader(
codec: Codec,
sender: tokio::sync::mpsc::Sender<StdinEvent>,
mut reader: TransportReader,
) {
thread::spawn(move || {
loop {
match codec.read_message(&mut reader) {
Ok(None) => {
let _ = sender.blocking_send(StdinEvent::Closed);
break;
}
Ok(Some(bytes)) => match codec.decode::<IncomingMessage>(&bytes) {
Ok(msg) => {
if sender.blocking_send(StdinEvent::Message(msg)).is_err() {
return;
}
}
Err(e) => {
let warning = format!("parse error: {e}");
if sender.blocking_send(StdinEvent::Warning(warning)).is_err() {
return;
}
}
},
Err(e) => {
let _ = sender.blocking_send(StdinEvent::Warning(format!("read error: {e}")));
let _ = sender.blocking_send(StdinEvent::Closed);
break;
}
}
}
});
}