pub struct Runtime { /* private fields */ }Expand description
A running Byteflow runtime: worker pool + timer thread over one shared
Chunk.
Owns M:N scheduling for flows (spawn, yield, sleep, mailboxes,
FlowCap resolution, supervised restarts). Optional trace JIT when built
with feature = "jit" and enabled in RuntimeConfig::jit.
Construct with Runtime::new (no natives) or
Runtime::with_natives when the chunk uses CallNative /
crate::std_native_table.
Implementations§
Source§impl Runtime
impl Runtime
Sourcepub fn new(chunk: Chunk) -> Result<Self, SpawnError>
pub fn new(chunk: Chunk) -> Result<Self, SpawnError>
Convenience constructor for chunks that never call out through
Opcode::CallNative. Equivalent to
Runtime::with_natives(chunk, NativeTable::empty()).
Returns SpawnError instead of panicking: verify failures and OS
thread-spawn refusals are category-A errors (see
docs::error_model).
Sourcepub fn with_natives(
chunk: Chunk,
natives: Arc<NativeTable>,
) -> Result<Self, SpawnError>
pub fn with_natives( chunk: Chunk, natives: Arc<NativeTable>, ) -> Result<Self, SpawnError>
Construct a runtime whose flows can call into natives via
Opcode::CallNative — the host FFI boundary.
Sourcepub fn with_config(
chunk: Chunk,
config: RuntimeConfig,
) -> Result<Self, SpawnError>
pub fn with_config( chunk: Chunk, config: RuntimeConfig, ) -> Result<Self, SpawnError>
Examples found in repository?
22fn run_once(jit: bool) -> Result<(), Box<dyn std::error::Error>> {
23 let chunk = loop_chunk(1_000_000);
24 let rt = Runtime::with_config(
25 chunk,
26 RuntimeConfig {
27 workers: 1,
28 quantum: 50_000_000,
29 jit: JitConfig {
30 enabled: jit,
31 hot_threshold: 1,
32 },
33 ..Default::default()
34 },
35 )?;
36 let start = Instant::now();
37 let outcome = rt.spawn(0, &[])?.join();
38 let snapshot = rt.metrics();
39 rt.shutdown();
40 let elapsed = start.elapsed();
41 match outcome {
42 byteflow::FlowOutcome::Completed(Value::Int(n)) => {
43 println!("jit={jit} result={n} elapsed={elapsed:?} metrics={snapshot}");
44 Ok(())
45 }
46 other => Err(format!("unexpected outcome: {other:?}").into()),
47 }
48}More examples
15fn main() {
16 let rt = match Runtime::with_config(
17 samples::boom(),
18 RuntimeConfig {
19 workers: 1,
20 quantum: 1_000,
21 mailbox: byteflow::MailboxConfig::DEFAULT,
22 ..Default::default()
23 },
24 ) {
25 Ok(rt) => rt,
26 Err(e) => {
27 eprintln!("runtime: {e}");
28 std::process::exit(1);
29 }
30 };
31 let Some(boom) = rt.function_index("boom") else {
32 eprintln!("missing boom");
33 std::process::exit(1);
34 };
35 let sup = match Supervisor::with_config(
36 rt.spawner(),
37 SupervisorConfig {
38 max_restarts: 2,
39 max_period: Duration::from_secs(5),
40 strategy: byteflow::RestartStrategy::OneForOne,
41 },
42 ) {
43 Ok(s) => s,
44 Err(e) => {
45 eprintln!("supervisor: {e}");
46 std::process::exit(1);
47 }
48 };
49 if let Err(e) = sup.start_child(ChildSpec::new("boom", boom).restart(RestartPolicy::OnFailure))
50 {
51 eprintln!("start_child: {e}");
52 std::process::exit(1);
53 }
54
55 let start = std::time::Instant::now();
56 while !sup.intensity_exceeded() {
57 if start.elapsed() > Duration::from_secs(2) {
58 eprintln!("supervisor did not hit intensity in time");
59 std::process::exit(1);
60 }
61 thread::sleep(Duration::from_millis(5));
62 }
63
64 let metrics = rt.metrics();
65 println!("intensity exceeded after {} failures", metrics.processes_failed);
66 println!("{metrics}");
67 sup.shutdown();
68 rt.shutdown();
69}20fn main() {
21 let n: u32 = match std::env::args().nth(1) {
22 Some(s) => match s.parse() {
23 Ok(v) => v,
24 Err(_) => 50_000,
25 },
26 None => 50_000,
27 };
28
29 let workers = match std::thread::available_parallelism() {
30 Ok(p) => p.get(),
31 Err(_) => 1,
32 };
33
34 let rt = match Runtime::with_config(
35 trivial_chunk(),
36 RuntimeConfig {
37 workers,
38 quantum: 10_000,
39 mailbox: byteflow::MailboxConfig::DEFAULT,
40 ..Default::default()
41 },
42 ) {
43 Ok(rt) => rt,
44 Err(e) => {
45 eprintln!("runtime: {e}");
46 std::process::exit(1);
47 }
48 };
49 let Some(worker_fn) = rt.function_index("worker") else {
50 eprintln!("missing worker");
51 std::process::exit(1);
52 };
53
54 let start = Instant::now();
55 let mut handles = Vec::with_capacity(n as usize);
56 for _ in 0..n {
57 match rt.spawn(worker_fn, &[]) {
58 Ok(h) => handles.push(h),
59 Err(e) => {
60 eprintln!("spawn: {e}");
61 std::process::exit(1);
62 }
63 }
64 }
65 let mut ok = 0u32;
66 for h in handles {
67 if matches!(h.join(), FlowOutcome::Completed(Value::Int(1))) {
68 ok += 1;
69 }
70 }
71 let elapsed = start.elapsed();
72 let metrics = rt.metrics();
73 rt.shutdown();
74
75 let secs = elapsed.as_secs_f64().max(1e-9);
76 println!("processes={ok}/{n}");
77 println!("workers={workers}");
78 println!("elapsed_ms={:.2}", elapsed.as_secs_f64() * 1000.0);
79 println!("spawns_per_sec={:.0}", ok as f64 / secs);
80 println!("{metrics}");
81}Sourcepub fn with_std_natives_and_config(
chunk: Chunk,
config: RuntimeConfig,
) -> Result<Self, SpawnError>
pub fn with_std_natives_and_config( chunk: Chunk, config: RuntimeConfig, ) -> Result<Self, SpawnError>
Like Runtime::with_config but wires crate::std_native_table_with
using RuntimeConfig::output for the print native.
Sourcepub fn with_natives_and_config(
chunk: Chunk,
natives: Arc<NativeTable>,
config: RuntimeConfig,
) -> Result<Self, SpawnError>
pub fn with_natives_and_config( chunk: Chunk, natives: Arc<NativeTable>, config: RuntimeConfig, ) -> Result<Self, SpawnError>
Verify chunk, spawn the worker pool + timer thread, and return a
live Runtime.
Failures here mean the runtime was never started (no orphan
threads): either the bytecode is invalid
(SpawnError::VerifyFailed) or the OS refused a thread
(SpawnError::ThreadSpawnFailed).
Examples found in repository?
9fn main() {
10 let chunk = samples::ping_pong();
11 let rt = match Runtime::with_natives_and_config(
12 chunk,
13 std_native_table(),
14 RuntimeConfig {
15 workers: 1,
16 quantum: 10_000,
17 mailbox: byteflow::MailboxConfig::DEFAULT,
18 ..Default::default()
19 },
20 ) {
21 Ok(rt) => rt,
22 Err(e) => {
23 eprintln!("runtime: {e}");
24 std::process::exit(1);
25 }
26 };
27 let Some(main) = rt.function_index("main") else {
28 eprintln!("missing main");
29 std::process::exit(1);
30 };
31 let handle = match rt.spawn(main, &[]) {
32 Ok(h) => h,
33 Err(e) => {
34 eprintln!("spawn: {e}");
35 std::process::exit(1);
36 }
37 };
38 let outcome = handle.join();
39 let metrics = rt.metrics();
40 rt.shutdown();
41
42 match outcome {
43 FlowOutcome::Completed(Value::Int(2)) => {
44 println!("pong replied 2 (Atomic Hop)");
45 println!("{metrics}");
46 }
47 other => {
48 eprintln!("unexpected {other:?}");
49 std::process::exit(1);
50 }
51 }
52}More examples
21fn main() {
22 let chunk = samples::atomic_request_reply();
23 let rt = match Runtime::with_natives_and_config(
24 chunk,
25 std_native_table(),
26 RuntimeConfig {
27 workers: 2,
28 quantum: 10_000,
29 mailbox: byteflow::MailboxConfig::DEFAULT,
30 ..Default::default()
31 },
32 ) {
33 Ok(rt) => rt,
34 Err(e) => {
35 eprintln!("runtime: {e}");
36 std::process::exit(1);
37 }
38 };
39 let Some(main) = rt.function_index("main") else {
40 eprintln!("missing main");
41 std::process::exit(1);
42 };
43 let handle = match rt.spawn(main, &[]) {
44 Ok(h) => h,
45 Err(e) => {
46 eprintln!("spawn: {e}");
47 std::process::exit(1);
48 }
49 };
50 let outcome = handle.join();
51 let metrics = rt.metrics();
52 rt.shutdown();
53
54 match outcome {
55 FlowOutcome::Completed(Value::Int(42)) => {
56 println!("atomic request-reply ok: payload=42");
57 println!("{metrics}");
58 }
59 other => {
60 eprintln!("unexpected {other:?}");
61 std::process::exit(1);
62 }
63 }
64}Sourcepub fn spawn(
&self,
function: u32,
args: &[Value],
) -> Result<FlowHandle, SpawnError>
pub fn spawn( &self, function: u32, args: &[Value], ) -> Result<FlowHandle, SpawnError>
Spawn a top-level flow starting at function in this runtime’s
chunk, returning a FlowHandle the caller can .join().
Returns SpawnError::BadFunction if function is out of range.
Examples found in repository?
22fn run_once(jit: bool) -> Result<(), Box<dyn std::error::Error>> {
23 let chunk = loop_chunk(1_000_000);
24 let rt = Runtime::with_config(
25 chunk,
26 RuntimeConfig {
27 workers: 1,
28 quantum: 50_000_000,
29 jit: JitConfig {
30 enabled: jit,
31 hot_threshold: 1,
32 },
33 ..Default::default()
34 },
35 )?;
36 let start = Instant::now();
37 let outcome = rt.spawn(0, &[])?.join();
38 let snapshot = rt.metrics();
39 rt.shutdown();
40 let elapsed = start.elapsed();
41 match outcome {
42 byteflow::FlowOutcome::Completed(Value::Int(n)) => {
43 println!("jit={jit} result={n} elapsed={elapsed:?} metrics={snapshot}");
44 Ok(())
45 }
46 other => Err(format!("unexpected outcome: {other:?}").into()),
47 }
48}More examples
9fn main() {
10 let chunk = samples::ping_pong();
11 let rt = match Runtime::with_natives_and_config(
12 chunk,
13 std_native_table(),
14 RuntimeConfig {
15 workers: 1,
16 quantum: 10_000,
17 mailbox: byteflow::MailboxConfig::DEFAULT,
18 ..Default::default()
19 },
20 ) {
21 Ok(rt) => rt,
22 Err(e) => {
23 eprintln!("runtime: {e}");
24 std::process::exit(1);
25 }
26 };
27 let Some(main) = rt.function_index("main") else {
28 eprintln!("missing main");
29 std::process::exit(1);
30 };
31 let handle = match rt.spawn(main, &[]) {
32 Ok(h) => h,
33 Err(e) => {
34 eprintln!("spawn: {e}");
35 std::process::exit(1);
36 }
37 };
38 let outcome = handle.join();
39 let metrics = rt.metrics();
40 rt.shutdown();
41
42 match outcome {
43 FlowOutcome::Completed(Value::Int(2)) => {
44 println!("pong replied 2 (Atomic Hop)");
45 println!("{metrics}");
46 }
47 other => {
48 eprintln!("unexpected {other:?}");
49 std::process::exit(1);
50 }
51 }
52}21fn main() {
22 let chunk = samples::atomic_request_reply();
23 let rt = match Runtime::with_natives_and_config(
24 chunk,
25 std_native_table(),
26 RuntimeConfig {
27 workers: 2,
28 quantum: 10_000,
29 mailbox: byteflow::MailboxConfig::DEFAULT,
30 ..Default::default()
31 },
32 ) {
33 Ok(rt) => rt,
34 Err(e) => {
35 eprintln!("runtime: {e}");
36 std::process::exit(1);
37 }
38 };
39 let Some(main) = rt.function_index("main") else {
40 eprintln!("missing main");
41 std::process::exit(1);
42 };
43 let handle = match rt.spawn(main, &[]) {
44 Ok(h) => h,
45 Err(e) => {
46 eprintln!("spawn: {e}");
47 std::process::exit(1);
48 }
49 };
50 let outcome = handle.join();
51 let metrics = rt.metrics();
52 rt.shutdown();
53
54 match outcome {
55 FlowOutcome::Completed(Value::Int(42)) => {
56 println!("atomic request-reply ok: payload=42");
57 println!("{metrics}");
58 }
59 other => {
60 eprintln!("unexpected {other:?}");
61 std::process::exit(1);
62 }
63 }
64}20fn main() {
21 let n: u32 = match std::env::args().nth(1) {
22 Some(s) => match s.parse() {
23 Ok(v) => v,
24 Err(_) => 50_000,
25 },
26 None => 50_000,
27 };
28
29 let workers = match std::thread::available_parallelism() {
30 Ok(p) => p.get(),
31 Err(_) => 1,
32 };
33
34 let rt = match Runtime::with_config(
35 trivial_chunk(),
36 RuntimeConfig {
37 workers,
38 quantum: 10_000,
39 mailbox: byteflow::MailboxConfig::DEFAULT,
40 ..Default::default()
41 },
42 ) {
43 Ok(rt) => rt,
44 Err(e) => {
45 eprintln!("runtime: {e}");
46 std::process::exit(1);
47 }
48 };
49 let Some(worker_fn) = rt.function_index("worker") else {
50 eprintln!("missing worker");
51 std::process::exit(1);
52 };
53
54 let start = Instant::now();
55 let mut handles = Vec::with_capacity(n as usize);
56 for _ in 0..n {
57 match rt.spawn(worker_fn, &[]) {
58 Ok(h) => handles.push(h),
59 Err(e) => {
60 eprintln!("spawn: {e}");
61 std::process::exit(1);
62 }
63 }
64 }
65 let mut ok = 0u32;
66 for h in handles {
67 if matches!(h.join(), FlowOutcome::Completed(Value::Int(1))) {
68 ok += 1;
69 }
70 }
71 let elapsed = start.elapsed();
72 let metrics = rt.metrics();
73 rt.shutdown();
74
75 let secs = elapsed.as_secs_f64().max(1e-9);
76 println!("processes={ok}/{n}");
77 println!("workers={workers}");
78 println!("elapsed_ms={:.2}", elapsed.as_secs_f64() * 1000.0);
79 println!("spawns_per_sec={:.0}", ok as f64 / secs);
80 println!("{metrics}");
81}Sourcepub fn spawner(&self) -> RuntimeSpawner
pub fn spawner(&self) -> RuntimeSpawner
A cheap, Send + Sync handle that can spawn processes into this
runtime from any thread, independent of Runtime’s own lifetime
bookkeeping (worker JoinHandles). Used by super::supervisor::Supervisor.
Examples found in repository?
15fn main() {
16 let rt = match Runtime::with_config(
17 samples::boom(),
18 RuntimeConfig {
19 workers: 1,
20 quantum: 1_000,
21 mailbox: byteflow::MailboxConfig::DEFAULT,
22 ..Default::default()
23 },
24 ) {
25 Ok(rt) => rt,
26 Err(e) => {
27 eprintln!("runtime: {e}");
28 std::process::exit(1);
29 }
30 };
31 let Some(boom) = rt.function_index("boom") else {
32 eprintln!("missing boom");
33 std::process::exit(1);
34 };
35 let sup = match Supervisor::with_config(
36 rt.spawner(),
37 SupervisorConfig {
38 max_restarts: 2,
39 max_period: Duration::from_secs(5),
40 strategy: byteflow::RestartStrategy::OneForOne,
41 },
42 ) {
43 Ok(s) => s,
44 Err(e) => {
45 eprintln!("supervisor: {e}");
46 std::process::exit(1);
47 }
48 };
49 if let Err(e) = sup.start_child(ChildSpec::new("boom", boom).restart(RestartPolicy::OnFailure))
50 {
51 eprintln!("start_child: {e}");
52 std::process::exit(1);
53 }
54
55 let start = std::time::Instant::now();
56 while !sup.intensity_exceeded() {
57 if start.elapsed() > Duration::from_secs(2) {
58 eprintln!("supervisor did not hit intensity in time");
59 std::process::exit(1);
60 }
61 thread::sleep(Duration::from_millis(5));
62 }
63
64 let metrics = rt.metrics();
65 println!("intensity exceeded after {} failures", metrics.processes_failed);
66 println!("{metrics}");
67 sup.shutdown();
68 rt.shutdown();
69}Sourcepub fn supervisor(&self) -> Result<Supervisor, SpawnError>
pub fn supervisor(&self) -> Result<Supervisor, SpawnError>
A super::supervisor::Supervisor bound to this runtime, ready to
take supervised children (design notes §15).
Sourcepub fn function_index(&self, name: &str) -> Option<u32>
pub fn function_index(&self, name: &str) -> Option<u32>
Look up a function by name in the runtime’s chunk — convenience for
callers that built their chunk with crate::Program
and don’t want to thread raw indices through their own code.
Examples found in repository?
9fn main() {
10 let chunk = samples::ping_pong();
11 let rt = match Runtime::with_natives_and_config(
12 chunk,
13 std_native_table(),
14 RuntimeConfig {
15 workers: 1,
16 quantum: 10_000,
17 mailbox: byteflow::MailboxConfig::DEFAULT,
18 ..Default::default()
19 },
20 ) {
21 Ok(rt) => rt,
22 Err(e) => {
23 eprintln!("runtime: {e}");
24 std::process::exit(1);
25 }
26 };
27 let Some(main) = rt.function_index("main") else {
28 eprintln!("missing main");
29 std::process::exit(1);
30 };
31 let handle = match rt.spawn(main, &[]) {
32 Ok(h) => h,
33 Err(e) => {
34 eprintln!("spawn: {e}");
35 std::process::exit(1);
36 }
37 };
38 let outcome = handle.join();
39 let metrics = rt.metrics();
40 rt.shutdown();
41
42 match outcome {
43 FlowOutcome::Completed(Value::Int(2)) => {
44 println!("pong replied 2 (Atomic Hop)");
45 println!("{metrics}");
46 }
47 other => {
48 eprintln!("unexpected {other:?}");
49 std::process::exit(1);
50 }
51 }
52}More examples
21fn main() {
22 let chunk = samples::atomic_request_reply();
23 let rt = match Runtime::with_natives_and_config(
24 chunk,
25 std_native_table(),
26 RuntimeConfig {
27 workers: 2,
28 quantum: 10_000,
29 mailbox: byteflow::MailboxConfig::DEFAULT,
30 ..Default::default()
31 },
32 ) {
33 Ok(rt) => rt,
34 Err(e) => {
35 eprintln!("runtime: {e}");
36 std::process::exit(1);
37 }
38 };
39 let Some(main) = rt.function_index("main") else {
40 eprintln!("missing main");
41 std::process::exit(1);
42 };
43 let handle = match rt.spawn(main, &[]) {
44 Ok(h) => h,
45 Err(e) => {
46 eprintln!("spawn: {e}");
47 std::process::exit(1);
48 }
49 };
50 let outcome = handle.join();
51 let metrics = rt.metrics();
52 rt.shutdown();
53
54 match outcome {
55 FlowOutcome::Completed(Value::Int(42)) => {
56 println!("atomic request-reply ok: payload=42");
57 println!("{metrics}");
58 }
59 other => {
60 eprintln!("unexpected {other:?}");
61 std::process::exit(1);
62 }
63 }
64}15fn main() {
16 let rt = match Runtime::with_config(
17 samples::boom(),
18 RuntimeConfig {
19 workers: 1,
20 quantum: 1_000,
21 mailbox: byteflow::MailboxConfig::DEFAULT,
22 ..Default::default()
23 },
24 ) {
25 Ok(rt) => rt,
26 Err(e) => {
27 eprintln!("runtime: {e}");
28 std::process::exit(1);
29 }
30 };
31 let Some(boom) = rt.function_index("boom") else {
32 eprintln!("missing boom");
33 std::process::exit(1);
34 };
35 let sup = match Supervisor::with_config(
36 rt.spawner(),
37 SupervisorConfig {
38 max_restarts: 2,
39 max_period: Duration::from_secs(5),
40 strategy: byteflow::RestartStrategy::OneForOne,
41 },
42 ) {
43 Ok(s) => s,
44 Err(e) => {
45 eprintln!("supervisor: {e}");
46 std::process::exit(1);
47 }
48 };
49 if let Err(e) = sup.start_child(ChildSpec::new("boom", boom).restart(RestartPolicy::OnFailure))
50 {
51 eprintln!("start_child: {e}");
52 std::process::exit(1);
53 }
54
55 let start = std::time::Instant::now();
56 while !sup.intensity_exceeded() {
57 if start.elapsed() > Duration::from_secs(2) {
58 eprintln!("supervisor did not hit intensity in time");
59 std::process::exit(1);
60 }
61 thread::sleep(Duration::from_millis(5));
62 }
63
64 let metrics = rt.metrics();
65 println!("intensity exceeded after {} failures", metrics.processes_failed);
66 println!("{metrics}");
67 sup.shutdown();
68 rt.shutdown();
69}20fn main() {
21 let n: u32 = match std::env::args().nth(1) {
22 Some(s) => match s.parse() {
23 Ok(v) => v,
24 Err(_) => 50_000,
25 },
26 None => 50_000,
27 };
28
29 let workers = match std::thread::available_parallelism() {
30 Ok(p) => p.get(),
31 Err(_) => 1,
32 };
33
34 let rt = match Runtime::with_config(
35 trivial_chunk(),
36 RuntimeConfig {
37 workers,
38 quantum: 10_000,
39 mailbox: byteflow::MailboxConfig::DEFAULT,
40 ..Default::default()
41 },
42 ) {
43 Ok(rt) => rt,
44 Err(e) => {
45 eprintln!("runtime: {e}");
46 std::process::exit(1);
47 }
48 };
49 let Some(worker_fn) = rt.function_index("worker") else {
50 eprintln!("missing worker");
51 std::process::exit(1);
52 };
53
54 let start = Instant::now();
55 let mut handles = Vec::with_capacity(n as usize);
56 for _ in 0..n {
57 match rt.spawn(worker_fn, &[]) {
58 Ok(h) => handles.push(h),
59 Err(e) => {
60 eprintln!("spawn: {e}");
61 std::process::exit(1);
62 }
63 }
64 }
65 let mut ok = 0u32;
66 for h in handles {
67 if matches!(h.join(), FlowOutcome::Completed(Value::Int(1))) {
68 ok += 1;
69 }
70 }
71 let elapsed = start.elapsed();
72 let metrics = rt.metrics();
73 rt.shutdown();
74
75 let secs = elapsed.as_secs_f64().max(1e-9);
76 println!("processes={ok}/{n}");
77 println!("workers={workers}");
78 println!("elapsed_ms={:.2}", elapsed.as_secs_f64() * 1000.0);
79 println!("spawns_per_sec={:.0}", ok as f64 / secs);
80 println!("{metrics}");
81}Sourcepub fn metrics(&self) -> RuntimeMetricsSnapshot
pub fn metrics(&self) -> RuntimeMetricsSnapshot
Examples found in repository?
22fn run_once(jit: bool) -> Result<(), Box<dyn std::error::Error>> {
23 let chunk = loop_chunk(1_000_000);
24 let rt = Runtime::with_config(
25 chunk,
26 RuntimeConfig {
27 workers: 1,
28 quantum: 50_000_000,
29 jit: JitConfig {
30 enabled: jit,
31 hot_threshold: 1,
32 },
33 ..Default::default()
34 },
35 )?;
36 let start = Instant::now();
37 let outcome = rt.spawn(0, &[])?.join();
38 let snapshot = rt.metrics();
39 rt.shutdown();
40 let elapsed = start.elapsed();
41 match outcome {
42 byteflow::FlowOutcome::Completed(Value::Int(n)) => {
43 println!("jit={jit} result={n} elapsed={elapsed:?} metrics={snapshot}");
44 Ok(())
45 }
46 other => Err(format!("unexpected outcome: {other:?}").into()),
47 }
48}More examples
9fn main() {
10 let chunk = samples::ping_pong();
11 let rt = match Runtime::with_natives_and_config(
12 chunk,
13 std_native_table(),
14 RuntimeConfig {
15 workers: 1,
16 quantum: 10_000,
17 mailbox: byteflow::MailboxConfig::DEFAULT,
18 ..Default::default()
19 },
20 ) {
21 Ok(rt) => rt,
22 Err(e) => {
23 eprintln!("runtime: {e}");
24 std::process::exit(1);
25 }
26 };
27 let Some(main) = rt.function_index("main") else {
28 eprintln!("missing main");
29 std::process::exit(1);
30 };
31 let handle = match rt.spawn(main, &[]) {
32 Ok(h) => h,
33 Err(e) => {
34 eprintln!("spawn: {e}");
35 std::process::exit(1);
36 }
37 };
38 let outcome = handle.join();
39 let metrics = rt.metrics();
40 rt.shutdown();
41
42 match outcome {
43 FlowOutcome::Completed(Value::Int(2)) => {
44 println!("pong replied 2 (Atomic Hop)");
45 println!("{metrics}");
46 }
47 other => {
48 eprintln!("unexpected {other:?}");
49 std::process::exit(1);
50 }
51 }
52}21fn main() {
22 let chunk = samples::atomic_request_reply();
23 let rt = match Runtime::with_natives_and_config(
24 chunk,
25 std_native_table(),
26 RuntimeConfig {
27 workers: 2,
28 quantum: 10_000,
29 mailbox: byteflow::MailboxConfig::DEFAULT,
30 ..Default::default()
31 },
32 ) {
33 Ok(rt) => rt,
34 Err(e) => {
35 eprintln!("runtime: {e}");
36 std::process::exit(1);
37 }
38 };
39 let Some(main) = rt.function_index("main") else {
40 eprintln!("missing main");
41 std::process::exit(1);
42 };
43 let handle = match rt.spawn(main, &[]) {
44 Ok(h) => h,
45 Err(e) => {
46 eprintln!("spawn: {e}");
47 std::process::exit(1);
48 }
49 };
50 let outcome = handle.join();
51 let metrics = rt.metrics();
52 rt.shutdown();
53
54 match outcome {
55 FlowOutcome::Completed(Value::Int(42)) => {
56 println!("atomic request-reply ok: payload=42");
57 println!("{metrics}");
58 }
59 other => {
60 eprintln!("unexpected {other:?}");
61 std::process::exit(1);
62 }
63 }
64}15fn main() {
16 let rt = match Runtime::with_config(
17 samples::boom(),
18 RuntimeConfig {
19 workers: 1,
20 quantum: 1_000,
21 mailbox: byteflow::MailboxConfig::DEFAULT,
22 ..Default::default()
23 },
24 ) {
25 Ok(rt) => rt,
26 Err(e) => {
27 eprintln!("runtime: {e}");
28 std::process::exit(1);
29 }
30 };
31 let Some(boom) = rt.function_index("boom") else {
32 eprintln!("missing boom");
33 std::process::exit(1);
34 };
35 let sup = match Supervisor::with_config(
36 rt.spawner(),
37 SupervisorConfig {
38 max_restarts: 2,
39 max_period: Duration::from_secs(5),
40 strategy: byteflow::RestartStrategy::OneForOne,
41 },
42 ) {
43 Ok(s) => s,
44 Err(e) => {
45 eprintln!("supervisor: {e}");
46 std::process::exit(1);
47 }
48 };
49 if let Err(e) = sup.start_child(ChildSpec::new("boom", boom).restart(RestartPolicy::OnFailure))
50 {
51 eprintln!("start_child: {e}");
52 std::process::exit(1);
53 }
54
55 let start = std::time::Instant::now();
56 while !sup.intensity_exceeded() {
57 if start.elapsed() > Duration::from_secs(2) {
58 eprintln!("supervisor did not hit intensity in time");
59 std::process::exit(1);
60 }
61 thread::sleep(Duration::from_millis(5));
62 }
63
64 let metrics = rt.metrics();
65 println!("intensity exceeded after {} failures", metrics.processes_failed);
66 println!("{metrics}");
67 sup.shutdown();
68 rt.shutdown();
69}20fn main() {
21 let n: u32 = match std::env::args().nth(1) {
22 Some(s) => match s.parse() {
23 Ok(v) => v,
24 Err(_) => 50_000,
25 },
26 None => 50_000,
27 };
28
29 let workers = match std::thread::available_parallelism() {
30 Ok(p) => p.get(),
31 Err(_) => 1,
32 };
33
34 let rt = match Runtime::with_config(
35 trivial_chunk(),
36 RuntimeConfig {
37 workers,
38 quantum: 10_000,
39 mailbox: byteflow::MailboxConfig::DEFAULT,
40 ..Default::default()
41 },
42 ) {
43 Ok(rt) => rt,
44 Err(e) => {
45 eprintln!("runtime: {e}");
46 std::process::exit(1);
47 }
48 };
49 let Some(worker_fn) = rt.function_index("worker") else {
50 eprintln!("missing worker");
51 std::process::exit(1);
52 };
53
54 let start = Instant::now();
55 let mut handles = Vec::with_capacity(n as usize);
56 for _ in 0..n {
57 match rt.spawn(worker_fn, &[]) {
58 Ok(h) => handles.push(h),
59 Err(e) => {
60 eprintln!("spawn: {e}");
61 std::process::exit(1);
62 }
63 }
64 }
65 let mut ok = 0u32;
66 for h in handles {
67 if matches!(h.join(), FlowOutcome::Completed(Value::Int(1))) {
68 ok += 1;
69 }
70 }
71 let elapsed = start.elapsed();
72 let metrics = rt.metrics();
73 rt.shutdown();
74
75 let secs = elapsed.as_secs_f64().max(1e-9);
76 println!("processes={ok}/{n}");
77 println!("workers={workers}");
78 println!("elapsed_ms={:.2}", elapsed.as_secs_f64() * 1000.0);
79 println!("spawns_per_sec={:.0}", ok as f64 / secs);
80 println!("{metrics}");
81}Sourcepub fn live_flows(&self) -> usize
pub fn live_flows(&self) -> usize
Number of flows currently registered in the directory — i.e. alive (running, ready, sleeping, or waiting), not counting ones that have already completed or failed.
pub fn worker_count(&self) -> usize
Sourcepub fn send(&self, target: FlowId, message: Value) -> Result<(), SendError>
pub fn send(&self, target: FlowId, message: Value) -> Result<(), SendError>
Deliver an Atomic Hop (Value::Message) to target from the
embedder (not from bytecode).
§Host trust boundary
This path takes a FlowId directly — no Cap required. The host
is trusted; bytecode must use Value::Cap via Opcode::Send /
Ask. Host-injected messages are not re-stamped (sender /
reply_cap stay as built). Bare scalars are rejected
(SendError::NotAHop).
Sourcepub fn mint_cap(&self, flow: FlowId) -> Result<CapId, LifecycleError>
pub fn mint_cap(&self, flow: FlowId) -> Result<CapId, LifecycleError>
Mint a SEND|ASK Cap for a live flow (host equivalent of SelfPid).
Sourcepub fn grant_cap(
&self,
holder: FlowId,
target: FlowId,
) -> Result<CapId, LifecycleError>
pub fn grant_cap( &self, holder: FlowId, target: FlowId, ) -> Result<CapId, LifecycleError>
Mint a Cap held by holder targeting target (host introduction).
Sourcepub fn monitor(
&self,
owner: FlowId,
target: FlowId,
) -> Result<MonitorRef, LifecycleError>
pub fn monitor( &self, owner: FlowId, target: FlowId, ) -> Result<MonitorRef, LifecycleError>
Watch target; when it exits, owner receives a crate::TAG_SYS_DOWN hop.
Both flows must be live. owner == target is [LifecycleError::SelfRelation].
Sourcepub fn demonitor(
&self,
owner: FlowId,
monitor: MonitorRef,
) -> Result<(), LifecycleError>
pub fn demonitor( &self, owner: FlowId, monitor: MonitorRef, ) -> Result<(), LifecycleError>
Drop monitor if owner still owns it.
Sourcepub fn link(&self, a: FlowId, b: FlowId) -> Result<LinkId, LifecycleError>
pub fn link(&self, a: FlowId, b: FlowId) -> Result<LinkId, LifecycleError>
Bidirectional link. Abnormal exit of either side kills the peer.
Sourcepub fn unlink(&self, owner: FlowId, link: LinkId) -> Result<(), LifecycleError>
pub fn unlink(&self, owner: FlowId, link: LinkId) -> Result<(), LifecycleError>
Drop link if owner is one of the endpoints.
Sourcepub fn register_name(
&self,
name: &str,
cap: CapId,
) -> Result<(), LifecycleError>
pub fn register_name( &self, name: &str, cap: CapId, ) -> Result<(), LifecycleError>
Bind name to a live Cap (address). Names are swept when that flow exits.
Sourcepub fn whereis(&self, name: &str) -> Result<Option<CapId>, LifecycleError>
pub fn whereis(&self, name: &str) -> Result<Option<CapId>, LifecycleError>
Look up a registered Cap, or None if the name is free / was swept.
Sourcepub fn kill(&self, id: FlowId) -> Result<(), LifecycleError>
pub fn kill(&self, id: FlowId) -> Result<(), LifecycleError>
Cooperative abort. Parked flows finalize immediately; a running flow
dies at the next quantum with super::monitor::FlowExitReason::Killed.
Sourcepub fn mint_admin_cap(&self, holder: FlowId) -> Result<CapId, LifecycleError>
pub fn mint_admin_cap(&self, holder: FlowId) -> Result<CapId, LifecycleError>
Mint a scheduler ADMIN cap for a live flow (host introduction).
Sourcepub fn admin_kill(
&self,
holder: FlowId,
cap: CapId,
target: FlowId,
) -> Result<(), LifecycleError>
pub fn admin_kill( &self, holder: FlowId, cap: CapId, target: FlowId, ) -> Result<(), LifecycleError>
Kill target only if holder presents a live ADMIN scheduler cap.
Sourcepub fn admin_top_up_cpu(
&self,
holder: FlowId,
cap: CapId,
target: FlowId,
extra: i64,
) -> Result<(), LifecycleError>
pub fn admin_top_up_cpu( &self, holder: FlowId, cap: CapId, target: FlowId, extra: i64, ) -> Result<(), LifecycleError>
Top up target’s CPU budget. Requires a live ADMIN scheduler cap.
Sourcepub fn unregister_name(&self, name: &str) -> Result<bool, LifecycleError>
pub fn unregister_name(&self, name: &str) -> Result<bool, LifecycleError>
Remove a name without waiting for the flow to exit. false if unknown.
Sourcepub fn reload_chunk(&mut self, chunk: Chunk) -> Result<(), SpawnError>
pub fn reload_chunk(&mut self, chunk: Chunk) -> Result<(), SpawnError>
Replace the image used by new host Self::spawn calls and
drop JIT traces. Live flows and bytecode Spawn keep the parent’s
existing Vm chunk.
Sourcepub fn shutdown(self)
pub fn shutdown(self)
Stop accepting new scheduling work and join every worker + the timer thread. Processes that are mid-quantum are allowed to reach their next natural suspension point; this does not forcibly abort running bytecode (there is no safe way to do that to an OS thread mid-instruction — see design notes §11 on why preemption here is cooperative/budgeted rather than signal-based).
Examples found in repository?
22fn run_once(jit: bool) -> Result<(), Box<dyn std::error::Error>> {
23 let chunk = loop_chunk(1_000_000);
24 let rt = Runtime::with_config(
25 chunk,
26 RuntimeConfig {
27 workers: 1,
28 quantum: 50_000_000,
29 jit: JitConfig {
30 enabled: jit,
31 hot_threshold: 1,
32 },
33 ..Default::default()
34 },
35 )?;
36 let start = Instant::now();
37 let outcome = rt.spawn(0, &[])?.join();
38 let snapshot = rt.metrics();
39 rt.shutdown();
40 let elapsed = start.elapsed();
41 match outcome {
42 byteflow::FlowOutcome::Completed(Value::Int(n)) => {
43 println!("jit={jit} result={n} elapsed={elapsed:?} metrics={snapshot}");
44 Ok(())
45 }
46 other => Err(format!("unexpected outcome: {other:?}").into()),
47 }
48}More examples
9fn main() {
10 let chunk = samples::ping_pong();
11 let rt = match Runtime::with_natives_and_config(
12 chunk,
13 std_native_table(),
14 RuntimeConfig {
15 workers: 1,
16 quantum: 10_000,
17 mailbox: byteflow::MailboxConfig::DEFAULT,
18 ..Default::default()
19 },
20 ) {
21 Ok(rt) => rt,
22 Err(e) => {
23 eprintln!("runtime: {e}");
24 std::process::exit(1);
25 }
26 };
27 let Some(main) = rt.function_index("main") else {
28 eprintln!("missing main");
29 std::process::exit(1);
30 };
31 let handle = match rt.spawn(main, &[]) {
32 Ok(h) => h,
33 Err(e) => {
34 eprintln!("spawn: {e}");
35 std::process::exit(1);
36 }
37 };
38 let outcome = handle.join();
39 let metrics = rt.metrics();
40 rt.shutdown();
41
42 match outcome {
43 FlowOutcome::Completed(Value::Int(2)) => {
44 println!("pong replied 2 (Atomic Hop)");
45 println!("{metrics}");
46 }
47 other => {
48 eprintln!("unexpected {other:?}");
49 std::process::exit(1);
50 }
51 }
52}21fn main() {
22 let chunk = samples::atomic_request_reply();
23 let rt = match Runtime::with_natives_and_config(
24 chunk,
25 std_native_table(),
26 RuntimeConfig {
27 workers: 2,
28 quantum: 10_000,
29 mailbox: byteflow::MailboxConfig::DEFAULT,
30 ..Default::default()
31 },
32 ) {
33 Ok(rt) => rt,
34 Err(e) => {
35 eprintln!("runtime: {e}");
36 std::process::exit(1);
37 }
38 };
39 let Some(main) = rt.function_index("main") else {
40 eprintln!("missing main");
41 std::process::exit(1);
42 };
43 let handle = match rt.spawn(main, &[]) {
44 Ok(h) => h,
45 Err(e) => {
46 eprintln!("spawn: {e}");
47 std::process::exit(1);
48 }
49 };
50 let outcome = handle.join();
51 let metrics = rt.metrics();
52 rt.shutdown();
53
54 match outcome {
55 FlowOutcome::Completed(Value::Int(42)) => {
56 println!("atomic request-reply ok: payload=42");
57 println!("{metrics}");
58 }
59 other => {
60 eprintln!("unexpected {other:?}");
61 std::process::exit(1);
62 }
63 }
64}15fn main() {
16 let rt = match Runtime::with_config(
17 samples::boom(),
18 RuntimeConfig {
19 workers: 1,
20 quantum: 1_000,
21 mailbox: byteflow::MailboxConfig::DEFAULT,
22 ..Default::default()
23 },
24 ) {
25 Ok(rt) => rt,
26 Err(e) => {
27 eprintln!("runtime: {e}");
28 std::process::exit(1);
29 }
30 };
31 let Some(boom) = rt.function_index("boom") else {
32 eprintln!("missing boom");
33 std::process::exit(1);
34 };
35 let sup = match Supervisor::with_config(
36 rt.spawner(),
37 SupervisorConfig {
38 max_restarts: 2,
39 max_period: Duration::from_secs(5),
40 strategy: byteflow::RestartStrategy::OneForOne,
41 },
42 ) {
43 Ok(s) => s,
44 Err(e) => {
45 eprintln!("supervisor: {e}");
46 std::process::exit(1);
47 }
48 };
49 if let Err(e) = sup.start_child(ChildSpec::new("boom", boom).restart(RestartPolicy::OnFailure))
50 {
51 eprintln!("start_child: {e}");
52 std::process::exit(1);
53 }
54
55 let start = std::time::Instant::now();
56 while !sup.intensity_exceeded() {
57 if start.elapsed() > Duration::from_secs(2) {
58 eprintln!("supervisor did not hit intensity in time");
59 std::process::exit(1);
60 }
61 thread::sleep(Duration::from_millis(5));
62 }
63
64 let metrics = rt.metrics();
65 println!("intensity exceeded after {} failures", metrics.processes_failed);
66 println!("{metrics}");
67 sup.shutdown();
68 rt.shutdown();
69}20fn main() {
21 let n: u32 = match std::env::args().nth(1) {
22 Some(s) => match s.parse() {
23 Ok(v) => v,
24 Err(_) => 50_000,
25 },
26 None => 50_000,
27 };
28
29 let workers = match std::thread::available_parallelism() {
30 Ok(p) => p.get(),
31 Err(_) => 1,
32 };
33
34 let rt = match Runtime::with_config(
35 trivial_chunk(),
36 RuntimeConfig {
37 workers,
38 quantum: 10_000,
39 mailbox: byteflow::MailboxConfig::DEFAULT,
40 ..Default::default()
41 },
42 ) {
43 Ok(rt) => rt,
44 Err(e) => {
45 eprintln!("runtime: {e}");
46 std::process::exit(1);
47 }
48 };
49 let Some(worker_fn) = rt.function_index("worker") else {
50 eprintln!("missing worker");
51 std::process::exit(1);
52 };
53
54 let start = Instant::now();
55 let mut handles = Vec::with_capacity(n as usize);
56 for _ in 0..n {
57 match rt.spawn(worker_fn, &[]) {
58 Ok(h) => handles.push(h),
59 Err(e) => {
60 eprintln!("spawn: {e}");
61 std::process::exit(1);
62 }
63 }
64 }
65 let mut ok = 0u32;
66 for h in handles {
67 if matches!(h.join(), FlowOutcome::Completed(Value::Int(1))) {
68 ok += 1;
69 }
70 }
71 let elapsed = start.elapsed();
72 let metrics = rt.metrics();
73 rt.shutdown();
74
75 let secs = elapsed.as_secs_f64().max(1e-9);
76 println!("processes={ok}/{n}");
77 println!("workers={workers}");
78 println!("elapsed_ms={:.2}", elapsed.as_secs_f64() * 1000.0);
79 println!("spawns_per_sec={:.0}", ok as f64 / secs);
80 println!("{metrics}");
81}