use mumu::{
parser::interpreter::Interpreter,
parser::types::{Value, FunctionValue},
parser::interpreter::apply_one_function_value,
};
use std::ffi::c_void;
use std::{
sync::{
Arc,
Mutex,
mpsc::{channel, Receiver, Sender},
atomic::{AtomicUsize, Ordering},
},
collections::HashMap,
};
use lazy_static::lazy_static;
use indexmap::IndexMap;
use whoami;
pub static ACTIVE_TASKS: AtomicUsize = AtomicUsize::new(0);
#[cfg(unix)]
fn get_uid() -> i32 {
use nix::unistd::getuid;
getuid().as_raw() as i32
}
#[cfg(not(unix))]
fn get_uid() -> i32 {
0
}
#[allow(dead_code)]
enum ProcessMessage {
Ok(usize, Value),
Err(usize, String),
}
struct ProcessTask {
callback: Box<FunctionValue>,
done: bool,
}
struct ProcessManager {
_next_id: usize,
tasks: HashMap<usize, ProcessTask>,
_tx: Sender<ProcessMessage>,
rx: Receiver<ProcessMessage>,
active_count: AtomicUsize,
}
impl ProcessManager {
fn new() -> Self {
let (tx, rx) = channel();
Self {
_next_id: 0,
tasks: HashMap::new(),
_tx: tx,
rx,
active_count: AtomicUsize::new(0),
}
}
fn poll_events(&mut self, interp: &mut Interpreter) {
while let Ok(msg) = self.rx.try_recv() {
match msg {
ProcessMessage::Ok(id, data_val) => {
if let Some(t) = self.tasks.get_mut(&id) {
if !t.done {
t.done = true;
let _ = apply_one_function_value(interp, t.callback.clone(), data_val);
}
}
}
ProcessMessage::Err(id, err_str) => {
if let Some(t) = self.tasks.get_mut(&id) {
if !t.done {
t.done = true;
let mut map = IndexMap::new();
map.insert("error".to_string(), Value::SingleString(err_str));
let final_val = Value::KeyedArray(map);
let _ = apply_one_function_value(interp, t.callback.clone(), final_val);
}
}
}
}
}
let before = self.tasks.len();
self.tasks.retain(|_, t| !t.done);
let removed = before.saturating_sub(self.tasks.len());
if removed > 0 {
self.active_count.fetch_sub(removed, Ordering::SeqCst);
}
}
fn count_tasks(&self) -> usize {
self.active_count.load(Ordering::SeqCst)
}
}
#[allow(dead_code)]
fn process_spawn_bridge(
_interp: &mut Interpreter,
_args: Vec<Value>
) -> Result<Value, String> {
Ok(Value::Bool(true))
}
fn process_info_bridge(
interp: &mut Interpreter,
mut args: Vec<Value>
) -> Result<Value, String> {
if args.len() != 1 {
return Err(format!("process:info => expected 1 argument => callback, got {}", args.len()));
}
let callback_val = args.remove(0);
let callback_func = match callback_val {
Value::Function(fb) => fb,
other => return Err(format!("process:info => first arg must be function, got {:?}", other)),
};
let mut map = IndexMap::new();
let pid_u32 = std::process::id();
map.insert("pid".to_string(), Value::Int(pid_u32 as i32));
map.insert("uid".to_string(), Value::Int(get_uid()));
map.insert("username".to_string(), Value::SingleString(whoami::username()));
map.insert("binary_name".to_string(), Value::SingleString("mumu".into()));
map.insert("event_loop_len".to_string(), Value::Int(42));
let info_val = Value::KeyedArray(map);
let _ = apply_one_function_value(interp, callback_func, info_val)?;
Ok(Value::Bool(true))
}
fn process_check_tasks_bridge(
interp: &mut Interpreter,
_args: Vec<Value>
) -> Result<Value, String> {
let mut mgr = PROCESS_MANAGER.lock().unwrap();
mgr.poll_events(interp);
let count = mgr.count_tasks();
Ok(Value::Int(count as i32))
}
lazy_static! {
static ref PROCESS_MANAGER: Mutex<ProcessManager> = Mutex::new(ProcessManager::new());
}
#[export_name = "Cargo_lock"]
pub unsafe extern "C" fn cargo_lock(
interp_ptr: *mut c_void,
_extra_str: *const c_void,
) -> i32 {
if interp_ptr.is_null() {
return 1;
}
let interp_ref = &mut *(interp_ptr as *mut Interpreter);
let info_fn = Arc::new(Mutex::new(process_info_bridge));
interp_ref.register_dynamic_function("process:info", info_fn);
interp_ref.set_variable(
"process:info",
Value::Function(Box::new(FunctionValue::Named("process:info".to_string())))
);
let check_fn = Arc::new(Mutex::new(process_check_tasks_bridge));
interp_ref.register_dynamic_function("process:check_tasks", check_fn);
interp_ref.set_variable(
"process:check_tasks",
Value::Function(Box::new(FunctionValue::Named("process:check_tasks".to_string())))
);
let poller = Arc::new(Mutex::new(move |interp: &mut Interpreter| {
let mut mgr = PROCESS_MANAGER.lock().unwrap();
mgr.poll_events(interp);
mgr.count_tasks()
}));
interp_ref.add_poller(poller);
0
}