use crate::host::{is_callable, with_host, JsObj};
use fusevm::Value;
use indexmap::IndexMap;
use std::collections::{HashMap, VecDeque};
use std::io::{self, Write};
pub const METHODS: &[&str] = &[
"createInterface",
"clearLine",
"clearScreenDown",
"cursorTo",
"moveCursor",
"emitKeypressEvents",
];
pub const INTERFACE_METHODS: &[&str] = &[
"question",
"write",
"close",
"pause",
"resume",
"prompt",
"setPrompt",
"getPrompt",
"on",
"once",
"addListener",
"prependListener",
"removeListener",
"off",
"removeAllListeners",
"@@asyncIterator",
];
pub const PROMISES_METHODS: &[&str] = &["createInterface"];
pub fn call(method: &str, args: &[Value]) -> Option<Result<Value, String>> {
if method.starts_with("@@") {
return internal_call(method, args);
}
Some(match method {
"createInterface" => Ok(create_interface(args, false)),
"cursorTo" => {
let x = super::arg_num(args, 1);
let y = args.get(2).filter(|v| !matches!(v, Value::Undef));
let seq = match y {
Some(yv) => format!(
"\x1b[{};{}H",
with_host(|h| h.to_number(yv)) as i64 + 1,
x as i64 + 1
),
None => format!("\x1b[{}G", x as i64 + 1),
};
write_stdout(&seq);
Ok(Value::Bool(true))
}
"moveCursor" => {
let dx = super::arg_num(args, 1) as i64;
let dy = super::arg_num(args, 2) as i64;
let mut seq = String::new();
if dx > 0 {
seq.push_str(&format!("\x1b[{dx}C"));
} else if dx < 0 {
seq.push_str(&format!("\x1b[{}D", -dx));
}
if dy > 0 {
seq.push_str(&format!("\x1b[{dy}B"));
} else if dy < 0 {
seq.push_str(&format!("\x1b[{}A", -dy));
}
write_stdout(&seq);
Ok(Value::Bool(true))
}
"clearLine" => {
let dir = super::arg_num(args, 1);
let seq = if dir < 0.0 {
"\x1b[1K"
} else if dir > 0.0 {
"\x1b[0K"
} else {
"\x1b[2K"
};
write_stdout(seq);
Ok(Value::Bool(true))
}
"clearScreenDown" => {
write_stdout("\x1b[0J");
Ok(Value::Bool(true))
}
"emitKeypressEvents" => Ok(Value::Undef),
_ => return None,
})
}
pub fn construct(args: &[Value]) -> Result<Value, String> {
Ok(create_interface(args, false))
}
pub fn constant(name: &str) -> Option<Value> {
match name {
"Interface" => Some(with_host(|h| h.alloc(JsObj::Builtin("Interface".into())))),
_ => None,
}
}
pub fn promises_call(method: &str, args: &[Value]) -> Option<Result<Value, String>> {
match method {
"createInterface" => Some(Ok(create_interface(args, true))),
_ => None,
}
}
fn create_interface(args: &[Value], promises: bool) -> Value {
let id = NEXT_ID.with(|n| {
let id = n.get();
n.set(id + 1);
id
});
let (input, output) = match args.first() {
Some(o) if opt_prop(o, "input").is_some() => (
opt_prop(o, "input").unwrap_or(Value::Undef),
opt_prop(o, "output").unwrap_or(Value::Undef),
),
_ => (
args.first().cloned().unwrap_or(Value::Undef),
args.get(1).cloned().unwrap_or(Value::Undef),
),
};
let iface = with_host(|h| {
let listeners = h.new_object(IndexMap::new());
let prompt = h.new_str("> ");
let mut m = IndexMap::new();
m.insert("@@native".into(), h.new_str("Interface"));
m.insert("@@input".into(), input.clone());
m.insert("@@output".into(), output);
m.insert("@@prompt".into(), prompt);
m.insert("@@listeners".into(), listeners);
m.insert("@@rlid".into(), Value::Float(id as f64));
if promises {
let flag = h.new_str("1");
m.insert("@@promises".into(), flag);
}
h.new_object(m)
});
if matches!(input, Value::Obj(_)) && !is_stdin(&input) {
attach_stream(&iface, &input);
}
iface
}
pub fn instance_call(recv: &Value, method: &str, args: Vec<Value>) -> Result<Value, String> {
match method {
"question" => {
let query = with_host(|h| args.first().map(|v| h.str_of(v)).unwrap_or_default());
write_stdout(&query);
if with_state(recv, |s| s.flowing) {
if read_hidden(recv, "@@promises") == "1" {
let (promise, id) = with_host(|h| {
let p = h.new_promise();
let id = h.promise_id(&p).unwrap_or(0);
(p, id)
});
with_state(recv, |s| s.question = Some(QuestionReply::Promise(id)));
return Ok(promise);
}
let cb = args
.iter()
.rev()
.find(|v| with_host(|h| is_callable(h, v)))
.cloned();
if let Some(cb) = cb {
with_state(recv, |s| s.question = Some(QuestionReply::Callback(cb)));
}
return Ok(Value::Undef);
}
let line = read_line();
if read_hidden(recv, "@@promises") == "1" {
let line_val = with_host(|h| h.new_str(line));
return crate::builtins::promise_resolve_pub(line_val);
}
let cb = args
.iter()
.rev()
.find(|v| with_host(|h| is_callable(h, v)))
.cloned();
if let Some(cb) = cb {
let line_val = with_host(|h| h.new_str(line));
crate::host::invoke(&cb, vec![line_val], None)?;
}
Ok(Value::Undef)
}
"write" => {
let data = with_host(|h| args.first().map(|v| h.str_of(v)).unwrap_or_default());
write_output(recv, &data);
Ok(Value::Undef)
}
"prompt" => {
let p = read_hidden(recv, "@@prompt");
write_output(recv, &p);
Ok(Value::Undef)
}
"setPrompt" => {
let p = with_host(|h| args.first().map(|v| h.str_of(v)).unwrap_or_default());
with_host(|h| {
let pv = h.new_str(p);
if let Some(JsObj::Object(m)) = h.get_mut(recv) {
m.insert("@@prompt".into(), pv);
}
});
Ok(Value::Undef)
}
"getPrompt" => {
let prompt = read_hidden(recv, "@@prompt");
Ok(with_host(|h| h.new_str(prompt)))
}
"on" | "once" | "addListener" | "prependListener" => {
if let (Some(ev), Some(cb)) = (args.first(), args.get(1)) {
let event = with_host(|h| h.str_of(ev));
store_listener(
recv,
&event,
cb.clone(),
method == "once",
method == "prependListener",
);
if event == "line" || event == "close" {
ensure_flowing(recv);
}
}
Ok(recv.clone())
}
"removeListener" | "off" => {
if let (Some(ev), Some(cb)) = (args.first(), args.get(1)) {
let event = with_host(|h| h.str_of(ev));
remove_listener(recv, &event, cb);
}
Ok(recv.clone())
}
"removeAllListeners" => {
let event = args
.first()
.filter(|v| !matches!(v, Value::Undef))
.map(|v| with_host(|h| h.str_of(v)));
clear_listeners(recv, event.as_deref());
Ok(recv.clone())
}
"close" => {
close(recv)?;
Ok(Value::Undef)
}
"@@asyncIterator" => Ok(async_iterator(recv)),
"pause" | "resume" => Ok(recv.clone()),
_ => Err(crate::host::type_error(&format!(
"{method} is not a function"
))),
}
}
fn read_hidden(recv: &Value, key: &str) -> String {
with_host(|h| match h.get(recv) {
Some(JsObj::Object(p)) => p.get(key).map(|v| h.str_of(v)).unwrap_or_default(),
_ => String::new(),
})
}
fn listeners_obj(recv: &Value) -> Option<Value> {
opt_prop(recv, "@@listeners")
}
fn store_listener(recv: &Value, event: &str, cb: Value, once: bool, prepend: bool) {
let Some(listeners) = listeners_obj(recv) else {
return;
};
with_host(|h| {
let entry = h.new_array(vec![cb, Value::Bool(once)]);
let arr = match h.get(&listeners) {
Some(JsObj::Object(p)) => p.get(event).cloned(),
_ => None,
};
let arr = arr.filter(|a| matches!(h.get(a), Some(JsObj::Array(_))));
match arr {
Some(a) => {
if let Some(JsObj::Array(items)) = h.get_mut(&a) {
if prepend {
items.insert(0, entry);
} else {
items.push(entry);
}
}
}
None => {
let a = h.new_array(vec![entry]);
if let Some(JsObj::Object(p)) = h.get_mut(&listeners) {
p.insert(event.to_string(), a);
}
}
}
});
}
fn listener_entries(recv: &Value, event: &str) -> Vec<(Value, bool)> {
let Some(listeners) = listeners_obj(recv) else {
return Vec::new();
};
with_host(|h| {
let arr = match h.get(&listeners) {
Some(JsObj::Object(p)) => p.get(event).cloned(),
_ => None,
};
let Some(Some(JsObj::Array(items))) = arr.map(|a| h.get(&a).cloned()) else {
return Vec::new();
};
items
.iter()
.filter_map(|e| match h.get(e) {
Some(JsObj::Array(pair)) if pair.len() == 2 => {
Some((pair[0].clone(), h.truthy(&pair[1])))
}
_ => None,
})
.collect()
})
}
fn remove_listener(recv: &Value, event: &str, cb: &Value) {
let Some(listeners) = listeners_obj(recv) else {
return;
};
with_host(|h| {
let arr = match h.get(&listeners) {
Some(JsObj::Object(p)) => p.get(event).cloned(),
_ => None,
};
let Some(arr) = arr else { return };
let pos = match h.get(&arr) {
Some(JsObj::Array(items)) => items.iter().position(|e| match h.get(e) {
Some(JsObj::Array(pair)) => pair.first().is_some_and(|f| h.strict_eq(f, cb)),
_ => false,
}),
_ => None,
};
if let (Some(i), Some(JsObj::Array(items))) = (pos, h.get_mut(&arr)) {
items.remove(i);
}
});
}
fn clear_listeners(recv: &Value, event: Option<&str>) {
let Some(listeners) = listeners_obj(recv) else {
return;
};
with_host(|h| {
if let Some(JsObj::Object(p)) = h.get_mut(&listeners) {
match event {
Some(e) => {
p.shift_remove(e);
}
None => p.clear(),
}
}
});
}
fn read_line() -> String {
let mut line = String::new();
let _ = io::stdin().read_line(&mut line);
while line.ends_with('\n') || line.ends_with('\r') {
line.pop();
}
line
}
fn write_stdout(s: &str) {
let mut out = io::stdout();
let _ = out.write_all(s.as_bytes());
let _ = out.flush();
}
fn write_output(recv: &Value, s: &str) {
let out = opt_prop(recv, "@@output").unwrap_or(Value::Undef);
if matches!(out, Value::Obj(_)) {
let payload = with_host(|h| h.new_str(s.to_string()));
if crate::host::call_method(&out, "write", vec![payload]).is_ok() {
return;
}
}
write_stdout(s);
}
fn opt_prop(v: &Value, key: &str) -> Option<Value> {
with_host(|h| match h.get(v) {
Some(JsObj::Object(p)) => p.get(key).cloned(),
_ => None,
})
}
#[derive(Default)]
struct LineState {
pending: Vec<u8>,
after_cr: bool,
to_emit: VecDeque<String>,
emitting: bool,
ended: bool,
closed: bool,
flowing: bool,
iterating: bool,
buffered: VecDeque<String>,
waiters: VecDeque<u32>,
question: Option<QuestionReply>,
}
enum QuestionReply {
Callback(Value),
Promise(u32),
}
thread_local! {
static LINES: std::cell::RefCell<HashMap<u64, LineState>> =
std::cell::RefCell::new(HashMap::new());
static NEXT_ID: std::cell::Cell<u64> = const { std::cell::Cell::new(1) };
}
fn rl_id(recv: &Value) -> u64 {
opt_prop(recv, "@@rlid")
.map(|v| with_host(|h| h.to_number(&v)) as u64)
.unwrap_or(0)
}
fn with_state<R>(recv: &Value, f: impl FnOnce(&mut LineState) -> R) -> R {
let id = rl_id(recv);
LINES.with(|m| f(m.borrow_mut().entry(id).or_default()))
}
fn is_stdin(input: &Value) -> bool {
with_host(|h| match h.get(input) {
Some(JsObj::Object(p)) => {
p.get("@@native").map(|v| h.str_of(v)).as_deref() == Some("WriteStream")
&& p.get("fd").map(|v| h.to_number(v)) == Some(0.0)
}
_ => false,
})
}
fn ensure_flowing(recv: &Value) {
let input = opt_prop(recv, "@@input").unwrap_or(Value::Undef);
if !is_stdin(&input) {
return;
}
let start = with_state(recv, |s| {
!std::mem::replace(&mut s.flowing, true) && !s.closed
});
if start {
attach_stream(recv, &input);
}
}
fn on_data(recv: &Value, bytes: &[u8]) -> Result<(), String> {
if with_state(recv, |s| s.closed) {
return Ok(());
}
let start = with_state(recv, |s| {
let mut i = 0;
if s.after_cr && bytes.first() == Some(&b'\n') {
i = 1;
}
s.after_cr = false;
while i < bytes.len() {
match bytes[i] {
b'\n' => {
let line = String::from_utf8_lossy(&s.pending).into_owned();
s.pending.clear();
s.to_emit.push_back(line);
}
b'\r' => {
let line = String::from_utf8_lossy(&s.pending).into_owned();
s.pending.clear();
s.to_emit.push_back(line);
match bytes.get(i + 1) {
Some(b'\n') => i += 1,
None => s.after_cr = true,
_ => {}
}
}
b => s.pending.push(b),
}
i += 1;
}
!s.to_emit.is_empty() && !std::mem::replace(&mut s.emitting, true)
});
if start {
emit_lines(recv)?;
}
Ok(())
}
fn on_end(recv: &Value) -> Result<(), String> {
let (start, close_now) = with_state(recv, |s| {
s.ended = true;
if !s.pending.is_empty() {
let line = String::from_utf8_lossy(&s.pending).into_owned();
s.pending.clear();
s.to_emit.push_back(line);
}
let start = !s.to_emit.is_empty() && !std::mem::replace(&mut s.emitting, true);
(start, s.to_emit.is_empty() && !s.emitting)
});
if start {
emit_lines(recv)?;
} else if close_now {
close(recv)?;
}
Ok(())
}
fn emit_lines(recv: &Value) -> Result<(), String> {
while let Some(line) = with_state(recv, |s| s.to_emit.pop_front()) {
if let Err(e) = deliver(recv, line) {
with_state(recv, |s| s.emitting = false);
return Err(e);
}
}
let close_now = with_state(recv, |s| {
s.emitting = false;
s.ended
});
if close_now {
close(recv)?;
}
Ok(())
}
fn deliver(recv: &Value, line: String) -> Result<(), String> {
match with_state(recv, |s| s.question.take()) {
Some(QuestionReply::Callback(cb)) => {
let v = with_host(|h| h.new_str(line));
crate::host::invoke(&cb, vec![v], None)?;
return Ok(());
}
Some(QuestionReply::Promise(id)) => {
let v = with_host(|h| h.new_str(line));
crate::host::resolve_promise_val(id, v);
return Ok(());
}
None => {}
}
let v = with_host(|h| h.new_str(line.clone()));
emit(recv, "line", vec![v.clone()])?;
let waiter = with_state(recv, |s| {
if !s.iterating {
return None;
}
let w = s.waiters.pop_front();
if w.is_none() {
s.buffered.push_back(line);
}
w
});
if let Some(id) = waiter {
crate::host::resolve_promise_val(id, iter_result(v, false));
}
Ok(())
}
fn close(recv: &Value) -> Result<(), String> {
let waiters = with_state(recv, |s| {
if std::mem::replace(&mut s.closed, true) {
return None;
}
Some(std::mem::take(&mut s.waiters))
});
let Some(waiters) = waiters else {
return Ok(());
};
let input = opt_prop(recv, "@@input").unwrap_or(Value::Undef);
if is_stdin(&input) {
crate::host::call_method(&input, "pause", Vec::new())?;
}
emit(recv, "close", Vec::new())?;
for id in waiters {
crate::host::resolve_promise_val(id, iter_result(Value::Undef, true));
}
Ok(())
}
fn emit(recv: &Value, event: &str, args: Vec<Value>) -> Result<(), String> {
let entries = listener_entries(recv, event);
for (cb, once) in entries {
if once {
remove_listener(recv, event, &cb);
}
crate::host::invoke(&cb, args.clone(), Some(recv.clone()))?;
}
Ok(())
}
fn iter_result(value: Value, done: bool) -> Value {
with_host(|h| {
let mut m = IndexMap::new();
m.insert("value".into(), value);
m.insert("done".into(), Value::Bool(done));
h.new_object(m)
})
}
fn internal_call(method: &str, args: &[Value]) -> Option<Result<Value, String>> {
let recv = args.first().cloned().unwrap_or(Value::Undef);
let r = match method {
"@@feed" => {
let chunk = args.get(1).cloned().unwrap_or(Value::Undef);
let bytes = super::buffer::view_bytes(&chunk)
.unwrap_or_else(|| with_host(|h| h.str_of(&chunk)).into_bytes());
on_data(&recv, &bytes)
}
"@@end" => on_end(&recv),
"@@close" => close(&recv),
_ => return None,
};
Some(r.map(|_| Value::Undef))
}
fn attach_stream(recv: &Value, input: &Value) {
with_state(recv, |s| s.flowing = true);
let has_on = crate::builtins::get_property(input, "on")
.map(|f| with_host(|h| is_callable(h, &f)))
.unwrap_or(false);
if !has_on {
return;
}
for (event, internal) in [("data", "readline.@@feed"), ("end", "readline.@@end")] {
let (name, cb) = with_host(|h| {
let target = h.alloc(JsObj::Builtin(internal.into()));
let undef = Value::Undef;
let cb = h.alloc(JsObj::BoundFunc {
target,
this: undef,
args: vec![recv.clone()],
});
(h.new_str(event), cb)
});
let _ = crate::host::call_method(input, "on", vec![name, cb]);
}
}
pub const ITERATOR_METHODS: &[&str] = &["next", "return", "@@asyncIterator"];
fn async_iterator(recv: &Value) -> Value {
with_state(recv, |s| s.iterating = true);
ensure_flowing(recv);
with_host(|h| {
let mut m = IndexMap::new();
m.insert("@@native".into(), h.new_str("ReadlineIterator"));
m.insert("@@iface".into(), recv.clone());
h.new_object(m)
})
}
pub fn iterator_call(recv: &Value, method: &str) -> Result<Value, String> {
let iface = opt_prop(recv, "@@iface").unwrap_or(Value::Undef);
match method {
"@@asyncIterator" => Ok(recv.clone()),
"next" => {
let ready = with_state(&iface, |s| match s.buffered.pop_front() {
Some(line) => Some(Some(line)),
None if s.closed => Some(None),
None => None,
});
match ready {
Some(Some(line)) => {
let v = with_host(|h| h.new_str(line));
crate::builtins::promise_resolve_pub(iter_result(v, false))
}
Some(None) => crate::builtins::promise_resolve_pub(iter_result(Value::Undef, true)),
None => {
let (promise, id) = with_host(|h| {
let p = h.new_promise();
let id = h.promise_id(&p).unwrap_or(0);
(p, id)
});
with_state(&iface, |s| s.waiters.push_back(id));
Ok(promise)
}
}
}
"return" => {
with_host(|h| {
let cb = h.alloc(JsObj::Builtin("readline.@@close".into()));
h.queue_nexttick(cb, vec![iface.clone()]);
});
crate::builtins::promise_resolve_pub(iter_result(Value::Undef, true))
}
_ => Err(crate::host::type_error(&format!(
"{method} is not a function"
))),
}
}