runner API
Luna-Flow/mare_mark/runner 负责测量循环。你描述一个基准测试用例,把它编译为经过验证的计划,然后用种子、环境快照、事件接收器和经过验证的协议 run 它。运行器在计时之前用判定器验证每个实现,校准批次大小,以平衡区组进行测量,并发出原始事件。每一步的理由见 runner 设计。
源码:src/runner/runner.mbt、src/runner/bench_spec.mbt。
import {
"Luna-Flow/mare_mark/model",
"Luna-Flow/mare_mark/event",
"Luna-Flow/mare_mark/runner",
"moonbitlang/async",
}
run 和 execute_operation 是 async 的。请从 async fn main 或 async test 中调用它们;驱动它们的 moonbitlang/async 运行时可用于 native、JS 和 wasm 目标,但不可用于 wasm-gc。子进程工作者需要 native 目标。
以下类型参数在整个包中反复出现:
| 参数 | 含义 |
|---|---|
Scale | 数据集的规模参数,例如 Int |
Input | 夹具为一个数据集物化出的值 |
Prepared | 实现运行所用的值,由夹具的 prepare 产生 |
Expected | 参考判定器计算出的值 |
Output | 实现返回的值 |
Context | 从一次操作传递到下一次操作的状态(无状态时为 Unit) |
State, SinkValue | OutputSink 的累加器和最终值 |
运行计划
run
run 执行一个已编译的计划并返回运行摘要。
pub async fn[Scale, Input, Prepared, Expected, Output, Context, State, SinkValue] run(ValidatedBenchPlan[Scale, Input, Prepared, Expected, Output, Context, State, SinkValue], RunContext) -> @model.RunSummary
对于每个规模,运行器按顺序把输入物化一次,用判定器验证每个实现,对每个实现进行预热和校准,然后以平衡顺序测量 exploratory_samples 和 confirmatory_samples 个区组。每次验证、失败、校准和观测都在发生时发送给接收器;摘要最后发送,接收器的 finish 返回的位置字符串成为返回的摘要中的 artifact_location。
返回的 RunSummary 具有 complete == true、各事件计数、验证统计(passed_count、failed_count、unsupported_count、expected_difference_count)以及环境。它的 run_id 为 protocol_identity(protocol) + ":" + case_id;它标识的是协议和用例,而不是某一次执行,因此请在环境的来源信息中放入唯一 id。
验证失败不会中止测量:该实现仍然会被计时,失败也会被报告,这样报告就能把不匹配显示在它所否定的系列旁边。
fn environment() -> @model.EnvironmentSnapshot {
@model.EnvironmentSnapshot::new(
@model.SemanticEnvironment::new(@model.ExecutionTarget::Native, "moonc 0.10", "", "i32"),
@model.PerformanceEnvironment::new("native", "laptop", "default", 1, "monotonic"),
@model.ProvenanceEnvironment::new("macos", "host", "2026-10-08T00:00:00Z", "HEAD", "run-1"),
)
}
async test "run a plan" {
let square = @runner.Implementation::stateless("square", "1", (x : Int) => {
@model.OperationResult::completed(x * x, ())
})
let plan = @runner.single_step("square", [2, 3])
.with_immutable_input(context => context.dataset_key.scale, x => x.to_string())
.compare([square])
.against_equal(x => x * x, (expected, actual) => expected == actual)
.compile()
.unwrap()
let memory = @event.InMemorySink::new()
let summary = @runner.run(
plan,
@runner.RunContext::new(
environment(),
memory.as_sink(),
42UL,
@runner.ProtocolPreset::QuickCheck.validated(),
),
)
inspect(summary.run_id, content="mmkp_1:1:3:1:square")
inspect(summary.passed_count, content="2")
inspect(summary.observation_count, content="8")
inspect(summary.artifact_location.unwrap(), content="memory://run/mmkp_1:1:3:1:square")
}
使用 QuickCheck 时,每个规模有一个探索性区组和三个验证性区组,因此两个规模和一个实现会产生八个观测。
RunContext
RunContext 携带一次运行除计划以外所需的一切。
pub struct RunContext {
seed : UInt64
environment : @model.EnvironmentSnapshot
sink : @event.ObservationSink
protocol : ValidatedProtocol
}
pub fn RunContext::new(@model.EnvironmentSnapshot, @event.ObservationSink, UInt64, ValidatedProtocol) -> Self
注意 RunContext::new 的参数顺序:环境、接收器、种子、协议。种子通过 GenerationContext.seed 原样传给夹具,并混入区组顺序(参见 balanced_order)。
balanced_order
balanced_order 返回某个区组所用的实现索引的循环轮换。
pub fn balanced_order(Int, Int) -> Array[Int]
balanced_order(k, b) 为 。在 OrderPolicy::BalancedBlocks(order_seed) 下,运行器使用区组偏移 ,其中 ;在 FixedOrder 下使用恒等排列。在任意 个连续区组中,每个实现恰好在每个位置出现一次。
test "rotated block order" {
debug_inspect(@runner.balanced_order(3, 0), content="[0, 1, 2]")
debug_inspect(@runner.balanced_order(3, 1), content="[1, 2, 0]")
debug_inspect(@runner.balanced_order(3, 5), content="[2, 0, 1]")
}
协议
ValidatedProtocol
ValidatedProtocol 是通过了 validate_protocol 的 RunProtocol。
pub struct ValidatedProtocol {
protocol : @model.RunProtocol
}
它只能从 validate_protocol 或 ProtocolPreset 获得,因此计划永远不会以负的样本数或空的迭代范围运行。
validate_protocol
validate_protocol 检查协议,返回一个 ValidatedProtocol 或发现的全部违规。
pub fn validate_protocol(@model.RunProtocol) -> Result[ValidatedProtocol, Array[ProtocolConfigError]]
| 规则 | 错误 |
|---|---|
warmup_iterations >= 0 | InvalidInteger("warmup_iterations", n) |
warmup_time_us 设置时须 >= 0 | InvalidDuration("warmup_time_us", t) |
exploratory_samples >= 0 | InvalidInteger("exploratory_samples", n) |
confirmatory_samples > 0 | InvalidInteger("confirmatory_samples", n) |
target_batch_time_us > 0 | InvalidDuration("target_batch_time_us", t) |
max_sample_time_us > 0 | InvalidDuration("max_sample_time_us", t) |
target_batch_time_us <= max_sample_time_us | InvalidDuration("target_batch_time_us", t) |
0 < min_batch_iterations <= max_batch_iterations | InvalidIterationRange(min, max) |
practical_delta_pct >= 0 | InvalidPracticalDelta(pct) |
所有规则都会被检查;错误按此顺序一并返回。
test "protocol validation reports every problem" {
let protocol = @model.RunProtocol::new(
@model.ExperimentDesign::FixedDatasetRepeatedMeasurements,
-1,
None,
@model.CalibrationProtocol::new(1000.0, 10, 5, 10000.0, @model.BatchPolicy::PerImplementation),
1.0,
@model.OrderPolicy::BalancedBlocks(7UL),
@model.OutlierPolicy::ReportOnly,
@model.ValidationCoverage::EveryDataset,
2,
0,
)
guard @runner.validate_protocol(protocol) is Err(errors) else { fail("expected errors") }
inspect(errors.length(), content="3")
inspect(errors[0] is InvalidInteger("warmup_iterations", -1), content="true")
inspect(errors[1] is InvalidInteger("confirmatory_samples", 0), content="true")
inspect(errors[2] is InvalidIterationRange(10, 5), content="true")
}
ProtocolConfigError
ProtocolConfigError 指明一条被违反的协议规则。
pub(all) enum ProtocolConfigError {
InvalidInteger(String, Int)
InvalidDuration(String, Double)
InvalidIterationRange(Int, Int)
InvalidPracticalDelta(Double)
}
String 负载是字段名;数值是被拒绝的值。
ProtocolPreset
ProtocolPreset 指明一个现成的协议。
pub(all) enum ProtocolPreset {
QuickCheck
Development
RegressionGate
Custom(ValidatedProtocol)
}
pub fn ProtocolPreset::validated(Self) -> ValidatedProtocol
validated 返回预设的协议(对于 Custom,返回其包装的协议)。预设包括:
| 字段 | QuickCheck | Development | RegressionGate |
|---|---|---|---|
experiment_design | FixedDatasetRepeatedMeasurements | FixedDatasetRepeatedMeasurements | HierarchicalDatasetsAndRepeats |
warmup_iterations | 1 | 3 | 10 |
warmup_time_us | None | Some(5000.0) | Some(10000.0) |
target_batch_time_us | 1000 | 5000 | 10000 |
| 批次迭代次数 | 1 到 1000 | 1 到 10000 | 5 到 10000 |
max_sample_time_us | 100000 | 250000 | 1000000 |
batch_policy | PerImplementation | PerImplementation | PerImplementation |
practical_delta_pct | 1.0 | 1.0 | 0.5 |
order_policy | BalancedBlocks(1) | BalancedBlocks(1) | BalancedBlocks(1) |
outlier_policy | ReportOnly | ReportOnly | TukeyFence |
validation_coverage | ConfirmatoryOnly | EveryDataset | EveryMeasurement |
exploratory_samples | 1 | 3 | 5 |
confirmatory_samples | 3 | 10 | 20 |
运行器依据预热、校准、顺序和样本字段行事。experiment_design、outlier_policy、validation_coverage 和 practical_delta_pct 是记录下来的意图:验证总是在计时之前对每个数据集运行一次,离群值过滤和决策在 stats 中进行。
实现
Implementation
Implementation 是基准测试中的一个参赛者。
pub struct Implementation[Prepared, Output, Context] {
id : String
version : String
initial_context : () -> Context
execution : ExecutionMode[Prepared, Output, Context]
synchronize : () -> Unit
}
id 在事件中为它命名,在一个用例内必须唯一且非空;version 会被复制到每个观测和验证中。initial_context 作为每个验证序列和每个测量批次的起点。synchronize 在时钟启动前和停止后立即调用;可用它等待排队中的异步工作(GPU 流、线程池)。
Implementation::stateless
Implementation::stateless 包装一个作用于准备好的输入、不需要上下文的函数。
pub fn[Prepared, Output] Implementation::stateless(String, String, (Prepared) -> @model.OperationResult[Output, Unit]) -> Self[Prepared, Output, Unit]
函数结果中的 next_context 会被忽略并替换为 Some(()),因此无状态操作永远不会提前结束序列。
Implementation::in_process
Implementation::in_process 构建一个在当前进程中运行并传递上下文的实现。
pub fn[Prepared, Output, Context] Implementation::in_process(String, String, () -> Context, (Prepared, Context) -> @model.OperationResult[Output, Context], synchronize? : () -> Unit) -> Self[Prepared, Output, Context]
返回 next_context = Some(c) 以继续使用 c;返回 None 以结束序列(在计时批次中,这会把该观测标记为无效)。
test "a stateful implementation" {
let counter = @runner.Implementation::in_process("counter", "1", () => 0, (step : Int, total : Int) => {
@model.OperationResult::completed(total + step, total + step)
})
inspect(counter.id, content="counter")
inspect((counter.initial_context)(), content="0")
}
Implementation::worker
Implementation::worker 构建一个在单独进程中运行每次操作的实现。
pub fn[Prepared, Output, Context] Implementation::worker(String, String, () -> Context, WorkerSpec[Prepared, Output, Context], synchronize? : () -> Unit) -> Self[Prepared, Output, Context]
把它用于可能崩溃、挂起或破坏内存的代码。崩溃会变成 Aborted 结果,而不会终止基准测试。
ExecutionMode
ExecutionMode 表示操作在哪里运行。
pub enum ExecutionMode[Prepared, Output, Context] {
InProcess((Prepared, Context) -> @model.OperationResult[Output, Context])
Subprocess(WorkerSpec[Prepared, Output, Context])
}
它在包外是只读的;请通过 Implementation 的构造函数构建它。
WorkerSpec
WorkerSpec 描述如何以子进程方式运行一次操作。
pub struct WorkerSpec[Prepared, Output, Context] {
command : String
build_arguments : (Prepared, Context) -> Array[String]
decode : (String) -> @model.OperationResult[Output, Context]
timeout_ms : Int
cwd : String?
}
pub fn[Prepared, Output, Context] WorkerSpec::new(String, (Prepared, Context) -> Array[String], (String) -> @model.OperationResult[Output, Context], timeout_ms? : Int, cwd? : String) -> Self[Prepared, Output, Context]
build_arguments 生成参数向量,decode 把工作进程的 stdout 转换为结果,timeout_ms 默认为 5000,cwd 默认为当前目录。
fn echo_worker() -> @runner.Implementation[Int, String, Unit] {
@runner.Implementation::worker(
"echo",
"1",
() => (),
@runner.WorkerSpec::new(
"sh",
(n, _) => ["-c", "echo " + n.to_string()],
stdout => @model.OperationResult::completed(stdout.trim().to_owned(), ()),
timeout_ms=1000,
),
)
}
test "a worker is described, not started" {
inspect(echo_worker().id, content="echo")
}
execute_operation
execute_operation 运行实现的一次操作。
pub async fn[Prepared, Output, Context] execute_operation(Implementation[Prepared, Output, Context], Prepared, Context) -> @model.OperationResult[Output, Context]
对于 InProcess,它调用该函数。对于 Subprocess,它启动命令,捕获 stdout 和 stderr,并按如下方式映射结果:
| 工作进程结果 | 结果 |
|---|---|
| 退出码为 0 | decode(stdout) 返回的结果和上下文,并附上 stdout、stderr 和退出码 |
| 退出码 | Aborted(code, stderr),没有下一个上下文 |
在 timeout_ms 内未退出 | Timeout(timeout_ms, "worker exceeded timeout");进程被终止 |
| 进程无法启动 | ParseFailure(error text) |
| 不在 native 目标上 | Aborted(-1, "subprocess workers require the native target") |
OutputSink
OutputSink 折叠一个批次的输出,使编译器无法丢弃这些工作。
pub struct OutputSink[Output, State, SinkValue] {
initial : () -> State
fold : (State, Output) -> State
finish : (State) -> SinkValue
}
pub fn[Output, State, SinkValue] OutputSink::new(() -> State, (State, Output) -> State, (State) -> SinkValue) -> Self[Output, State, SinkValue]
pub fn[Output] OutputSink::keep_last() -> Self[Output, Output?, Output?]
fold 在计时区域内、每次操作之后运行;finish 在时钟停止后运行,其值被传给 @bench.Bench::keep。keep_last 保留最后一个输出。当输出必须被消费时,请配合一个廉价的校验和使用 OutputSink::new:
test "a checksum sink" {
let sink : @runner.OutputSink[Int, Int, Int] = @runner.OutputSink::new(
() => 0,
(state, output) => state ^ output,
state => state,
)
let folded = [3, 5, 6].fold(init=(sink.initial)(), (state, x) => (sink.fold)(state, x))
inspect((sink.finish)(folded), content="0")
}
描述用例
single_step
single_step 为每个输入只有一次操作的用例启动简短构建器。
pub fn[Scale] single_step(String, Array[Scale]) -> SingleStepCase[Scale]
构建器链为 single_step(id, scales) .with_immutable_input(generate, fingerprint) .compare(implementations) .against_equal(reference, comparator) .compile()。每一步都会复制它接收到的数组。
Case
Case 是构建器入口点的命名空间。
pub(all) enum Case {
Case
}
pub fn[Scale] Case::single_step(String, Array[Scale]) -> SingleStepCase[Scale]
Case::single_step(id, scales) 与 single_step(id, scales) 相同。
SingleStepCase
SingleStepCase 是一个已有 id 和规模、但还没有输入的用例。
pub struct SingleStepCase[Scale] {
id : String
scales : Array[Scale]
}
pub fn[Scale] SingleStepCase::compile(Self[Scale]) -> Result[Unit, Array[BenchConfigError]]
pub fn[Scale, Input] SingleStepCase::with_immutable_input(Self[Scale], (@model.GenerationContext[Scale]) -> Input, (Input) -> String) -> ImmutableSingleStepCase[Scale, Input]
compile 在这一阶段总是失败,并列出缺少的内容(MissingFixture、EmptyImplementations、MissingOracle、MissingOutputSink,适用时还有 EmptyCaseId 和 EmptyScales)。with_immutable_input 添加一个名为 id + "-fixture"、版本为 "1"、用 @fixture.Fixture::immutable 构建的不可变夹具。
ImmutableSingleStepCase
ImmutableSingleStepCase 是带有不可变输入的用例。
pub struct ImmutableSingleStepCase[Scale, Input] {
id : String
scales : Array[Scale]
fixture : @fixture.Fixture[Scale, Input, Input]
}
pub fn[Scale, Input, Output] ImmutableSingleStepCase::compare(Self[Scale, Input], Array[Implementation[Input, Output, Unit]]) -> ComparedSingleStepCase[Scale, Input, Output]
compare 添加要比较的无状态实现。
ComparedSingleStepCase
ComparedSingleStepCase 是带有输入和实现的用例。
pub struct ComparedSingleStepCase[Scale, Input, Output] {
id : String
scales : Array[Scale]
fixture : @fixture.Fixture[Scale, Input, Input]
implementations : Array[Implementation[Input, Output, Unit]]
}
pub fn[Scale, Input, Expected, Output] ComparedSingleStepCase::against_equal(Self[Scale, Input, Output], (Input) -> Expected, (Expected, Output) -> Bool) -> BenchSpec[Scale, Input, Input, Expected, Output, Unit, Output?, Output?]
against_equal(reference, comparator) 添加一个名为 id + "-reference"、用 @experiment.ReferenceOracle::equal 构建的参考判定器,一个 OutputSink::keep_last 接收器,序列长度 1,以及占位用的文本函数("<scale>"、"<input>"、"<output>")。它的重放规格带有空命令,唯一参数是实现 id。当你需要在事件和重放产物中使用真实文本时,请使用 BenchSpec::advanced。
BenchSpec
BenchSpec 是一份完整但尚未验证的用例描述。
pub struct BenchSpec[Scale, Input, Prepared, Expected, Output, Context, State, SinkValue] {
case : DifferentialCase[Scale, Input, Prepared, Expected, Output, Context, State, SinkValue]
}
BenchSpec::advanced
BenchSpec::advanced 由用例的全部组成部分构建用例。
pub fn[Scale, Input, Prepared, Expected, Output, Context, State, SinkValue] BenchSpec::advanced(String, @fixture.Fixture[Scale, Input, Prepared], Array[Implementation[Prepared, Output, Context]], OutputSink[Output, State, SinkValue], @experiment.OracleSpec[Input, Expected, Output, Context], Array[Scale], (Scale) -> String, Int, (Input, Int) -> @model.CaseDescriptor, (Input) -> String, (Output) -> String, (Context) -> String, (Input, String) -> @model.ReplaySpec, shrinker? : @experiment.Shrinker[Input]) -> Self[Scale, Input, Prepared, Expected, Output, Context, State, SinkValue]
| 参数 | 作用 |
|---|---|
id | 用例 id,复制到每个事件中 |
fixture | 物化、克隆、准备和重置输入 |
implementations | 参赛者(数组会被复制) |
output_sink | 在计时批次内折叠输出 |
oracle | 参考验证和/或关系验证 |
scales | 每个规模一个数据集,按此顺序(会被复制) |
scale_text | 验证事件中规模的文本 |
sequence_length | 每个实现和数据集验证的操作次数 |
describe | 第 步的操作、操作数、上下文和舍入方式,用于证据 |
input_text, output_text, context_text | 证据和失败产物中使用的文本形式 |
replay | 复现输入和实现 id 的命令 |
shrinker | 可选;最小化失败的输入 |
BenchSpec::with_sequence_length
with_sequence_length 返回规格的一个副本,带有另一个序列长度。
pub fn[Scale, Input, Prepared, Expected, Output, Context, State, SinkValue] BenchSpec::with_sequence_length(Self[Scale, Input, Prepared, Expected, Output, Context, State, SinkValue], Int) -> Self[Scale, Input, Prepared, Expected, Output, Context, State, SinkValue]
BenchSpec::compile
compile 验证规格,返回一个计划或全部配置错误。
pub fn[Scale, Input, Prepared, Expected, Output, Context, State, SinkValue] BenchSpec::compile(Self[Scale, Input, Prepared, Expected, Output, Context, State, SinkValue]) -> Result[ValidatedBenchPlan[Scale, Input, Prepared, Expected, Output, Context, State, SinkValue], Array[BenchConfigError]]
它依次检查:非空的用例 id、至少一个规模、至少一个实现、正的序列长度,以及非空且唯一的实现 id。
test "compile collects every configuration error" {
let anonymous = @runner.Implementation::stateless("", "1", (x : Int) => {
@model.OperationResult::completed(x, ())
})
let result = @runner.single_step("", ([] : Array[Int]))
.with_immutable_input(context => context.dataset_key.scale, x => x.to_string())
.compare([anonymous])
.against_equal(x => x, (expected, actual) => expected == actual)
.compile()
guard result is Err(errors) else { fail("expected errors") }
inspect(errors.length(), content="3")
inspect(errors[0] is EmptyCaseId, content="true")
inspect(errors[1] is EmptyScales, content="true")
inspect(errors[2] is EmptyImplementationId(0), content="true")
}
BenchConfigError
BenchConfigError 指明用例描述中的一个问题。
pub(all) enum BenchConfigError {
MissingFixture
MissingOracle
MissingOutputSink
MissingSerializer(String)
EmptyCaseId
EmptyImplementations
EmptyScales
EmptyImplementationId(Int)
DuplicateImplementationId(String)
InvalidSequenceLength(Int)
}
EmptyImplementationId 携带实现的索引,DuplicateImplementationId 携带重复的 id,InvalidSequenceLength 携带被拒绝的长度。Missing* 构造器来自 SingleStepCase::compile;MissingSerializer 是保留项,当前代码不会产生它。
DifferentialCase
DifferentialCase 是规格或计划背后的记录。
pub struct DifferentialCase[Scale, Input, Prepared, Expected, Output, Context, State, SinkValue] {
id : String
fixture : @fixture.Fixture[Scale, Input, Prepared]
implementations : Array[Implementation[Prepared, Output, Context]]
output_sink : OutputSink[Output, State, SinkValue]
oracle : @experiment.OracleSpec[Input, Expected, Output, Context]
scales : Array[Scale]
scale_text : (Scale) -> String
sequence_length : Int
describe : (Input, Int) -> @model.CaseDescriptor
input_text : (Input) -> String
output_text : (Output) -> String
context_text : (Context) -> String
shrinker : @experiment.Shrinker[Input]?
replay : (Input, String) -> @model.ReplaySpec
}
它的字段与 BenchSpec::advanced 的参数一一对应。
ValidatedBenchPlan
ValidatedBenchPlan 是通过了 BenchSpec::compile 的用例。
pub struct ValidatedBenchPlan[Scale, Input, Prepared, Expected, Output, Context, State, SinkValue] {
case : DifferentialCase[Scale, Input, Prepared, Expected, Output, Context, State, SinkValue]
}
它只能由 compile 创建,而 run 只接受它。