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, SinkValueOutputSink 的累加器和最终值

运行计划

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) 为 [(0+b) mod k,(1+b) mod k,…,(k−1+b) mod k][(0 + b) \bmod k, (1 + b) \bmod k, \dots, (k-1+b) \bmod k]。在 OrderPolicy::BalancedBlocks(order_seed) 下,运行器使用区组偏移 b+ob + o,其中 o=(order_seed⊕run_seed) mod ko = (\text{order\_seed} \oplus \text{run\_seed}) \bmod k;在 FixedOrder 下使用恒等排列。在任意 kk 个连续区组中,每个实现恰好在每个位置出现一次。

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 >= 0InvalidInteger("warmup_iterations", n)
warmup_time_us 设置时须 >= 0InvalidDuration("warmup_time_us", t)
exploratory_samples >= 0InvalidInteger("exploratory_samples", n)
confirmatory_samples > 0InvalidInteger("confirmatory_samples", n)
target_batch_time_us > 0InvalidDuration("target_batch_time_us", t)
max_sample_time_us > 0InvalidDuration("max_sample_time_us", t)
target_batch_time_us <= max_sample_time_usInvalidDuration("target_batch_time_us", t)
0 < min_batch_iterations <= max_batch_iterationsInvalidIterationRange(min, max)
practical_delta_pct >= 0InvalidPracticalDelta(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,返回其包装的协议)。预设包括:

字段QuickCheckDevelopmentRegressionGate
experiment_designFixedDatasetRepeatedMeasurementsFixedDatasetRepeatedMeasurementsHierarchicalDatasetsAndRepeats
warmup_iterations1310
warmup_time_usNoneSome(5000.0)Some(10000.0)
target_batch_time_us1000500010000
批次迭代次数1 到 10001 到 100005 到 10000
max_sample_time_us1000002500001000000
batch_policyPerImplementationPerImplementationPerImplementation
practical_delta_pct1.01.00.5
order_policyBalancedBlocks(1)BalancedBlocks(1)BalancedBlocks(1)
outlier_policyReportOnlyReportOnlyTukeyFence
validation_coverageConfirmatoryOnlyEveryDatasetEveryMeasurement
exploratory_samples135
confirmatory_samples31020

运行器依据预热、校准、顺序和样本字段行事。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,并按如下方式映射结果:

工作进程结果结果
退出码为 0decode(stdout) 返回的结果和上下文,并附上 stdout、stderr 和退出码
退出码 ≠0\ne 0Aborted(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第 ii 步的操作、操作数、上下文和舍入方式,用于证据
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 只接受它。