1use 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
156fn 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
221fn 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 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#[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 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 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 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 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 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}