RX-MBT

    Rx-Rust 的 MoonBit 复刻:基于 Observer 模式的响应式编程库(Observable / Subject / Scheduler / 操作符),零外部依赖。

    reactive
    rx
    observable
    observer
    stream
    Download zip
    Author
    Version
    0.1.0
    License
    Apache-2.0
    Last updated
    last month
    Downloads
    6
    # RX-MBT

    Rx-Rust 的 MoonBit 复刻:基于 Observer 模式的响应式编程库,纯 MoonBit 实现,零外部依赖。

    #作为三方依赖使用

    本仓库是标准 MoonBit 库模块,库代码在 rx/ 子包中,可通过两种方式引入:

    #方式一:git 依赖(无需发布到中心仓库)

    # 在你的项目里执行(仓库地址替换为实际 git 地址) moon add git+https://<你的仓库地址>.git

    然后在依赖包的 moon.pkg 中声明并起别名:

    import { "vicTop-cw/RX-MBT/rx" @rx, }

    #方式二:发布到 MoonBit 中心仓库

    moon publish # 之后任何项目可用: moon add vicTop-cw/RX-MBT

    #使用示例(外部项目视角)

    // moon.pkg 中已 import "vicTop-cw/RX-MBT/rx" @rx
    fn main {
    // 链式调用:map -> filter -> collect
    let result = @rx.Observable::from_iter([1, 2, 3, 4, 5])
    .map(fn(x) { x * 2 })
    .filter(fn(x) { x > 4 })
    let _ = result.subscribe_on_next(fn(v) { println(v) })

    // Subject 多播
    let subject : @rx.PublishSubject[Int, Unit] = @rx.PublishSubject::new()
    let _ = subject.subscribe(@rx.Observer::on_next_only(fn(v) { println("收到: \{v}") }))
    subject.on_next(42)
    subject.on_completed()
    }

    本仓库自带的 cmd/main 就是一个完整可运行的示例:moon run cmd/main

    #启动与测试

    moon check # 类型检查(零警告零错误) moon test # 运行全部测试

    完整的 API 签名与语义说明见 DOC.md(按分类整理的全部公开 API)。

    #版本与测试环境

    本项目在以下环境验证通过:

    项目
    MoonBit 工具链moon 0.1.20260724 (5f1406a 2026-07-24)moon version 输出)
    操作系统Windows 11(PowerShell 7 / Windows PowerShell 5.1 均可运行脚本)
    编译目标wasmpreferred_target = "wasm",见 moon.mod
    测试后端moon test 默认 wasm 后端
    库版本vicTop-cw/RX-MBT@0.1.0moon.modversion 字段)
    测试数量377(基本 + 高压 + Haskell API 回归,全部通过,见 TEST_REPORT.md
    基准测试test-suit/bench#skip 默认跳过,moon test -g "bench" 运行)

    工具链版本以 moon version 输出为准。要求 moon >= 0.1.20260724(本项目使用了较新的 moon.mod / moon.pkg 语法与 Feature flags:rr_moon_mod, rr_moon_pkg)。

    #测试报告

    运行报告生成脚本(PowerShell,无需额外依赖):

    powershell -ExecutionPolicy Bypass -File scripts/gen_test_report.ps1

    产物:

    文件说明
    TEST_REPORT.md测试环境、汇总统计、按包分类的完整测试清单
    TEST_REPORT.json机器可读报告(供 CI / 对比脚本消费)

    报告包含 moon 版本、生成时间、总测试数/通过/失败,以及每个测试包的全部测试名清单。

    测试按分类组织在 test-suit/ 下(每个子目录独立成包,黑盒引用 @rx):

    目录覆盖
    test-suit/createof / from_iter / empty / never / error / repeat / range / from_callable / start
    test-suit/transformmap / filter_map / flat_map / scan / concat_map / switch_map / group_by / curry_map
    test-suit/filterfilter / take / skip / first / last / distinct / skip_until / take_until / contains
    test-suit/aggregationcount / reduce / collect / sum / average / min / max / mean / median / variance / rolling_* / to_map / to_set
    test-suit/compositestart_with / concat / merge / zip / combine_latest / with_latest_from / amb / iif / sequence_equal
    test-suit/math_conditionsdrop_none / fill_none / clamp / abs / every / some / all / find / is_empty
    test-suit/corerun / debug / Subject / Scheduler / cum_* 等基础行为
    test-suit/itorItor / Node 可控制迭代器:send 插队 / set_pause / resume_iter / restart 重放 / stop / history 策略
    test-suit/haskellHaskell 迭代器 API(Data.List / Data.Foldable 语义):iterate / scanl1 / foldr1 / map_accum_l / inits / tails / union / insert / lines / words 等
    test-suit/stress高压测试:10 万元素链、多订阅者、Subject 高压、递归重试
    test-suit/bench基准测试(#skip 默认跳过,moon test -g bench 运行)

    基准测试以控制台 JSON 输出结果(每个用例一行),可与 Rx-Rust 的 benches/bench.rs 结果对比。如需保存:

    moon test -p test-suit/bench --include-skipped > moon_bench.json

    #核心抽象

    类型说明
    Observable[T, E]可观察对象(冷流),包装订阅函数,每次 subscribe 重新执行
    Observer[T, E]观察者,承载 on_next / on_error / on_completed 三个回调
    Subscription订阅句柄,dispose() 幂等取消订阅
    Clock虚拟时钟,时间类操作符基于它实现,advance_time(ms) 推进
    Node[T]惰性链表节点:持有值 + 懒加载源 + fallback,from_iter 构建
    Itor[T]可控制迭代器(vools 风格):支持 send 插队 / 暂停 / 重放 / 停止 / 历史策略

    #Itor / Node(可控制迭代器)

    复刻 vools 的 Itor 相关 API,包装任意可重复迭代的源(() -> Iter[T]):

    API说明
    Node::new(val, next?) / Node::from_iter(iter)构造节点 / 从迭代器构建惰性链表
    Node::val() / Node::next()取值 / 取下一节点(懒加载展开)
    Node::to_itor() / Node::to_iter()节点链表 → Itor / Iter
    Itor::new(source) / from_array / from_iter构造(数组源可重复迭代,副本独立)
    Itor::next()取下一个值(自动初始化源)
    Itor::send(node, jump_when?) / send_value(v) / send_values(arr)紧急插队 / 条件插队
    Itor::set_pause() / resume_iter() / stop()暂停 / 恢复 / 终止
    Itor::restart()回到 Pending,从历史头部重放,重放后继续源
    Itor::history_strategy(fn) / set_history_max(n)历史保留策略(-1 全保留 / n 条窗口)
    Itor::state() / is_pending / is_iterring / is_paused / is_stopped状态查询
    Itor::copy() / call()返回独立副本(共享源工厂)
    Itor::do_iter(f, pre_f?, sub_f?)副作用遍历(支持前置 / 后置钩子)

    取值顺序:历史重放(restart 后)→ 紧急插队 → 待处理源值 → 源迭代器。

    #创建类工厂

    from_iter / of / repeat / empty / never / error / range / from_range / from_range_with_step / from_callable / start / interval / timer

    #操作符

    转换类map / filter_map / flat_map / flat_map_latest / scan / concat_map / switch_map / group_by / curry_map

    过滤类filter / take / skip / first / last / take_while / skip_while / skip_n_events / take_n_events / skip_last / take_last / element_at / distinct / distinct_by / distinct_until_changed / distinct_until_changed_by / skip_until / take_until / contains / includes / sort / top_k / bottom_k / drop_none / fill_none / clamp

    聚合类count / reduce / collect / to_list / sum / average / minimum / maximum / min / max / mean / median / variance / std / quantile / n_unique / arg_min / arg_max / cum_prod / cum_mean / cum_sum / cum_min / cum_max / rolling_sum / rolling_mean / rolling_min / rolling_max / rolling_count / to_map / to_set / abs

    组合类start_with / concat / merge / zip / combine_latest / with_latest_from / amb / end_with / iif / sequence_equal

    条件类every / all / some / find / find_index / is_empty

    错误处理类retry / retry_indefinitely / retry_when / retry_with_backoff / retry_with_backoff_on / catch_error / on_error_return / on_error_resume_next / circuit_breaker / circuit_breaker_on(背压操作符 backpressure_* 保留签名、同步模型下透传)

    时间类(基于虚拟时钟):delay / debounce / throttle / throttle_first / rate_limit / timeout / sample / timestamp(均提供 _on(clock) 变体)

    工具类tap / tap_on_next / default_if_empty / ignore_elements / pairwise / pairwise_with_buffer / buffer_count / buffer_time / window / switch / sample_first / skip_until_data / take_until_data / run / debug

    #Subject(多播主题)

    PublishSubject / BehaviorSubject / ReplaySubjectwith_buffer_size 限制重放窗口),均提供 subscribe / on_next / on_error / on_completed / as_observable

    #调度器

    同步语义 + 虚拟时钟等价实现(MoonBit wasm 目标无线程):

    CurrentThreadScheduler / ImmediateScheduler / AsyncScheduler / ThreadPoolScheduler

    每个提供 now() 返回虚拟时间(毫秒)、schedule(task) 立即同步执行;ThreadPoolScheduler 额外提供 get_num_threads()

    #示例

    // 链式调用:map -> filter -> collect
    let result = Observable::from_iter([1, 2, 3, 4, 5])
    .map(fn(x) { x * 2 })
    .filter(fn(x) { x > 4 })
    .collect()
    let _ = result.subscribe_on_next(fn(v) { println(v) })

    // Subject 多播
    let subject : PublishSubject[Int, Unit] = PublishSubject::new()
    let _ = subject.subscribe_on_next(fn(v) { println("收到: \{v}") })
    subject.on_next(42)
    subject.on_completed()

    #与 Rx-Rust 的差异

    • 监控类 API(文件/文件夹监听、watch_file / watch_folder 等)未移植
    • 时间类操作符基于虚拟时钟而非真实定时器,测试时用 advance_time(ms) 推进
    • 背压与线程池调度器在同步模型下为签名占位,语义等价透传/立即执行