summaryrefslogtreecommitdiff
path: root/src/LivingVillage.Headless/BatchOutput.fs
diff options
context:
space:
mode:
Diffstat (limited to 'src/LivingVillage.Headless/BatchOutput.fs')
-rw-r--r--src/LivingVillage.Headless/BatchOutput.fs75
1 files changed, 75 insertions, 0 deletions
diff --git a/src/LivingVillage.Headless/BatchOutput.fs b/src/LivingVillage.Headless/BatchOutput.fs
new file mode 100644
index 0000000..aeaef52
--- /dev/null
+++ b/src/LivingVillage.Headless/BatchOutput.fs
@@ -0,0 +1,75 @@
+namespace LivingVillage.Headless
+
+open System.Threading.Tasks
+
+/// M3 批次的输出格式与流式并行驱动。
+///
+/// 全部是纯函数 / 无模拟依赖的驱动,便于单测;真实运行由 `Program.runBatch`
+/// 传入 `evalWorld` 与 `Console.Out`。格式约定:
+/// - `world=... ` 判据行保持原前缀与字段,行尾追加 `elapsed_s=<秒>`;
+/// - 每完成一个世界追加一行 `[done k/N] elapsed=…s world=…`;
+/// - `batch_summary` 保留旧字段,行尾追加 `workers=` 与 `wall_s=`(总墙钟)。
+module BatchOutput =
+
+ /// `--cost-probe` 默认天数阶梯(单世界逐段实测)。
+ let defaultCostLadder : int64 list = [ 1L; 5L; 10L; 20L; 50L; 100L ]
+
+ /// world 结果行:前缀不变,行尾追加耗时字段。
+ let worldLine (baseLine: string) (elapsedSeconds: float) : string =
+ sprintf "%s elapsed_s=%.1f" baseLine elapsedSeconds
+
+ /// 逐世界完成进度行。
+ let doneLine (completed: int) (total: int) (elapsedSeconds: float) (world: int) : string =
+ sprintf "[done %d/%d] elapsed=%.1fs world=%d" completed total elapsedSeconds world
+
+ /// batch_summary:旧字段原样保留,追加 workers 与 wall_s(总墙钟,append 不删旧字段)。
+ let summaryLine
+ (worlds: int)
+ (passed: int)
+ (failed: int)
+ (oldPassed: int)
+ (ratioMin: float)
+ (ratioMax: float)
+ (giniMin: float)
+ (giniMax: float)
+ (workers: int)
+ (elapsedSeconds: float) : string =
+ sprintf
+ "batch_summary worlds=%d passed=%d failed=%d old_passed=%d/%d ratio_min=%.3f ratio_max=%.3f gini_min=%.3f gini_max=%.3f elapsed_s=%.1f workers=%d wall_s=%.1f"
+ worlds passed failed oldPassed worlds ratioMin ratioMax giniMin giniMax elapsedSeconds workers elapsedSeconds
+
+ /// --cost-probe 单段输出:ticks/s、墙钟、分配与 GC 全部来自真实测量。
+ let costLine
+ (days: int64)
+ (ticks: int64)
+ (elapsedSeconds: float)
+ (allocatedBytes: int64)
+ (gen0: int64)
+ (gen1: int64)
+ (gen2: int64) : string =
+ let ticksPerSecond = if elapsedSeconds > 0.0 then float ticks / elapsedSeconds else 0.0
+ sprintf
+ "cost days=%d ticks=%d elapsed_s=%.3f ticks_per_s=%.1f allocated_bytes=%d gen0=%d gen1=%d gen2=%d"
+ days ticks elapsedSeconds ticksPerSecond allocatedBytes gen0 gen1 gen2
+
+ /// 并行跑 total 个工作,**完成即回调** `emit world result completedCount`。
+ /// 返回按 world 序号存放的结果数组;`work` 必须是仅依赖 k 的确定性纯世界计算。
+ let runStreaming (total: int) (workers: int) (work: int -> 'T) (emit: int -> 'T -> int -> unit) : 'T[] =
+ let results = Array.zeroCreate total
+ let gate = obj ()
+ let mutable completed = 0
+ let options = ParallelOptions(MaxDegreeOfParallelism = max 1 workers)
+ Parallel.For(
+ 0,
+ total,
+ options,
+ fun k ->
+ let result = work k
+ let completedCount =
+ lock gate (fun () ->
+ results.[k] <- result
+ completed <- completed + 1
+ completed)
+ emit k result completedCount)
+ |> ignore
+ results