Skip to main content

byteflow/
samples.rs

1//! Built-in demo chunks assembled with [`crate::Program`].
2//!
3//! Use these as runnable specs of the messaging contract (and as regression
4//! tests). Prefer copying a sample over inventing hop register layouts from
5//! scratch.
6//!
7//! | Sample | Shows |
8//! |--------|--------|
9//! | [`add_forty_two`] | Scalar VM path (no natives) |
10//! | [`ping_pong`] | Cap spawn + Atomic Hop round-trip |
11//! | [`atomic_request_reply`] | Tagged REQ/REP + `print` |
12//! | [`selective_receive`] | `ReceiveMatch` FIFO skip |
13//! | [`ask_reply`] | `Ask` RPC hop |
14//! | [`forged_sender_send`] / [`forged_sender_ask`] | S1: forged `make_msg` sender dies |
15//! | [`boom`] | Immediate trap (supervisor demos) |
16//!
17//! Hop samples require [`crate::std_native_table`].
18
19use crate::{Chunk, Program};
20
21/// Native indices (must match [`crate::std_native_map`]).
22const N_PRINT: u32 = 0;
23const N_MAKE_MSG: u32 = 2;
24const N_MSG_SENDER: u32 = 3;
25const N_MSG_REQUEST_ID: u32 = 4;
26const N_MSG_TAG: u32 = 5;
27const N_MSG_PAYLOAD: u32 = 6;
28const N_MSG_REPLY_CAP: u32 = 7;
29
30/// Protocol tags for Atomic Hop samples.
31pub const TAG_REQ: i32 = 1;
32pub const TAG_REP: i32 = 2;
33pub const TAG_PING: i32 = 10;
34pub const TAG_PONG: i32 = 11;
35/// Decoy hop for [`selective_receive`] — must be skipped by `ReceiveMatch`.
36pub const TAG_JUNK: i32 = 99;
37
38/// `41 + 1; return` — the 60-second sanity chunk.
39pub fn add_forty_two() -> Chunk {
40    let mut p = Program::new("add-forty-two");
41    p.function("main", 0, |f| {
42        let a = f.load_i32(41);
43        let b = f.load_i32(1);
44        let sum = f.add(a, b);
45        f.return_(sum);
46    });
47    p.build()
48}
49
50/// Two flows, one **Atomic Hop** round-trip: `main` sends a `Message` to
51/// `pong`, `pong` replies with payload+1 via `msg_reply_cap`, `main` returns
52/// that payload (`2`).
53pub fn ping_pong() -> Chunk {
54    let mut p = Program::new("ping-pong");
55    let pong =     p.function("pong", 0, |f| {
56        let msg = f.receive();
57        let reply_cap = f.native1_from(msg, N_MSG_REPLY_CAP);
58        let req_id = f.native1_from(msg, N_MSG_REQUEST_ID);
59        let payload = f.native1_from(msg, N_MSG_PAYLOAD);
60        f.add_imm(payload, 1);
61        let self_cap = f.self_cap();
62        let reply = f.make_msg(N_MAKE_MSG, self_cap, req_id, TAG_PONG, payload);
63        f.send(reply_cap, reply);
64        f.exit(reply);
65    });
66    p.function("main", 0, |f| {
67        let self_cap = f.self_cap();
68        let child = f.spawn(pong, 0);
69        let req_id = f.load_i32(1);
70        let payload = f.load_i32(1);
71        let req = f.make_msg(N_MAKE_MSG, self_cap, req_id, TAG_PING, payload);
72        f.send(child, req);
73        let reply = f.receive();
74        let out = f.native1_from(reply, N_MSG_PAYLOAD);
75        f.return_(out);
76    });
77    p.build()
78}
79
80/// Atomic request-reply with [`crate::Value::Message`] (one envelope per hop).
81pub fn atomic_request_reply() -> Chunk {
82    let mut p = Program::new("atomic-request-reply");
83    let server = p.function("server", 0, |f| {
84        let msg = f.receive();
85        f.native1_on(msg, N_PRINT);
86        let reply_cap = f.native1_from(msg, N_MSG_REPLY_CAP);
87        let req_id = f.native1_from(msg, N_MSG_REQUEST_ID);
88        let _tag = f.native1_from(msg, N_MSG_TAG);
89        let payload = f.native1_from(msg, N_MSG_PAYLOAD);
90        f.add_imm(payload, 1);
91        let self_cap = f.self_cap();
92        let reply = f.make_msg(N_MAKE_MSG, self_cap, req_id, TAG_REP, payload);
93        f.native1_on(reply, N_PRINT);
94        f.send(reply_cap, reply);
95        f.exit(reply);
96    });
97    p.function("main", 0, |f| {
98        let self_cap = f.self_cap();
99        let server_cap = f.spawn(server, 0);
100        let req_id = f.load_i32(1);
101        let payload = f.load_i32(41);
102        let req = f.make_msg(N_MAKE_MSG, self_cap, req_id, TAG_REQ, payload);
103        f.native1_on(req, N_PRINT);
104        f.send(server_cap, req);
105        let reply = f.receive();
106        f.native1_on(reply, N_PRINT);
107        let out = f.native1_from(reply, N_MSG_PAYLOAD);
108        f.return_(out);
109    });
110    p.build()
111}
112
113/// Selective Atomic Hop: server waits for `TAG_REQ` while a `TAG_JUNK` hop
114/// sits ahead in the mailbox (FIFO skip, not drop).
115pub fn selective_receive() -> Chunk {
116    let mut p = Program::new("selective-receive");
117    let server = p.function("server", 0, |f| {
118        let msg = f.receive_match_imm(TAG_REQ as u16);
119        let reply_cap = f.native1_from(msg, N_MSG_REPLY_CAP);
120        let req_id = f.native1_from(msg, N_MSG_REQUEST_ID);
121        let payload = f.native1_from(msg, N_MSG_PAYLOAD);
122        f.add_imm(payload, 1);
123        let self_cap = f.self_cap();
124        let reply = f.make_msg(N_MAKE_MSG, self_cap, req_id, TAG_REP, payload);
125        f.send(reply_cap, reply);
126        let junk = f.receive();
127        let tag = f.native1_from(junk, N_MSG_TAG);
128        let is_junk = f.eq_imm(tag, TAG_JUNK);
129        let trap_lbl = f.label();
130        f.branch_if_falsy(is_junk, trap_lbl);
131        f.exit(reply);
132        f.bind(trap_lbl);
133        f.trap(2);
134    });
135    p.function("main", 0, |f| {
136        let self_cap = f.self_cap();
137        let server_cap = f.spawn(server, 0);
138        let req_id = f.load_i32(1);
139        let zero = f.load_i32(0);
140        let junk = f.make_msg(N_MAKE_MSG, self_cap, req_id, TAG_JUNK, zero);
141        f.send(server_cap, junk);
142        let payload = f.load_i32(41);
143        let req = f.make_msg(N_MAKE_MSG, self_cap, req_id, TAG_REQ, payload);
144        f.send(server_cap, req);
145        let reply = f.receive_match_imm(TAG_REP as u16);
146        let out = f.native1_from(reply, N_MSG_PAYLOAD);
147        f.return_(out);
148    });
149    p.build()
150}
151
152/// Atomic request/reply via `Ask` (RPC hop).
153pub fn ask_reply() -> Chunk {
154    let mut p = Program::new("ask-reply");
155    let server = p.function("server", 0, |f| {
156        let msg = f.receive_match_imm(TAG_REQ as u16);
157        let reply_cap = f.native1_from(msg, N_MSG_REPLY_CAP);
158        let req_id = f.native1_from(msg, N_MSG_REQUEST_ID);
159        let payload = f.native1_from(msg, N_MSG_PAYLOAD);
160        f.add_imm(payload, 1);
161        let self_cap = f.self_cap();
162        let reply = f.make_msg(N_MAKE_MSG, self_cap, req_id, TAG_REP, payload);
163        f.send(reply_cap, reply);
164        f.exit(reply);
165    });
166    p.function("main", 0, |f| {
167        let self_cap = f.self_cap();
168        let server_cap = f.spawn(server, 0);
169        let req_id = f.load_i32(1);
170        let payload = f.load_i32(41);
171        let req = f.make_msg(N_MAKE_MSG, self_cap, req_id, TAG_REQ, payload);
172        let reply = f.ask(server_cap, req);
173        let out = f.native1_from(reply, N_MSG_PAYLOAD);
174        f.return_(out);
175    });
176    p.build()
177}
178
179/// Security regression: forged `make_msg` sender must not survive `Send`.
180pub fn forged_sender_send() -> Chunk {
181    let mut p = Program::new("forged-sender-send");
182    let server = p.function("server", 0, |f| {
183        let msg = f.receive();
184        let reply_cap = f.native1_from(msg, N_MSG_REPLY_CAP);
185        let req_id = f.native1_from(msg, N_MSG_REQUEST_ID);
186        let sender = f.native1_from(msg, N_MSG_SENDER);
187        let self_cap = f.self_cap();
188        let reply = f.make_msg(N_MAKE_MSG, self_cap, req_id, TAG_REP, sender);
189        f.send(reply_cap, reply);
190        f.exit(reply);
191    });
192    p.function("main", 0, |f| {
193        let server_cap = f.spawn(server, 0);
194        let forged = f.load_i32(999);
195        let req_id = f.load_i32(1);
196        let zero = f.load_i32(0);
197        let req = f.make_msg(N_MAKE_MSG, forged, req_id, TAG_REQ, zero);
198        f.send(server_cap, req);
199        let reply = f.receive();
200        let out = f.native1_from(reply, N_MSG_PAYLOAD);
201        f.return_(out);
202    });
203    p.build()
204}
205
206/// Same security property as [`forged_sender_send`], via `Ask`.
207pub fn forged_sender_ask() -> Chunk {
208    let mut p = Program::new("forged-sender-ask");
209    let server = p.function("server", 0, |f| {
210        let msg = f.receive_match_imm(TAG_REQ as u16);
211        let reply_cap = f.native1_from(msg, N_MSG_REPLY_CAP);
212        let req_id = f.native1_from(msg, N_MSG_REQUEST_ID);
213        let sender = f.native1_from(msg, N_MSG_SENDER);
214        let self_cap = f.self_cap();
215        let reply = f.make_msg(N_MAKE_MSG, self_cap, req_id, TAG_REP, sender);
216        f.send(reply_cap, reply);
217        f.exit(reply);
218    });
219    p.function("main", 0, |f| {
220        let server_cap = f.spawn(server, 0);
221        let forged = f.load_i32(999);
222        let req_id = f.load_i32(1);
223        let zero = f.load_i32(0);
224        let req = f.make_msg(N_MAKE_MSG, forged, req_id, TAG_REQ, zero);
225        let reply = f.ask(server_cap, req);
226        let out = f.native1_from(reply, N_MSG_PAYLOAD);
227        f.return_(out);
228    });
229    p.build()
230}
231
232/// Immediate `Trap` — used to show [`crate::Supervisor`] restart.
233pub fn boom() -> Chunk {
234    let mut p = Program::new("boom");
235    p.function("boom", 0, |f| f.trap(1));
236    p.build()
237}
238
239#[cfg(test)]
240mod tests {
241    use super::*;
242    use crate::{
243        decode, encode, std_native_table, verify, FlowOutcome, Runtime, RuntimeConfig, Value,
244    };
245
246    fn tiny(chunk: Chunk) -> Result<Runtime, crate::SpawnError> {
247        Runtime::with_config(
248            chunk,
249            RuntimeConfig {
250                workers: 1,
251                quantum: 10_000,
252                mailbox: crate::MailboxConfig::DEFAULT,
253                ..Default::default()
254            },
255        )
256    }
257
258    fn tiny_natives(chunk: Chunk) -> Result<Runtime, crate::SpawnError> {
259        Runtime::with_natives_and_config(
260            chunk,
261            std_native_table(),
262            RuntimeConfig {
263                workers: 1,
264                quantum: 10_000,
265                mailbox: crate::MailboxConfig::DEFAULT,
266                ..Default::default()
267            },
268        )
269    }
270
271    #[test]
272    fn add_forty_two_joins_42() -> Result<(), Box<dyn std::error::Error>> {
273        let rt = tiny(add_forty_two())?;
274        let idx = rt.function_index("main").ok_or("main")?;
275        let outcome = rt.spawn(idx, &[])?.join();
276        rt.shutdown();
277        assert!(matches!(outcome, FlowOutcome::Completed(Value::Int(42))));
278        Ok(())
279    }
280
281    #[test]
282    fn ping_pong_joins_2() -> Result<(), Box<dyn std::error::Error>> {
283        let chunk = ping_pong();
284        assert!(verify(&chunk).is_ok());
285        let bytes = encode(&chunk);
286        let chunk = decode(&bytes)?;
287        let rt = tiny_natives(chunk)?;
288        let idx = rt.function_index("main").ok_or("main")?;
289        let outcome = rt.spawn(idx, &[])?.join();
290        let sent = rt.metrics().messages_sent;
291        rt.shutdown();
292        assert!(matches!(outcome, FlowOutcome::Completed(Value::Int(2))));
293        assert!(sent >= 2);
294        Ok(())
295    }
296
297    #[test]
298    fn atomic_request_reply_joins_42() -> Result<(), Box<dyn std::error::Error>> {
299        let chunk = atomic_request_reply();
300        assert!(verify(&chunk).is_ok());
301        let bytes = encode(&chunk);
302        let chunk = decode(&bytes)?;
303        let rt = tiny_natives(chunk)?;
304        let idx = rt.function_index("main").ok_or("main")?;
305        let outcome = rt.spawn(idx, &[])?.join();
306        let sent = rt.metrics().messages_sent;
307        rt.shutdown();
308        assert!(
309            matches!(outcome, FlowOutcome::Completed(Value::Int(42))),
310            "got {outcome:?}"
311        );
312        assert!(sent >= 1);
313        Ok(())
314    }
315
316    #[test]
317    fn selective_receive_skips_junk_tag() -> Result<(), Box<dyn std::error::Error>> {
318        let chunk = selective_receive();
319        assert!(verify(&chunk).is_ok());
320        let bytes = encode(&chunk);
321        let chunk = decode(&bytes)?;
322        let rt = tiny_natives(chunk)?;
323        let idx = rt.function_index("main").ok_or("main")?;
324        let outcome = rt.spawn(idx, &[])?.join();
325        rt.shutdown();
326        assert!(
327            matches!(outcome, FlowOutcome::Completed(Value::Int(42))),
328            "got {outcome:?}"
329        );
330        Ok(())
331    }
332
333    #[test]
334    fn ask_reply_joins_42() -> Result<(), Box<dyn std::error::Error>> {
335        let chunk = ask_reply();
336        assert!(verify(&chunk).is_ok());
337        let bytes = encode(&chunk);
338        let chunk = decode(&bytes)?;
339        let rt = tiny_natives(chunk)?;
340        let idx = rt.function_index("main").ok_or("main")?;
341        let outcome = rt.spawn(idx, &[])?.join();
342        let sent = rt.metrics().messages_sent;
343        rt.shutdown();
344        assert!(
345            matches!(outcome, FlowOutcome::Completed(Value::Int(42))),
346            "got {outcome:?}"
347        );
348        assert!(sent >= 2);
349        Ok(())
350    }
351
352    #[test]
353    fn send_overwrites_forged_sender() -> Result<(), Box<dyn std::error::Error>> {
354        let chunk = forged_sender_send();
355        assert!(verify(&chunk).is_ok());
356        let rt = tiny_natives(chunk)?;
357        let idx = rt.function_index("main").ok_or("main")?;
358        let outcome = rt.spawn(idx, &[])?.join();
359        rt.shutdown();
360        match outcome {
361            FlowOutcome::Completed(Value::Int(n)) => {
362                assert_ne!(n, 999, "forged make_msg sender must not survive Send");
363                assert!(n >= 1, "authenticated sender must be a live flow id");
364                Ok(())
365            }
366            other => Err(format!("expected Completed(Int), got {other:?}").into()),
367        }
368    }
369
370    #[test]
371    fn ask_overwrites_forged_request_sender() -> Result<(), Box<dyn std::error::Error>> {
372        let chunk = forged_sender_ask();
373        assert!(verify(&chunk).is_ok());
374        let rt = tiny_natives(chunk)?;
375        let idx = rt.function_index("main").ok_or("main")?;
376        let outcome = rt.spawn(idx, &[])?.join();
377        rt.shutdown();
378        match outcome {
379            FlowOutcome::Completed(Value::Int(n)) => {
380                assert_ne!(n, 999, "forged make_msg sender must not survive Ask");
381                assert!(n >= 1, "authenticated sender must be a live flow id");
382                Ok(())
383            }
384            other => Err(format!("expected Completed(Int), got {other:?}").into()),
385        }
386    }
387
388    #[test]
389    fn send_scalar_target_traps() -> Result<(), Box<dyn std::error::Error>> {
390        let mut p = Program::new("bad-cap-target");
391        p.function("main", 0, |f| {
392            let bad_cap = f.load_i32(99);
393            let sender = f.load_i32(0);
394            let req_id = f.load_i32(1);
395            let payload = f.load_i32(1);
396            let msg = f.make_msg(N_MAKE_MSG, sender, req_id, TAG_PING, payload);
397            f.send(bad_cap, msg);
398            f.return_(msg);
399        });
400        let rt = tiny_natives(p.build())?;
401        let outcome = rt.spawn(0, &[])?.join();
402        rt.shutdown();
403        assert!(
404            matches!(outcome, FlowOutcome::Failed(_)),
405            "non-Cap Send target must fail, got {outcome:?}"
406        );
407        Ok(())
408    }
409
410    #[test]
411    fn send_scalar_is_not_an_atomic_hop() -> Result<(), Box<dyn std::error::Error>> {
412        let mut p = Program::new("bad-hop");
413        p.function("main", 0, |f| {
414            let cap = f.self_cap();
415            let scalar = f.load_i32(99);
416            f.send(cap, scalar);
417            f.return_(scalar);
418        });
419        let rt = tiny(p.build())?;
420        let outcome = rt.spawn(0, &[])?.join();
421        rt.shutdown();
422        assert!(
423            matches!(outcome, FlowOutcome::Failed(_)),
424            "scalar Send must trap, got {outcome:?}"
425        );
426        Ok(())
427    }
428}