use std::sync::Arc;
use crate::Result;
use async_broadcast::{broadcast, Receiver};
use tokio::task::JoinHandle;
use crazyflie_link::Packet;
use flume as channel;
use futures::{lock::Mutex, Stream, StreamExt};
pub struct Console {
stream_broadcast_receiver: Receiver<String>,
console_buffer: Arc<Mutex<String>>,
line_broadcast_receiver: Receiver<String>,
console_lines: Arc<Mutex<Vec<String>>>,
_console_task: JoinHandle<()>,
}
impl Console {
pub(crate) async fn new(
downlink: channel::Receiver<Packet>,
) -> Result<Self> {
let (mut stream_broadcast, stream_broadcast_receiver) = broadcast(1000);
let console_buffer: Arc<Mutex<String>> = Default::default();
let (mut line_broadcast, line_broadcast_receiver) = broadcast(1000);
stream_broadcast.set_overflow(true);
line_broadcast.set_overflow(true);
let console_lines: Arc<Mutex<Vec<String>>> = Default::default();
let buffer = console_buffer.clone();
let lines = console_lines.clone();
let _console_task = tokio::spawn(async move {
let mut line_buffer = String::new();
while let Ok(pk) = downlink.recv_async().await {
let text = String::from_utf8_lossy(pk.get_data());
buffer.lock().await.push_str(&text);
let _ = stream_broadcast.broadcast(text.clone().into_owned()).await;
line_buffer.push_str(&text);
if let Some((line, rest)) = line_buffer.clone().split_once("\n") {
line_buffer = rest.to_owned();
lines.lock().await.push(line.to_owned().clone());
let _ = line_broadcast.broadcast(line.to_owned()).await;
}
}
});
Ok(Self {
stream_broadcast_receiver,
console_buffer,
line_broadcast_receiver,
console_lines,
_console_task,
})
}
pub async fn stream(&self) -> impl Stream<Item = String> + use<> {
let buffer = self.console_buffer.lock().await;
let history_buffer = buffer.clone();
let history_stream = futures::stream::once(async { history_buffer }).boxed();
history_stream.chain(self.stream_broadcast_receiver.new_receiver())
}
pub async fn stream_no_history(&self) -> impl Stream<Item = String> + use<> {
self.stream_broadcast_receiver.new_receiver()
}
pub async fn line_stream(&self) -> impl Stream<Item = String> + use<> {
let lines = self.console_lines.lock().await;
let history_lines = lines.clone();
let history_stream = futures::stream::iter(history_lines.into_iter()).boxed();
history_stream.chain(self.line_broadcast_receiver.new_receiver())
}
pub async fn line_stream_no_history(&self) -> impl Stream<Item = String> + use<> {
self.line_broadcast_receiver.new_receiver()
}
}