event API
Luna-Flow/mare_mark/event 接收运行器发出的内容。ObservationSink 是由五个回调组成的记录;该包提供内存接收器、带缓冲的 JSONL 接收器、流式 JSONL 接收器以及一个扇出组合子。JSONL 格式是报告和重放所读取的审计记录。参见 event 设计。
import {
"Luna-Flow/mare_mark/model",
"Luna-Flow/mare_mark/event",
}
接收器
ObservationSink
ObservationSink 是运行器与存储之间的接口。
pub struct ObservationSink {
emit_observation : (@model.Observation) -> Unit
emit_validation : (@model.Validation) -> Unit
emit_failure : (@model.ValidationFailure) -> Unit
emit_calibration : (@model.CalibrationEvent) -> Unit
finish : (@model.RunSummary) -> String
}
pub fn ObservationSink::new((@model.Observation) -> Unit, (@model.Validation) -> Unit, (@model.CalibrationEvent) -> Unit, (@model.RunSummary) -> String, emit_failure? : (@model.ValidationFailure) -> Unit) -> Self
finish 在结束时接收一次运行摘要,并返回一个位置字符串,运行器将其存入 RunSummary.artifact_location。注意 new 的参数顺序:观测、验证、校准、结束,以及可选的失败回调(默认忽略失败)。
test "a counting sink" {
let count = Ref(0)
let sink = @event.ObservationSink::new(
_ => count.val += 1,
_ => (),
_ => (),
summary => "counted://" + summary.run_id,
)
(sink.emit_observation)(
@model.Observation::new("c", "a", "1", 0, 0, 0, Confirmatory, 1.5, 10, Kept, ExcludedFromMeasurement, true),
)
inspect(count.val, content="1")
inspect((sink.finish)(@model.RunSummary::new("r", 1, 0, 0, true, None)), content="counted://r")
}
InMemorySink
InMemorySink 把每个事件保存在数组中。
pub struct InMemorySink {
observations : Array[@model.Observation]
validations : Array[@model.Validation]
failures : Array[@model.ValidationFailure]
calibrations : Array[@model.CalibrationEvent]
}
pub fn InMemorySink::new() -> Self
pub fn InMemorySink::as_sink(Self) -> ObservationSink
as_sink 向数组追加;它的 finish 返回 "memory://run/" + run_id,且不保存摘要。
JsonlSink
JsonlSink 把每个事件保存为一行 JSON。
pub struct JsonlSink {
lines : Array[String]
}
pub fn JsonlSink::new() -> Self
pub fn JsonlSink::as_sink(Self) -> ObservationSink
pub fn JsonlSink::to_jsonl(Self) -> String
as_sink 为每个事件(包括摘要)追加一行;它的 finish 返回 "jsonl://memory/" + run_id。to_jsonl 用 "\n" 连接各行,末尾不带换行符。
test "JSONL lines" {
let jsonl = @event.JsonlSink::new()
let sink = jsonl.as_sink()
(sink.emit_calibration)(@model.CalibrationEvent::new("a", 0, 64, 1012.5, 1000.0, 1))
inspect(
jsonl.to_jsonl(),
content="{\"artifact_version\":\"mmka_1\",\"type\":\"calibration\",\"implementation\":\"a\",\"dataset_id\":0,\"batch_iterations\":64,\"elapsed_us\":1012.5,\"target_elapsed_us\":1000,\"retries\":1}",
)
}
streaming_jsonl
streaming_jsonl 返回一个立即写出每个事件行的接收器。
pub fn streaming_jsonl((String) -> Unit, String) -> ObservationSink
第一个参数接收不带换行符的每一行;第二个参数是 finish 返回的位置。可用它在运行过程中向文件或管道追加内容,这样即使崩溃,也会留下流的一个有效前缀。
test "stream lines as they happen" {
let written : Array[String] = []
let sink = @event.streaming_jsonl(line => written.push(line), "file://events.jsonl")
let location = (sink.finish)(@model.RunSummary::new("r", 0, 0, 0, true, None))
inspect(location, content="file://events.jsonl")
inspect(written[0].contains("\"type\":\"summary\""), content="true")
}
tee
tee 把每个事件发送到两个接收器。
pub fn tee(ObservationSink, ObservationSink) -> ObservationSink
事件先发往左侧接收器。finish 调用两者,并返回右侧接收器的位置。
JSONL 记录
每一行都是一个带有 "artifact_version": "mmka_1" 和 "type" 的对象:
type | 字段 |
|---|---|
observation | case, implementation, implementation_version, dataset_id, repetition_id, block_id, phase, elapsed_us, iterations, batch_sink, valid |
validation | status、implementation、oracle、scale;带证据时还有 case、dataset_id、step_id、operation、operands、context、rounding、expected、actual、expected_kind、actual_kind、expected_flags、actual_flags、trap、stderr、exit_code(存在时)、fingerprint、implementation_version、replay_command、replay_arguments、replay_timeout_ms |
validation_failure | 验证的全部字段,外加 seed(十进制字符串)、original_fingerprint、minimal_fingerprint、shrink_path、minimal_input |
calibration | implementation, dataset_id, batch_iterations, elapsed_us, target_elapsed_us, retries |
summary | run_id、observation_count、validation_count、calibration_count、complete、passed_count、failed_count、unsupported_count、expected_difference_count,以及已知时带有 semantic、performance 和 provenance 对象的 environment |
status 取 valid、invalid、skipped、expected_difference、unsupported、infrastructure_failure 之一;状态的原因字符串不会写出。phase 为 exploratory 或 confirmatory;batch_sink 为 kept 或 discarded:<reason>。观测的 setup_timing 不会写出。