Compare commits
3 commits
ca34b2f778
...
6cb5044d06
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
6cb5044d06 | ||
|
|
0601e61290 | ||
| 85c7c557ee |
16 changed files with 1652 additions and 88 deletions
1
.gitignore
vendored
1
.gitignore
vendored
|
|
@ -1,2 +1,3 @@
|
||||||
/target
|
/target
|
||||||
.vscode/
|
.vscode/
|
||||||
|
.venv
|
||||||
106
Cargo.lock
generated
106
Cargo.lock
generated
|
|
@ -203,12 +203,27 @@ dependencies = [
|
||||||
"zerocopy",
|
"zerocopy",
|
||||||
]
|
]
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "heck"
|
||||||
|
version = "0.5.0"
|
||||||
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
|
checksum = "2304e00983f87ffb38b55b444b5e3b60a884b5d30c0fca7d82fe33449bbe55ea"
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "hermit-abi"
|
name = "hermit-abi"
|
||||||
version = "0.5.2"
|
version = "0.5.2"
|
||||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
checksum = "fc0fef456e4baa96da950455cd02c081ca953b141298e41db3fc7e36b1da849c"
|
checksum = "fc0fef456e4baa96da950455cd02c081ca953b141298e41db3fc7e36b1da849c"
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "indoc"
|
||||||
|
version = "2.0.7"
|
||||||
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
|
checksum = "79cf5c93f93228cf8efb3ba362535fb11199ac548a09ce117c9b1adc3030d706"
|
||||||
|
dependencies = [
|
||||||
|
"rustversion",
|
||||||
|
]
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "is-terminal"
|
name = "is-terminal"
|
||||||
version = "0.4.17"
|
version = "0.4.17"
|
||||||
|
|
@ -257,6 +272,15 @@ version = "2.7.6"
|
||||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
checksum = "f52b00d39961fc5b2736ea853c9cc86238e165017a493d1d5c8eac6bdc4cc273"
|
checksum = "f52b00d39961fc5b2736ea853c9cc86238e165017a493d1d5c8eac6bdc4cc273"
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "memoffset"
|
||||||
|
version = "0.9.1"
|
||||||
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
|
checksum = "488016bfae457b036d996092f6cb448677611ce4449e970ceaf42695203f218a"
|
||||||
|
dependencies = [
|
||||||
|
"autocfg",
|
||||||
|
]
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "num-traits"
|
name = "num-traits"
|
||||||
version = "0.2.19"
|
version = "0.2.19"
|
||||||
|
|
@ -306,6 +330,12 @@ dependencies = [
|
||||||
"plotters-backend",
|
"plotters-backend",
|
||||||
]
|
]
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "portable-atomic"
|
||||||
|
version = "1.13.1"
|
||||||
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
|
checksum = "c33a9471896f1c69cecef8d20cbe2f7accd12527ce60845ff44c153bb2a21b49"
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "proc-macro2"
|
name = "proc-macro2"
|
||||||
version = "1.0.106"
|
version = "1.0.106"
|
||||||
|
|
@ -315,6 +345,69 @@ dependencies = [
|
||||||
"unicode-ident",
|
"unicode-ident",
|
||||||
]
|
]
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "pyo3"
|
||||||
|
version = "0.23.5"
|
||||||
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
|
checksum = "7778bffd85cf38175ac1f545509665d0b9b92a198ca7941f131f85f7a4f9a872"
|
||||||
|
dependencies = [
|
||||||
|
"cfg-if",
|
||||||
|
"indoc",
|
||||||
|
"libc",
|
||||||
|
"memoffset",
|
||||||
|
"once_cell",
|
||||||
|
"portable-atomic",
|
||||||
|
"pyo3-build-config",
|
||||||
|
"pyo3-ffi",
|
||||||
|
"pyo3-macros",
|
||||||
|
"unindent",
|
||||||
|
]
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "pyo3-build-config"
|
||||||
|
version = "0.23.5"
|
||||||
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
|
checksum = "94f6cbe86ef3bf18998d9df6e0f3fc1050a8c5efa409bf712e661a4366e010fb"
|
||||||
|
dependencies = [
|
||||||
|
"once_cell",
|
||||||
|
"target-lexicon",
|
||||||
|
]
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "pyo3-ffi"
|
||||||
|
version = "0.23.5"
|
||||||
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
|
checksum = "e9f1b4c431c0bb1c8fb0a338709859eed0d030ff6daa34368d3b152a63dfdd8d"
|
||||||
|
dependencies = [
|
||||||
|
"libc",
|
||||||
|
"pyo3-build-config",
|
||||||
|
]
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "pyo3-macros"
|
||||||
|
version = "0.23.5"
|
||||||
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
|
checksum = "fbc2201328f63c4710f68abdf653c89d8dbc2858b88c5d88b0ff38a75288a9da"
|
||||||
|
dependencies = [
|
||||||
|
"proc-macro2",
|
||||||
|
"pyo3-macros-backend",
|
||||||
|
"quote",
|
||||||
|
"syn",
|
||||||
|
]
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "pyo3-macros-backend"
|
||||||
|
version = "0.23.5"
|
||||||
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
|
checksum = "fca6726ad0f3da9c9de093d6f116a93c1a38e417ed73bf138472cf4064f72028"
|
||||||
|
dependencies = [
|
||||||
|
"heck",
|
||||||
|
"proc-macro2",
|
||||||
|
"pyo3-build-config",
|
||||||
|
"quote",
|
||||||
|
"syn",
|
||||||
|
]
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "quote"
|
name = "quote"
|
||||||
version = "1.0.44"
|
version = "1.0.44"
|
||||||
|
|
@ -439,6 +532,7 @@ dependencies = [
|
||||||
"crossbeam-queue",
|
"crossbeam-queue",
|
||||||
"crossbeam-utils",
|
"crossbeam-utils",
|
||||||
"getrandom",
|
"getrandom",
|
||||||
|
"pyo3",
|
||||||
]
|
]
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
|
|
@ -452,6 +546,12 @@ dependencies = [
|
||||||
"unicode-ident",
|
"unicode-ident",
|
||||||
]
|
]
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "target-lexicon"
|
||||||
|
version = "0.12.16"
|
||||||
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
|
checksum = "61c41af27dd6d1e27b1b16b489db798443478cef1f06a660c96db617ba5de3b1"
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "tinytemplate"
|
name = "tinytemplate"
|
||||||
version = "1.2.1"
|
version = "1.2.1"
|
||||||
|
|
@ -468,6 +568,12 @@ version = "1.0.22"
|
||||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
checksum = "9312f7c4f6ff9069b165498234ce8be658059c6728633667c526e27dc2cf1df5"
|
checksum = "9312f7c4f6ff9069b165498234ce8be658059c6728633667c526e27dc2cf1df5"
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "unindent"
|
||||||
|
version = "0.2.4"
|
||||||
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
|
checksum = "7264e107f553ccae879d21fbea1d6724ac785e8c3bfc762137959b5802826ef3"
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "walkdir"
|
name = "walkdir"
|
||||||
version = "2.5.0"
|
version = "2.5.0"
|
||||||
|
|
|
||||||
|
|
@ -11,12 +11,13 @@ crate-type = ["cdylib", "rlib"]
|
||||||
default = ["getrandom"]
|
default = ["getrandom"]
|
||||||
getrandom = ["dep:getrandom"]
|
getrandom = ["dep:getrandom"]
|
||||||
no_random = [] # compile without access to a source of randomness
|
no_random = [] # compile without access to a source of randomness
|
||||||
stress = [] # Enable stress tests
|
python = ["dep:pyo3"]
|
||||||
|
|
||||||
[dependencies]
|
[dependencies]
|
||||||
getrandom = { version = "0.2", optional = true }
|
getrandom = { version = "0.2", optional = true }
|
||||||
crossbeam-queue = "0.3.12"
|
crossbeam-queue = "0.3.12"
|
||||||
crossbeam-utils = "0.8.21"
|
crossbeam-utils = "0.8.21"
|
||||||
|
pyo3 = { version = "0.23", features = ["extension-module"], optional = true }
|
||||||
|
|
||||||
[dev-dependencies]
|
[dev-dependencies]
|
||||||
criterion = { version = "0.5", features = ["html_reports"] }
|
criterion = { version = "0.5", features = ["html_reports"] }
|
||||||
|
|
@ -24,3 +25,7 @@ criterion = { version = "0.5", features = ["html_reports"] }
|
||||||
[[bench]]
|
[[bench]]
|
||||||
name = "runtime_benchmarks"
|
name = "runtime_benchmarks"
|
||||||
harness = false
|
harness = false
|
||||||
|
|
||||||
|
[[bench]]
|
||||||
|
name = "worker_benchmarks"
|
||||||
|
harness = false
|
||||||
|
|
|
||||||
153
benches/worker_benchmarks.rs
Normal file
153
benches/worker_benchmarks.rs
Normal file
|
|
@ -0,0 +1,153 @@
|
||||||
|
use criterion::{
|
||||||
|
criterion_group, criterion_main, BenchmarkId, Criterion, Throughput,
|
||||||
|
};
|
||||||
|
use swactor::worker::Mailbox;
|
||||||
|
|
||||||
|
// ---------------------------------------------------------------------------
|
||||||
|
// Push throughput
|
||||||
|
// ---------------------------------------------------------------------------
|
||||||
|
|
||||||
|
fn mailbox_push(c: &mut Criterion) {
|
||||||
|
let mut group = c.benchmark_group("mailbox_push");
|
||||||
|
for n in [100, 1_000, 10_000] {
|
||||||
|
group.throughput(Throughput::Elements(n as u64));
|
||||||
|
group.bench_with_input(BenchmarkId::from_parameter(n), &n, |b, &n| {
|
||||||
|
b.iter(|| {
|
||||||
|
let mut mb: Mailbox<u64> = Mailbox::new(n);
|
||||||
|
for i in 0..n {
|
||||||
|
mb.push(i as u64);
|
||||||
|
}
|
||||||
|
});
|
||||||
|
});
|
||||||
|
}
|
||||||
|
group.finish();
|
||||||
|
}
|
||||||
|
|
||||||
|
// ---------------------------------------------------------------------------
|
||||||
|
// Pop throughput
|
||||||
|
// ---------------------------------------------------------------------------
|
||||||
|
|
||||||
|
fn mailbox_pop(c: &mut Criterion) {
|
||||||
|
let mut group = c.benchmark_group("mailbox_pop");
|
||||||
|
for n in [100, 1_000, 10_000] {
|
||||||
|
group.throughput(Throughput::Elements(n as u64));
|
||||||
|
group.bench_with_input(BenchmarkId::from_parameter(n), &n, |b, &n| {
|
||||||
|
b.iter_batched(
|
||||||
|
|| {
|
||||||
|
let mut mb: Mailbox<u64> = Mailbox::new(n);
|
||||||
|
for i in 0..n {
|
||||||
|
mb.push(i as u64);
|
||||||
|
}
|
||||||
|
mb
|
||||||
|
},
|
||||||
|
|mut mb| {
|
||||||
|
for _ in 0..n {
|
||||||
|
std::hint::black_box(mb.pop());
|
||||||
|
}
|
||||||
|
},
|
||||||
|
criterion::BatchSize::SmallInput,
|
||||||
|
);
|
||||||
|
});
|
||||||
|
}
|
||||||
|
group.finish();
|
||||||
|
}
|
||||||
|
|
||||||
|
// ---------------------------------------------------------------------------
|
||||||
|
// Interleaved push+pop
|
||||||
|
// ---------------------------------------------------------------------------
|
||||||
|
|
||||||
|
fn mailbox_interleaved(c: &mut Criterion) {
|
||||||
|
let mut group = c.benchmark_group("mailbox_interleaved");
|
||||||
|
for n in [100, 1_000, 10_000] {
|
||||||
|
group.throughput(Throughput::Elements(n as u64 * 2));
|
||||||
|
group.bench_with_input(BenchmarkId::from_parameter(n), &n, |b, &n| {
|
||||||
|
b.iter(|| {
|
||||||
|
let mut mb: Mailbox<u64> = Mailbox::new(n);
|
||||||
|
for i in 0..n {
|
||||||
|
mb.push(i as u64);
|
||||||
|
std::hint::black_box(mb.pop());
|
||||||
|
}
|
||||||
|
});
|
||||||
|
});
|
||||||
|
}
|
||||||
|
group.finish();
|
||||||
|
}
|
||||||
|
|
||||||
|
// ---------------------------------------------------------------------------
|
||||||
|
// drain_count O(1) verification
|
||||||
|
// ---------------------------------------------------------------------------
|
||||||
|
|
||||||
|
fn mailbox_drain_count(c: &mut Criterion) {
|
||||||
|
let mut group = c.benchmark_group("mailbox_drain_count");
|
||||||
|
|
||||||
|
// Below waterlevel
|
||||||
|
group.bench_function("below", |b| {
|
||||||
|
let mut mb: Mailbox<u64> = Mailbox::new(100);
|
||||||
|
for i in 0..50 {
|
||||||
|
mb.push(i);
|
||||||
|
}
|
||||||
|
b.iter(|| std::hint::black_box(mb.drain_count()));
|
||||||
|
});
|
||||||
|
|
||||||
|
// At waterlevel
|
||||||
|
group.bench_function("at", |b| {
|
||||||
|
let mut mb: Mailbox<u64> = Mailbox::new(100);
|
||||||
|
for i in 0..100 {
|
||||||
|
mb.push(i);
|
||||||
|
}
|
||||||
|
b.iter(|| std::hint::black_box(mb.drain_count()));
|
||||||
|
});
|
||||||
|
|
||||||
|
// Above waterlevel
|
||||||
|
group.bench_function("above", |b| {
|
||||||
|
let mut mb: Mailbox<u64> = Mailbox::new(100);
|
||||||
|
for i in 0..500 {
|
||||||
|
mb.push(i);
|
||||||
|
}
|
||||||
|
b.iter(|| std::hint::black_box(mb.drain_count()));
|
||||||
|
});
|
||||||
|
|
||||||
|
group.finish();
|
||||||
|
}
|
||||||
|
|
||||||
|
// ---------------------------------------------------------------------------
|
||||||
|
// Simulated actor tick: drain_count + pop N
|
||||||
|
// ---------------------------------------------------------------------------
|
||||||
|
|
||||||
|
fn mailbox_actor_tick(c: &mut Criterion) {
|
||||||
|
let mut group = c.benchmark_group("mailbox_actor_tick");
|
||||||
|
|
||||||
|
for (wl, fill) in [(10, 5), (10, 10), (10, 50), (100, 200)] {
|
||||||
|
let param = format!("wl={wl},fill={fill}");
|
||||||
|
group.bench_function(BenchmarkId::from_parameter(¶m), |b| {
|
||||||
|
b.iter_batched(
|
||||||
|
|| {
|
||||||
|
let mut mb: Mailbox<u64> = Mailbox::new(wl);
|
||||||
|
for i in 0..fill {
|
||||||
|
mb.push(i as u64);
|
||||||
|
}
|
||||||
|
mb
|
||||||
|
},
|
||||||
|
|mut mb| {
|
||||||
|
let n = mb.drain_count();
|
||||||
|
for _ in 0..n {
|
||||||
|
std::hint::black_box(mb.pop());
|
||||||
|
}
|
||||||
|
},
|
||||||
|
criterion::BatchSize::SmallInput,
|
||||||
|
);
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
group.finish();
|
||||||
|
}
|
||||||
|
|
||||||
|
criterion_group!(
|
||||||
|
benches,
|
||||||
|
mailbox_push,
|
||||||
|
mailbox_pop,
|
||||||
|
mailbox_interleaved,
|
||||||
|
mailbox_drain_count,
|
||||||
|
mailbox_actor_tick,
|
||||||
|
);
|
||||||
|
criterion_main!(benches);
|
||||||
37
examples/hello_async.py
Normal file
37
examples/hello_async.py
Normal file
|
|
@ -0,0 +1,37 @@
|
||||||
|
"""Async hello world: runtime runs in background threads, driven from asyncio."""
|
||||||
|
|
||||||
|
import asyncio
|
||||||
|
from swactor import Runtime, RuntimeConfig
|
||||||
|
|
||||||
|
|
||||||
|
async def recv(inbox, timeout=1.0):
|
||||||
|
"""Poll an inbox until a message arrives."""
|
||||||
|
while timeout > 0:
|
||||||
|
msg = inbox.try_recv()
|
||||||
|
if msg is not None:
|
||||||
|
return msg
|
||||||
|
await asyncio.sleep(0.01)
|
||||||
|
timeout -= 0.01
|
||||||
|
return None
|
||||||
|
|
||||||
|
|
||||||
|
async def main():
|
||||||
|
rt = Runtime(RuntimeConfig(num_threads=2))
|
||||||
|
|
||||||
|
def echo(ctx, msg):
|
||||||
|
ctx.send(msg["reply_to"], f"hello, {msg['name']}!")
|
||||||
|
|
||||||
|
addr = rt.spawn(echo)
|
||||||
|
inbox = rt.inbox()
|
||||||
|
handle = rt.run()
|
||||||
|
|
||||||
|
for name in ["alice", "bob", "charlie"]:
|
||||||
|
handle.send(addr, {"name": name, "reply_to": inbox.addr})
|
||||||
|
reply = await recv(inbox)
|
||||||
|
print(reply)
|
||||||
|
|
||||||
|
handle.shutdown()
|
||||||
|
handle.join()
|
||||||
|
|
||||||
|
|
||||||
|
asyncio.run(main())
|
||||||
13
examples/hello_single_thread.py
Normal file
13
examples/hello_single_thread.py
Normal file
|
|
@ -0,0 +1,13 @@
|
||||||
|
"""Hello world: spawn an echo actor, send a message, get it back."""
|
||||||
|
|
||||||
|
from swactor import Runtime
|
||||||
|
|
||||||
|
def echo(ctx, msg):
|
||||||
|
ctx.send(msg["reply_to"], f"hello, {msg['name']}!")
|
||||||
|
|
||||||
|
rt = Runtime()
|
||||||
|
addr = rt.spawn(echo)
|
||||||
|
inbox = rt.inbox()
|
||||||
|
rt.send(addr, {"name": "world", "reply_to": inbox.addr})
|
||||||
|
rt.tick()
|
||||||
|
print(inbox.try_recv())
|
||||||
11
pyproject.toml
Normal file
11
pyproject.toml
Normal file
|
|
@ -0,0 +1,11 @@
|
||||||
|
[build-system]
|
||||||
|
requires = ["maturin>=1.7,<2"]
|
||||||
|
build-backend = "maturin"
|
||||||
|
|
||||||
|
[project]
|
||||||
|
name = "swactor"
|
||||||
|
version = "0.1.0"
|
||||||
|
requires-python = ">=3.9"
|
||||||
|
|
||||||
|
[tool.maturin]
|
||||||
|
features = ["python"]
|
||||||
51
src/actor.rs
51
src/actor.rs
|
|
@ -1,6 +1,6 @@
|
||||||
use std::any::Any;
|
use std::any::Any;
|
||||||
|
|
||||||
use crate::{get_random, runtime::{ContextInner, Ctx}, worker::Mailbox};
|
use crate::runtime::Ctx;
|
||||||
|
|
||||||
/// The primary trait defining data that can be passed to and from actor processes
|
/// The primary trait defining data that can be passed to and from actor processes
|
||||||
pub trait Message: 'static + Sized + Clone + Send + Sync {}
|
pub trait Message: 'static + Sized + Clone + Send + Sync {}
|
||||||
|
|
@ -20,63 +20,32 @@ pub struct ActorAddress(pub [u8; 32]);
|
||||||
impl ActorAddress {
|
impl ActorAddress {
|
||||||
pub fn new_random() -> Self {
|
pub fn new_random() -> Self {
|
||||||
let mut bytes = [0u8; 32];
|
let mut bytes = [0u8; 32];
|
||||||
get_random(&mut bytes);
|
crate::get_random(&mut bytes);
|
||||||
Self(bytes)
|
Self(bytes)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// The actor process as represented in the Runtime, with the actor state stored with its mailbox.
|
/// The actor process as represented in the Runtime — thin wrapper around user state.
|
||||||
pub(crate) struct Actor<A>
|
pub(crate) struct Actor<A: ActorInterface>(A);
|
||||||
where
|
|
||||||
A: ActorInterface,
|
|
||||||
{
|
|
||||||
addr: ActorAddress,
|
|
||||||
mailbox: Mailbox<A::Incoming>,
|
|
||||||
inner: A,
|
|
||||||
}
|
|
||||||
|
|
||||||
impl<A: ActorInterface> Actor<A> {
|
impl<A: ActorInterface> Actor<A> {
|
||||||
pub(crate) fn new(addr: ActorAddress, mailbox: Mailbox<A::Incoming>, inner: A) -> Self {
|
pub(crate) fn new(inner: A) -> Self {
|
||||||
Self {
|
Self(inner)
|
||||||
addr,
|
|
||||||
mailbox,
|
|
||||||
inner,
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Trait for type-erased actors
|
/// Trait for type-erased actors — single-message handler.
|
||||||
pub(crate) trait AnyActor: Send {
|
pub(crate) trait AnyActor: Send {
|
||||||
/// Tick the actor, processing pending messages. Returns `true` if any work was done.
|
fn handle_any(&mut self, ctx: &Ctx, msg: Box<dyn Any + Send>);
|
||||||
fn tick(&mut self, inner: &dyn ContextInner) -> bool;
|
|
||||||
/// Deliver a type-erased message into this actor's mailbox.
|
|
||||||
/// Returns `true` if the downcast succeeded.
|
|
||||||
fn deliver(&mut self, msg: Box<dyn Any + Send>) -> bool;
|
|
||||||
}
|
}
|
||||||
|
|
||||||
impl<A> AnyActor for Actor<A>
|
impl<A> AnyActor for Actor<A>
|
||||||
where
|
where
|
||||||
A: ActorInterface,
|
A: ActorInterface,
|
||||||
{
|
{
|
||||||
fn tick(&mut self, inner: &dyn ContextInner) -> bool {
|
fn handle_any(&mut self, ctx: &Ctx, msg: Box<dyn Any + Send>) {
|
||||||
let n = self.mailbox.drain_count();
|
|
||||||
if n > 0 {
|
|
||||||
let ctx = Ctx::new(inner, self.addr);
|
|
||||||
for _ in 0..n {
|
|
||||||
if let Some(msg) = self.mailbox.pop() {
|
|
||||||
self.inner.handle(&ctx, msg);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
n > 0
|
|
||||||
}
|
|
||||||
|
|
||||||
fn deliver(&mut self, msg: Box<dyn Any + Send>) -> bool {
|
|
||||||
if let Ok(typed) = msg.downcast::<A::Incoming>() {
|
if let Ok(typed) = msg.downcast::<A::Incoming>() {
|
||||||
self.mailbox.push(*typed);
|
self.0.handle(ctx, *typed);
|
||||||
true
|
|
||||||
} else {
|
|
||||||
false
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -11,6 +11,15 @@ pub mod config;
|
||||||
|
|
||||||
pub mod runtime;
|
pub mod runtime;
|
||||||
|
|
||||||
|
#[cfg(feature = "python")]
|
||||||
|
mod python;
|
||||||
|
|
||||||
|
#[cfg(feature = "python")]
|
||||||
|
#[pyo3::pymodule]
|
||||||
|
fn swactor(m: &pyo3::Bound<'_, pyo3::types::PyModule>) -> pyo3::PyResult<()> {
|
||||||
|
python::register(m)
|
||||||
|
}
|
||||||
|
|
||||||
#[cfg(feature = "getrandom")]
|
#[cfg(feature = "getrandom")]
|
||||||
pub(crate) fn get_random(buf: &mut [u8]) {
|
pub(crate) fn get_random(buf: &mut [u8]) {
|
||||||
getrandom::getrandom(buf).unwrap()
|
getrandom::getrandom(buf).unwrap()
|
||||||
|
|
|
||||||
455
src/python.rs
Normal file
455
src/python.rs
Normal file
|
|
@ -0,0 +1,455 @@
|
||||||
|
use std::any::Any;
|
||||||
|
use std::cell::RefCell;
|
||||||
|
|
||||||
|
use pyo3::prelude::*;
|
||||||
|
use pyo3::types::PyModule;
|
||||||
|
|
||||||
|
use crate::actor::{Actor, ActorAddress, ActorInterface, AnyActor};
|
||||||
|
use crate::config::{BackoffPolicy, RuntimeConfig};
|
||||||
|
use crate::runtime::{Ctx, Inbox, Runtime, RuntimeHandle};
|
||||||
|
use crate::worker::Mailbox;
|
||||||
|
use crate::Error;
|
||||||
|
|
||||||
|
// ─── PyMsg newtype ───────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
/// Newtype around `PyObject` that implements `Clone + Send + Sync`.
|
||||||
|
///
|
||||||
|
/// `Py<PyAny>` in pyo3 0.23 doesn't implement `Clone` by default.
|
||||||
|
/// We implement it by acquiring the GIL to bump the refcount.
|
||||||
|
/// `Send + Sync` are safe because `Py<T>` is a reference-counted
|
||||||
|
/// pointer to a Python object protected by the GIL.
|
||||||
|
#[derive(Debug)]
|
||||||
|
struct PyMsg(PyObject);
|
||||||
|
|
||||||
|
impl Clone for PyMsg {
|
||||||
|
fn clone(&self) -> Self {
|
||||||
|
Python::with_gil(|py| PyMsg(self.0.clone_ref(py)))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Safety: Py<PyAny> is Send + Sync — access is serialized by the GIL.
|
||||||
|
unsafe impl Send for PyMsg {}
|
||||||
|
unsafe impl Sync for PyMsg {}
|
||||||
|
|
||||||
|
impl PyMsg {
|
||||||
|
fn into_inner(self) -> PyObject {
|
||||||
|
self.0
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// ─── Helpers ─────────────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
fn to_py_err(e: Error) -> PyErr {
|
||||||
|
pyo3::exceptions::PyRuntimeError::new_err(e.to_string())
|
||||||
|
}
|
||||||
|
|
||||||
|
// ─── PyActorAddress ──────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
#[pyclass(name = "ActorAddress")]
|
||||||
|
#[derive(Clone)]
|
||||||
|
pub struct PyActorAddress {
|
||||||
|
inner: ActorAddress,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[pymethods]
|
||||||
|
impl PyActorAddress {
|
||||||
|
fn hex(&self) -> String {
|
||||||
|
self.inner
|
||||||
|
.0
|
||||||
|
.iter()
|
||||||
|
.map(|b| format!("{b:02x}"))
|
||||||
|
.collect()
|
||||||
|
}
|
||||||
|
|
||||||
|
fn to_bytes(&self) -> Vec<u8> {
|
||||||
|
self.inner.0.to_vec()
|
||||||
|
}
|
||||||
|
|
||||||
|
fn __repr__(&self) -> String {
|
||||||
|
let hex = self.hex();
|
||||||
|
format!("ActorAddress({hex})")
|
||||||
|
}
|
||||||
|
|
||||||
|
fn __eq__(&self, other: &PyActorAddress) -> bool {
|
||||||
|
self.inner == other.inner
|
||||||
|
}
|
||||||
|
|
||||||
|
fn __hash__(&self) -> u64 {
|
||||||
|
use std::hash::{Hash, Hasher};
|
||||||
|
let mut hasher = std::collections::hash_map::DefaultHasher::new();
|
||||||
|
self.inner.hash(&mut hasher);
|
||||||
|
hasher.finish()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl From<ActorAddress> for PyActorAddress {
|
||||||
|
fn from(inner: ActorAddress) -> Self {
|
||||||
|
Self { inner }
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// ─── Effects / PyCtx ─────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
enum Effect {
|
||||||
|
Send {
|
||||||
|
addr: ActorAddress,
|
||||||
|
msg: PyObject,
|
||||||
|
},
|
||||||
|
Spawn {
|
||||||
|
addr: ActorAddress,
|
||||||
|
handler: PyObject,
|
||||||
|
},
|
||||||
|
}
|
||||||
|
|
||||||
|
#[pyclass(name = "Ctx", unsendable)]
|
||||||
|
pub struct PyCtx {
|
||||||
|
self_addr: ActorAddress,
|
||||||
|
effects: RefCell<Vec<Effect>>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl PyCtx {
|
||||||
|
fn new(self_addr: ActorAddress) -> Self {
|
||||||
|
Self {
|
||||||
|
self_addr,
|
||||||
|
effects: RefCell::new(Vec::new()),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn take_effects(&self) -> Vec<Effect> {
|
||||||
|
self.effects.borrow_mut().drain(..).collect()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[pymethods]
|
||||||
|
impl PyCtx {
|
||||||
|
#[getter]
|
||||||
|
fn self_addr(&self) -> PyActorAddress {
|
||||||
|
PyActorAddress::from(self.self_addr)
|
||||||
|
}
|
||||||
|
|
||||||
|
fn send(&self, addr: &PyActorAddress, msg: PyObject) {
|
||||||
|
self.effects.borrow_mut().push(Effect::Send {
|
||||||
|
addr: addr.inner,
|
||||||
|
msg,
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
fn spawn(&self, handler: PyObject) -> PyActorAddress {
|
||||||
|
let addr = ActorAddress::new_random();
|
||||||
|
self.effects.borrow_mut().push(Effect::Spawn {
|
||||||
|
addr,
|
||||||
|
handler,
|
||||||
|
});
|
||||||
|
PyActorAddress::from(addr)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// ─── PyActor ─────────────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
struct PyActor {
|
||||||
|
handler: PyObject,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl PyActor {
|
||||||
|
fn new(handler: PyObject) -> Self {
|
||||||
|
Self { handler }
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl ActorInterface for PyActor {
|
||||||
|
type Incoming = PyMsg;
|
||||||
|
type Response = PyMsg;
|
||||||
|
|
||||||
|
fn handle(&mut self, ctx: &Ctx, msg: Self::Incoming) {
|
||||||
|
let py_ctx = PyCtx::new(ctx.self_addr());
|
||||||
|
|
||||||
|
let call_result = Python::with_gil(|py| {
|
||||||
|
let ctx_bound = Bound::new(py, py_ctx)?;
|
||||||
|
self.handler
|
||||||
|
.call1(py, (&ctx_bound, msg.into_inner()))?;
|
||||||
|
let ctx_ref = ctx_bound.borrow();
|
||||||
|
Ok::<Vec<Effect>, PyErr>(ctx_ref.take_effects())
|
||||||
|
});
|
||||||
|
|
||||||
|
match call_result {
|
||||||
|
Ok(effects) => {
|
||||||
|
for effect in effects {
|
||||||
|
match effect {
|
||||||
|
Effect::Send { addr, msg } => {
|
||||||
|
let _ = ctx.raw_inner().send_via_queue(
|
||||||
|
addr,
|
||||||
|
Box::new(PyMsg(msg)) as Box<dyn Any + Send>,
|
||||||
|
);
|
||||||
|
}
|
||||||
|
Effect::Spawn { addr, handler } => {
|
||||||
|
let waterlevel = ctx.raw_inner().mailbox_waterlevel();
|
||||||
|
let actor = PyActor::new(handler);
|
||||||
|
let actor = Actor::new(addr, Mailbox::new(waterlevel), actor);
|
||||||
|
let boxed: Box<dyn AnyActor> = Box::new(actor);
|
||||||
|
let _ = ctx.raw_inner().spawn_any(addr, boxed);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
Err(e) => {
|
||||||
|
Python::with_gil(|py| {
|
||||||
|
e.print(py);
|
||||||
|
});
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// ─── PyInbox ─────────────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
#[pyclass(name = "Inbox")]
|
||||||
|
pub struct PyInbox {
|
||||||
|
inner: Inbox<PyMsg>,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[pymethods]
|
||||||
|
impl PyInbox {
|
||||||
|
#[getter]
|
||||||
|
fn addr(&self) -> PyActorAddress {
|
||||||
|
PyActorAddress::from(*self.inner.addr())
|
||||||
|
}
|
||||||
|
|
||||||
|
fn try_recv(&self) -> Option<PyObject> {
|
||||||
|
self.inner.try_recv().map(|m| m.into_inner())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// ─── PyRuntimeConfig ─────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
#[pyclass(name = "RuntimeConfig")]
|
||||||
|
#[derive(Clone)]
|
||||||
|
pub struct PyRuntimeConfig {
|
||||||
|
#[pyo3(get, set)]
|
||||||
|
num_threads: usize,
|
||||||
|
#[pyo3(get, set)]
|
||||||
|
max_actors: usize,
|
||||||
|
#[pyo3(get, set)]
|
||||||
|
actor_max_messages: usize,
|
||||||
|
#[pyo3(get, set)]
|
||||||
|
mailbox_waterlevel: usize,
|
||||||
|
#[pyo3(get, set)]
|
||||||
|
spin_threshold: u32,
|
||||||
|
#[pyo3(get, set)]
|
||||||
|
yield_threshold: u32,
|
||||||
|
#[pyo3(get, set)]
|
||||||
|
sleep_increment_us: u64,
|
||||||
|
#[pyo3(get, set)]
|
||||||
|
sleep_max_us: u64,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[pymethods]
|
||||||
|
impl PyRuntimeConfig {
|
||||||
|
#[new]
|
||||||
|
#[pyo3(signature = (
|
||||||
|
*,
|
||||||
|
num_threads = 1,
|
||||||
|
max_actors = 1_000,
|
||||||
|
actor_max_messages = 1_000,
|
||||||
|
mailbox_waterlevel = 10,
|
||||||
|
spin_threshold = 64,
|
||||||
|
yield_threshold = 256,
|
||||||
|
sleep_increment_us = 50,
|
||||||
|
sleep_max_us = 1_000,
|
||||||
|
))]
|
||||||
|
fn new(
|
||||||
|
num_threads: usize,
|
||||||
|
max_actors: usize,
|
||||||
|
actor_max_messages: usize,
|
||||||
|
mailbox_waterlevel: usize,
|
||||||
|
spin_threshold: u32,
|
||||||
|
yield_threshold: u32,
|
||||||
|
sleep_increment_us: u64,
|
||||||
|
sleep_max_us: u64,
|
||||||
|
) -> Self {
|
||||||
|
Self {
|
||||||
|
num_threads,
|
||||||
|
max_actors,
|
||||||
|
actor_max_messages,
|
||||||
|
mailbox_waterlevel,
|
||||||
|
spin_threshold,
|
||||||
|
yield_threshold,
|
||||||
|
sleep_increment_us,
|
||||||
|
sleep_max_us,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl From<PyRuntimeConfig> for RuntimeConfig {
|
||||||
|
fn from(py: PyRuntimeConfig) -> Self {
|
||||||
|
RuntimeConfig {
|
||||||
|
num_threads: py.num_threads,
|
||||||
|
max_actors: py.max_actors,
|
||||||
|
actor_max_messages: py.actor_max_messages,
|
||||||
|
mailbox_waterlevel: py.mailbox_waterlevel,
|
||||||
|
backoff_policy: BackoffPolicy {
|
||||||
|
spin_threshold: py.spin_threshold,
|
||||||
|
yield_threshold: py.yield_threshold,
|
||||||
|
sleep_increment_us: py.sleep_increment_us,
|
||||||
|
sleep_max_us: py.sleep_max_us,
|
||||||
|
},
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// ─── PyRuntime ───────────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
#[pyclass(name = "Runtime")]
|
||||||
|
pub struct PyRuntime {
|
||||||
|
inner: Option<Runtime>,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[pymethods]
|
||||||
|
impl PyRuntime {
|
||||||
|
#[new]
|
||||||
|
#[pyo3(signature = (config=None))]
|
||||||
|
fn new(config: Option<PyRuntimeConfig>) -> Self {
|
||||||
|
let config: RuntimeConfig = match config {
|
||||||
|
Some(c) => c.into(),
|
||||||
|
None => RuntimeConfig::default(),
|
||||||
|
};
|
||||||
|
Self {
|
||||||
|
inner: Some(Runtime::new(config)),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn spawn(&self, handler: PyObject) -> PyResult<PyActorAddress> {
|
||||||
|
let rt = self
|
||||||
|
.inner
|
||||||
|
.as_ref()
|
||||||
|
.ok_or_else(|| pyo3::exceptions::PyRuntimeError::new_err("Runtime consumed by run()"))?;
|
||||||
|
let actor = PyActor::new(handler);
|
||||||
|
let addr = rt.spawn(actor).map_err(to_py_err)?;
|
||||||
|
Ok(PyActorAddress::from(addr))
|
||||||
|
}
|
||||||
|
|
||||||
|
fn send(&self, addr: &PyActorAddress, msg: PyObject) -> PyResult<()> {
|
||||||
|
let rt = self
|
||||||
|
.inner
|
||||||
|
.as_ref()
|
||||||
|
.ok_or_else(|| pyo3::exceptions::PyRuntimeError::new_err("Runtime consumed by run()"))?;
|
||||||
|
rt.send_to(addr.inner, PyMsg(msg)).map_err(to_py_err)
|
||||||
|
}
|
||||||
|
|
||||||
|
fn inbox(&self) -> PyResult<PyInbox> {
|
||||||
|
let rt = self
|
||||||
|
.inner
|
||||||
|
.as_ref()
|
||||||
|
.ok_or_else(|| pyo3::exceptions::PyRuntimeError::new_err("Runtime consumed by run()"))?;
|
||||||
|
let inbox: Inbox<PyMsg> = rt.new_inbox().map_err(to_py_err)?;
|
||||||
|
Ok(PyInbox { inner: inbox })
|
||||||
|
}
|
||||||
|
|
||||||
|
fn tick(&self) -> PyResult<()> {
|
||||||
|
let rt = self
|
||||||
|
.inner
|
||||||
|
.as_ref()
|
||||||
|
.ok_or_else(|| pyo3::exceptions::PyRuntimeError::new_err("Runtime consumed by run()"))?;
|
||||||
|
rt.tick();
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
fn run(&mut self, py: Python<'_>) -> PyResult<PyRuntimeHandle> {
|
||||||
|
let rt = self
|
||||||
|
.inner
|
||||||
|
.take()
|
||||||
|
.ok_or_else(|| pyo3::exceptions::PyRuntimeError::new_err("Runtime consumed by run()"))?;
|
||||||
|
let handle = py.allow_threads(|| rt.run().map_err(to_py_err))?;
|
||||||
|
Ok(PyRuntimeHandle {
|
||||||
|
inner: Some(handle),
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
fn shutdown(&self) -> PyResult<()> {
|
||||||
|
let rt = self
|
||||||
|
.inner
|
||||||
|
.as_ref()
|
||||||
|
.ok_or_else(|| pyo3::exceptions::PyRuntimeError::new_err("Runtime consumed by run()"))?;
|
||||||
|
rt.shutdown();
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// ─── PyRuntimeHandle ─────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
#[pyclass(name = "RuntimeHandle")]
|
||||||
|
pub struct PyRuntimeHandle {
|
||||||
|
inner: Option<RuntimeHandle>,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[pymethods]
|
||||||
|
impl PyRuntimeHandle {
|
||||||
|
fn spawn(&self, handler: PyObject) -> PyResult<PyActorAddress> {
|
||||||
|
let handle = self
|
||||||
|
.inner
|
||||||
|
.as_ref()
|
||||||
|
.ok_or_else(|| {
|
||||||
|
pyo3::exceptions::PyRuntimeError::new_err("RuntimeHandle consumed by join()")
|
||||||
|
})?;
|
||||||
|
let actor = PyActor::new(handler);
|
||||||
|
let addr = handle.runtime.spawn(actor).map_err(to_py_err)?;
|
||||||
|
Ok(PyActorAddress::from(addr))
|
||||||
|
}
|
||||||
|
|
||||||
|
fn send(&self, addr: &PyActorAddress, msg: PyObject) -> PyResult<()> {
|
||||||
|
let handle = self
|
||||||
|
.inner
|
||||||
|
.as_ref()
|
||||||
|
.ok_or_else(|| {
|
||||||
|
pyo3::exceptions::PyRuntimeError::new_err("RuntimeHandle consumed by join()")
|
||||||
|
})?;
|
||||||
|
handle
|
||||||
|
.runtime
|
||||||
|
.send_to(addr.inner, PyMsg(msg))
|
||||||
|
.map_err(to_py_err)
|
||||||
|
}
|
||||||
|
|
||||||
|
fn inbox(&self) -> PyResult<PyInbox> {
|
||||||
|
let handle = self
|
||||||
|
.inner
|
||||||
|
.as_ref()
|
||||||
|
.ok_or_else(|| {
|
||||||
|
pyo3::exceptions::PyRuntimeError::new_err("RuntimeHandle consumed by join()")
|
||||||
|
})?;
|
||||||
|
let inbox: Inbox<PyMsg> = handle.runtime.new_inbox().map_err(to_py_err)?;
|
||||||
|
Ok(PyInbox { inner: inbox })
|
||||||
|
}
|
||||||
|
|
||||||
|
fn shutdown(&self) -> PyResult<()> {
|
||||||
|
let handle = self
|
||||||
|
.inner
|
||||||
|
.as_ref()
|
||||||
|
.ok_or_else(|| {
|
||||||
|
pyo3::exceptions::PyRuntimeError::new_err("RuntimeHandle consumed by join()")
|
||||||
|
})?;
|
||||||
|
handle.shutdown();
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
fn join(&mut self, py: Python<'_>) -> PyResult<()> {
|
||||||
|
let handle = self
|
||||||
|
.inner
|
||||||
|
.take()
|
||||||
|
.ok_or_else(|| {
|
||||||
|
pyo3::exceptions::PyRuntimeError::new_err("RuntimeHandle consumed by join()")
|
||||||
|
})?;
|
||||||
|
py.allow_threads(|| handle.join());
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// ─── Module registration ─────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
pub fn register(m: &Bound<'_, PyModule>) -> PyResult<()> {
|
||||||
|
m.add_class::<PyActorAddress>()?;
|
||||||
|
m.add_class::<PyCtx>()?;
|
||||||
|
m.add_class::<PyInbox>()?;
|
||||||
|
m.add_class::<PyRuntimeConfig>()?;
|
||||||
|
m.add_class::<PyRuntime>()?;
|
||||||
|
m.add_class::<PyRuntimeHandle>()?;
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
@ -10,7 +10,6 @@ use crate::address_map::{AddressMap, Placement, WorkerId};
|
||||||
use crate::channel::{Receiver, Sender};
|
use crate::channel::{Receiver, Sender};
|
||||||
// Re-export config types so existing code using `runtime::RuntimeConfig` still works
|
// Re-export config types so existing code using `runtime::RuntimeConfig` still works
|
||||||
pub use crate::config::{BackoffPolicy, RuntimeConfig};
|
pub use crate::config::{BackoffPolicy, RuntimeConfig};
|
||||||
use crate::worker::Mailbox;
|
|
||||||
use crate::worker::{TickContext, Worker};
|
use crate::worker::{TickContext, Worker};
|
||||||
use crate::Error;
|
use crate::Error;
|
||||||
|
|
||||||
|
|
@ -64,6 +63,11 @@ impl<'a> Ctx<'a> {
|
||||||
Self { inner, self_addr }
|
Self { inner, self_addr }
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[cfg(feature = "python")]
|
||||||
|
pub(crate) fn raw_inner(&self) -> &dyn ContextInner {
|
||||||
|
self.inner
|
||||||
|
}
|
||||||
|
|
||||||
/// Returns the address of the actor currently being ticked.
|
/// Returns the address of the actor currently being ticked.
|
||||||
pub fn self_addr(&self) -> ActorAddress {
|
pub fn self_addr(&self) -> ActorAddress {
|
||||||
self.self_addr
|
self.self_addr
|
||||||
|
|
@ -77,9 +81,7 @@ impl<'a> Ctx<'a> {
|
||||||
/// Spawn a new actor, returning its address.
|
/// Spawn a new actor, returning its address.
|
||||||
pub fn spawn<A: ActorInterface>(&self, actor: A) -> Result<ActorAddress, Error> {
|
pub fn spawn<A: ActorInterface>(&self, actor: A) -> Result<ActorAddress, Error> {
|
||||||
let addr = ActorAddress::new_random();
|
let addr = ActorAddress::new_random();
|
||||||
let waterlevel = self.inner.mailbox_waterlevel();
|
let boxed: Box<dyn AnyActor> = Box::new(Actor::new(actor));
|
||||||
let actor = Actor::new(addr, Mailbox::new(waterlevel), actor);
|
|
||||||
let boxed: Box<dyn AnyActor> = Box::new(actor);
|
|
||||||
self.inner.spawn_any(addr, boxed)?;
|
self.inner.spawn_any(addr, boxed)?;
|
||||||
Ok(addr)
|
Ok(addr)
|
||||||
}
|
}
|
||||||
|
|
@ -187,8 +189,7 @@ impl Runtime {
|
||||||
let addr = ActorAddress::new_random();
|
let addr = ActorAddress::new_random();
|
||||||
let worker_id = self.placement.next_worker();
|
let worker_id = self.placement.next_worker();
|
||||||
self.address_map.insert(addr, worker_id);
|
self.address_map.insert(addr, worker_id);
|
||||||
let actor = Actor::new(addr, Mailbox::new(self.config.mailbox_waterlevel), actor);
|
let boxed: Box<dyn AnyActor> = Box::new(Actor::new(actor));
|
||||||
let boxed: Box<dyn AnyActor> = Box::new(actor);
|
|
||||||
self.spawn_txs[worker_id.as_usize()]
|
self.spawn_txs[worker_id.as_usize()]
|
||||||
.try_send((addr, boxed))
|
.try_send((addr, boxed))
|
||||||
.map_err(|_| Error::from("Runtime error: spawn queue full"))?;
|
.map_err(|_| Error::from("Runtime error: spawn queue full"))?;
|
||||||
|
|
@ -233,25 +234,22 @@ impl Runtime {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Spawn worker threads and start processing, returning a set of handles and
|
/// Spawn worker threads and start processing, returning a handle
|
||||||
/// a Runtime object to interface with.
|
/// to interact with the runtime and join the threads later.
|
||||||
///
|
///
|
||||||
/// ### WARN:
|
/// Works in both single-threaded and multi-threaded configurations.
|
||||||
/// ##### Returns an error if the configuration is set as single threaded
|
/// In single-threaded mode, one background thread is spawned.
|
||||||
/// `config.num_threads == 1`
|
|
||||||
pub fn run(mut self) -> Result<RuntimeHandle, Error> {
|
pub fn run(mut self) -> Result<RuntimeHandle, Error> {
|
||||||
if self.config.num_threads < 2 {
|
|
||||||
return Err(Error::from(
|
|
||||||
"Runtime error: cannot call `Runtime::run()` from a single-threaded context.",
|
|
||||||
));
|
|
||||||
}
|
|
||||||
|
|
||||||
self.is_running.store(true, Ordering::Release);
|
self.is_running.store(true, Ordering::Release);
|
||||||
|
|
||||||
let workers = self
|
let mut workers: Vec<Worker> = Vec::new();
|
||||||
.pending_workers
|
|
||||||
.take()
|
if let Some(w) = self.single_worker.take() {
|
||||||
.expect("Workers must be present for multi-threaded runtime");
|
workers.push(w.into_inner());
|
||||||
|
}
|
||||||
|
if let Some(ws) = self.pending_workers.take() {
|
||||||
|
workers.extend(ws);
|
||||||
|
}
|
||||||
|
|
||||||
let rt = Arc::new(self);
|
let rt = Arc::new(self);
|
||||||
let mut handles: Vec<JoinHandle<()>> = Vec::with_capacity(workers.len());
|
let mut handles: Vec<JoinHandle<()>> = Vec::with_capacity(workers.len());
|
||||||
|
|
|
||||||
|
|
@ -8,17 +8,17 @@ use crate::actor::{ActorAddress, AnyActor, Message};
|
||||||
use crate::address_map::{AddressMap, Placement, WorkerId};
|
use crate::address_map::{AddressMap, Placement, WorkerId};
|
||||||
use crate::channel::{Receiver, Sender};
|
use crate::channel::{Receiver, Sender};
|
||||||
use crate::config::{BackoffPolicy, RuntimeConfig};
|
use crate::config::{BackoffPolicy, RuntimeConfig};
|
||||||
use crate::runtime::{ContextInner, Envelope, InboxRegistry};
|
use crate::runtime::{ContextInner, Ctx, Envelope, InboxRegistry};
|
||||||
use crate::Error;
|
use crate::Error;
|
||||||
|
|
||||||
/// Shared state passed to tick_once — single thin pointer avoids register spill.
|
/// Shared state passed to tick_once — single thin pointer avoids register spill.
|
||||||
pub(crate) struct TickContext<'a> {
|
pub(crate) struct TickContext<'a> {
|
||||||
pub address_map: &'a AddressMap,
|
pub(crate) address_map: &'a AddressMap,
|
||||||
pub transfer_txs: &'a [Sender<Envelope>],
|
pub(crate) transfer_txs: &'a [Sender<Envelope>],
|
||||||
pub spawn_txs: &'a [Sender<(ActorAddress, Box<dyn AnyActor>)>],
|
pub(crate) spawn_txs: &'a [Sender<(ActorAddress, Box<dyn AnyActor>)>],
|
||||||
pub placement: &'a Placement,
|
pub(crate) placement: &'a Placement,
|
||||||
pub inbox_registry: &'a InboxRegistry,
|
pub(crate) inbox_registry: &'a InboxRegistry,
|
||||||
pub config: &'a RuntimeConfig,
|
pub(crate) config: &'a RuntimeConfig,
|
||||||
}
|
}
|
||||||
|
|
||||||
/// A worker owns a set of actors and runs them in a loop.
|
/// A worker owns a set of actors and runs them in a loop.
|
||||||
|
|
@ -30,7 +30,7 @@ pub(crate) struct Worker {
|
||||||
}
|
}
|
||||||
|
|
||||||
impl Worker {
|
impl Worker {
|
||||||
pub fn new(
|
pub(crate) fn new(
|
||||||
id: WorkerId,
|
id: WorkerId,
|
||||||
transfer_rx: Receiver<Envelope>,
|
transfer_rx: Receiver<Envelope>,
|
||||||
spawn_rx: Receiver<(ActorAddress, Box<dyn AnyActor>)>,
|
spawn_rx: Receiver<(ActorAddress, Box<dyn AnyActor>)>,
|
||||||
|
|
@ -44,7 +44,7 @@ impl Worker {
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Run one iteration of the worker loop. Returns `true` if any work was done.
|
/// Run one iteration of the worker loop. Returns `true` if any work was done.
|
||||||
pub fn tick_once(&mut self, tc: &TickContext) -> bool {
|
pub(crate) fn tick_once(&mut self, tc: &TickContext) -> bool {
|
||||||
let mut did_work = false;
|
let mut did_work = false;
|
||||||
|
|
||||||
// 1. Drain spawn queue → add actors to pool
|
// 1. Drain spawn queue → add actors to pool
|
||||||
|
|
@ -166,9 +166,25 @@ impl ContextInner for WorkerContext<'_> {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Per-worker actor storage.
|
/// How many messages to process this tick:
|
||||||
|
/// - `len < waterlevel` → process all (`len`)
|
||||||
|
/// - `len >= waterlevel` → process half (`len >> 1`)
|
||||||
|
pub(crate) fn drain_count(len: usize, waterlevel: usize) -> usize {
|
||||||
|
if len < waterlevel {
|
||||||
|
len
|
||||||
|
} else {
|
||||||
|
len >> 1
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
struct ActorSlot {
|
||||||
|
mailbox: VecDeque<Box<dyn Any + Send>>,
|
||||||
|
actor: Box<dyn AnyActor>,
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Per-worker actor storage. Owns per-actor mailboxes.
|
||||||
pub(crate) struct ActorPool {
|
pub(crate) struct ActorPool {
|
||||||
actors: HashMap<ActorAddress, Box<dyn AnyActor>>,
|
actors: HashMap<ActorAddress, ActorSlot>,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl ActorPool {
|
impl ActorPool {
|
||||||
|
|
@ -179,18 +195,22 @@ impl ActorPool {
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn insert(&mut self, addr: ActorAddress, actor: Box<dyn AnyActor>) {
|
pub fn insert(&mut self, addr: ActorAddress, actor: Box<dyn AnyActor>) {
|
||||||
self.actors.insert(addr, actor);
|
self.actors.insert(addr, ActorSlot {
|
||||||
|
mailbox: VecDeque::new(),
|
||||||
|
actor,
|
||||||
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn remove(&mut self, addr: &ActorAddress) -> Option<Box<dyn AnyActor>> {
|
pub fn remove(&mut self, addr: &ActorAddress) -> Option<Box<dyn AnyActor>> {
|
||||||
self.actors.remove(addr)
|
self.actors.remove(addr).map(|slot| slot.actor)
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Deliver a type-erased message to the actor at `addr`.
|
/// Deliver a type-erased message to the actor at `addr`.
|
||||||
/// Returns `true` if the actor was found and the message type matched.
|
/// Returns `true` if the actor exists (message is queued; type check deferred to tick).
|
||||||
pub fn deliver(&mut self, addr: &ActorAddress, msg: Box<dyn Any + Send>) -> bool {
|
pub fn deliver(&mut self, addr: &ActorAddress, msg: Box<dyn Any + Send>) -> bool {
|
||||||
if let Some(actor) = self.actors.get_mut(addr) {
|
if let Some(slot) = self.actors.get_mut(addr) {
|
||||||
actor.deliver(msg)
|
slot.mailbox.push_back(msg);
|
||||||
|
true
|
||||||
} else {
|
} else {
|
||||||
false
|
false
|
||||||
}
|
}
|
||||||
|
|
@ -199,8 +219,16 @@ impl ActorPool {
|
||||||
/// Tick all actors in the pool. Returns `true` if any actor processed messages.
|
/// Tick all actors in the pool. Returns `true` if any actor processed messages.
|
||||||
pub fn tick_all(&mut self, inner: &dyn ContextInner) -> bool {
|
pub fn tick_all(&mut self, inner: &dyn ContextInner) -> bool {
|
||||||
let mut did_work = false;
|
let mut did_work = false;
|
||||||
for actor in self.actors.values_mut() {
|
for (&addr, slot) in self.actors.iter_mut() {
|
||||||
if actor.tick(inner) {
|
let len = slot.mailbox.len();
|
||||||
|
let n = drain_count(len, inner.mailbox_waterlevel());
|
||||||
|
if n > 0 {
|
||||||
|
let ctx = Ctx::new(inner, addr);
|
||||||
|
for _ in 0..n {
|
||||||
|
if let Some(msg) = slot.mailbox.pop_front() {
|
||||||
|
slot.actor.handle_any(&ctx, msg);
|
||||||
|
}
|
||||||
|
}
|
||||||
did_work = true;
|
did_work = true;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -247,12 +275,9 @@ impl<M: Message> Mailbox<M> {
|
||||||
/// - `len < waterlevel` → process all (`len`)
|
/// - `len < waterlevel` → process all (`len`)
|
||||||
/// - `len >= waterlevel` → process half (`len >> 1`)
|
/// - `len >= waterlevel` → process half (`len >> 1`)
|
||||||
pub fn drain_count(&self) -> usize {
|
pub fn drain_count(&self) -> usize {
|
||||||
let len = self.queue.len();
|
drain_count(self.queue.len(), self.waterlevel)
|
||||||
if len < self.waterlevel {
|
|
||||||
len
|
|
||||||
} else {
|
|
||||||
len >> 1
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
mod tests;
|
||||||
612
src/worker/tests.rs
Normal file
612
src/worker/tests.rs
Normal file
|
|
@ -0,0 +1,612 @@
|
||||||
|
use std::any::Any;
|
||||||
|
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
|
||||||
|
use std::sync::Arc;
|
||||||
|
use std::thread;
|
||||||
|
|
||||||
|
use crate::actor::{ActorAddress, AnyActor};
|
||||||
|
use crate::address_map::{AddressMap, Placement, WorkerId};
|
||||||
|
use crate::channel::Receiver;
|
||||||
|
use crate::config::{BackoffPolicy, RuntimeConfig};
|
||||||
|
use crate::runtime::{ContextInner, Ctx, Envelope, InboxRegistry};
|
||||||
|
use super::{ActorPool, Mailbox, TickContext, Worker};
|
||||||
|
use crate::Error;
|
||||||
|
|
||||||
|
// ── Test helpers ────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
#[derive(Clone, Debug, PartialEq)]
|
||||||
|
struct TestMsg(u64);
|
||||||
|
|
||||||
|
/// A minimal actor that counts how many messages it handled.
|
||||||
|
struct CounterActor {
|
||||||
|
counter: Arc<AtomicUsize>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl AnyActor for CounterActor {
|
||||||
|
fn handle_any(&mut self, _ctx: &Ctx, msg: Box<dyn Any + Send>) {
|
||||||
|
if msg.downcast::<u64>().is_ok() {
|
||||||
|
self.counter.fetch_add(1, Ordering::Relaxed);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn make_test_actor(
|
||||||
|
_addr: ActorAddress,
|
||||||
|
) -> (Box<dyn AnyActor>, Arc<AtomicUsize>) {
|
||||||
|
let counter = Arc::new(AtomicUsize::new(0));
|
||||||
|
let actor = CounterActor {
|
||||||
|
counter: counter.clone(),
|
||||||
|
};
|
||||||
|
(Box::new(actor), counter)
|
||||||
|
}
|
||||||
|
|
||||||
|
/// No-op ContextInner for ActorPool tests.
|
||||||
|
struct StubContextInner {
|
||||||
|
waterlevel: usize,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl ContextInner for StubContextInner {
|
||||||
|
fn send_any(&self, _addr: ActorAddress, _msg: Box<dyn Any + Send>) -> Result<(), Error> {
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
fn spawn_any(&self, _addr: ActorAddress, _actor: Box<dyn AnyActor>) -> Result<(), Error> {
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
fn mailbox_waterlevel(&self) -> usize {
|
||||||
|
self.waterlevel
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn make_addr(id: u8) -> ActorAddress {
|
||||||
|
let mut bytes = [0u8; 32];
|
||||||
|
bytes[0] = id;
|
||||||
|
ActorAddress(bytes)
|
||||||
|
}
|
||||||
|
|
||||||
|
// ── Mailbox tests ───────────────────────────────────────────────────
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn mailbox_new_is_empty() {
|
||||||
|
let mb: Mailbox<u64> = Mailbox::new(10);
|
||||||
|
assert_eq!(mb.len(), 0);
|
||||||
|
assert!(mb.is_empty());
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn mailbox_push_increments_len() {
|
||||||
|
let mut mb: Mailbox<u64> = Mailbox::new(10);
|
||||||
|
mb.push(1);
|
||||||
|
assert_eq!(mb.len(), 1);
|
||||||
|
mb.push(2);
|
||||||
|
assert_eq!(mb.len(), 2);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn mailbox_pop_returns_fifo_order() {
|
||||||
|
let mut mb: Mailbox<u64> = Mailbox::new(10);
|
||||||
|
mb.push(10);
|
||||||
|
mb.push(20);
|
||||||
|
mb.push(30);
|
||||||
|
assert_eq!(mb.pop(), Some(10));
|
||||||
|
assert_eq!(mb.pop(), Some(20));
|
||||||
|
assert_eq!(mb.pop(), Some(30));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn mailbox_pop_empty_returns_none() {
|
||||||
|
let mut mb: Mailbox<u64> = Mailbox::new(10);
|
||||||
|
assert_eq!(mb.pop(), None);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn mailbox_pop_drains_to_empty() {
|
||||||
|
let mut mb: Mailbox<u64> = Mailbox::new(10);
|
||||||
|
mb.push(1);
|
||||||
|
mb.push(2);
|
||||||
|
mb.pop();
|
||||||
|
mb.pop();
|
||||||
|
assert!(mb.is_empty());
|
||||||
|
assert_eq!(mb.len(), 0);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn mailbox_push_pop_interleaved() {
|
||||||
|
let mut mb: Mailbox<u64> = Mailbox::new(10);
|
||||||
|
mb.push(1);
|
||||||
|
mb.push(2);
|
||||||
|
assert_eq!(mb.pop(), Some(1));
|
||||||
|
mb.push(3);
|
||||||
|
assert_eq!(mb.pop(), Some(2));
|
||||||
|
assert_eq!(mb.pop(), Some(3));
|
||||||
|
assert_eq!(mb.pop(), None);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn drain_count_empty() {
|
||||||
|
let mb: Mailbox<u64> = Mailbox::new(10);
|
||||||
|
assert_eq!(mb.drain_count(), 0);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn drain_count_below_waterlevel() {
|
||||||
|
let mut mb: Mailbox<u64> = Mailbox::new(10);
|
||||||
|
for i in 0..5 {
|
||||||
|
mb.push(i);
|
||||||
|
}
|
||||||
|
assert_eq!(mb.drain_count(), 5);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn drain_count_at_waterlevel_minus_one() {
|
||||||
|
let mut mb: Mailbox<u64> = Mailbox::new(10);
|
||||||
|
for i in 0..9 {
|
||||||
|
mb.push(i);
|
||||||
|
}
|
||||||
|
// len=9 < waterlevel=10 → returns len
|
||||||
|
assert_eq!(mb.drain_count(), 9);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn drain_count_at_waterlevel() {
|
||||||
|
let mut mb: Mailbox<u64> = Mailbox::new(10);
|
||||||
|
for i in 0..10 {
|
||||||
|
mb.push(i);
|
||||||
|
}
|
||||||
|
// len=10 >= waterlevel=10 → returns len >> 1 = 5
|
||||||
|
assert_eq!(mb.drain_count(), 5);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn drain_count_at_waterlevel_plus_one() {
|
||||||
|
let mut mb: Mailbox<u64> = Mailbox::new(10);
|
||||||
|
for i in 0..11 {
|
||||||
|
mb.push(i);
|
||||||
|
}
|
||||||
|
// len=11 >= waterlevel=10 → returns 11 >> 1 = 5
|
||||||
|
assert_eq!(mb.drain_count(), 5);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn drain_count_well_above_waterlevel() {
|
||||||
|
let mut mb: Mailbox<u64> = Mailbox::new(10);
|
||||||
|
for i in 0..100 {
|
||||||
|
mb.push(i);
|
||||||
|
}
|
||||||
|
// len=100 >= waterlevel=10 → returns 100 >> 1 = 50
|
||||||
|
assert_eq!(mb.drain_count(), 50);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn drain_count_odd_len_truncates() {
|
||||||
|
let mut mb: Mailbox<u64> = Mailbox::new(1);
|
||||||
|
for i in 0..7 {
|
||||||
|
mb.push(i);
|
||||||
|
}
|
||||||
|
// len=7 >= waterlevel=1 → returns 7 >> 1 = 3
|
||||||
|
assert_eq!(mb.drain_count(), 3);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn drain_count_waterlevel_one() {
|
||||||
|
let mut mb: Mailbox<u64> = Mailbox::new(1);
|
||||||
|
mb.push(42);
|
||||||
|
// len=1 >= waterlevel=1 → returns 1 >> 1 = 0
|
||||||
|
assert_eq!(mb.drain_count(), 0);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn drain_count_waterlevel_zero() {
|
||||||
|
// waterlevel=0 means len >= 0 is always true → always half-drain
|
||||||
|
let mb: Mailbox<u64> = Mailbox::new(0);
|
||||||
|
assert_eq!(mb.drain_count(), 0); // empty: 0 >> 1 = 0
|
||||||
|
|
||||||
|
let mut mb2: Mailbox<u64> = Mailbox::new(0);
|
||||||
|
mb2.push(1);
|
||||||
|
// len=1 >= waterlevel=0 → returns 1 >> 1 = 0
|
||||||
|
assert_eq!(mb2.drain_count(), 0);
|
||||||
|
|
||||||
|
let mut mb3: Mailbox<u64> = Mailbox::new(0);
|
||||||
|
mb3.push(1);
|
||||||
|
mb3.push(2);
|
||||||
|
// len=2 >= waterlevel=0 → returns 2 >> 1 = 1
|
||||||
|
assert_eq!(mb3.drain_count(), 1);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn drain_count_does_not_mutate() {
|
||||||
|
let mut mb: Mailbox<u64> = Mailbox::new(10);
|
||||||
|
for i in 0..5 {
|
||||||
|
mb.push(i);
|
||||||
|
}
|
||||||
|
let dc1 = mb.drain_count();
|
||||||
|
let dc2 = mb.drain_count();
|
||||||
|
assert_eq!(dc1, dc2);
|
||||||
|
assert_eq!(mb.len(), 5);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn drain_count_large_waterlevel() {
|
||||||
|
let mut mb: Mailbox<u64> = Mailbox::new(usize::MAX);
|
||||||
|
for i in 0..100 {
|
||||||
|
mb.push(i);
|
||||||
|
}
|
||||||
|
// len=100 < waterlevel=usize::MAX → always full drain
|
||||||
|
assert_eq!(mb.drain_count(), 100);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn mailbox_with_struct_messages() {
|
||||||
|
let mut mb: Mailbox<TestMsg> = Mailbox::new(10);
|
||||||
|
mb.push(TestMsg(1));
|
||||||
|
mb.push(TestMsg(2));
|
||||||
|
assert_eq!(mb.pop(), Some(TestMsg(1)));
|
||||||
|
assert_eq!(mb.pop(), Some(TestMsg(2)));
|
||||||
|
assert!(mb.is_empty());
|
||||||
|
}
|
||||||
|
|
||||||
|
// ── ActorPool tests ─────────────────────────────────────────────────
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn pool_new_is_empty() {
|
||||||
|
let pool = ActorPool::new();
|
||||||
|
assert_eq!(pool.len(), 0);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn pool_insert_increments_len() {
|
||||||
|
let mut pool = ActorPool::new();
|
||||||
|
let addr = make_addr(1);
|
||||||
|
let (actor, _) = make_test_actor(addr);
|
||||||
|
pool.insert(addr, actor);
|
||||||
|
assert_eq!(pool.len(), 1);
|
||||||
|
|
||||||
|
let addr2 = make_addr(2);
|
||||||
|
let (actor2, _) = make_test_actor(addr2);
|
||||||
|
pool.insert(addr2, actor2);
|
||||||
|
assert_eq!(pool.len(), 2);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn pool_remove_returns_actor() {
|
||||||
|
let mut pool = ActorPool::new();
|
||||||
|
let addr = make_addr(1);
|
||||||
|
let (actor, _) = make_test_actor(addr);
|
||||||
|
pool.insert(addr, actor);
|
||||||
|
assert!(pool.remove(&addr).is_some());
|
||||||
|
assert_eq!(pool.len(), 0);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn pool_remove_unknown_returns_none() {
|
||||||
|
let mut pool = ActorPool::new();
|
||||||
|
let addr = make_addr(99);
|
||||||
|
assert!(pool.remove(&addr).is_none());
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn pool_deliver_correct_type() {
|
||||||
|
let mut pool = ActorPool::new();
|
||||||
|
let addr = make_addr(1);
|
||||||
|
let (actor, _) = make_test_actor(addr);
|
||||||
|
pool.insert(addr, actor);
|
||||||
|
|
||||||
|
let msg: Box<dyn Any + Send> = Box::new(42u64);
|
||||||
|
assert!(pool.deliver(&addr, msg));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn pool_deliver_unknown_addr() {
|
||||||
|
let mut pool = ActorPool::new();
|
||||||
|
let addr = make_addr(99);
|
||||||
|
let msg: Box<dyn Any + Send> = Box::new(42u64);
|
||||||
|
assert!(!pool.deliver(&addr, msg));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn pool_deliver_wrong_type() {
|
||||||
|
let mut pool = ActorPool::new();
|
||||||
|
let addr = make_addr(1);
|
||||||
|
let (actor, _) = make_test_actor(addr);
|
||||||
|
pool.insert(addr, actor);
|
||||||
|
|
||||||
|
// Actor expects u64, we send String — queued (type check deferred to tick)
|
||||||
|
let msg: Box<dyn Any + Send> = Box::new("wrong type".to_string());
|
||||||
|
assert!(pool.deliver(&addr, msg));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn pool_tick_all_processes_messages() {
|
||||||
|
let mut pool = ActorPool::new();
|
||||||
|
let addr = make_addr(1);
|
||||||
|
let (actor, counter) = make_test_actor(addr);
|
||||||
|
pool.insert(addr, actor);
|
||||||
|
|
||||||
|
// Deliver 3 messages
|
||||||
|
pool.deliver(&addr, Box::new(1u64));
|
||||||
|
pool.deliver(&addr, Box::new(2u64));
|
||||||
|
pool.deliver(&addr, Box::new(3u64));
|
||||||
|
|
||||||
|
let stub = StubContextInner { waterlevel: 100 };
|
||||||
|
let did_work = pool.tick_all(&stub);
|
||||||
|
assert!(did_work);
|
||||||
|
assert_eq!(counter.load(Ordering::Relaxed), 3);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn pool_tick_all_empty_returns_false() {
|
||||||
|
let mut pool = ActorPool::new();
|
||||||
|
let addr = make_addr(1);
|
||||||
|
let (actor, _) = make_test_actor(addr);
|
||||||
|
pool.insert(addr, actor);
|
||||||
|
|
||||||
|
// No messages delivered
|
||||||
|
let stub = StubContextInner { waterlevel: 100 };
|
||||||
|
let did_work = pool.tick_all(&stub);
|
||||||
|
assert!(!did_work);
|
||||||
|
}
|
||||||
|
|
||||||
|
// ── Worker tests ────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn worker_tick_once_no_work() {
|
||||||
|
let transfer_rx = Receiver::<Envelope>::new(64);
|
||||||
|
let spawn_rx = Receiver::<(ActorAddress, Box<dyn AnyActor>)>::new(64);
|
||||||
|
|
||||||
|
let transfer_tx = transfer_rx.new_sender();
|
||||||
|
let spawn_tx = spawn_rx.new_sender();
|
||||||
|
|
||||||
|
let mut worker = Worker::new(WorkerId(0), transfer_rx, spawn_rx);
|
||||||
|
|
||||||
|
let address_map = AddressMap::new();
|
||||||
|
let placement = Placement::new(1);
|
||||||
|
let inbox_registry = InboxRegistry::new();
|
||||||
|
let config = RuntimeConfig::default();
|
||||||
|
|
||||||
|
let tc = TickContext {
|
||||||
|
address_map: &address_map,
|
||||||
|
transfer_txs: &[transfer_tx],
|
||||||
|
spawn_txs: &[spawn_tx],
|
||||||
|
placement: &placement,
|
||||||
|
inbox_registry: &inbox_registry,
|
||||||
|
config: &config,
|
||||||
|
};
|
||||||
|
|
||||||
|
assert!(!worker.tick_once(&tc));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn worker_tick_once_drains_spawns() {
|
||||||
|
let transfer_rx = Receiver::<Envelope>::new(64);
|
||||||
|
let spawn_rx = Receiver::<(ActorAddress, Box<dyn AnyActor>)>::new(64);
|
||||||
|
|
||||||
|
let transfer_tx = transfer_rx.new_sender();
|
||||||
|
let spawn_tx = spawn_rx.new_sender();
|
||||||
|
|
||||||
|
let addr = make_addr(1);
|
||||||
|
let (actor, _) = make_test_actor(addr);
|
||||||
|
spawn_tx.try_send((addr, actor)).ok().unwrap();
|
||||||
|
|
||||||
|
let mut worker = Worker::new(WorkerId(0), transfer_rx, spawn_rx);
|
||||||
|
|
||||||
|
let address_map = AddressMap::new();
|
||||||
|
let placement = Placement::new(1);
|
||||||
|
let inbox_registry = InboxRegistry::new();
|
||||||
|
let config = RuntimeConfig::default();
|
||||||
|
|
||||||
|
let tc = TickContext {
|
||||||
|
address_map: &address_map,
|
||||||
|
transfer_txs: &[transfer_tx],
|
||||||
|
spawn_txs: &[spawn_tx],
|
||||||
|
placement: &placement,
|
||||||
|
inbox_registry: &inbox_registry,
|
||||||
|
config: &config,
|
||||||
|
};
|
||||||
|
|
||||||
|
assert!(worker.tick_once(&tc));
|
||||||
|
assert_eq!(worker.pool.len(), 1);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn worker_tick_once_drains_transfers() {
|
||||||
|
let transfer_rx = Receiver::<Envelope>::new(64);
|
||||||
|
let spawn_rx = Receiver::<(ActorAddress, Box<dyn AnyActor>)>::new(64);
|
||||||
|
|
||||||
|
let transfer_tx = transfer_rx.new_sender();
|
||||||
|
let transfer_tx2 = transfer_rx.new_sender();
|
||||||
|
let spawn_tx = spawn_rx.new_sender();
|
||||||
|
let spawn_tx2 = spawn_rx.new_sender();
|
||||||
|
|
||||||
|
let addr = make_addr(1);
|
||||||
|
let (actor, counter) = make_test_actor(addr);
|
||||||
|
|
||||||
|
let mut worker = Worker::new(WorkerId(0), transfer_rx, spawn_rx);
|
||||||
|
|
||||||
|
// Spawn the actor first
|
||||||
|
spawn_tx.try_send((addr, actor)).ok().unwrap();
|
||||||
|
|
||||||
|
let address_map = AddressMap::new();
|
||||||
|
let placement = Placement::new(1);
|
||||||
|
let inbox_registry = InboxRegistry::new();
|
||||||
|
let config = RuntimeConfig::default();
|
||||||
|
|
||||||
|
let tc = TickContext {
|
||||||
|
address_map: &address_map,
|
||||||
|
transfer_txs: &[transfer_tx2],
|
||||||
|
spawn_txs: &[spawn_tx2],
|
||||||
|
placement: &placement,
|
||||||
|
inbox_registry: &inbox_registry,
|
||||||
|
config: &config,
|
||||||
|
};
|
||||||
|
|
||||||
|
worker.tick_once(&tc); // drain spawns
|
||||||
|
|
||||||
|
// Now send a transfer envelope
|
||||||
|
let envelope = Envelope::new(addr, Box::new(42u64));
|
||||||
|
transfer_tx.try_send(envelope).ok().unwrap();
|
||||||
|
|
||||||
|
// Tick again to drain transfers + tick actors
|
||||||
|
let did_work = worker.tick_once(&tc);
|
||||||
|
assert!(did_work);
|
||||||
|
assert_eq!(counter.load(Ordering::Relaxed), 1);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn worker_tick_once_processes_messages() {
|
||||||
|
let transfer_rx = Receiver::<Envelope>::new(64);
|
||||||
|
let spawn_rx = Receiver::<(ActorAddress, Box<dyn AnyActor>)>::new(64);
|
||||||
|
|
||||||
|
let transfer_tx = transfer_rx.new_sender();
|
||||||
|
let transfer_tx2 = transfer_rx.new_sender();
|
||||||
|
let spawn_tx = spawn_rx.new_sender();
|
||||||
|
let spawn_tx2 = spawn_rx.new_sender();
|
||||||
|
|
||||||
|
let addr = make_addr(1);
|
||||||
|
let (actor, counter) = make_test_actor(addr);
|
||||||
|
|
||||||
|
let mut worker = Worker::new(WorkerId(0), transfer_rx, spawn_rx);
|
||||||
|
|
||||||
|
spawn_tx.try_send((addr, actor)).ok().unwrap();
|
||||||
|
|
||||||
|
let address_map = AddressMap::new();
|
||||||
|
let placement = Placement::new(1);
|
||||||
|
let inbox_registry = InboxRegistry::new();
|
||||||
|
let config = RuntimeConfig::default();
|
||||||
|
|
||||||
|
let tc = TickContext {
|
||||||
|
address_map: &address_map,
|
||||||
|
transfer_txs: &[transfer_tx2],
|
||||||
|
spawn_txs: &[spawn_tx2],
|
||||||
|
placement: &placement,
|
||||||
|
inbox_registry: &inbox_registry,
|
||||||
|
config: &config,
|
||||||
|
};
|
||||||
|
|
||||||
|
worker.tick_once(&tc); // spawn
|
||||||
|
|
||||||
|
// Deliver multiple messages
|
||||||
|
for i in 0..5u64 {
|
||||||
|
transfer_tx
|
||||||
|
.try_send(Envelope::new(addr, Box::new(i)))
|
||||||
|
.ok()
|
||||||
|
.unwrap();
|
||||||
|
}
|
||||||
|
|
||||||
|
worker.tick_once(&tc); // transfer + tick
|
||||||
|
assert_eq!(counter.load(Ordering::Relaxed), 5);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn worker_tick_once_multiple_spawns() {
|
||||||
|
let transfer_rx = Receiver::<Envelope>::new(64);
|
||||||
|
let spawn_rx = Receiver::<(ActorAddress, Box<dyn AnyActor>)>::new(64);
|
||||||
|
|
||||||
|
let transfer_tx = transfer_rx.new_sender();
|
||||||
|
let spawn_tx = spawn_rx.new_sender();
|
||||||
|
|
||||||
|
for i in 0..5u8 {
|
||||||
|
let addr = make_addr(i);
|
||||||
|
let (actor, _) = make_test_actor(addr);
|
||||||
|
spawn_tx.try_send((addr, actor)).ok().unwrap();
|
||||||
|
}
|
||||||
|
|
||||||
|
let mut worker = Worker::new(WorkerId(0), transfer_rx, spawn_rx);
|
||||||
|
|
||||||
|
let address_map = AddressMap::new();
|
||||||
|
let placement = Placement::new(1);
|
||||||
|
let inbox_registry = InboxRegistry::new();
|
||||||
|
let config = RuntimeConfig::default();
|
||||||
|
|
||||||
|
let tc = TickContext {
|
||||||
|
address_map: &address_map,
|
||||||
|
transfer_txs: &[transfer_tx],
|
||||||
|
spawn_txs: &[spawn_tx],
|
||||||
|
placement: &placement,
|
||||||
|
inbox_registry: &inbox_registry,
|
||||||
|
config: &config,
|
||||||
|
};
|
||||||
|
|
||||||
|
assert!(worker.tick_once(&tc));
|
||||||
|
assert_eq!(worker.pool.len(), 5);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn worker_tick_once_wrong_type_no_panic() {
|
||||||
|
let transfer_rx = Receiver::<Envelope>::new(64);
|
||||||
|
let spawn_rx = Receiver::<(ActorAddress, Box<dyn AnyActor>)>::new(64);
|
||||||
|
|
||||||
|
let transfer_tx = transfer_rx.new_sender();
|
||||||
|
let transfer_tx2 = transfer_rx.new_sender();
|
||||||
|
let spawn_tx = spawn_rx.new_sender();
|
||||||
|
let spawn_tx2 = spawn_rx.new_sender();
|
||||||
|
|
||||||
|
let addr = make_addr(1);
|
||||||
|
let (actor, counter) = make_test_actor(addr);
|
||||||
|
|
||||||
|
let mut worker = Worker::new(WorkerId(0), transfer_rx, spawn_rx);
|
||||||
|
|
||||||
|
spawn_tx.try_send((addr, actor)).ok().unwrap();
|
||||||
|
|
||||||
|
let address_map = AddressMap::new();
|
||||||
|
let placement = Placement::new(1);
|
||||||
|
let inbox_registry = InboxRegistry::new();
|
||||||
|
let config = RuntimeConfig::default();
|
||||||
|
|
||||||
|
let tc = TickContext {
|
||||||
|
address_map: &address_map,
|
||||||
|
transfer_txs: &[transfer_tx2],
|
||||||
|
spawn_txs: &[spawn_tx2],
|
||||||
|
placement: &placement,
|
||||||
|
inbox_registry: &inbox_registry,
|
||||||
|
config: &config,
|
||||||
|
};
|
||||||
|
|
||||||
|
worker.tick_once(&tc); // spawn
|
||||||
|
|
||||||
|
// Send wrong type (String instead of u64) — should not panic
|
||||||
|
let bad_envelope = Envelope::new(addr, Box::new("wrong".to_string()));
|
||||||
|
transfer_tx.try_send(bad_envelope).ok().unwrap();
|
||||||
|
|
||||||
|
// Send correct type after
|
||||||
|
let good_envelope = Envelope::new(addr, Box::new(99u64));
|
||||||
|
transfer_tx.try_send(good_envelope).ok().unwrap();
|
||||||
|
|
||||||
|
worker.tick_once(&tc); // transfer + tick
|
||||||
|
// The correct message should still be processed
|
||||||
|
assert_eq!(counter.load(Ordering::Relaxed), 1);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn worker_run_stops_on_signal() {
|
||||||
|
let transfer_rx = Receiver::<Envelope>::new(64);
|
||||||
|
let spawn_rx = Receiver::<(ActorAddress, Box<dyn AnyActor>)>::new(64);
|
||||||
|
|
||||||
|
let transfer_tx = transfer_rx.new_sender();
|
||||||
|
let spawn_tx = spawn_rx.new_sender();
|
||||||
|
|
||||||
|
let mut worker = Worker::new(WorkerId(0), transfer_rx, spawn_rx);
|
||||||
|
let is_running = AtomicBool::new(false); // start as false → should exit immediately
|
||||||
|
let backoff = BackoffPolicy::default();
|
||||||
|
|
||||||
|
let address_map = AddressMap::new();
|
||||||
|
let placement = Placement::new(1);
|
||||||
|
let inbox_registry = InboxRegistry::new();
|
||||||
|
let config = RuntimeConfig::default();
|
||||||
|
|
||||||
|
let tc = TickContext {
|
||||||
|
address_map: &address_map,
|
||||||
|
transfer_txs: &[transfer_tx],
|
||||||
|
spawn_txs: &[spawn_tx],
|
||||||
|
placement: &placement,
|
||||||
|
inbox_registry: &inbox_registry,
|
||||||
|
config: &config,
|
||||||
|
};
|
||||||
|
|
||||||
|
// Run in a scoped thread to verify it actually terminates
|
||||||
|
thread::scope(|s| {
|
||||||
|
s.spawn(|| {
|
||||||
|
worker.run(&tc, &is_running, &backoff);
|
||||||
|
});
|
||||||
|
});
|
||||||
|
// If we get here, the thread exited — test passes
|
||||||
|
}
|
||||||
162
tests/test_python.py
Normal file
162
tests/test_python.py
Normal file
|
|
@ -0,0 +1,162 @@
|
||||||
|
"""Tests for swactor Python bindings."""
|
||||||
|
|
||||||
|
import unittest
|
||||||
|
from swactor import Runtime, RuntimeConfig, ActorAddress
|
||||||
|
|
||||||
|
|
||||||
|
class TestActorAddress(unittest.TestCase):
|
||||||
|
def test_repr(self):
|
||||||
|
rt = Runtime()
|
||||||
|
addr = rt.spawn(lambda ctx, msg: None)
|
||||||
|
r = repr(addr)
|
||||||
|
self.assertTrue(r.startswith("ActorAddress("))
|
||||||
|
self.assertTrue(r.endswith(")"))
|
||||||
|
# hex string should be 64 chars (32 bytes)
|
||||||
|
hex_part = r[len("ActorAddress("):-1]
|
||||||
|
self.assertEqual(len(hex_part), 64)
|
||||||
|
|
||||||
|
def test_hex(self):
|
||||||
|
rt = Runtime()
|
||||||
|
addr = rt.spawn(lambda ctx, msg: None)
|
||||||
|
self.assertEqual(len(addr.hex()), 64)
|
||||||
|
|
||||||
|
def test_to_bytes(self):
|
||||||
|
rt = Runtime()
|
||||||
|
addr = rt.spawn(lambda ctx, msg: None)
|
||||||
|
self.assertEqual(len(addr.to_bytes()), 32)
|
||||||
|
|
||||||
|
def test_equality(self):
|
||||||
|
rt = Runtime()
|
||||||
|
addr = rt.spawn(lambda ctx, msg: None)
|
||||||
|
# Same address object should be equal to itself
|
||||||
|
self.assertEqual(addr, addr)
|
||||||
|
|
||||||
|
def test_hashable(self):
|
||||||
|
rt = Runtime()
|
||||||
|
addr1 = rt.spawn(lambda ctx, msg: None)
|
||||||
|
addr2 = rt.spawn(lambda ctx, msg: None)
|
||||||
|
s = {addr1, addr2}
|
||||||
|
self.assertEqual(len(s), 2)
|
||||||
|
s.add(addr1)
|
||||||
|
self.assertEqual(len(s), 2)
|
||||||
|
|
||||||
|
|
||||||
|
class TestRuntimeConfig(unittest.TestCase):
|
||||||
|
def test_defaults(self):
|
||||||
|
cfg = RuntimeConfig()
|
||||||
|
self.assertEqual(cfg.num_threads, 1)
|
||||||
|
self.assertEqual(cfg.max_actors, 1000)
|
||||||
|
self.assertEqual(cfg.actor_max_messages, 1000)
|
||||||
|
self.assertEqual(cfg.mailbox_waterlevel, 10)
|
||||||
|
self.assertEqual(cfg.spin_threshold, 64)
|
||||||
|
self.assertEqual(cfg.yield_threshold, 256)
|
||||||
|
self.assertEqual(cfg.sleep_increment_us, 50)
|
||||||
|
self.assertEqual(cfg.sleep_max_us, 1000)
|
||||||
|
|
||||||
|
def test_custom(self):
|
||||||
|
cfg = RuntimeConfig(num_threads=4, max_actors=500)
|
||||||
|
self.assertEqual(cfg.num_threads, 4)
|
||||||
|
self.assertEqual(cfg.max_actors, 500)
|
||||||
|
|
||||||
|
|
||||||
|
class TestSingleThreaded(unittest.TestCase):
|
||||||
|
def test_echo(self):
|
||||||
|
"""Spawn an echo actor, send a message, tick, and recv."""
|
||||||
|
rt = Runtime()
|
||||||
|
|
||||||
|
def echo(ctx, msg):
|
||||||
|
ctx.send(msg["reply_to"], msg["payload"])
|
||||||
|
|
||||||
|
addr = rt.spawn(echo)
|
||||||
|
inbox = rt.inbox()
|
||||||
|
rt.send(addr, {"payload": "hello", "reply_to": inbox.addr})
|
||||||
|
rt.tick()
|
||||||
|
result = inbox.try_recv()
|
||||||
|
self.assertEqual(result, "hello")
|
||||||
|
|
||||||
|
def test_spawn_from_handler(self):
|
||||||
|
"""Actor spawns a child and forwards work to it."""
|
||||||
|
rt = Runtime()
|
||||||
|
|
||||||
|
def child(ctx, msg):
|
||||||
|
ctx.send(msg["reply_to"], "from_child")
|
||||||
|
|
||||||
|
def parent(ctx, msg):
|
||||||
|
c = ctx.spawn(child)
|
||||||
|
ctx.send(c, {"reply_to": msg["reply_to"]})
|
||||||
|
|
||||||
|
addr = rt.spawn(parent)
|
||||||
|
inbox = rt.inbox()
|
||||||
|
rt.send(addr, {"reply_to": inbox.addr})
|
||||||
|
# First tick: parent runs, spawns child, sends to child
|
||||||
|
rt.tick()
|
||||||
|
# Second tick: child runs, sends to inbox
|
||||||
|
rt.tick()
|
||||||
|
result = inbox.try_recv()
|
||||||
|
self.assertEqual(result, "from_child")
|
||||||
|
|
||||||
|
def test_stateful_actor(self):
|
||||||
|
"""Callable class maintains state across messages."""
|
||||||
|
rt = Runtime()
|
||||||
|
|
||||||
|
class Counter:
|
||||||
|
def __init__(self):
|
||||||
|
self.n = 0
|
||||||
|
|
||||||
|
def __call__(self, ctx, msg):
|
||||||
|
self.n += 1
|
||||||
|
ctx.send(msg["reply_to"], self.n)
|
||||||
|
|
||||||
|
addr = rt.spawn(Counter())
|
||||||
|
inbox = rt.inbox()
|
||||||
|
rt.send(addr, {"reply_to": inbox.addr})
|
||||||
|
rt.send(addr, {"reply_to": inbox.addr})
|
||||||
|
rt.tick()
|
||||||
|
self.assertEqual(inbox.try_recv(), 1)
|
||||||
|
self.assertEqual(inbox.try_recv(), 2)
|
||||||
|
|
||||||
|
def test_no_message_returns_none(self):
|
||||||
|
rt = Runtime()
|
||||||
|
inbox = rt.inbox()
|
||||||
|
self.assertIsNone(inbox.try_recv())
|
||||||
|
|
||||||
|
|
||||||
|
class TestMultiThreaded(unittest.TestCase):
|
||||||
|
def test_run_shutdown_join(self):
|
||||||
|
"""Multi-threaded runtime can spawn, send, and receive."""
|
||||||
|
import time
|
||||||
|
|
||||||
|
rt = Runtime(RuntimeConfig(num_threads=2))
|
||||||
|
|
||||||
|
def echo(ctx, msg):
|
||||||
|
ctx.send(msg["reply_to"], msg["payload"])
|
||||||
|
|
||||||
|
addr = rt.spawn(echo)
|
||||||
|
inbox = rt.inbox()
|
||||||
|
handle = rt.run()
|
||||||
|
handle.send(addr, {"payload": "mt_hello", "reply_to": inbox.addr})
|
||||||
|
|
||||||
|
# Poll for result
|
||||||
|
result = None
|
||||||
|
for _ in range(100):
|
||||||
|
result = inbox.try_recv()
|
||||||
|
if result is not None:
|
||||||
|
break
|
||||||
|
time.sleep(0.01)
|
||||||
|
self.assertEqual(result, "mt_hello")
|
||||||
|
|
||||||
|
handle.shutdown()
|
||||||
|
handle.join()
|
||||||
|
|
||||||
|
def test_run_consumes_runtime(self):
|
||||||
|
"""After run(), tick() should raise."""
|
||||||
|
rt = Runtime(RuntimeConfig(num_threads=2))
|
||||||
|
handle = rt.run()
|
||||||
|
with self.assertRaises(RuntimeError):
|
||||||
|
rt.tick()
|
||||||
|
handle.shutdown()
|
||||||
|
handle.join()
|
||||||
|
|
||||||
|
|
||||||
|
if __name__ == "__main__":
|
||||||
|
unittest.main()
|
||||||
8
uv.lock
Normal file
8
uv.lock
Normal file
|
|
@ -0,0 +1,8 @@
|
||||||
|
version = 1
|
||||||
|
revision = 3
|
||||||
|
requires-python = ">=3.9"
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "swactor"
|
||||||
|
version = "0.1.0"
|
||||||
|
source = { editable = "." }
|
||||||
Loading…
Reference in a new issue