1use serde::{Deserialize, Serialize};
47use std::io::{Read, Write};
48use std::net::TcpStream;
49use std::path::PathBuf;
50use std::sync::atomic::{AtomicBool, Ordering};
51use std::sync::Arc;
52use std::time::{Duration, Instant};
53
54pub mod frame_kind {
56 pub const AGENT_HELLO: u8 = 0x10;
58 pub const AGENT_HELLO_ACK: u8 = 0x11;
60 pub const FIRE: u8 = 0x12;
62 pub const METRICS: u8 = 0x13;
64 pub const TERMINATE: u8 = 0x14;
66 pub const TERM_ACK: u8 = 0x15;
68 pub const SHUTDOWN: u8 = 0x16;
70 pub const STATUS: u8 = 0x17;
72 pub const STATUS_RESP: u8 = 0x18;
74 pub const EVAL: u8 = 0x19;
76 pub const EVAL_RESULT: u8 = 0x1A;
78 pub const AGENT_AUTH: u8 = 0x1B;
82 pub const ERROR: u8 = 0xFF;
84}
85
86pub const AGENT_PROTO_VERSION: u32 = 2;
90
91#[derive(Debug, Clone, Serialize, Deserialize, Default)]
93pub struct AgentConfig {
94 #[serde(default)]
96 pub controller: ControllerConfig,
97 #[serde(default)]
99 pub limits: LimitsConfig,
100 #[serde(default)]
102 pub agent: AgentIdentity,
103}
104#[derive(Debug, Clone, Serialize, Deserialize)]
106pub struct ControllerConfig {
107 #[serde(default = "default_host")]
109 pub host: String,
110 #[serde(default = "default_port")]
112 pub port: u16,
113}
114
115fn default_host() -> String {
116 "localhost".to_string()
117}
118fn default_port() -> u16 {
119 9999
120}
121
122impl Default for ControllerConfig {
123 fn default() -> Self {
124 Self {
125 host: default_host(),
126 port: default_port(),
127 }
128 }
129}
130#[derive(Debug, Clone, Serialize, Deserialize)]
132pub struct LimitsConfig {
133 #[serde(default = "default_max_temp")]
135 pub max_temp: u32,
136 #[serde(default = "default_max_duration")]
138 pub max_duration: u64,
139}
140
141fn default_max_temp() -> u32 {
142 85
143}
144fn default_max_duration() -> u64 {
145 3600
146}
147
148impl Default for LimitsConfig {
149 fn default() -> Self {
150 Self {
151 max_temp: default_max_temp(),
152 max_duration: default_max_duration(),
153 }
154 }
155}
156#[derive(Debug, Clone, Serialize, Deserialize, Default)]
158pub struct AgentIdentity {
159 #[serde(default)]
161 pub name: Option<String>,
162}
163
164#[derive(Debug, Clone, Serialize, Deserialize)]
166pub struct AgentHello {
167 pub proto_version: u32,
169 pub stryke_version: String,
171 pub hostname: String,
173 pub cores: usize,
175 pub memory_bytes: u64,
177 pub agent_name: Option<String>,
179}
180
181#[derive(Debug, Clone, Serialize, Deserialize)]
183pub struct AgentHelloAck {
184 pub session_id: u64,
186 pub accepted: bool,
188 pub message: String,
190}
191
192#[derive(Debug, Clone, Serialize, Deserialize)]
196pub struct AgentAuth {
197 pub token: String,
199}
200
201#[derive(Debug, Clone, Serialize, Deserialize)]
203pub struct FireCommand {
204 pub workload: WorkloadType,
206 pub duration_secs: f64,
208 pub intensity: f64, }
210#[derive(Debug, Clone, Serialize, Deserialize)]
212pub enum WorkloadType {
213 Cpu,
215 Memory {
216 bytes: u64,
217 },
218 Io {
219 dir: String,
220 iterations: u64,
221 },
222 Combined,
224 Custom {
225 code: String,
226 },
227}
228
229#[derive(Debug, Clone, Serialize, Deserialize)]
231pub struct AgentMetrics {
232 pub cpu_percent: f64,
234 pub memory_used: u64,
236 pub hashes_per_sec: u64,
238 pub elapsed_secs: f64,
240 pub state: AgentState,
242}
243#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq)]
245pub enum AgentState {
246 Idle,
248 Armed,
250 Firing,
252 Terminated,
254}
255
256#[derive(Debug, Clone, Serialize, Deserialize)]
258pub struct TermAck {
259 pub total_hashes: u64,
261 pub total_duration: f64,
263 pub peak_cpu: f64,
265}
266
267#[derive(Debug, Clone, Serialize, Deserialize)]
273pub struct EvalCommand {
274 pub code: String,
276}
277
278#[derive(Debug, Clone, Serialize, Deserialize)]
281pub struct EvalResult {
282 pub ok: bool,
284 pub output: String,
286}
287
288pub fn handle_eval_frame<W: Write>(
292 stream: &mut W,
293 interp: &mut crate::vm_helper::VMHelper,
294 payload: &[u8],
295) -> std::io::Result<()> {
296 let result = match bincode::deserialize::<EvalCommand>(payload) {
297 Ok(cmd) => match crate::parse_and_run_string(&cmd.code, interp) {
298 Ok(v) => EvalResult {
299 ok: true,
300 output: v.to_string(),
301 },
302 Err(e) => EvalResult {
303 ok: false,
304 output: format!("{}", e),
305 },
306 },
307 Err(e) => EvalResult {
308 ok: false,
309 output: format!("malformed EVAL frame: {}", e),
310 },
311 };
312 let bytes = bincode::serialize(&result).expect("serialize EvalResult");
313 write_frame(stream, frame_kind::EVAL_RESULT, &bytes)
314}
315
316pub fn read_frame<R: Read>(r: &mut R) -> std::io::Result<(u8, Vec<u8>)> {
318 let mut len_buf = [0u8; 8];
319 r.read_exact(&mut len_buf)?;
320 let len = u64::from_le_bytes(len_buf) as usize;
321 if len < 1 {
322 return Err(std::io::Error::new(
323 std::io::ErrorKind::InvalidData,
324 "empty frame",
325 ));
326 }
327 let mut payload = vec![0u8; len];
328 r.read_exact(&mut payload)?;
329 let kind = payload[0];
330 Ok((kind, payload[1..].to_vec()))
331}
332
333pub fn write_frame<W: Write>(w: &mut W, kind: u8, payload: &[u8]) -> std::io::Result<()> {
335 let total_len = 1 + payload.len();
336 w.write_all(&(total_len as u64).to_le_bytes())?;
337 w.write_all(&[kind])?;
338 w.write_all(payload)?;
339 w.flush()
340}
341
342pub fn default_config_path() -> PathBuf {
344 dirs::config_dir()
345 .unwrap_or_else(|| PathBuf::from("."))
346 .join("stryke")
347 .join("agent.toml")
348}
349
350pub fn load_config(path: Option<&str>) -> AgentConfig {
352 let config_path = path.map(PathBuf::from).unwrap_or_else(default_config_path);
353
354 if config_path.exists() {
355 match std::fs::read_to_string(&config_path) {
356 Ok(content) => match toml::from_str(&content) {
357 Ok(config) => {
358 eprintln!("stryke agent: loaded config from {}", config_path.display());
359 return config;
360 }
361 Err(e) => {
362 eprintln!(
363 "stryke agent: config parse error {}: {}",
364 config_path.display(),
365 e
366 );
367 }
368 },
369 Err(e) => {
370 eprintln!("stryke agent: cannot read {}: {}", config_path.display(), e);
371 }
372 }
373 }
374
375 eprintln!("stryke agent: using default config (controller=localhost:9999)");
376 AgentConfig::default()
377}
378
379fn get_hostname() -> String {
381 hostname::get()
382 .map(|h| h.to_string_lossy().to_string())
383 .unwrap_or_else(|_| "unknown".to_string())
384}
385
386fn get_cores() -> usize {
388 std::thread::available_parallelism()
389 .map(|p| p.get())
390 .unwrap_or(1)
391}
392
393fn get_memory() -> u64 {
395 16 * 1024 * 1024 * 1024 }
399
400fn run_workload(
402 workload: &WorkloadType,
403 duration_secs: f64,
404 terminate: Arc<AtomicBool>,
405) -> (u64, f64) {
406 use sha2::{Digest, Sha256};
407 use std::sync::atomic::AtomicU64;
408
409 let start = Instant::now();
410 let duration = Duration::from_secs_f64(duration_secs);
411 let num_cores = std::thread::available_parallelism()
412 .map(|p| p.get())
413 .unwrap_or(1);
414
415 match workload {
416 WorkloadType::Cpu | WorkloadType::Combined => {
417 let total_hashes = AtomicU64::new(0);
418
419 std::thread::scope(|s| {
420 for _ in 0..num_cores {
421 let term = Arc::clone(&terminate);
422 let counter = &total_hashes;
423 s.spawn(move || {
424 let mut local_count: u64 = 0;
425 let mut data = [0u8; 64];
426
427 while start.elapsed() < duration && !term.load(Ordering::Relaxed) {
428 for _ in 0..1000 {
429 let hash = Sha256::digest(data);
430 data[..32].copy_from_slice(&hash);
431 local_count += 1;
432 }
433 }
434
435 counter.fetch_add(local_count, Ordering::Relaxed);
436 });
437 }
438 });
439
440 (
441 total_hashes.load(Ordering::Relaxed),
442 start.elapsed().as_secs_f64(),
443 )
444 }
445 WorkloadType::Memory { bytes } => {
446 let bytes_per_core = *bytes as usize / num_cores;
447
448 std::thread::scope(|s| {
449 for core_id in 0..num_cores {
450 let term = Arc::clone(&terminate);
451 s.spawn(move || {
452 if term.load(Ordering::Relaxed) {
453 return;
454 }
455 let mut buf: Vec<u8> = vec![0u8; bytes_per_core];
456 for i in (0..bytes_per_core).step_by(4096) {
457 if term.load(Ordering::Relaxed) {
458 break;
459 }
460 buf[i] = ((i + core_id) & 0xff) as u8;
461 }
462 std::hint::black_box(&buf);
463 });
464 }
465 });
466
467 (*bytes, start.elapsed().as_secs_f64())
468 }
469 WorkloadType::Io { dir, iterations } => {
470 use std::fs;
471 use std::io::Write as IoWrite;
472
473 let total_bytes = AtomicU64::new(0);
474 let iters_per_core = *iterations as usize / num_cores;
475
476 std::thread::scope(|s| {
477 for core_id in 0..num_cores {
478 let term = Arc::clone(&terminate);
479 let counter = &total_bytes;
480 let dir = dir.clone();
481 s.spawn(move || {
482 let io_data = vec![0xABu8; 1_000_000];
483 for i in 0..iters_per_core {
484 if term.load(Ordering::Relaxed) {
485 break;
486 }
487 let path = format!("{}/stryke_stress_{}_{}", dir, core_id, i);
488 if let Ok(mut f) = fs::File::create(&path) {
489 let _ = f.write_all(&io_data);
490 }
491 let _ = fs::read(&path);
492 let _ = fs::remove_file(&path);
493 counter.fetch_add(io_data.len() as u64, Ordering::Relaxed);
494 }
495 });
496 }
497 });
498
499 (
500 total_bytes.load(Ordering::Relaxed),
501 start.elapsed().as_secs_f64(),
502 )
503 }
504 WorkloadType::Custom { code: _ } => {
505 (0, start.elapsed().as_secs_f64())
507 }
508 }
509}
510
511pub fn run_agent(config_path: Option<&str>) -> i32 {
513 run_agent_with_overrides(config_path, None, None)
514}
515
516pub fn run_agent_with_explicit(host: &str, port: u16, name: Option<&str>) -> i32 {
521 let config = AgentConfig {
522 controller: ControllerConfig {
523 host: host.to_string(),
524 port,
525 },
526 limits: LimitsConfig::default(),
527 agent: AgentIdentity {
528 name: name.map(|s| s.to_string()),
529 },
530 };
531 run_agent_with_config(config)
532}
533
534pub fn run_agent_with_overrides(
536 config_path: Option<&str>,
537 controller_override: Option<&str>,
538 port_override: Option<u16>,
539) -> i32 {
540 let mut config = load_config(config_path);
541
542 if let Some(host) = controller_override {
543 config.controller.host = host.to_string();
544 }
545 if let Some(port) = port_override {
546 config.controller.port = port;
547 }
548
549 run_agent_with_config(config)
550}
551
552fn run_agent_with_config(config: AgentConfig) -> i32 {
557 let addr = format!("{}:{}", config.controller.host, config.controller.port);
558
559 eprintln!("stryke agent: connecting to controller at {}", addr);
560
561 let mut stream = match TcpStream::connect(&addr) {
562 Ok(s) => s,
563 Err(e) => {
564 eprintln!("stryke agent: connection failed: {}", e);
565 return 1;
566 }
567 };
568
569 let _ = stream.set_read_timeout(Some(Duration::from_millis(100)));
571
572 let hello = AgentHello {
574 proto_version: AGENT_PROTO_VERSION,
575 stryke_version: env!("CARGO_PKG_VERSION").to_string(),
576 hostname: get_hostname(),
577 cores: get_cores(),
578 memory_bytes: get_memory(),
579 agent_name: config.agent.name.clone(),
580 };
581
582 let hello_bytes = bincode::serialize(&hello).expect("serialize hello");
583 if let Err(e) = write_frame(&mut stream, frame_kind::AGENT_HELLO, &hello_bytes) {
584 eprintln!("stryke agent: failed to send hello: {}", e);
585 return 1;
586 }
587
588 if let Ok(token) = std::env::var("STRYKE_AGENT_TOKEN") {
594 if !token.is_empty() {
595 let auth = AgentAuth { token };
596 let auth_bytes = bincode::serialize(&auth).expect("serialize auth");
597 let _ = write_frame(&mut stream, frame_kind::AGENT_AUTH, &auth_bytes);
598 }
599 }
600
601 let (kind, payload) = match read_frame(&mut stream) {
603 Ok(f) => f,
604 Err(e) => {
605 eprintln!("stryke agent: failed to read hello ack: {}", e);
606 return 1;
607 }
608 };
609
610 if kind != frame_kind::AGENT_HELLO_ACK {
611 eprintln!("stryke agent: unexpected frame kind: {}", kind);
612 return 1;
613 }
614
615 let ack: AgentHelloAck = match bincode::deserialize(&payload) {
616 Ok(a) => a,
617 Err(e) => {
618 eprintln!("stryke agent: failed to parse hello ack: {}", e);
619 return 1;
620 }
621 };
622
623 if !ack.accepted {
624 eprintln!("stryke agent: rejected by controller: {}", ack.message);
625 return 1;
626 }
627
628 eprintln!(
629 "stryke agent: connected (session_id={}, cores={}, hostname={})",
630 ack.session_id,
631 get_cores(),
632 get_hostname()
633 );
634 eprintln!("stryke agent: awaiting commands...");
635
636 let _ = stream.set_read_timeout(None);
638
639 let terminate = Arc::new(AtomicBool::new(false));
641 #[allow(unused_assignments)]
642 let mut state = AgentState::Idle;
643 let mut interp = crate::vm_helper::VMHelper::new();
646 let mut session_start: Option<Instant> = None;
647 let mut total_hashes: u64 = 0;
648 let mut peak_cpu: f64 = 0.0;
649
650 loop {
651 let (kind, payload) = match read_frame(&mut stream) {
652 Ok(f) => f,
653 Err(e) => {
654 if e.kind() == std::io::ErrorKind::UnexpectedEof {
655 eprintln!("stryke agent: controller disconnected");
656 } else {
657 eprintln!("stryke agent: read error: {}", e);
658 }
659 break;
660 }
661 };
662
663 match kind {
664 frame_kind::FIRE => {
665 let cmd: FireCommand = match bincode::deserialize(&payload) {
666 Ok(c) => c,
667 Err(e) => {
668 eprintln!("stryke agent: invalid FIRE command: {}", e);
669 continue;
670 }
671 };
672
673 eprintln!(
674 "stryke agent: FIRE received (duration={}s, intensity={})",
675 cmd.duration_secs, cmd.intensity
676 );
677
678 #[allow(unused_assignments)]
679 {
680 state = AgentState::Firing;
681 }
682 session_start = Some(Instant::now());
683 terminate.store(false, Ordering::Relaxed);
684
685 let term_clone = Arc::clone(&terminate);
687 let workload = cmd.workload.clone();
688 let duration = cmd.duration_secs;
689
690 let handle =
691 std::thread::spawn(move || run_workload(&workload, duration, term_clone));
692
693 let (hashes, elapsed) = handle.join().unwrap_or((0, 0.0));
695 total_hashes += hashes;
696
697 let metrics = AgentMetrics {
699 cpu_percent: 100.0, memory_used: 0,
701 hashes_per_sec: if elapsed > 0.0 {
702 (hashes as f64 / elapsed) as u64
703 } else {
704 0
705 },
706 elapsed_secs: elapsed,
707 state: AgentState::Idle,
708 };
709
710 let metrics_bytes = bincode::serialize(&metrics).expect("serialize metrics");
711 let _ = write_frame(&mut stream, frame_kind::METRICS, &metrics_bytes);
712
713 state = AgentState::Idle;
714 eprintln!(
715 "stryke agent: workload complete ({} hashes in {:.2}s)",
716 hashes, elapsed
717 );
718 }
719
720 frame_kind::TERMINATE => {
721 eprintln!("stryke agent: TERMINATE received");
722 terminate.store(true, Ordering::Relaxed);
723
724 let elapsed = session_start
725 .map(|s| s.elapsed().as_secs_f64())
726 .unwrap_or(0.0);
727 let term_ack = TermAck {
728 total_hashes,
729 total_duration: elapsed,
730 peak_cpu,
731 };
732
733 let ack_bytes = bincode::serialize(&term_ack).expect("serialize term_ack");
734 let _ = write_frame(&mut stream, frame_kind::TERM_ACK, &ack_bytes);
735
736 state = AgentState::Idle;
737 total_hashes = 0;
738 peak_cpu = 0.0;
739 session_start = None;
740 }
741
742 frame_kind::STATUS => {
743 let metrics = AgentMetrics {
744 cpu_percent: if state == AgentState::Firing {
745 100.0
746 } else {
747 0.0
748 },
749 memory_used: 0,
750 hashes_per_sec: 0,
751 elapsed_secs: session_start
752 .map(|s| s.elapsed().as_secs_f64())
753 .unwrap_or(0.0),
754 state,
755 };
756
757 let metrics_bytes = bincode::serialize(&metrics).expect("serialize metrics");
758 let _ = write_frame(&mut stream, frame_kind::STATUS_RESP, &metrics_bytes);
759 }
760
761 frame_kind::EVAL => {
762 eprintln!("stryke agent: EVAL received ({} bytes)", payload.len());
763 if let Err(e) = handle_eval_frame(&mut stream, &mut interp, &payload) {
764 eprintln!("stryke agent: failed to write EVAL_RESULT: {}", e);
765 }
766 }
767
768 frame_kind::SHUTDOWN => {
769 eprintln!("stryke agent: SHUTDOWN received, exiting");
770 terminate.store(true, Ordering::Relaxed);
771 break;
772 }
773
774 _ => {
775 eprintln!("stryke agent: unknown frame kind: {}", kind);
776 }
777 }
778 }
779
780 eprintln!("stryke agent: disconnected");
781 0
782}
783
784pub fn print_help() {
786 println!("stryke agent — Distributed load testing agent");
787 println!();
788 println!("USAGE:");
789 println!(" stryke agent [OPTIONS]");
790 println!();
791 println!("OPTIONS:");
792 println!(" -c, --config PATH Config file (default: ~/.config/stryke/agent.toml)");
793 println!(" --controller HOST Controller address (overrides config)");
794 println!(" --port PORT Controller port (overrides config)");
795 println!(" --help Print this help");
796 println!();
797 println!("CONFIG FILE:");
798 println!(" ~/.config/stryke/agent.toml");
799 println!();
800 println!(" [controller]");
801 println!(" host = \"controller.example.com\"");
802 println!(" port = 9999");
803 println!();
804 println!(" [limits]");
805 println!(" max_temp = 85");
806 println!(" max_duration = 3600");
807 println!();
808 println!(" [agent]");
809 println!(" name = \"node-01\"");
810 println!();
811 println!("EXAMPLE:");
812 println!(" stryke agent # use config file");
813 println!(" stryke agent --controller 10.0.0.1 # connect to specific host");
814}
815
816#[cfg(test)]
817mod tests {
818 use super::*;
819 use std::io::Cursor;
820 use std::net::{TcpListener, TcpStream};
821 use std::thread;
822
823 #[test]
831 fn agent_hello_bincode_roundtrip_preserves_all_fields() {
832 let hello = AgentHello {
833 proto_version: AGENT_PROTO_VERSION,
834 stryke_version: "0.14.30".to_string(),
835 hostname: "lab-node-07.example.com".to_string(),
836 cores: 32,
837 memory_bytes: 256u64 * 1024 * 1024 * 1024, agent_name: Some("gpu-worker-α".to_string()), };
840 let bytes = bincode::serialize(&hello).expect("serialize Some-name");
841 let back: AgentHello = bincode::deserialize(&bytes).expect("deserialize");
842 assert_eq!(back.proto_version, AGENT_PROTO_VERSION);
843 assert_eq!(back.stryke_version, "0.14.30");
844 assert_eq!(back.hostname, "lab-node-07.example.com");
845 assert_eq!(back.cores, 32);
846 assert_eq!(back.memory_bytes, 256u64 * 1024 * 1024 * 1024);
847 assert_eq!(back.agent_name.as_deref(), Some("gpu-worker-α"));
848
849 let anon = AgentHello {
851 proto_version: AGENT_PROTO_VERSION,
852 stryke_version: "0.14.30".to_string(),
853 hostname: "anon".to_string(),
854 cores: 1,
855 memory_bytes: 0,
856 agent_name: None,
857 };
858 let bytes = bincode::serialize(&anon).unwrap();
859 let back: AgentHello = bincode::deserialize(&bytes).unwrap();
860 assert!(back.agent_name.is_none(), "None must round-trip as None");
861 }
862
863 #[test]
867 fn agent_hello_ack_bincode_roundtrip_preserves_all_fields() {
868 let ack = AgentHelloAck {
869 session_id: u64::MAX,
870 accepted: false,
871 message: "incompatible proto_version: agent=1 controller=2".to_string(),
872 };
873 let bytes = bincode::serialize(&ack).expect("serialize ack");
874 let back: AgentHelloAck = bincode::deserialize(&bytes).expect("deserialize ack");
875 assert_eq!(back.session_id, u64::MAX);
876 assert!(!back.accepted);
877 assert_eq!(
878 back.message,
879 "incompatible proto_version: agent=1 controller=2"
880 );
881 }
882
883 #[test]
886 fn eval_command_result_bincode_roundtrip() {
887 let cmd = EvalCommand {
888 code: "1 + 2".to_string(),
889 };
890 let bytes = bincode::serialize(&cmd).unwrap();
891 let back: EvalCommand = bincode::deserialize(&bytes).unwrap();
892 assert_eq!(back.code, "1 + 2");
893
894 let res = EvalResult {
895 ok: false,
896 output: "die at -e line 1".to_string(),
897 };
898 let bytes = bincode::serialize(&res).unwrap();
899 let back: EvalResult = bincode::deserialize(&bytes).unwrap();
900 assert!(!back.ok);
901 assert_eq!(back.output, "die at -e line 1");
902 }
903
904 #[test]
908 fn handle_eval_frame_success_writes_eval_result() {
909 let mut interp = crate::vm_helper::VMHelper::new();
910 let cmd = EvalCommand {
911 code: "21 * 2".to_string(),
912 };
913 let payload = bincode::serialize(&cmd).unwrap();
914
915 let mut out = Vec::new();
916 handle_eval_frame(&mut out, &mut interp, &payload).expect("write EVAL_RESULT");
917
918 let mut cur = Cursor::new(out);
919 let (kind, body) = read_frame(&mut cur).expect("read back");
920 assert_eq!(kind, frame_kind::EVAL_RESULT);
921 let r: EvalResult = bincode::deserialize(&body).unwrap();
922 assert!(r.ok, "expected ok=true, got {:?}", r);
923 assert_eq!(r.output, "42");
924 }
925
926 #[test]
928 fn handle_eval_frame_error_writes_eval_result_with_ok_false() {
929 let mut interp = crate::vm_helper::VMHelper::new();
930 let cmd = EvalCommand {
931 code: "this is not valid stryke @@@".to_string(),
932 };
933 let payload = bincode::serialize(&cmd).unwrap();
934
935 let mut out = Vec::new();
936 handle_eval_frame(&mut out, &mut interp, &payload).expect("write");
937
938 let mut cur = Cursor::new(out);
939 let (kind, body) = read_frame(&mut cur).expect("read back");
940 assert_eq!(kind, frame_kind::EVAL_RESULT);
941 let r: EvalResult = bincode::deserialize(&body).unwrap();
942 assert!(!r.ok, "expected ok=false on parse failure, got {:?}", r);
943 assert!(!r.output.is_empty(), "error output must not be empty");
944 }
945
946 #[test]
955 fn successive_eval_frames_share_persistent_vm_state() {
956 let mut interp = crate::vm_helper::VMHelper::new();
957
958 let cmd1 = EvalCommand {
960 code: "$main::counter = 100; sub bump { $main::counter + 7 } $main::counter"
961 .to_string(),
962 };
963 let mut out1 = Vec::new();
964 handle_eval_frame(&mut out1, &mut interp, &bincode::serialize(&cmd1).unwrap()).unwrap();
965 let (_, body1) = read_frame(&mut Cursor::new(out1)).unwrap();
966 let r1: EvalResult = bincode::deserialize(&body1).unwrap();
967 assert!(r1.ok, "frame 1 must succeed, got {:?}", r1);
968 assert_eq!(r1.output, "100");
969
970 let cmd2 = EvalCommand {
972 code: "bump()".to_string(),
973 };
974 let mut out2 = Vec::new();
975 handle_eval_frame(&mut out2, &mut interp, &bincode::serialize(&cmd2).unwrap()).unwrap();
976 let (_, body2) = read_frame(&mut Cursor::new(out2)).unwrap();
977 let r2: EvalResult = bincode::deserialize(&body2).unwrap();
978 assert!(r2.ok, "frame 2 must succeed, got {:?}", r2);
979 assert_eq!(
980 r2.output, "107",
981 "frame 2 must call the sub defined in frame 1 and see $main::counter"
982 );
983 }
984
985 #[test]
990 fn tcp_loopback_eval_roundtrip() {
991 let listener = TcpListener::bind("127.0.0.1:0").expect("bind loopback");
992 let addr = listener.local_addr().unwrap();
993
994 let agent_handle = thread::spawn(move || {
995 let (mut server, _) = listener.accept().expect("accept");
996 let mut interp = crate::vm_helper::VMHelper::new();
997 let (kind, payload) = read_frame(&mut server).expect("read EVAL");
998 assert_eq!(kind, frame_kind::EVAL);
999 handle_eval_frame(&mut server, &mut interp, &payload).expect("reply");
1000 });
1001
1002 let mut client = TcpStream::connect(addr).expect("connect");
1003 let cmd = EvalCommand {
1004 code: "join(\",\", 1:5)".to_string(),
1005 };
1006 let body = bincode::serialize(&cmd).unwrap();
1007 write_frame(&mut client, frame_kind::EVAL, &body).expect("send EVAL");
1008
1009 let (kind, payload) = read_frame(&mut client).expect("read EVAL_RESULT");
1010 assert_eq!(kind, frame_kind::EVAL_RESULT);
1011 let r: EvalResult = bincode::deserialize(&payload).unwrap();
1012 assert!(r.ok, "expected ok=true, got {:?}", r);
1013 assert_eq!(r.output, "1,2,3,4,5");
1014
1015 agent_handle.join().expect("agent thread");
1016 }
1017}