use crate::host::*;
use super::text_output::*;
use futures::prelude::*;
use futures::channel::mpsc;
use futures::executor;
use futures::{pin_mut};
use std::str;
use std::thread;
use std::io::{BufRead};
use serde::*;
pub static STDIN_PROGRAM: StaticSubProgramId = StaticSubProgramId::called("flo_scene::stdin");
#[derive(Clone, Debug, PartialEq, Eq)]
#[derive(Serialize, Deserialize)]
pub enum TextInput {
RequestCharacter(SubProgramId),
RequestLine(SubProgramId),
PromptRequestLine(Vec<TextOutput>, SubProgramId),
}
#[derive(Clone, PartialEq, PartialOrd, Ord, Eq, Hash, Debug)]
#[derive(Serialize, Deserialize)]
pub enum TextInputResult {
Characters(String),
Eof,
}
impl SceneMessage for TextInputResult {
#[inline]
fn message_type_name() -> String { "flo_scene::TextInputResult".into() }
}
impl SceneMessage for TextInput {
fn default_target() -> StreamTarget { (*STDIN_PROGRAM).into() }
#[inline]
fn message_type_name() -> String { "flo_scene::TextInput".into() }
}
pub async fn text_input_subprogram(source: impl 'static + Send + BufRead, messages: impl Stream<Item=TextInput>, context: SceneContext) {
use std::mem;
let mut text_output = context.send(()).unwrap();
let (send_request, recv_request) = mpsc::channel::<TextInput>(0);
let (send_result, recv_result) = mpsc::channel::<(SubProgramId, TextInputResult)>(0);
let monitor_thread = thread::spawn(move || executor::block_on(async move {
use TextInput::*;
let mut recv_request = recv_request;
let mut source = source;
let mut send_result = send_result;
while let Some(input_request) = recv_request.next().await {
match input_request {
RequestCharacter(target) => {
let mut bytes = vec![];
let result = loop {
let pos = bytes.len();
bytes.push(0);
let read_err = source.read(&mut bytes[pos..pos+1]);
match read_err {
Err(err) => { break Err(err); },
Ok(0) => { bytes.pop(); continue; }
Ok(_) => { }
}
let utf8_error = str::from_utf8(&bytes);
match utf8_error {
Ok(chr) => { break Ok(chr.to_string()); }
Err(err) => {
if err.error_len().is_some() {
break Ok("\u{fffd}".to_string());
}
}
}
};
let result_is_err = result.is_err();
let send_err = match result {
Ok(chr) => send_result.send((target, TextInputResult::Characters(chr))).await,
Err(_) => send_result.send((target, TextInputResult::Eof)).await,
};
if result_is_err || send_err.is_err() {
break;
}
}
RequestLine(target) | PromptRequestLine(_, target) => {
let mut line = String::new();
let read_err = source.read_line(&mut line);
if line.ends_with('\n') {
line.remove(line.len()-1);
if line.ends_with('\r') {
line.remove(line.len()-1);
}
}
let send_err = match read_err {
Ok(0) => { send_result.send((target, TextInputResult::Eof)).await.ok(); Err(()) },
Ok(_) => send_result.send((target, TextInputResult::Characters(line))).await.map_err(|_| ()),
Err(_) => send_result.send((target, TextInputResult::Eof)).await.map_err(|_| ()),
};
if read_err.is_err() || send_err.is_err() {
break;
}
}
}
}
}));
pin_mut!(messages);
let mut send_request = send_request;
let mut recv_result = recv_result;
while let Some(request) = messages.next().await {
use TextInput::*;
let target = match &request {
RequestCharacter(target) | RequestLine(target) | PromptRequestLine(_, target) => *target
};
let request = match request {
PromptRequestLine(prompt, target) => {
for prompt_output in prompt {
text_output.send(prompt_output).await.ok();
}
RequestLine(target)
}
other => other,
};
let result = {
if send_request.send(request).await.is_ok() {
if let Some((target, message)) = recv_result.next().await {
if let Ok(mut target) = context.send(target) {
target.send(message).await.ok();
}
Ok(())
} else {
Err(())
}
} else {
Err(())
}
};
if result.is_err() {
if let Ok(mut target) = context.send(target) {
target.send(TextInputResult::Eof).await.ok();
}
}
}
send_request.close_channel();
recv_result.close();
mem::drop(send_request);
mem::drop(recv_result);
monitor_thread.join().ok();
}