1use crate::{Chunk, Program};
25use crate::natives::std_native;
26
27const N_PRINT: u32 = std_native::PRINT;
29const N_MAKE_MSG: u32 = std_native::MAKE_MSG;
30
31pub const TAG_REQ: i32 = 1;
33pub const TAG_REP: i32 = 2;
34pub const TAG_PING: i32 = 10;
35pub const TAG_PONG: i32 = 11;
36pub const TAG_JUNK: i32 = 99;
38
39pub fn add_forty_two() -> Chunk {
41 let mut p = Program::new("add-forty-two");
42 p.function("main", 0, |f| {
43 let a = f.load_i32(41);
44 let b = f.load_i32(1);
45 let sum = f.add(a, b);
46 f.return_(sum);
47 });
48 p.build()
49}
50
51pub fn ping_pong() -> Chunk {
55 let mut p = Program::new("ping-pong");
56 let pong = p.function("pong", 0, |f| {
57 let msg = f.receive();
58 let payload = f.hop_payload(msg);
59 f.add_imm(payload, 1);
60 f.send_reply(msg, TAG_PONG, payload);
61 f.exit(payload);
62 });
63 p.function("main", 0, |f| {
64 let child = f.spawn(pong, 0);
65 let req_id = f.load_i32(1);
66 let payload = f.load_i32(1);
67 let req = f.hop(req_id, TAG_PING, payload);
68 f.send(child, req);
69 let reply = f.receive();
70 let out = f.hop_payload(reply);
71 f.return_(out);
72 });
73 p.build()
74}
75
76pub fn atomic_request_reply() -> Chunk {
78 let mut p = Program::new("atomic-request-reply");
79 let server = p.function("server", 0, |f| {
80 let msg = f.receive();
81 f.native1_on(msg, N_PRINT);
82 let payload = f.hop_payload(msg);
83 f.add_imm(payload, 1);
84 f.send_reply(msg, TAG_REP, payload);
85 f.exit(payload);
86 });
87 p.function("main", 0, |f| {
88 let server_cap = f.spawn(server, 0);
89 let req_id = f.load_i32(1);
90 let payload = f.load_i32(41);
91 let req = f.hop(req_id, TAG_REQ, payload);
92 f.native1_on(req, N_PRINT);
93 f.send(server_cap, req);
94 let reply = f.receive();
95 f.native1_on(reply, N_PRINT);
96 let out = f.hop_payload(reply);
97 f.return_(out);
98 });
99 p.build()
100}
101
102pub fn selective_receive() -> Chunk {
105 let mut p = Program::new("selective-receive");
106 let server = p.function("server", 0, |f| {
107 let msg = f.receive_match_imm(TAG_REQ as u16);
108 let payload = f.hop_payload(msg);
109 f.add_imm(payload, 1);
110 f.send_reply(msg, TAG_REP, payload);
111 let junk = f.receive();
112 let tag = f.hop_tag(junk);
113 let is_junk = f.eq_imm(tag, TAG_JUNK);
114 let trap_lbl = f.label();
115 f.branch_if_falsy(is_junk, trap_lbl);
116 f.exit(payload);
117 f.bind(trap_lbl);
118 f.trap(2);
119 });
120 p.function("main", 0, |f| {
121 let server_cap = f.spawn(server, 0);
122 let req_id = f.load_i32(1);
123 let zero = f.load_i32(0);
124 let junk = f.hop(req_id, TAG_JUNK, zero);
125 f.send(server_cap, junk);
126 let payload = f.load_i32(41);
127 let req = f.hop(req_id, TAG_REQ, payload);
128 f.send(server_cap, req);
129 let reply = f.receive_match_imm(TAG_REP as u16);
130 let out = f.hop_payload(reply);
131 f.return_(out);
132 });
133 p.build()
134}
135
136pub fn ask_reply() -> Chunk {
138 let mut p = Program::new("ask-reply");
139 let server = p.function("server", 0, |f| {
140 let msg = f.receive_match_imm(TAG_REQ as u16);
141 let payload = f.hop_payload(msg);
142 f.add_imm(payload, 1);
143 f.send_reply(msg, TAG_REP, payload);
144 f.exit(payload);
145 });
146 p.function("main", 0, |f| {
147 let server_cap = f.spawn(server, 0);
148 let req_id = f.load_i32(1);
149 let payload = f.load_i32(41);
150 let req = f.hop(req_id, TAG_REQ, payload);
151 let reply = f.ask(server_cap, req);
152 let out = f.hop_payload(reply);
153 f.return_(out);
154 });
155 p.build()
156}
157
158pub fn ask_timeout_expires() -> Chunk {
160 let mut p = Program::new("ask-timeout");
161 let server = p.function("server", 0, |f| {
162 let _msg = f.receive();
163 let ms = f.load_i32(10_000);
164 f.sleep(ms);
165 let zero = f.load_i32(0);
166 f.return_(zero);
167 });
168 p.function("main", 0, |f| {
169 let server_cap = f.spawn(server, 0);
170 let req_id = f.load_i32(1);
171 let payload = f.load_i32(0);
172 let req = f.hop(req_id, TAG_REQ, payload);
173 let ms = f.load_i32(40);
174 let reply = f.ask_timeout(server_cap, req, ms);
175 f.return_(reply);
176 });
177 p.build()
178}
179
180pub fn ask_target_exits() -> Chunk {
182 let mut p = Program::new("ask-target-exits");
183 let server = p.function("server", 0, |f| {
184 let _msg = f.receive();
185 let z = f.load_i32(0);
186 f.return_(z);
187 });
188 p.function("main", 0, |f| {
189 let server_cap = f.spawn(server, 0);
190 let req_id = f.load_i32(1);
191 let payload = f.load_i32(0);
192 let req = f.hop(req_id, TAG_REQ, payload);
193 let reply = f.ask(server_cap, req);
194 let tag = f.hop_tag(reply);
195 f.return_(tag);
196 });
197 p.build()
198}
199
200pub fn server_loop() -> Chunk {
202 let mut p = Program::new("server-loop");
203 let server = p.function("server", 0, |f| {
204 let loop_lbl = f.label();
205 f.bind(loop_lbl);
206 let req = f.receive_match_imm(TAG_REQ as u16);
207 let payload = f.hop_payload(req);
208 f.add_imm(payload, 1);
209 f.send_reply(req, TAG_REP, payload);
210 f.jump(loop_lbl);
211 });
212 p.function("main", 0, |f| {
213 let server_cap = f.spawn(server, 0);
214 let req_id = f.load_i32(1);
215 let payload = f.load_i32(41);
216 let req = f.hop(req_id, TAG_REQ, payload);
217 f.send(server_cap, req);
218 let reply = f.receive_match_imm(TAG_REP as u16);
219 let out = f.hop_payload(reply);
220 f.return_(out);
221 });
222 p.build()
223}
224
225pub fn forged_sender_send() -> Chunk {
227 let mut p = Program::new("forged-sender-send");
228 let server = p.function("server", 0, |f| {
229 let msg = f.receive();
230 let sender = f.hop_sender(msg);
231 f.send_reply(msg, TAG_REP, sender);
232 f.exit(sender);
233 });
234 p.function("main", 0, |f| {
235 let server_cap = f.spawn(server, 0);
236 let forged = f.load_i32(999);
237 let req_id = f.load_i32(1);
238 let zero = f.load_i32(0);
239 let req = f.make_msg(N_MAKE_MSG, forged, req_id, TAG_REQ, zero);
240 f.send(server_cap, req);
241 let reply = f.receive();
242 let out = f.hop_payload(reply);
243 f.return_(out);
244 });
245 p.build()
246}
247
248pub fn forged_sender_ask() -> Chunk {
250 let mut p = Program::new("forged-sender-ask");
251 let server = p.function("server", 0, |f| {
252 let msg = f.receive_match_imm(TAG_REQ as u16);
253 let sender = f.hop_sender(msg);
254 f.send_reply(msg, TAG_REP, sender);
255 f.exit(sender);
256 });
257 p.function("main", 0, |f| {
258 let server_cap = f.spawn(server, 0);
259 let forged = f.load_i32(999);
260 let req_id = f.load_i32(1);
261 let zero = f.load_i32(0);
262 let req = f.make_msg(N_MAKE_MSG, forged, req_id, TAG_REQ, zero);
263 let reply = f.ask(server_cap, req);
264 let out = f.hop_payload(reply);
265 f.return_(out);
266 });
267 p.build()
268}
269
270pub fn monitor_down() -> Chunk {
272 let mut p = Program::new("monitor-down");
273 let child = p.function("child", 0, |f| {
274 let ms = f.load_i32(40);
275 f.sleep(ms);
276 let z = f.load_i32(0);
277 f.return_(z);
278 });
279 p.function("main", 0, |f| {
280 let cap = f.spawn(child, 0);
281 let mon = f.monitor(cap);
282 let msg = f.receive_match_imm(crate::TAG_SYS_DOWN);
283 let id = f.hop_request_id(msg);
284 let ok = f.eq(id, mon);
285 let trap = f.label();
286 f.branch_if_falsy(ok, trap);
287 let out = f.hop_payload(msg);
288 f.return_(out);
289 f.bind(trap);
290 f.trap(3);
291 });
292 p.build()
293}
294
295pub fn waiting_send() -> Chunk {
298 let mut p = Program::new("waiting-send");
299 p.function("server", 0, |f| {
300 let ms = f.load_i32(50);
301 f.sleep(ms);
302 let _first = f.receive();
303 let second = f.receive();
304 let out = f.hop_payload(second);
305 f.return_(out);
306 });
307 p.function("client", 1, |f| {
308 let server = f.reg(0);
309 let req_id = f.load_i32(1);
310 let one = f.load_i32(1);
311 let first = f.hop(req_id, TAG_REQ, one);
312 f.send(server, first);
313 let forty_two = f.load_i32(42);
314 let second = f.hop(req_id, TAG_REQ, forty_two);
315 f.send(server, second);
316 f.exit(second);
317 });
318 p.build()
319}
320
321pub fn boom() -> Chunk {
323 let mut p = Program::new("boom");
324 p.function("boom", 0, |f| f.trap(1));
325 p.build()
326}
327
328#[cfg(test)]
329mod tests {
330 use super::*;
331 use crate::{
332 decode, encode, std_native_table, verify, FlowOutcome, Runtime, RuntimeConfig, Value,
333 };
334
335 fn tiny(chunk: Chunk) -> Result<Runtime, crate::SpawnError> {
336 Runtime::with_config(
337 chunk,
338 RuntimeConfig {
339 workers: 1,
340 quantum: 10_000,
341 mailbox: crate::MailboxConfig::DEFAULT,
342 ..Default::default()
343 },
344 )
345 }
346
347 fn tiny_natives(chunk: Chunk) -> Result<Runtime, crate::SpawnError> {
348 Runtime::with_natives_and_config(
349 chunk,
350 std_native_table(),
351 RuntimeConfig {
352 workers: 1,
353 quantum: 10_000,
354 mailbox: crate::MailboxConfig::DEFAULT,
355 ..Default::default()
356 },
357 )
358 }
359
360 #[test]
361 fn add_forty_two_joins_42() -> Result<(), Box<dyn std::error::Error>> {
362 let rt = tiny(add_forty_two())?;
363 let idx = rt.function_index("main").ok_or("main")?;
364 let outcome = rt.spawn(idx, &[])?.join();
365 rt.shutdown();
366 assert!(matches!(outcome, FlowOutcome::Completed(Value::Int(42))));
367 Ok(())
368 }
369
370 #[test]
371 fn ping_pong_joins_2() -> Result<(), Box<dyn std::error::Error>> {
372 let chunk = ping_pong();
373 assert!(verify(&chunk).is_ok());
374 let bytes = encode(&chunk);
375 let chunk = decode(&bytes)?;
376 let rt = tiny_natives(chunk)?;
377 let idx = rt.function_index("main").ok_or("main")?;
378 let outcome = rt.spawn(idx, &[])?.join();
379 let sent = rt.metrics().messages_sent;
380 rt.shutdown();
381 assert!(matches!(outcome, FlowOutcome::Completed(Value::Int(2))));
382 assert!(sent >= 2);
383 Ok(())
384 }
385
386 #[test]
387 fn atomic_request_reply_joins_42() -> Result<(), Box<dyn std::error::Error>> {
388 let chunk = atomic_request_reply();
389 assert!(verify(&chunk).is_ok());
390 let bytes = encode(&chunk);
391 let chunk = decode(&bytes)?;
392 let rt = tiny_natives(chunk)?;
393 let idx = rt.function_index("main").ok_or("main")?;
394 let outcome = rt.spawn(idx, &[])?.join();
395 let sent = rt.metrics().messages_sent;
396 rt.shutdown();
397 assert!(
398 matches!(outcome, FlowOutcome::Completed(Value::Int(42))),
399 "got {outcome:?}"
400 );
401 assert!(sent >= 1);
402 Ok(())
403 }
404
405 #[test]
406 fn selective_receive_skips_junk_tag() -> Result<(), Box<dyn std::error::Error>> {
407 let chunk = selective_receive();
408 assert!(verify(&chunk).is_ok());
409 let bytes = encode(&chunk);
410 let chunk = decode(&bytes)?;
411 let rt = tiny_natives(chunk)?;
412 let idx = rt.function_index("main").ok_or("main")?;
413 let outcome = rt.spawn(idx, &[])?.join();
414 rt.shutdown();
415 assert!(
416 matches!(outcome, FlowOutcome::Completed(Value::Int(42))),
417 "got {outcome:?}"
418 );
419 Ok(())
420 }
421
422 #[test]
423 fn ask_reply_joins_42() -> Result<(), Box<dyn std::error::Error>> {
424 let chunk = ask_reply();
425 assert!(verify(&chunk).is_ok());
426 let bytes = encode(&chunk);
427 let chunk = decode(&bytes)?;
428 let rt = tiny_natives(chunk)?;
429 let idx = rt.function_index("main").ok_or("main")?;
430 let outcome = rt.spawn(idx, &[])?.join();
431 let sent = rt.metrics().messages_sent;
432 rt.shutdown();
433 assert!(
434 matches!(outcome, FlowOutcome::Completed(Value::Int(42))),
435 "got {outcome:?}"
436 );
437 assert!(sent >= 2);
438 Ok(())
439 }
440
441 #[test]
442 fn ask_timeout_writes_unit() -> Result<(), Box<dyn std::error::Error>> {
443 let chunk = ask_timeout_expires();
444 assert!(verify(&chunk).is_ok());
445 let rt = tiny_natives(chunk)?;
446 let idx = rt.function_index("main").ok_or("main")?;
447 let outcome = rt.spawn(idx, &[])?.join();
448 rt.shutdown();
449 assert!(
450 matches!(outcome, FlowOutcome::Completed(Value::Unit)),
451 "got {outcome:?}"
452 );
453 Ok(())
454 }
455
456 #[test]
457 fn ask_target_exit_writes_sys_exit() -> Result<(), Box<dyn std::error::Error>> {
458 let chunk = ask_target_exits();
459 assert!(verify(&chunk).is_ok());
460 let rt = tiny_natives(chunk)?;
461 let idx = rt.function_index("main").ok_or("main")?;
462 let outcome = rt.spawn(idx, &[])?.join();
463 rt.shutdown();
464 assert!(
465 matches!(
466 outcome,
467 FlowOutcome::Completed(Value::Int(n)) if n == i64::from(crate::TAG_SYS_EXIT)
468 ),
469 "got {outcome:?}"
470 );
471 Ok(())
472 }
473
474 #[test]
475 fn server_loop_joins_42() -> Result<(), Box<dyn std::error::Error>> {
476 let chunk = server_loop();
477 assert!(verify(&chunk).is_ok());
478 let rt = tiny_natives(chunk)?;
479 let idx = rt.function_index("main").ok_or("main")?;
480 let outcome = rt.spawn(idx, &[])?.join();
481 rt.shutdown();
482 assert!(
483 matches!(outcome, FlowOutcome::Completed(Value::Int(42))),
484 "got {outcome:?}"
485 );
486 Ok(())
487 }
488
489 #[test]
490 fn send_overwrites_forged_sender() -> Result<(), Box<dyn std::error::Error>> {
491 let chunk = forged_sender_send();
492 assert!(verify(&chunk).is_ok());
493 let rt = tiny_natives(chunk)?;
494 let idx = rt.function_index("main").ok_or("main")?;
495 let outcome = rt.spawn(idx, &[])?.join();
496 rt.shutdown();
497 match outcome {
498 FlowOutcome::Completed(Value::Int(n)) => {
499 assert_ne!(n, 999, "forged make_msg sender must not survive Send");
500 assert!(n >= 1, "authenticated sender must be a live flow id");
501 Ok(())
502 }
503 other => Err(format!("expected Completed(Int), got {other:?}").into()),
504 }
505 }
506
507 #[test]
508 fn ask_overwrites_forged_request_sender() -> Result<(), Box<dyn std::error::Error>> {
509 let chunk = forged_sender_ask();
510 assert!(verify(&chunk).is_ok());
511 let rt = tiny_natives(chunk)?;
512 let idx = rt.function_index("main").ok_or("main")?;
513 let outcome = rt.spawn(idx, &[])?.join();
514 rt.shutdown();
515 match outcome {
516 FlowOutcome::Completed(Value::Int(n)) => {
517 assert_ne!(n, 999, "forged make_msg sender must not survive Ask");
518 assert!(n >= 1, "authenticated sender must be a live flow id");
519 Ok(())
520 }
521 other => Err(format!("expected Completed(Int), got {other:?}").into()),
522 }
523 }
524
525 #[test]
526 fn send_scalar_target_traps() -> Result<(), Box<dyn std::error::Error>> {
527 let mut p = Program::new("bad-cap-target");
528 p.function("main", 0, |f| {
529 let bad_cap = f.load_i32(99);
530 let sender = f.load_i32(0);
531 let req_id = f.load_i32(1);
532 let payload = f.load_i32(1);
533 let msg = f.make_msg(N_MAKE_MSG, sender, req_id, TAG_PING, payload);
534 f.send(bad_cap, msg);
535 f.return_(msg);
536 });
537 let rt = tiny_natives(p.build())?;
538 let outcome = rt.spawn(0, &[])?.join();
539 rt.shutdown();
540 assert!(
541 matches!(outcome, FlowOutcome::Failed(_)),
542 "non-Cap Send target must fail, got {outcome:?}"
543 );
544 Ok(())
545 }
546
547 #[test]
548 fn send_scalar_is_not_an_atomic_hop() -> Result<(), Box<dyn std::error::Error>> {
549 let mut p = Program::new("bad-hop");
550 p.function("main", 0, |f| {
551 let cap = f.self_cap();
552 let scalar = f.load_i32(99);
553 f.send(cap, scalar);
554 f.return_(scalar);
555 });
556 let rt = tiny(p.build())?;
557 let outcome = rt.spawn(0, &[])?.join();
558 rt.shutdown();
559 assert!(
560 matches!(outcome, FlowOutcome::Failed(_)),
561 "scalar Send must trap, got {outcome:?}"
562 );
563 Ok(())
564 }
565
566 #[test]
567 fn monitor_down_joins_normal_reason() -> Result<(), Box<dyn std::error::Error>> {
568 let chunk = monitor_down();
569 assert!(verify(&chunk).is_ok());
570 let rt = tiny_natives(chunk)?;
571 let idx = rt.function_index("main").ok_or("main")?;
572 let outcome = rt.spawn(idx, &[])?.join();
573 rt.shutdown();
574 assert!(
575 matches!(outcome, FlowOutcome::Completed(Value::Int(0))),
576 "DOWN reason should be Normal (0), got {outcome:?}"
577 );
578 Ok(())
579 }
580
581 #[test]
582 fn waiting_send_second_hop_arrives() -> Result<(), Box<dyn std::error::Error>> {
583 let cap = crate::MailboxCapacity::new(1).ok_or("cap")?;
584 let rt = Runtime::with_natives_and_config(
585 waiting_send(),
586 std_native_table(),
587 RuntimeConfig {
588 workers: 1,
589 quantum: 10_000,
590 mailbox: crate::MailboxConfig::new(cap, crate::OverflowPolicy::Reject),
591 ..Default::default()
592 },
593 )?;
594 let server = rt.function_index("server").ok_or("server")?;
595 let client = rt.function_index("client").ok_or("client")?;
596 let server_h = rt.spawn(server, &[])?;
597 let server_cap = rt.mint_cap(server_h.id())?;
598 rt.spawn(client, &[Value::Cap(server_cap.as_u64())])?;
599 let outcome = server_h.join();
600 rt.shutdown();
601 assert!(
602 matches!(outcome, FlowOutcome::Completed(Value::Int(42))),
603 "second hop should be admitted after the first pop, got {outcome:?}"
604 );
605 Ok(())
606 }
607
608 #[test]
609 fn link_kills_peer_on_fault() -> Result<(), Box<dyn std::error::Error>> {
610 let mut p = Program::new("link-kill");
611 p.function("park", 0, |f| {
612 let _ = f.receive();
613 f.trap(9);
614 });
615 p.function("boom", 0, |f| f.trap(1));
616 let rt = tiny(p.build())?;
617 let park = rt.function_index("park").ok_or("park")?;
618 let boom = rt.function_index("boom").ok_or("boom")?;
619 let parked = rt.spawn(park, &[])?;
620 let killer = rt.spawn(boom, &[])?;
621 rt.link(parked.id(), killer.id())?;
622 let boom_out = killer.join();
623 let park_out = parked.join();
624 rt.shutdown();
625 assert!(matches!(boom_out, FlowOutcome::Failed(_)), "{boom_out:?}");
626 assert!(
627 matches!(park_out, FlowOutcome::Failed(_)),
628 "linked peer must die on fault, got {park_out:?}"
629 );
630 Ok(())
631 }
632
633 #[test]
634 fn linked_exit_down_carries_link_reason() -> Result<(), Box<dyn std::error::Error>> {
635 let mut p = Program::new("link-down-reason");
636 p.function("watcher", 0, |f| {
637 let msg = f.receive_match_imm(crate::TAG_SYS_DOWN);
638 let out = f.hop_payload(msg);
639 f.return_(out);
640 });
641 p.function("park", 0, |f| {
642 let _ = f.receive();
643 f.trap(9);
644 });
645 p.function("boom", 0, |f| f.trap(1));
646 let rt = tiny_natives(p.build())?;
647 let watcher = rt.spawn(rt.function_index("watcher").ok_or("watcher")?, &[])?;
648 let parked = rt.spawn(rt.function_index("park").ok_or("park")?, &[])?;
649 let killer = rt.spawn(rt.function_index("boom").ok_or("boom")?, &[])?;
650 rt.monitor(watcher.id(), parked.id())?;
651 rt.link(parked.id(), killer.id())?;
652 let _ = killer.join();
653 let watched = watcher.join();
654 let _ = parked.join();
655 rt.shutdown();
656 assert!(
657 matches!(
658 watched,
659 FlowOutcome::Completed(Value::Int(n)) if n == crate::FlowExitReason::Link.as_u64() as i64
660 ),
661 "DOWN payload must be Link, got {watched:?}"
662 );
663 Ok(())
664 }
665
666 #[test]
667 fn monitor_dead_owner_is_rejected() -> Result<(), Box<dyn std::error::Error>> {
668 let rt = tiny(add_forty_two())?;
669 let first = rt.spawn(0, &[])?;
670 let dead = first.id();
671 let done = first.join();
672 assert!(matches!(done, FlowOutcome::Completed(_)));
673 let live = rt.spawn(0, &[])?;
674 let err = rt.monitor(dead, live.id());
675 live.join();
676 rt.shutdown();
677 assert!(
678 matches!(err, Err(crate::LifecycleError::NoSuchFlow(_))),
679 "{err:?}"
680 );
681 Ok(())
682 }
683
684 #[test]
685 fn forged_cap_cannot_be_registered() -> Result<(), Box<dyn std::error::Error>> {
686 let rt = tiny(add_forty_two())?;
687 let err = rt.register_name("svc", crate::CapId(99_999));
688 rt.shutdown();
689 assert_eq!(err, Err(crate::LifecycleError::InvalidCapability));
690 Ok(())
691 }
692
693 #[test]
694 fn registry_clears_on_exit() -> Result<(), Box<dyn std::error::Error>> {
695 let mut p = Program::new("reg");
696 p.function("main", 0, |f| {
697 let ms = f.load_i32(80);
698 f.sleep(ms);
699 let z = f.load_i32(1);
700 f.return_(z);
701 });
702 let rt = tiny(p.build())?;
703 let h = rt.spawn(0, &[])?;
704 let cap = rt.mint_cap(h.id())?;
705 rt.register_name("svc", cap)?;
706 assert_eq!(rt.whereis("svc")?, Some(cap));
707 let _ = h.join();
708 assert_eq!(rt.whereis("svc")?, None);
709 rt.shutdown();
710 Ok(())
711 }
712}