Skip to main content

stryke/
remote_wire.rs

1//! Framed bincode over stdin/stdout for `stryke --remote-worker` (distributed `pmap_on`).
2//!
3//! ## Wire protocol
4//!
5//! Every message is a length-prefixed frame: `[u64 LE length][u8 kind][bincode payload]`.
6//! The single-byte `kind` discriminator lets future revisions add message types without
7//! breaking older workers — an unknown kind is a hard error so version skew is loud.
8//!
9//! ### Message flow (v3 — persistent session)
10//!
11//! ```text
12//! dispatcher                    worker
13//!     │                            │
14//!     │── HELLO ─────────────────►│   (proto version, build id)
15//!     │◄───────────── HELLO_ACK ──│   (worker stryke version, hostname)
16//!     │── SESSION_INIT ──────────►│   (subs prelude, block source, captured lexicals)
17//!     │◄────────── SESSION_ACK ───│   (or ERROR)
18//!     │── JOB(seq=0) ────────────►│   (item)
19//!     │◄────────── JOB_RESP(0) ───│
20//!     │── JOB(seq=1) ────────────►│
21//!     │◄────────── JOB_RESP(1) ───│
22//!     │           ...             │
23//!     │── SHUTDOWN ──────────────►│
24//!     │                            └─ exit 0
25//! ```
26//!
27//! Why this beats the basic v1 protocol: subs prelude + block source ship **once** per
28//! session instead of once per item, the parser+compiler runs once per worker instead of
29//! once per job, and one ssh handshake amortizes across the whole map.
30//!
31//! Dynamic [`serde_json::Value`] fields are embedded as JSON UTF-8 bytes inside the bincode
32//! envelope (v3+). Bincode cannot deserialize `Value` directly (`deserialize_any`); nested
33//! JSON keeps the on-wire type self-describing.
34
35use std::collections::HashMap;
36use std::io::{Read, Write};
37use std::process::{Command, Stdio};
38use std::sync::Arc;
39
40use serde::{Deserialize, Serialize};
41
42use crate::ast::Block;
43use crate::value::{StrykeSub, StrykeValue};
44use crate::vm_helper::{FlowOrError, VMHelper};
45
46/// Frame-kind discriminator. Stored as the first byte of every wire payload after the
47/// length prefix. Sub-byte values are reserved (anything outside the documented set is
48/// rejected with a clean error rather than silently misparsed).
49#[allow(dead_code)]
50pub mod frame_kind {
51    /// `HELLO` constant.
52    pub const HELLO: u8 = 0x01;
53    /// `HELLO_ACK` constant.
54    pub const HELLO_ACK: u8 = 0x02;
55    /// `SESSION_INIT` constant.
56    pub const SESSION_INIT: u8 = 0x03;
57    /// `SESSION_ACK` constant.
58    pub const SESSION_ACK: u8 = 0x04;
59    /// `JOB` constant.
60    pub const JOB: u8 = 0x05;
61    /// `JOB_RESP` constant.
62    pub const JOB_RESP: u8 = 0x06;
63    /// `SHUTDOWN` constant.
64    pub const SHUTDOWN: u8 = 0x07;
65    /// `ERROR` constant.
66    pub const ERROR: u8 = 0xFF;
67}
68
69/// Wire protocol version. Bumped whenever the layout of an existing message changes in a
70/// backwards-incompatible way. The HELLO handshake fails fast on version mismatch so
71/// dispatcher and worker never silently disagree on layout.
72pub const PROTO_VERSION: u32 = 3;
73
74mod json_value_bincode {
75    use serde::{Deserialize, Deserializer, Serialize, Serializer};
76    /// `serialize` — see implementation.
77    pub fn serialize<S>(value: &serde_json::Value, serializer: S) -> Result<S::Ok, S::Error>
78    where
79        S: Serializer,
80    {
81        let buf = serde_json::to_vec(value).map_err(serde::ser::Error::custom)?;
82        buf.serialize(serializer)
83    }
84    /// `deserialize` — see implementation.
85    pub fn deserialize<'de, D>(deserializer: D) -> Result<serde_json::Value, D::Error>
86    where
87        D: Deserializer<'de>,
88    {
89        let buf: Vec<u8> = Vec::deserialize(deserializer)?;
90        serde_json::from_slice(&buf).map_err(serde::de::Error::custom)
91    }
92}
93
94mod capture_json_bincode {
95    use serde::{de::Deserializer, ser::SerializeSeq, Deserialize, Serializer};
96    /// `serialize` — see implementation.
97    pub fn serialize<S>(v: &[(String, serde_json::Value)], serializer: S) -> Result<S::Ok, S::Error>
98    where
99        S: Serializer,
100    {
101        let mut seq = serializer.serialize_seq(Some(v.len()))?;
102        for (k, val) in v {
103            let enc = serde_json::to_vec(val).map_err(serde::ser::Error::custom)?;
104            seq.serialize_element(&(k, enc))?;
105        }
106        seq.end()
107    }
108    /// `deserialize` — see implementation.
109    pub fn deserialize<'de, D>(
110        deserializer: D,
111    ) -> Result<Vec<(String, serde_json::Value)>, D::Error>
112    where
113        D: Deserializer<'de>,
114    {
115        let raw: Vec<(String, Vec<u8>)> = Vec::deserialize(deserializer)?;
116        let mut out = Vec::with_capacity(raw.len());
117        for (k, enc) in raw {
118            let val = serde_json::from_slice(&enc).map_err(serde::de::Error::custom)?;
119            out.push((k, val));
120        }
121        Ok(out)
122    }
123}
124/// `HelloMsg` — see fields for layout.
125#[derive(Debug, Clone, Serialize, Deserialize)]
126pub struct HelloMsg {
127    /// `proto_version` field.
128    pub proto_version: u32,
129    /// `pe_version` field.
130    pub pe_version: String,
131}
132/// `HelloAck` — see fields for layout.
133#[derive(Debug, Clone, Serialize, Deserialize)]
134pub struct HelloAck {
135    /// `proto_version` field.
136    pub proto_version: u32,
137    /// `pe_version` field.
138    pub pe_version: String,
139    /// `hostname` field.
140    pub hostname: String,
141}
142
143/// Sent **once** per worker session. Carries everything that doesn't change between jobs:
144/// the user's named subs, the `pmap_on` block source, and the captured-lexical snapshot.
145#[derive(Debug, Clone, Serialize, Deserialize)]
146pub struct SessionInit {
147    /// `subs_prelude` field.
148    pub subs_prelude: String,
149    /// `block_src` field.
150    pub block_src: String,
151    /// `capture` field.
152    #[serde(with = "capture_json_bincode")]
153    pub capture: Vec<(String, serde_json::Value)>,
154}
155/// `SessionAck` — see fields for layout.
156#[derive(Debug, Clone, Serialize, Deserialize)]
157pub struct SessionAck {
158    /// `ok` field.
159    pub ok: bool,
160    /// `err_msg` field.
161    pub err_msg: String,
162}
163/// `JobMsg` — see fields for layout.
164#[derive(Debug, Clone, Serialize, Deserialize)]
165pub struct JobMsg {
166    /// `seq` field.
167    pub seq: u64,
168    /// `item` field.
169    #[serde(with = "json_value_bincode")]
170    pub item: serde_json::Value,
171}
172/// `JobRespMsg` — see fields for layout.
173#[derive(Debug, Clone, Serialize, Deserialize)]
174pub struct JobRespMsg {
175    /// `seq` field.
176    pub seq: u64,
177    /// `ok` field.
178    pub ok: bool,
179    /// `result` field.
180    #[serde(with = "json_value_bincode")]
181    pub result: serde_json::Value,
182    /// `err_msg` field.
183    pub err_msg: String,
184}
185
186/// Read a typed frame: returns `(kind, body)` where `body` is the bincode payload after
187/// the kind byte. Caller decides how to interpret based on `kind`.
188pub fn read_typed_frame<R: Read>(r: &mut R) -> std::io::Result<(u8, Vec<u8>)> {
189    let raw = read_framed(r)?;
190    if raw.is_empty() {
191        return Err(std::io::Error::new(
192            std::io::ErrorKind::InvalidData,
193            "remote frame: empty payload (missing kind byte)",
194        ));
195    }
196    let kind = raw[0];
197    Ok((kind, raw[1..].to_vec()))
198}
199
200/// Write a typed frame: prepends the `kind` byte to `payload` and writes one length-prefixed
201/// frame.
202pub fn write_typed_frame<W: Write>(w: &mut W, kind: u8, payload: &[u8]) -> std::io::Result<()> {
203    let mut framed = Vec::with_capacity(payload.len() + 1);
204    framed.push(kind);
205    framed.extend_from_slice(payload);
206    write_framed(w, &framed)
207}
208
209/// Bincode + write helper. The two-step `bincode::serialize` + `write_typed_frame` pattern
210/// is the same in every send site so it lives here once.
211pub fn send_msg<W: Write, T: Serialize>(w: &mut W, kind: u8, msg: &T) -> Result<(), String> {
212    let payload = bincode::serialize(msg).map_err(|e| format!("bincode encode: {e}"))?;
213    write_typed_frame(w, kind, &payload).map_err(|e| format!("write frame: {e}"))
214}
215
216/// Bincode + read helper. Returns the deserialized message and verifies the kind matches.
217pub fn recv_msg<R: Read, T: for<'de> Deserialize<'de>>(
218    r: &mut R,
219    expected_kind: u8,
220) -> Result<T, String> {
221    let (kind, body) = read_typed_frame(r).map_err(|e| format!("read frame: {e}"))?;
222    if kind != expected_kind {
223        return Err(format!(
224            "wire: expected frame kind {:#04x}, got {:#04x}",
225            expected_kind, kind
226        ));
227    }
228    bincode::deserialize(&body).map_err(|e| format!("bincode decode: {e}"))
229}
230
231/// One unit of work executed on a remote `stryke --remote-worker`.
232#[derive(Debug, Clone, Serialize, Deserialize)]
233pub struct RemoteJobV1 {
234    /// `seq` field.
235    pub seq: u64,
236    /// `subs_prelude` field.
237    pub subs_prelude: String,
238    /// `block_src` field.
239    pub block_src: String,
240    /// `capture` field.
241    #[serde(with = "capture_json_bincode")]
242    pub capture: Vec<(String, serde_json::Value)>,
243    /// `item` field.
244    #[serde(with = "json_value_bincode")]
245    pub item: serde_json::Value,
246}
247/// `RemoteRespV1` — see fields for layout.
248#[derive(Debug, Clone, Serialize, Deserialize)]
249pub struct RemoteRespV1 {
250    /// `seq` field.
251    pub seq: u64,
252    /// `ok` field.
253    pub ok: bool,
254    /// `result` field.
255    #[serde(with = "json_value_bincode")]
256    pub result: serde_json::Value,
257    /// `err_msg` field.
258    pub err_msg: String,
259}
260
261const MAX_FRAME: usize = 256 * 1024 * 1024;
262/// `write_framed` — see implementation.
263pub fn write_framed<W: Write>(w: &mut W, payload: &[u8]) -> std::io::Result<()> {
264    w.write_all(&(payload.len() as u64).to_le_bytes())?;
265    w.write_all(payload)?;
266    w.flush()?;
267    Ok(())
268}
269/// `read_framed` — see implementation.
270pub fn read_framed<R: Read>(r: &mut R) -> std::io::Result<Vec<u8>> {
271    let mut h = [0u8; 8];
272    r.read_exact(&mut h)?;
273    let n = u64::from_le_bytes(h) as usize;
274    if n > MAX_FRAME {
275        return Err(std::io::Error::new(
276            std::io::ErrorKind::InvalidData,
277            format!("remote frame too large: {n}"),
278        ));
279    }
280    let mut v = vec![0u8; n];
281    r.read_exact(&mut v)?;
282    Ok(v)
283}
284/// `encode_job` — see implementation.
285pub fn encode_job(job: &RemoteJobV1) -> Result<Vec<u8>, String> {
286    bincode::serialize(job).map_err(|e| e.to_string())
287}
288/// `decode_job` — see implementation.
289pub fn decode_job(bytes: &[u8]) -> Result<RemoteJobV1, String> {
290    bincode::deserialize(bytes).map_err(|e| e.to_string())
291}
292/// `encode_resp` — see implementation.
293pub fn encode_resp(resp: &RemoteRespV1) -> Result<Vec<u8>, String> {
294    bincode::serialize(resp).map_err(|e| e.to_string())
295}
296/// `decode_resp` — see implementation.
297pub fn decode_resp(bytes: &[u8]) -> Result<RemoteRespV1, String> {
298    bincode::deserialize(bytes).map_err(|e| e.to_string())
299}
300/// `perl_to_json_value` — see implementation.
301pub fn perl_to_json_value(v: &StrykeValue) -> Result<serde_json::Value, String> {
302    if v.is_undef() {
303        return Ok(serde_json::Value::Null);
304    }
305    if let Some(i) = v.as_integer() {
306        return Ok(serde_json::json!(i));
307    }
308    if let Some(f) = v.as_float() {
309        return Ok(serde_json::json!(f));
310    }
311    if v.is_string_like() {
312        return Ok(serde_json::Value::String(v.to_string()));
313    }
314    if let Some(a) = v.as_array_vec() {
315        let mut out = Vec::with_capacity(a.len());
316        for x in a {
317            out.push(perl_to_json_value(&x)?);
318        }
319        return Ok(serde_json::Value::Array(out));
320    }
321    // Arrayref / hashref carry the same shape as flat array / hash for
322    // JSON — there's no ref/value distinction over the wire. Without this
323    // branch a stage block that ends in `[ ... ]` (used by `~d>` to keep
324    // list shape across the worker's scalar-return boundary) would fail
325    // with "value not supported for remote pmap".
326    if let Some(ar) = v.as_array_ref() {
327        let guard = ar.read();
328        let mut out = Vec::with_capacity(guard.len());
329        for x in guard.iter() {
330            out.push(perl_to_json_value(x)?);
331        }
332        return Ok(serde_json::Value::Array(out));
333    }
334    if let Some(h) = v.as_hash_map() {
335        let mut m = serde_json::Map::new();
336        for (k, val) in h {
337            m.insert(k.clone(), perl_to_json_value(&val)?);
338        }
339        return Ok(serde_json::Value::Object(m));
340    }
341    if let Some(hr) = v.as_hash_ref() {
342        let guard = hr.read();
343        let mut m = serde_json::Map::new();
344        for (k, val) in guard.iter() {
345            m.insert(k.clone(), perl_to_json_value(val)?);
346        }
347        return Ok(serde_json::Value::Object(m));
348    }
349    Err(format!(
350        "value not supported for remote pmap (need null, bool/int/float/string/array/hash): {}",
351        v.type_name()
352    ))
353}
354/// `json_to_perl` — see implementation.
355pub fn json_to_perl(v: &serde_json::Value) -> Result<StrykeValue, String> {
356    Ok(match v {
357        serde_json::Value::Null => StrykeValue::UNDEF,
358        serde_json::Value::Bool(b) => StrykeValue::integer(if *b { 1 } else { 0 }),
359        serde_json::Value::Number(n) => {
360            if let Some(i) = n.as_i64() {
361                StrykeValue::integer(i)
362            } else if let Some(u) = n.as_u64() {
363                StrykeValue::integer(u as i64)
364            } else {
365                StrykeValue::float(n.as_f64().unwrap_or(0.0))
366            }
367        }
368        serde_json::Value::String(s) => StrykeValue::string(s.clone()),
369        serde_json::Value::Array(a) => {
370            let mut items = Vec::with_capacity(a.len());
371            for x in a {
372                items.push(json_to_perl(x)?);
373            }
374            StrykeValue::array(items)
375        }
376        serde_json::Value::Object(o) => {
377            let mut map = indexmap::IndexMap::new();
378            for (k, val) in o {
379                map.insert(k.clone(), json_to_perl(val)?);
380            }
381            StrykeValue::hash(map)
382        }
383    })
384}
385/// `capture_entries_to_json` — see implementation.
386pub fn capture_entries_to_json(
387    entries: &[(String, StrykeValue)],
388) -> Result<Vec<(String, serde_json::Value)>, String> {
389    let mut out = Vec::with_capacity(entries.len());
390    for (k, v) in entries {
391        out.push((k.clone(), perl_to_json_value(v)?));
392    }
393    Ok(out)
394}
395/// `build_subs_prelude` — see implementation.
396pub fn build_subs_prelude(subs: &HashMap<String, Arc<StrykeSub>>) -> String {
397    let mut names: Vec<_> = subs.keys().cloned().collect();
398    names.sort();
399    let mut s = String::new();
400    for name in names {
401        let sub = &subs[&name];
402        if sub.closure_env.is_some() {
403            continue;
404        }
405        let sig = if !sub.params.is_empty() {
406            format!(
407                " ({})",
408                sub.params
409                    .iter()
410                    .map(crate::fmt::format_sub_sig_param)
411                    .collect::<Vec<_>>()
412                    .join(", ")
413            )
414        } else if let Some(ref p) = sub.prototype {
415            format!(" ({})", p)
416        } else {
417            String::new()
418        };
419        let body = crate::fmt::format_block(&sub.body);
420        s.push_str(&format!("fn {}{} {{\n{}\n}}\n", name, sig, body));
421    }
422    s
423}
424
425/// Run one job in-process (for tests / local debugging).
426pub fn run_job_local(job: &RemoteJobV1) -> RemoteRespV1 {
427    let mut interp = VMHelper::new();
428    let cap: Vec<(String, StrykeValue)> = match job
429        .capture
430        .iter()
431        .map(|(k, v)| json_to_perl(v).map(|pv| (k.clone(), pv)))
432        .collect()
433    {
434        Ok(c) => c,
435        Err(e) => {
436            return RemoteRespV1 {
437                seq: job.seq,
438                ok: false,
439                result: serde_json::Value::Null,
440                err_msg: e,
441            };
442        }
443    };
444    interp.scope_push_hook();
445    interp.scope.restore_capture(&cap);
446    let item_pv = match json_to_perl(&job.item) {
447        Ok(v) => v,
448        Err(e) => {
449            interp.scope_pop_hook();
450            return RemoteRespV1 {
451                seq: job.seq,
452                ok: false,
453                result: serde_json::Value::Null,
454                err_msg: e,
455            };
456        }
457    };
458    interp.scope.set_topic(item_pv);
459    let full_src = format!("{}\n{}", job.subs_prelude, job.block_src);
460    let prog = match crate::parse(&full_src) {
461        Ok(p) => p,
462        Err(e) => {
463            interp.scope_pop_hook();
464            return RemoteRespV1 {
465                seq: job.seq,
466                ok: false,
467                result: serde_json::Value::Null,
468                err_msg: e.message,
469            };
470        }
471    };
472    let block: Block = prog.statements;
473    let r = match interp.exec_block_smart(&block) {
474        Ok(v) => v,
475        Err(e) => {
476            interp.scope_pop_hook();
477            let msg = match e {
478                FlowOrError::Error(stryke) => stryke.to_string(),
479                FlowOrError::Flow(f) => format!("unexpected control flow: {:?}", f),
480            };
481            return RemoteRespV1 {
482                seq: job.seq,
483                ok: false,
484                result: serde_json::Value::Null,
485                err_msg: msg,
486            };
487        }
488    };
489    interp.scope_pop_hook();
490    match perl_to_json_value(&r) {
491        Ok(j) => RemoteRespV1 {
492            seq: job.seq,
493            ok: true,
494            result: j,
495            err_msg: String::new(),
496        },
497        Err(e) => RemoteRespV1 {
498            seq: job.seq,
499            ok: false,
500            result: serde_json::Value::Null,
501            err_msg: e,
502        },
503    }
504}
505
506/// Persistent v3 worker session: handles many jobs over a single stdin/stdout pair, with
507/// one Interpreter and one parsed block shared across the whole session.
508///
509/// Protocol order: HELLO → HELLO_ACK → SESSION_INIT → SESSION_ACK → JOB / JOB_RESP loop
510/// → SHUTDOWN → exit. Any wire error or unknown frame kind causes a clean non-zero exit so
511/// the dispatcher can re-route in-flight jobs to a different slot.
512///
513/// Why this beats the basic v1 [`run_remote_worker_stdio`]: subs prelude + block source
514/// ship **once** per session instead of per-item, parser+compiler runs once per worker,
515/// and one ssh handshake amortizes across the whole map.
516pub fn run_remote_worker_session() -> i32 {
517    let stdin = std::io::stdin();
518    let mut stdin = stdin.lock();
519    let mut stdout = std::io::stdout();
520
521    // 1. HELLO handshake. Dispatcher sends first; we reply with our build info.
522    let hello: HelloMsg = match recv_msg(&mut stdin, frame_kind::HELLO) {
523        Ok(h) => h,
524        Err(e) => {
525            let _ = writeln!(std::io::stderr(), "remote-worker: hello: {e}");
526            return 1;
527        }
528    };
529    if hello.proto_version != PROTO_VERSION {
530        let _ = writeln!(
531            std::io::stderr(),
532            "remote-worker: proto version mismatch (dispatcher {} vs worker {})",
533            hello.proto_version,
534            PROTO_VERSION
535        );
536        return 1;
537    }
538    let ack = HelloAck {
539        proto_version: PROTO_VERSION,
540        pe_version: env!("CARGO_PKG_VERSION").to_string(),
541        hostname: hostname_or_unknown(),
542    };
543    if let Err(e) = send_msg(&mut stdout, frame_kind::HELLO_ACK, &ack) {
544        let _ = writeln!(std::io::stderr(), "remote-worker: hello ack: {e}");
545        return 1;
546    }
547
548    // 2. SESSION_INIT: subs prelude + block source + captured lexicals.
549    let init: SessionInit = match recv_msg(&mut stdin, frame_kind::SESSION_INIT) {
550        Ok(i) => i,
551        Err(e) => {
552            let _ = writeln!(std::io::stderr(), "remote-worker: session init: {e}");
553            return 1;
554        }
555    };
556
557    // Parse subs prelude ONCE so they're registered for every JOB; parse block ONCE so we
558    // can hand the same `Block` to `exec_block_smart` per item without re-parsing.
559    let mut interp = VMHelper::new();
560    let prelude_program = match crate::parse(&init.subs_prelude) {
561        Ok(p) => p,
562        Err(e) => {
563            let nack = SessionAck {
564                ok: false,
565                err_msg: format!("parse subs prelude: {}", e.message),
566            };
567            let _ = send_msg(&mut stdout, frame_kind::SESSION_ACK, &nack);
568            return 2;
569        }
570    };
571    let block_program = match crate::parse(&init.block_src) {
572        Ok(p) => p,
573        Err(e) => {
574            let nack = SessionAck {
575                ok: false,
576                err_msg: format!("parse block: {}", e.message),
577            };
578            let _ = send_msg(&mut stdout, frame_kind::SESSION_ACK, &nack);
579            return 2;
580        }
581    };
582
583    // Restore captured lexicals once per session — they don't change across jobs.
584    let cap_pv: Vec<(String, StrykeValue)> = match init
585        .capture
586        .iter()
587        .map(|(k, v)| json_to_perl(v).map(|pv| (k.clone(), pv)))
588        .collect()
589    {
590        Ok(c) => c,
591        Err(e) => {
592            let nack = SessionAck {
593                ok: false,
594                err_msg: format!("decode capture: {e}"),
595            };
596            let _ = send_msg(&mut stdout, frame_kind::SESSION_ACK, &nack);
597            return 2;
598        }
599    };
600    interp.scope_push_hook();
601    interp.scope.restore_capture(&cap_pv);
602
603    // Run the prelude (sub decls) once. After this every JOB has the user's named subs in
604    // scope without re-parsing or re-executing the prelude per item.
605    if let Err(e) = interp.execute(&prelude_program) {
606        let nack = SessionAck {
607            ok: false,
608            err_msg: format!("session prelude: {e}"),
609        };
610        let _ = send_msg(&mut stdout, frame_kind::SESSION_ACK, &nack);
611        return 2;
612    }
613
614    let ack = SessionAck {
615        ok: true,
616        err_msg: String::new(),
617    };
618    if let Err(e) = send_msg(&mut stdout, frame_kind::SESSION_ACK, &ack) {
619        let _ = writeln!(std::io::stderr(), "remote-worker: session ack: {e}");
620        return 1;
621    }
622
623    let block: Block = block_program.statements;
624
625    // 3. JOB loop. Each iteration sets `$_ = item`, re-evaluates the cached block, and
626    // sends back the result. The Interpreter is reused — sub registrations, package state,
627    // anything mutated by SESSION_INIT persists across jobs.
628    loop {
629        let (kind, body) = match read_typed_frame(&mut stdin) {
630            Ok(p) => p,
631            Err(e) if e.kind() == std::io::ErrorKind::UnexpectedEof => return 0,
632            Err(e) => {
633                let _ = writeln!(std::io::stderr(), "remote-worker: read job: {e}");
634                return 1;
635            }
636        };
637        match kind {
638            frame_kind::JOB => {
639                let job: JobMsg = match bincode::deserialize(&body) {
640                    Ok(j) => j,
641                    Err(e) => {
642                        let resp = JobRespMsg {
643                            seq: 0,
644                            ok: false,
645                            result: serde_json::Value::Null,
646                            err_msg: format!("decode job: {e}"),
647                        };
648                        let _ = send_msg(&mut stdout, frame_kind::JOB_RESP, &resp);
649                        continue;
650                    }
651                };
652                let resp = run_one_session_job(&mut interp, &block, &job);
653                if let Err(e) = send_msg(&mut stdout, frame_kind::JOB_RESP, &resp) {
654                    let _ = writeln!(std::io::stderr(), "remote-worker: write resp: {e}");
655                    return 1;
656                }
657            }
658            frame_kind::SHUTDOWN => return 0,
659            other => {
660                let _ = writeln!(
661                    std::io::stderr(),
662                    "remote-worker: unexpected frame kind {:#04x} in JOB loop",
663                    other
664                );
665                return 1;
666            }
667        }
668    }
669}
670
671/// Run one JOB inside an active session. Sets `$_` to the item, evaluates the cached block,
672/// returns the JSON-marshalled result. Preserves Interpreter state across jobs so anything
673/// the prelude installed (named subs, package vars) stays live.
674fn run_one_session_job(interp: &mut VMHelper, block: &Block, job: &JobMsg) -> JobRespMsg {
675    let item_pv = match json_to_perl(&job.item) {
676        Ok(v) => v,
677        Err(e) => {
678            return JobRespMsg {
679                seq: job.seq,
680                ok: false,
681                result: serde_json::Value::Null,
682                err_msg: e,
683            };
684        }
685    };
686    interp.scope.set_topic(item_pv);
687    let r = match interp.exec_block_smart(block) {
688        Ok(v) => v,
689        Err(FlowOrError::Error(stryke)) => {
690            return JobRespMsg {
691                seq: job.seq,
692                ok: false,
693                result: serde_json::Value::Null,
694                err_msg: stryke.to_string(),
695            };
696        }
697        Err(FlowOrError::Flow(f)) => {
698            return JobRespMsg {
699                seq: job.seq,
700                ok: false,
701                result: serde_json::Value::Null,
702                err_msg: format!("unexpected control flow: {:?}", f),
703            };
704        }
705    };
706    match perl_to_json_value(&r) {
707        Ok(j) => JobRespMsg {
708            seq: job.seq,
709            ok: true,
710            result: j,
711            err_msg: String::new(),
712        },
713        Err(e) => JobRespMsg {
714            seq: job.seq,
715            ok: false,
716            result: serde_json::Value::Null,
717            err_msg: e,
718        },
719    }
720}
721
722fn hostname_or_unknown() -> String {
723    std::env::var("HOSTNAME").unwrap_or_else(|_| {
724        std::process::Command::new("hostname")
725            .output()
726            .ok()
727            .and_then(|o| String::from_utf8(o.stdout).ok())
728            .map(|s| s.trim().to_string())
729            .unwrap_or_else(|| "unknown".to_string())
730    })
731}
732
733/// stdin/stdout worker loop: one framed request → one framed response, then exit 0.
734pub fn run_remote_worker_stdio() -> i32 {
735    let stdin = std::io::stdin();
736    let mut stdin = stdin.lock();
737    let mut stdout = std::io::stdout();
738    let payload = match read_framed(&mut stdin) {
739        Ok(p) => p,
740        Err(e) => {
741            let _ = writeln!(std::io::stderr(), "remote-worker: read frame: {e}");
742            return 1;
743        }
744    };
745    let job = match decode_job(&payload) {
746        Ok(j) => j,
747        Err(e) => {
748            let _ = writeln!(std::io::stderr(), "remote-worker: decode job: {e}");
749            return 1;
750        }
751    };
752    let resp = run_job_local(&job);
753    let out = match encode_resp(&resp) {
754        Ok(b) => b,
755        Err(e) => {
756            let _ = writeln!(std::io::stderr(), "remote-worker: encode resp: {e}");
757            return 1;
758        }
759    };
760    if let Err(e) = write_framed(&mut stdout, &out) {
761        let _ = writeln!(std::io::stderr(), "remote-worker: write frame: {e}");
762        return 1;
763    }
764    if resp.ok {
765        0
766    } else {
767        let _ = writeln!(std::io::stderr(), "remote-worker: {}", resp.err_msg);
768        2
769    }
770}
771/// `ssh_invoke_remote_worker` — see implementation.
772pub fn ssh_invoke_remote_worker(
773    host: &str,
774    pe_bin: &str,
775    job: &RemoteJobV1,
776) -> Result<RemoteRespV1, String> {
777    let payload = encode_job(job)?;
778    let mut child = Command::new("ssh")
779        .arg(host)
780        .arg(pe_bin)
781        .arg("--remote-worker")
782        .stdin(Stdio::piped())
783        .stdout(Stdio::piped())
784        .stderr(Stdio::piped())
785        .spawn()
786        .map_err(|e| format!("ssh: {e}"))?;
787    let mut stdin = child.stdin.take().ok_or_else(|| "ssh: stdin".to_string())?;
788    write_framed(&mut stdin, &payload).map_err(|e| format!("ssh stdin: {e}"))?;
789    drop(stdin);
790    let mut stdout = child
791        .stdout
792        .take()
793        .ok_or_else(|| "ssh: stdout".to_string())?;
794    let mut stderr = child
795        .stderr
796        .take()
797        .ok_or_else(|| "ssh: stderr".to_string())?;
798    let stderr_task = std::thread::spawn(move || {
799        let mut s = String::new();
800        let _ = stderr.read_to_string(&mut s);
801        s
802    });
803    let out_bytes = read_framed(&mut stdout).map_err(|e| format!("ssh read frame: {e}"))?;
804    let status = child.wait().map_err(|e| format!("ssh wait: {e}"))?;
805    let stderr_text = stderr_task.join().unwrap_or_default();
806    if !status.success() {
807        return Err(format!(
808            "ssh remote stryke exited {:?}: {}",
809            status.code(),
810            stderr_text.trim()
811        ));
812    }
813    decode_resp(&out_bytes).map_err(|e| {
814        format!(
815            "decode remote response: {e}; stderr: {}",
816            stderr_text.trim()
817        )
818    })
819}
820
821#[cfg(test)]
822mod tests {
823    use super::*;
824
825    #[test]
826    fn job_resp_msg_bincode_roundtrip() {
827        let msg = JobRespMsg {
828            seq: 1,
829            ok: true,
830            result: serde_json::json!(42i64),
831            err_msg: String::new(),
832        };
833        let bytes = bincode::serialize(&msg).unwrap();
834        let back: JobRespMsg = bincode::deserialize(&bytes).unwrap();
835        assert_eq!(back.seq, msg.seq);
836        assert_eq!(back.ok, msg.ok);
837        assert_eq!(back.result, msg.result);
838        assert_eq!(back.err_msg, msg.err_msg);
839    }
840
841    #[test]
842    fn local_roundtrip_doubles() {
843        let job = RemoteJobV1 {
844            seq: 0,
845            subs_prelude: String::new(),
846            block_src: "$_ * 2;".to_string(),
847            capture: vec![],
848            item: serde_json::json!(21),
849        };
850        let r = run_job_local(&job);
851        assert!(r.ok, "{}", r.err_msg);
852        assert_eq!(r.result, serde_json::json!(42));
853    }
854
855    // ─── framed I/O ────────────────────────────────────────────────────
856
857    #[test]
858    fn write_then_read_framed_roundtrips_payload() {
859        let payload = b"hello powerline".to_vec();
860        let mut buf = Vec::new();
861        write_framed(&mut buf, &payload).expect("write");
862        let read = read_framed(&mut buf.as_slice()).expect("read");
863        assert_eq!(read, payload);
864    }
865
866    #[test]
867    fn write_framed_emits_le_length_prefix() {
868        let mut buf = Vec::new();
869        write_framed(&mut buf, b"abc").unwrap();
870        // First 8 bytes = little-endian u64 length = 3.
871        assert_eq!(&buf[..8], &3u64.to_le_bytes());
872        assert_eq!(&buf[8..], b"abc");
873    }
874
875    #[test]
876    fn read_framed_rejects_oversized_frame() {
877        // Synthesise a header claiming MAX_FRAME+1 bytes follow.
878        let buf = ((MAX_FRAME + 1) as u64).to_le_bytes().to_vec();
879        // No body — we expect the size check to fire BEFORE the body read.
880        let err = read_framed(&mut buf.as_slice()).expect_err("oversize must fail");
881        assert_eq!(err.kind(), std::io::ErrorKind::InvalidData);
882    }
883
884    #[test]
885    fn read_framed_rejects_truncated_body() {
886        // Header says 100 bytes; we only ship 4.
887        let mut buf = 100u64.to_le_bytes().to_vec();
888        buf.extend_from_slice(b"shrt");
889        let err = read_framed(&mut buf.as_slice()).expect_err("truncated must fail");
890        assert_eq!(err.kind(), std::io::ErrorKind::UnexpectedEof);
891    }
892
893    #[test]
894    fn read_framed_zero_length_frame_is_empty_vec() {
895        let buf = 0u64.to_le_bytes().to_vec();
896        let body = read_framed(&mut buf.as_slice()).unwrap();
897        assert!(body.is_empty());
898    }
899
900    // ─── typed framed I/O (kind byte) ──────────────────────────────────
901
902    #[test]
903    fn write_typed_then_read_typed_preserves_kind_and_payload() {
904        let mut buf = Vec::new();
905        write_typed_frame(&mut buf, 0x42, b"hello").unwrap();
906        let (kind, body) = read_typed_frame(&mut buf.as_slice()).unwrap();
907        assert_eq!(kind, 0x42);
908        assert_eq!(body, b"hello");
909    }
910
911    #[test]
912    fn write_typed_frame_emits_kind_then_payload() {
913        let mut buf = Vec::new();
914        write_typed_frame(&mut buf, 0xAB, b"xyz").unwrap();
915        // Header (8 byte length) then kind (1) then payload (3) = 12 bytes.
916        assert_eq!(buf.len(), 12);
917        assert_eq!(&buf[..8], &4u64.to_le_bytes()); // 1 kind + 3 body
918        assert_eq!(buf[8], 0xAB);
919        assert_eq!(&buf[9..], b"xyz");
920    }
921
922    #[test]
923    fn read_typed_frame_rejects_empty_payload_missing_kind_byte() {
924        let buf = 0u64.to_le_bytes().to_vec(); // zero-length frame
925        let err = read_typed_frame(&mut buf.as_slice()).expect_err("empty kind must fail");
926        assert_eq!(err.kind(), std::io::ErrorKind::InvalidData);
927    }
928
929    // ─── send_msg / recv_msg helpers ───────────────────────────────────
930
931    #[test]
932    fn send_recv_msg_roundtrips_struct() {
933        let original = JobRespMsg {
934            seq: 42,
935            ok: true,
936            result: serde_json::json!({"foo": 1, "bar": [2, 3]}),
937            err_msg: String::new(),
938        };
939        let mut buf = Vec::new();
940        send_msg(&mut buf, 0x05, &original).expect("send");
941        let received: JobRespMsg = recv_msg(&mut buf.as_slice(), 0x05).expect("recv");
942        assert_eq!(received.seq, original.seq);
943        assert_eq!(received.ok, original.ok);
944        assert_eq!(received.result, original.result);
945        assert_eq!(received.err_msg, original.err_msg);
946    }
947
948    #[test]
949    fn recv_msg_with_wrong_kind_returns_descriptive_error() {
950        let mut buf = Vec::new();
951        send_msg(&mut buf, 0x01, &"hello".to_string()).unwrap();
952        let err: Result<String, _> = recv_msg(&mut buf.as_slice(), 0x99);
953        let msg = err.expect_err("wrong kind must fail");
954        assert!(
955            msg.contains("expected frame kind 0x99") && msg.contains("got 0x01"),
956            "unexpected message: {msg}"
957        );
958    }
959
960    // ─── encode_job / decode_job ──────────────────────────────────────
961
962    #[test]
963    fn encode_decode_job_roundtrips() {
964        let job = RemoteJobV1 {
965            seq: 7,
966            subs_prelude: "sub greet { p 'hi' }\n".into(),
967            block_src: "greet()".into(),
968            capture: vec![
969                ("x".into(), serde_json::json!(10)),
970                ("name".into(), serde_json::json!("bob")),
971            ],
972            item: serde_json::json!([1, 2, 3]),
973        };
974        let bytes = encode_job(&job).expect("encode");
975        let back = decode_job(&bytes).expect("decode");
976        assert_eq!(back.seq, job.seq);
977        assert_eq!(back.subs_prelude, job.subs_prelude);
978        assert_eq!(back.block_src, job.block_src);
979        assert_eq!(back.capture, job.capture);
980        assert_eq!(back.item, job.item);
981    }
982
983    #[test]
984    fn decode_job_rejects_garbage_bytes() {
985        let err = decode_job(b"this is not bincode").expect_err("garbage must fail");
986        assert!(!err.is_empty());
987    }
988
989    // ─── encode_resp / decode_resp ────────────────────────────────────
990
991    #[test]
992    fn encode_decode_resp_roundtrips_ok_case() {
993        let resp = RemoteRespV1 {
994            seq: 99,
995            ok: true,
996            result: serde_json::json!({"sum": 1234}),
997            err_msg: String::new(),
998        };
999        let bytes = encode_resp(&resp).expect("encode");
1000        let back = decode_resp(&bytes).expect("decode");
1001        assert_eq!(back.seq, resp.seq);
1002        assert_eq!(back.ok, resp.ok);
1003        assert_eq!(back.result, resp.result);
1004        assert!(back.err_msg.is_empty());
1005    }
1006
1007    #[test]
1008    fn encode_decode_resp_roundtrips_error_case() {
1009        let resp = RemoteRespV1 {
1010            seq: 5,
1011            ok: false,
1012            result: serde_json::json!(null),
1013            err_msg: "division by zero".into(),
1014        };
1015        let bytes = encode_resp(&resp).expect("encode");
1016        let back = decode_resp(&bytes).expect("decode");
1017        assert!(!back.ok);
1018        assert_eq!(back.err_msg, "division by zero");
1019    }
1020
1021    // ─── perl_to_json_value / json_to_perl ────────────────────────────
1022
1023    #[test]
1024    fn perl_to_json_handles_undef_int_str() {
1025        let undef = StrykeValue::UNDEF;
1026        let i = StrykeValue::integer(42);
1027        let s = StrykeValue::string("hello".to_string());
1028        assert_eq!(perl_to_json_value(&undef).unwrap(), serde_json::Value::Null);
1029        assert_eq!(perl_to_json_value(&i).unwrap(), serde_json::json!(42));
1030        assert_eq!(perl_to_json_value(&s).unwrap(), serde_json::json!("hello"));
1031    }
1032
1033    #[test]
1034    fn json_to_perl_round_trips_through_perl_to_json() {
1035        // Perl has no first-class bool — true/false collapse to 1/0 by
1036        // design. Skip bool inputs here; covered separately below.
1037        for j in [
1038            serde_json::json!(null),
1039            serde_json::json!(42),
1040            serde_json::json!(3.5),
1041            serde_json::json!("hello"),
1042            serde_json::json!([1, 2, 3]),
1043            serde_json::json!({"foo": "bar", "n": 7}),
1044        ] {
1045            let p = json_to_perl(&j).expect("json -> perl");
1046            let back = perl_to_json_value(&p).expect("perl -> json");
1047            assert_eq!(back, j, "roundtrip diverged for {j}");
1048        }
1049    }
1050
1051    #[test]
1052    fn json_to_perl_collapses_bool_to_int_per_perl_semantics() {
1053        // Documented Perl semantics — true/false are just 1/0.
1054        let t = json_to_perl(&serde_json::json!(true)).unwrap();
1055        let f = json_to_perl(&serde_json::json!(false)).unwrap();
1056        assert_eq!(perl_to_json_value(&t).unwrap(), serde_json::json!(1));
1057        assert_eq!(perl_to_json_value(&f).unwrap(), serde_json::json!(0));
1058    }
1059
1060    // ─── build_subs_prelude ────────────────────────────────────────────
1061
1062    #[test]
1063    fn build_subs_prelude_returns_empty_string_for_empty_map() {
1064        let subs = HashMap::new();
1065        let prelude = build_subs_prelude(&subs);
1066        assert!(prelude.is_empty());
1067    }
1068
1069    // ─── MAX_FRAME bound ───────────────────────────────────────────────
1070
1071    #[test]
1072    fn max_frame_is_256mib() {
1073        // Documented size cap — bumping it changes the wire-protocol's
1074        // memory ceiling on every worker. Pin the value.
1075        assert_eq!(MAX_FRAME, 256 * 1024 * 1024);
1076    }
1077}