Skip to main content

mumusys/
lib.rs

1// sys/src/lib.rs
2//
3// sys plugin for MuMu/Lava:
4//   - sys:command(cmdString, callback)
5//   - sys:timestamp_ms()
6//   - sys:timestamp_micro()
7//   - sys:command_stream(cmdString)  [NEW: streaming output as InkIterator]
8
9use mumu::{
10    parser::interpreter::{Interpreter, apply_one_function_value},
11    parser::types::{FunctionValue, Value, InkIteratorHandle, InkIteratorKind, LavaStream},
12};
13use std::ffi::c_void;
14
15use std::{
16    collections::HashMap,
17    process::{Command as SystemCommand, Stdio},
18    sync::{
19        atomic::{AtomicUsize, Ordering},
20        mpsc::{channel, Receiver, Sender, RecvTimeoutError},
21        Arc, Mutex,
22    },
23    thread,
24    time::{SystemTime, UNIX_EPOCH, Duration},
25    io::{BufRead, BufReader},
26};
27use lazy_static::lazy_static;
28use indexmap::IndexMap;
29
30lazy_static! {
31    static ref COMMAND_MANAGER: Mutex<CommandManager> = Mutex::new(CommandManager::new());
32}
33
34#[derive(Debug)]
35enum CommandMessage {
36    Ok(usize, Value),
37    Err(usize, String),
38}
39
40struct CommandTask {
41    callback: Box<FunctionValue>,
42    done: bool,
43}
44
45struct CommandManager {
46    next_id: usize,
47    tasks: HashMap<usize, CommandTask>,
48    tx: Sender<CommandMessage>,
49    rx: Receiver<CommandMessage>,
50    active_count: AtomicUsize,
51}
52
53impl CommandManager {
54    fn new() -> Self {
55        let (tx, rx) = channel();
56        Self {
57            next_id: 0,
58            tasks: HashMap::new(),
59            tx,
60            rx,
61            active_count: AtomicUsize::new(0),
62        }
63    }
64
65    fn add_command_task(&mut self, cmd_str: String, callback: Box<FunctionValue>) {
66        let task_id = self.next_id;
67        self.next_id += 1;
68
69        self.tasks.insert(
70            task_id,
71            CommandTask {
72                callback,
73                done: false,
74            },
75        );
76        self.active_count.fetch_add(1, Ordering::SeqCst);
77
78        let txc = self.tx.clone();
79        thread::spawn(move || {
80            let parts: Vec<&str> = cmd_str.split_whitespace().collect();
81            if parts.is_empty() {
82                let _ = txc.send(CommandMessage::Err(task_id, "Empty command".into()));
83                return;
84            }
85            let prog = parts[0];
86            let prog_args = &parts[1..];
87
88            match SystemCommand::new(prog).args(prog_args).output() {
89                Ok(out) => {
90                    let code_i32 = out.status.code().unwrap_or(-1);
91                    let success_bool = code_i32 == 0;
92
93                    let stdout_str = String::from_utf8_lossy(&out.stdout).into_owned();
94                    let stderr_str = String::from_utf8_lossy(&out.stderr).into_owned();
95
96                    let mut map = IndexMap::new();
97                    map.insert("command".into(), Value::SingleString(cmd_str));
98                    map.insert("stdout".into(), Value::SingleString(stdout_str));
99                    map.insert("stderr".into(), Value::SingleString(stderr_str));
100                    map.insert("exit".into(), Value::Int(code_i32));
101                    map.insert("success".into(), Value::Bool(success_bool));
102
103                    let val = Value::KeyedArray(map);
104                    let _ = txc.send(CommandMessage::Ok(task_id, val));
105                }
106                Err(e) => {
107                    let _ = txc.send(CommandMessage::Err(task_id, format!("Command error: {e}")));
108                }
109            }
110        });
111    }
112
113    fn poll_events(&mut self, interp: &mut Interpreter) {
114        while let Ok(msg) = self.rx.try_recv() {
115            match msg {
116                CommandMessage::Ok(id, data_val) => {
117                    if let Some(t) = self.tasks.get_mut(&id) {
118                        if !t.done {
119                            t.done = true;
120                            let _ = apply_one_function_value(interp, t.callback.clone(), data_val);
121                        }
122                    }
123                }
124                CommandMessage::Err(id, err_str) => {
125                    if let Some(t) = self.tasks.get_mut(&id) {
126                        if !t.done {
127                            t.done = true;
128                            let mut map = IndexMap::new();
129                            map.insert("command".into(), Value::SingleString("<unknown>".into()));
130                            map.insert("stdout".into(), Value::SingleString(String::new()));
131                            map.insert("stderr".into(), Value::SingleString(err_str));
132                            map.insert("exit".into(), Value::Int(1));
133                            map.insert("success".into(), Value::Bool(false));
134                            let val = Value::KeyedArray(map);
135
136                            let _ = apply_one_function_value(interp, t.callback.clone(), val);
137                        }
138                    }
139                }
140            }
141        }
142
143        let before = self.tasks.len();
144        self.tasks.retain(|_, task| !task.done);
145        let removed = before.saturating_sub(self.tasks.len());
146        if removed > 0 {
147            self.active_count.fetch_sub(removed, Ordering::SeqCst);
148        }
149    }
150
151    fn count_active_tasks(&self) -> usize {
152        self.active_count.load(Ordering::SeqCst)
153    }
154}
155
156/* ─────────────────────────── bridge functions ─────────────────────────── */
157
158fn sys_command_bridge(
159    _interp: &mut Interpreter,
160    mut args: Vec<Value>,
161) -> Result<Value, String> {
162    if args.len() != 2 {
163        return Err(format!(
164            "sys:command expects (cmdString, callback); got {} arg(s)",
165            args.len()
166        ));
167    }
168
169    let cmd_str = match args.remove(0) {
170        Value::SingleString(s) => s,
171        Value::StrArray(ss) if ss.len() == 1 => ss[0].clone(),
172        _ => return Err("sys:command first arg must be a single string".into()),
173    };
174
175    let cb_func = match args.remove(0) {
176        Value::Function(fb) => fb,
177        _ => return Err("sys:command second arg must be a function".into()),
178    };
179
180    COMMAND_MANAGER.lock().unwrap().add_command_task(cmd_str, cb_func);
181
182    Ok(Value::Bool(true))
183}
184
185fn sys_timestamp_ms_bridge(
186    _interp: &mut Interpreter,
187    args: Vec<Value>,
188) -> Result<Value, String> {
189    if !args.is_empty() {
190        return Err(format!(
191            "sys:timestamp_ms expects 0 arguments, got {}",
192            args.len()
193        ));
194    }
195    let now = SystemTime::now();
196    let ms = now
197        .duration_since(UNIX_EPOCH)
198        .map_err(|e| format!("SystemTime error: {e}"))?
199        .as_millis();
200    Ok(Value::Long(ms as i64))
201}
202
203fn sys_timestamp_micro_bridge(
204    _interp: &mut Interpreter,
205    args: Vec<Value>,
206) -> Result<Value, String> {
207    if !args.is_empty() {
208        return Err(format!(
209            "sys:timestamp_micro expects 0 arguments, got {}",
210            args.len()
211        ));
212    }
213    let now = SystemTime::now();
214    let us = now
215        .duration_since(UNIX_EPOCH)
216        .map_err(|e| format!("SystemTime error: {e}"))?
217        .as_micros();
218    Ok(Value::Long(us as i64))
219}
220
221// -------------- NEW: sys:command_stream(cmdString) --------------
222
223fn sys_command_stream_bridge(
224    _interp: &mut Interpreter,
225    mut args: Vec<Value>,
226) -> Result<Value, String> {
227    if args.is_empty() {
228        return Err("sys:command_stream expects at least 1 argument".into());
229    }
230
231    let cmd_str = match args.remove(0) {
232        Value::SingleString(s) => s,
233        Value::StrArray(ss) if ss.len() == 1 => ss[0].clone(),
234        _ => return Err("sys:command_stream first arg must be a single string".into()),
235    };
236
237    // Use system shell for shell features, or split by whitespace for bare commands
238    // Here, for full generality, we run: sh -c "cmd_str"
239    let (tx, rx): (Sender<String>, Receiver<String>) = channel();
240
241    thread::spawn(move || {
242        let mut child = match SystemCommand::new("sh")
243            .arg("-c")
244            .arg(&cmd_str)
245            .stdout(Stdio::piped())
246            .spawn()
247        {
248            Ok(c) => c,
249            Err(e) => {
250                let _ = tx.send(format!("ERROR: failed to start: {e}"));
251                return;
252            }
253        };
254
255        if let Some(stdout) = child.stdout.take() {
256            let reader = BufReader::new(stdout);
257            for line in reader.lines() {
258                match line {
259                    Ok(l) => {
260                        if tx.send(l).is_err() {
261                            break;
262                        }
263                    }
264                    Err(e) => {
265                        let _ = tx.send(format!("ERROR: {e}"));
266                        break;
267                    }
268                }
269            }
270        }
271        let _ = child.wait();
272    });
273
274    let state = Arc::new(Mutex::new(CommandStreamState { rx, done: false }));
275
276    let handle = InkIteratorHandle {
277        kind: InkIteratorKind::Plugin(Arc::new(Mutex::new(Box::new(CommandStreamPlugin { state: state.clone() })))),
278    };
279
280    Ok(Value::InkIterator(handle))
281}
282
283struct CommandStreamState {
284    rx: Receiver<String>,
285    done: bool,
286}
287
288struct CommandStreamPlugin {
289    state: Arc<Mutex<CommandStreamState>>,
290}
291
292impl LavaStream for CommandStreamPlugin {
293    fn next_value(&mut self) -> Result<Value, String> {
294        let mut state = self.state.lock().unwrap();
295        if state.done {
296            return Err("NO_MORE_DATA".to_string());
297        }
298        match state.rx.recv_timeout(Duration::from_millis(50)) {
299            Ok(line) => Ok(Value::SingleString(line)),
300            Err(RecvTimeoutError::Timeout) => Err("NO_MORE_DATA".to_string()),
301            Err(_) => {
302                state.done = true;
303                Err("NO_MORE_DATA".to_string())
304            }
305        }
306    }
307}
308
309/* ───────────────────────── plugin entry point ───────────────────────── */
310
311#[export_name = "Cargo_lock"]
312pub unsafe extern "C" fn cargo_lock(
313    interp_ptr: *mut c_void,
314    _extra_str: *const c_void,
315) -> i32 {
316    if interp_ptr.is_null() {
317        return 1;
318    }
319    let interp = &mut *(interp_ptr as *mut Interpreter);
320
321    // sys:command
322    let cmd_fn = Arc::new(Mutex::new(sys_command_bridge));
323    interp.register_dynamic_function("sys:command", cmd_fn);
324    interp.set_variable(
325        "sys:command",
326        Value::Function(Box::new(FunctionValue::Named("sys:command".into()))),
327    );
328
329    // sys:timestamp_ms
330    let ms_fn = Arc::new(Mutex::new(sys_timestamp_ms_bridge));
331    interp.register_dynamic_function("sys:timestamp_ms", ms_fn);
332    interp.set_variable(
333        "sys:timestamp_ms",
334        Value::Function(Box::new(FunctionValue::Named("sys:timestamp_ms".into()))),
335    );
336
337    // sys:timestamp_micro
338    let micro_fn = Arc::new(Mutex::new(sys_timestamp_micro_bridge));
339    interp.register_dynamic_function("sys:timestamp_micro", micro_fn);
340    interp.set_variable(
341        "sys:timestamp_micro",
342        Value::Function(Box::new(FunctionValue::Named("sys:timestamp_micro".into()))),
343    );
344
345    // sys:command_stream
346    let stream_fn = Arc::new(Mutex::new(sys_command_stream_bridge));
347    interp.register_dynamic_function("sys:command_stream", stream_fn);
348    interp.set_variable(
349        "sys:command_stream",
350        Value::Function(Box::new(FunctionValue::Named("sys:command_stream".into()))),
351    );
352
353    // background poller (for sys:command)
354    let poller = Arc::new(Mutex::new(|interp: &mut Interpreter| {
355        let mut mgr = COMMAND_MANAGER.lock().unwrap();
356        mgr.poll_events(interp);
357        mgr.count_active_tasks()
358    }));
359    interp.add_poller(poller);
360
361    0
362}