A lightweight, high-performance Actor concurrency framework for MoonBit
Dependencies
// moon.mod
import {
"didiLjf/moonorbit@0.2.0"
}// 1. 定义消息类型
enum PingPongMsg {
Ping
Pong
}
// 2. 定义 Actor 行为 (Behavior)
async fn ping_pong_behavior(_context : Context, _state : Int, msg : PingPongMsg) -> Int raise {
@async.pause()
match msg {
Ping => {
println("Received Ping")
0
}
Pong => {
println("Received Pong")
0
}
}
}
// 3. 在异步上下文中启动系统与 Actor
async fn run_system() -> Unit {
@async.with_task_group((group) => {
// 初始化 Actor 系统
let system = ActorSystem::new("my_system", group)
// 生成 Actor 引用
let actor = system.spawn(ping_pong_behavior, 0)
// 发送消息
actor.send(Ping)
actor.send(Pong)
// 延迟并安全关闭系统
@async.sleep(50)
system.terminate()
})
}moon run --target native cmd/distributed_kvmoon run --target native cmd/work_pullingmoon run --target native cmd/benchmarksmoon fmt --checkmoon checkmoon test --target native --deny-warnmoon test --target js --deny-warn
moon test --target wasm-gc --deny-warn├── .github/workflows/ # CI 自动化工作流配置 (包含 MSVC 支持)
├── cmd/
│ ├── main/ # 默认运行主程序 (Ping-Pong 异步演示)
│ ├── distributed_kv/ # 复杂示例: 主从同步副本键值数据库
│ ├── work_pulling/ # 复杂示例: 拉取型 Worker 工作池
│ └── benchmarks/ # 环形 actor 消息吞吐性能基准测试
├── actor.mbt # Actor 核心定义与信箱包装
├── context.mbt # 运行上下文及子 Actor 级联生成
├── system.mbt # Actor 监管树核心控制流、核心执行循环
├── router.mbt # RoundRobin/Broadcast/Random 路由分发机制
├── supervisor.mbt # 监督策略与基础断言
├── timer.mbt # [v0.2.0] 定时器与周期调度器
├── stash.mbt # [v0.2.0] 状态自愈的二级消息暂存队列
├── fsm.mbt # [v0.2.0] 有限状态机 (FSM) 封装
├── pubsub.mbt # [v0.2.0] 发布订阅中介协调器
├── cluster.mbt # [v0.2.0] 集群网络、透明远程引用 RemoteRef 模拟
├── system_wbtest.mbt # 白盒测试套件 (涵盖容错、背压、路由等)
├── timer_wbtest.mbt # 定时器相关白盒测试
├── stash_wbtest.mbt # 暂存区相关白盒测试
├── fsm_wbtest.mbt # 状态机相关白盒测试
├── pubsub_wbtest.mbt # 发布订阅相关白盒测试
├── cluster_wbtest.mbt # 集群模拟相关白盒测试
├── LICENSE # OSI 认证开源许可证 (Apache-2.0)
└── moon.mod # 模块元数据定义type ActorControlfn[Msg] ActorRef::send_after(self : ActorRef[Msg], system : ActorSystem, delay : Int, msg : Msg) -> TimerSubscriptionfn[Msg] ActorRef::send_repeatedly(self : ActorRef[Msg], system : ActorSystem, initial_delay : Int, period : Int, msg : Msg) -> TimerSubscriptionpub enum ActorStatus {
Starting
Running
Restarting
Stopping
Stopped
Failed(String)
}fn[State, Msg] ActorSystem::spawn(self : ActorSystem, behavior : async (Context, State, Msg) -> State, initial_state : State, supervisor? : Supervisor, lifecycle? : LifecycleCallbacks[State], mailbox_kind? : Kind, parent_id? : Int?, deserializer? : (String) -> Msg) -> ActorRef[Msg]fn[State, Msg] Context::spawn(self : Context, behavior : async (Context, State, Msg) -> State, initial_state : State, supervisor? : Supervisor, lifecycle? : LifecycleCallbacks[State], mailbox_kind? : Kind) -> ActorRef[Msg]pub struct Fsm[State, Data] {
current_state : State
current_data : Data
timer_sub : TimerSubscription?
}fn[State, Data, Msg] Fsm::goto_with_timeout(self : Fsm[State, Data], next_state : State, next_data : Data, timeout : (Int, Msg), self_ref : ActorRef[Msg], system : ActorSystem) -> Unitpub struct NetworkSimulator {
nodes : Map[String, ActorSystem]
drop_probability : Int
latency : Int
seed : Int
}fn NetworkSimulator::register_node(self : NetworkSimulator, name : String, system : ActorSystem) -> Unitfn NetworkSimulator::send_remote(self : NetworkSimulator, target_node : String, target_actor_id : Int, serialized_msg : String) -> Unitpub struct PersistentActor[State, Event, Command] {
actor_id : Int
state : State
journal_ref : ActorRef[JournalMsg]
command_handler : (State, Command) -> Array[Event]
event_handler : (State, Event) -> State
serializer : (Event) -> String
deserializer : (String) -> Event
}fn[State, Event, Command] PersistentActor::new(actor_id : Int, initial_state : State, journal_ref : ActorRef[JournalMsg], command_handler : (State, Command) -> Array[Event], event_handler : (State, Event) -> State, serializer : (Event) -> String, deserializer : (String) -> Event) -> PersistentActor[State, Event, Command]async fn[State, Event, Command] PersistentActor::process_command(self : PersistentActor[State, Event, Command], _context : Context, command : Command) -> Unitasync fn[State, Event, Command] PersistentActor::recover(self : PersistentActor[State, Event, Command], _context : Context) -> Unitfn[State, Event, Command] PersistentActor::state(self : PersistentActor[State, Event, Command]) -> Statepub struct RemoteRef[Msg] {
target_node : String
target_actor_id : Int
simulator : NetworkSimulator
serializer : (Msg) -> String
}fn[Msg] RemoteRef::new(target_node : String, target_actor_id : Int, simulator : NetworkSimulator, serializer : (Msg) -> String) -> RemoteRef[Msg]pub struct Router[Msg] {
workers : Array[ActorRef[Msg]]
strategy : RoutingStrategy
next_index : Ref[Int]
}pub enum RoutingStrategy {
RoundRobin
Broadcast
Random
}pub enum SupervisionStrategy {
OneForOne
OneForAll
RestForOne
}pub enum SystemMessage {
Stop
Restart
ChildFailed(Int, Error)
}async fn journal_behavior(_context : Context, state : EventJournal, msg : JournalMsg) -> EventJournalasync fn[PubMsg] pubsub_behavior(_context : Context, state : PubSubState[PubMsg], msg : PubSubMsg[PubMsg]) -> PubSubState[PubMsg]A lightweight, high-performance Actor concurrency framework for MoonBit
Dependencies