From 43dab862cf5881b5fb7eb47c34aa55e02b63663f Mon Sep 17 00:00:00 2001 From: Joel Teply Date: Tue, 6 Oct 2026 08:15:04 -0500 Subject: [PATCH 1/4] server /train: a per-turn slowdown receipt beside the cycle average max_slowdown_ppm bounds the cycle's AVERAGE; Joel's bar is per turn (Cormac on #36). A turn that lands at a window's start decodes at the in-window rate for the whole window, which the average never shows. Serving now hands every finished turn to the trainer (its generation span and decode steps): a turn that overlapped no training window is her per-turn baseline, one that overlapped a window is a slowdown sample against it. A turn's rate is per slot, so it is measured against other turns, never against the lane's total rate. Turns under 16 decode steps time the scheduler, not her, and are skipped. A run's end closes its last window, so later turns never read as overlapped by training that no longer runs. /train status: turns_clean, turns_overlapped, turn_tps_clean, turn_slowdown_ppm_p50/p95/max. Measured on Qwen3.5-0.8B Q8_0, Metal (M5, beside a live 27B lane), 2 slots, T = 10%: - run A: client throughput fell 20% overall (the average bound about holds), but overlapped turns were slowed p50 28% / p95 86% (engine), p95 80% (client per-request rates). - run B: the engine's no-window rate read 39 tok/s against ~120 true, so the bound barely yielded; overlapped turns p50 54% / p95 89% (engine), 81% / 95% (client). Not reproduced in run A; this GPU's baseline moved 45-123 tok/s between runs with the 27B lane beside it. The receipt tracks the client within about 10 points at p95 in both, and shows what the average hides: a 6-12 s window nearly stalls the turns it overlaps. The cap on window duration (the walk's chunk size) is what makes the per-turn bar hold; this is its receipt. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01LoTjvf5j3Ez13g6k8mRkFo --- tools/server/server-context.cpp | 4 +++ tools/server/server-train.cpp | 53 +++++++++++++++++++++++++++++++++ tools/server/server-train.h | 15 ++++++++++ 3 files changed, 72 insertions(+) diff --git a/tools/server/server-context.cpp b/tools/server/server-context.cpp index a35e8d2a2b4..228a5e779c1 100644 --- a/tools/server/server-context.cpp +++ b/tools/server/server-context.cpp @@ -2020,6 +2020,10 @@ struct server_context_impl { } void send_final_response(server_slot & slot) { + // every finished turn reaches the trainer: her per-turn baseline, or a slowdown sample + if (trainer && slot.stats.t_gen_last > 0) { + trainer->on_turn(slot.stats.t_prompt_last, slot.stats.t_gen_last, (int64_t) slot.stats.n_gen_steps()); + } auto res = std::make_unique(); res->id = slot.task->id; diff --git a/tools/server/server-train.cpp b/tools/server/server-train.cpp index 99d8d2cbf1d..8ece3532829 100644 --- a/tools/server/server-train.cpp +++ b/tools/server/server-train.cpp @@ -465,6 +465,8 @@ json server_trainer::start(const json & body_in) { share_ppm.store(body.value("share_ppm", (int64_t) 250000)); max_slowdown_ppm.store(body.value("max_slowdown_ppm", (int64_t) 0)); busy_window_ms = busy_window_tokens = busy_yield_ms = busy_yield_tokens = 0; + window_spans.clear(); + turn_slowdown_ppm.clear(); rate_no_window = rate_estimate{}; rate_in_window = rate_estimate{}; last_no_window_sample = std::chrono::steady_clock::time_point{}; @@ -493,6 +495,29 @@ int64_t server_trainer::yield_for_slowdown(int64_t d_ms, int64_t max_slowdown_pp return y > 0 ? (int64_t) std::ceil(y) : 0; } +void server_trainer::on_turn(int64_t gen_start_us, int64_t gen_end_us, int64_t steps) { + if (steps < TURN_MIN_STEPS || gen_end_us <= gen_start_us) { + return; + } + std::lock_guard lock(mu); + int64_t overlap_us = 0; + for (const auto & [start, end] : window_spans) { + const int64_t e = end == 0 ? gen_end_us : end; // the window still running + overlap_us += std::max(0, std::min(e, gen_end_us) - std::max(start, gen_start_us)); + } + const double ms = (double) (gen_end_us - gen_start_us) / 1000.0; + if (overlap_us == 0) { + turn_rate_clean.add(ms, (double) steps); + turns_clean += 1; + return; + } + const double clean = turn_rate_clean.per_ms(); + if (clean > 0) { // no clean turn yet: nothing to measure her against, never a guessed sample + const double slowdown = std::max(0.0, 1.0 - ((double) steps / ms) / clean); + turn_slowdown_ppm.push_back((int64_t) std::llround(slowdown * 1e6)); + } +} + bool server_trainer::serving_busy() const { if (!yield_to_turns.load() || !serving.busy_slots) { return false; @@ -520,6 +545,9 @@ bool server_trainer::before_window(bool, void * user_data) { if (last_ms > 0) { // 0 = no window yet (the first call), never a sample self.window_ms_samples.push_back(last_ms); self.window_ms_max = std::max(self.window_ms_max, last_ms); + if (!self.window_spans.empty() && self.window_spans.back().second == 0) { + self.window_spans.back().second = ggml_time_us(); + } if (self.window_busy_start && working_t0) { // her rate WITH a window running self.busy_window_ms += last_ms; self.busy_window_tokens += tokens_t0 - self.window_tokens_start; @@ -597,6 +625,7 @@ bool server_trainer::before_window(bool, void * user_data) { self.windows_while_busy += 1; } self.window_started = std::chrono::steady_clock::now(); + self.window_spans.emplace_back(ggml_time_us(), 0); self.window_tokens_start = tokens_t1; self.window_busy_start = self.serving.busy_slots && self.serving.busy_slots() > 0; } @@ -679,6 +708,20 @@ json server_trainer::status() const { if (rate_in_window.ms > 0) { s["decode_tps_in_window_recent"] = rate_in_window.per_ms() * 1000.0; } + // the per-turn receipt: each turn that overlapped a window, against her clean turns + s["turns_clean"] = turns_clean; + s["turns_overlapped"] = (int64_t) turn_slowdown_ppm.size(); + if (turn_rate_clean.ms > 0) { + s["turn_tps_clean"] = turn_rate_clean.per_ms() * 1000.0; + } + if (!turn_slowdown_ppm.empty()) { + std::vector sorted = turn_slowdown_ppm; + std::sort(sorted.begin(), sorted.end()); + const auto pct = [&](double p) { return sorted[std::min(sorted.size() - 1, (size_t) (p * (double) sorted.size()))]; }; + s["turn_slowdown_ppm_p50"] = pct(0.50); + s["turn_slowdown_ppm_p95"] = pct(0.95); + s["turn_slowdown_ppm_max"] = sorted.back(); + } s["busy_window_ms"] = busy_window_ms; s["busy_yield_ms"] = busy_yield_ms; s["windows"] = windows.load(); @@ -707,6 +750,7 @@ void server_trainer::run(json req, examples_data ex) { pause_requested = false; paused = false; waiting_for_serving = false; + close_window_span(); running.store(false); }; @@ -952,5 +996,14 @@ void server_trainer::run(json req, examples_data ex) { state["state"] = "done"; state["adapter"] = out; } + close_window_span(); running.store(false); } + +void server_trainer::close_window_span() { + // a run that ends (done, failed, cancelled) ends its last window: an open span would read + // every later turn as overlapped by training that no longer runs + if (!window_spans.empty() && window_spans.back().second == 0) { + window_spans.back().second = ggml_time_us(); + } +} diff --git a/tools/server/server-train.h b/tools/server/server-train.h index d7e380b9a24..c0afe2e7666 100644 --- a/tools/server/server-train.h +++ b/tools/server/server-train.h @@ -102,6 +102,12 @@ class server_trainer { // Starts a run on a worker thread; refuses (ok=false) while one is running or on bad input. common_json start(const common_json & body_in); common_json status() const; + + // A finished turn, from serving: its generation began at gen_start_us and ended at + // gen_end_us (ggml_time_us) over `steps` decode steps. A turn that overlapped no training + // window is her per-turn baseline; one that overlapped a window is a per-turn slowdown + // sample against it (the bound is on the cycle's average, Joel's bar is per turn: Cormac). + void on_turn(int64_t gen_start_us, int64_t gen_end_us, int64_t steps); // Stops the running job at its next training window (no adapter is written); ok=false when // nothing is running. common_json cancel(); @@ -111,6 +117,7 @@ class server_trainer { private: void run(common_json req, examples_data ex); static bool before_window(bool train, void * user_data); + void close_window_span(); // under mu: a run's end ends its last window static void on_batch(bool train, ggml_opt_context_t, ggml_opt_dataset_t, ggml_opt_result_t, int64_t ibatch, int64_t ibatch_max, int64_t); @@ -157,6 +164,14 @@ class server_trainer { // price of the bound staying true. static constexpr int64_t PROBE_EVERY_MS = 30000; static constexpr int64_t PROBE_MS = 500; + // THE PER-TURN RECEIPT. Windows of this run as [start, end) in ggml_time_us (end 0 while + // one runs), every turn's overlap measured against them; a turn's rate is per slot, so its + // baseline is other turns' rate, never the lane's total rate. Guarded by mu. + static constexpr int64_t TURN_MIN_STEPS = 16; // fewer decode steps time the scheduler, not her + std::vector> window_spans; + rate_estimate turn_rate_clean; // turns that overlapped no window, across runs + std::vector turn_slowdown_ppm; // one per turn that overlapped a window, this run + int64_t turns_clean{0}; rate_estimate rate_no_window; // guarded by mu rate_estimate rate_in_window; // guarded by mu std::chrono::steady_clock::time_point last_no_window_sample{}; From 05a06928f81c13b9d41f81ff99fa83a8b5780439 Mon Sep 17 00:00:00 2001 From: Joel Teply Date: Tue, 6 Oct 2026 08:24:35 -0500 Subject: [PATCH 2/4] server /train: the per-turn receipt is bounded and says what a sample is (Cormac on #39) - window spans that ended before (now - the longest turn seen) can overlap no turn still to finish and are dropped, so a turn's overlap scan stays the size of the windows it could meet - slowdown samples live in a fixed histogram of 1% buckets (p50/p95 to the bucket, max exact) instead of a vector that grew with the run - a sample is the whole turn's slowdown: a turn a window partly overlapped reads diluted, and the status comment says so, so a p95 is never read as the slowdown inside a window - the call site states that t_prompt_last and t_gen_last are ggml_time_us timestamps (server-common.h), the decode span, not durations Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01LoTjvf5j3Ez13g6k8mRkFo --- tools/server/server-context.cpp | 2 ++ tools/server/server-train.cpp | 43 +++++++++++++++++++++++++-------- tools/server/server-train.h | 10 ++++++-- 3 files changed, 43 insertions(+), 12 deletions(-) diff --git a/tools/server/server-context.cpp b/tools/server/server-context.cpp index 228a5e779c1..04f384c4e1b 100644 --- a/tools/server/server-context.cpp +++ b/tools/server/server-context.cpp @@ -2021,6 +2021,8 @@ struct server_context_impl { void send_final_response(server_slot & slot) { // every finished turn reaches the trainer: her per-turn baseline, or a slowdown sample + // t_prompt_last / t_gen_last are ggml_time_us TIMESTAMPS (server-common.h): generation + // began at the prompt's last batch and ended at the last token if (trainer && slot.stats.t_gen_last > 0) { trainer->on_turn(slot.stats.t_prompt_last, slot.stats.t_gen_last, (int64_t) slot.stats.n_gen_steps()); } diff --git a/tools/server/server-train.cpp b/tools/server/server-train.cpp index 8ece3532829..f3869c7773e 100644 --- a/tools/server/server-train.cpp +++ b/tools/server/server-train.cpp @@ -466,7 +466,9 @@ json server_trainer::start(const json & body_in) { max_slowdown_ppm.store(body.value("max_slowdown_ppm", (int64_t) 0)); busy_window_ms = busy_window_tokens = busy_yield_ms = busy_yield_tokens = 0; window_spans.clear(); - turn_slowdown_ppm.clear(); + turn_slowdown_pct.fill(0); + turns_overlapped = 0; + turn_slowdown_ppm_max = 0; rate_no_window = rate_estimate{}; rate_in_window = rate_estimate{}; last_no_window_sample = std::chrono::steady_clock::time_point{}; @@ -500,6 +502,11 @@ void server_trainer::on_turn(int64_t gen_start_us, int64_t gen_end_us, int64_t s return; } std::lock_guard lock(mu); + longest_turn_us = std::max(longest_turn_us, gen_end_us - gen_start_us); + const int64_t horizon = gen_end_us - longest_turn_us; // no unfinished turn began before this + window_spans.erase(std::remove_if(window_spans.begin(), window_spans.end(), + [horizon](const std::pair & w) { return w.second != 0 && w.second < horizon; }), + window_spans.end()); int64_t overlap_us = 0; for (const auto & [start, end] : window_spans) { const int64_t e = end == 0 ? gen_end_us : end; // the window still running @@ -513,8 +520,13 @@ void server_trainer::on_turn(int64_t gen_start_us, int64_t gen_end_us, int64_t s } const double clean = turn_rate_clean.per_ms(); if (clean > 0) { // no clean turn yet: nothing to measure her against, never a guessed sample - const double slowdown = std::max(0.0, 1.0 - ((double) steps / ms) / clean); - turn_slowdown_ppm.push_back((int64_t) std::llround(slowdown * 1e6)); + // the WHOLE turn's rate: a turn a window overlapped for 20% of its span reads diluted, + // which is what she felt over that turn, not the in-window slowdown + const double slowdown = std::clamp(1.0 - ((double) steps / ms) / clean, 0.0, 1.0); + const int64_t ppm = (int64_t) std::llround(slowdown * 1e6); + turn_slowdown_pct[(size_t) std::llround(slowdown * 100)] += 1; + turns_overlapped += 1; + turn_slowdown_ppm_max = std::max(turn_slowdown_ppm_max, ppm); } } @@ -708,19 +720,30 @@ json server_trainer::status() const { if (rate_in_window.ms > 0) { s["decode_tps_in_window_recent"] = rate_in_window.per_ms() * 1000.0; } - // the per-turn receipt: each turn that overlapped a window, against her clean turns + // the per-turn receipt: each turn that overlapped a window, against her clean turns. A + // sample is the whole turn's slowdown, so a turn a window only partly overlapped reads + // diluted: a p95 here is what her turns felt, never the slowdown inside a window. p50/p95 + // to the 1% bucket (a fixed histogram), max exact. s["turns_clean"] = turns_clean; - s["turns_overlapped"] = (int64_t) turn_slowdown_ppm.size(); + s["turns_overlapped"] = turns_overlapped; if (turn_rate_clean.ms > 0) { s["turn_tps_clean"] = turn_rate_clean.per_ms() * 1000.0; } - if (!turn_slowdown_ppm.empty()) { - std::vector sorted = turn_slowdown_ppm; - std::sort(sorted.begin(), sorted.end()); - const auto pct = [&](double p) { return sorted[std::min(sorted.size() - 1, (size_t) (p * (double) sorted.size()))]; }; + if (turns_overlapped > 0) { + const auto pct = [&](double p) { + const int64_t rank = std::min(turns_overlapped - 1, (int64_t) (p * (double) turns_overlapped)); + int64_t seen = 0; + for (size_t b = 0; b < turn_slowdown_pct.size(); ++b) { + seen += turn_slowdown_pct[b]; + if (seen > rank) { + return (int64_t) b * 10000; + } + } + return (int64_t) 1000000; + }; s["turn_slowdown_ppm_p50"] = pct(0.50); s["turn_slowdown_ppm_p95"] = pct(0.95); - s["turn_slowdown_ppm_max"] = sorted.back(); + s["turn_slowdown_ppm_max"] = turn_slowdown_ppm_max; } s["busy_window_ms"] = busy_window_ms; s["busy_yield_ms"] = busy_yield_ms; diff --git a/tools/server/server-train.h b/tools/server/server-train.h index c0afe2e7666..8e1fc877c1e 100644 --- a/tools/server/server-train.h +++ b/tools/server/server-train.h @@ -45,6 +45,7 @@ #include "json.h" #include "ggml-opt.h" +#include #include #include #include @@ -168,9 +169,14 @@ class server_trainer { // one runs), every turn's overlap measured against them; a turn's rate is per slot, so its // baseline is other turns' rate, never the lane's total rate. Guarded by mu. static constexpr int64_t TURN_MIN_STEPS = 16; // fewer decode steps time the scheduler, not her + // Bounded: spans that ended before (now - the longest turn seen) can overlap no turn still + // to finish, and are dropped; samples live in a fixed histogram of 1% buckets. std::vector> window_spans; - rate_estimate turn_rate_clean; // turns that overlapped no window, across runs - std::vector turn_slowdown_ppm; // one per turn that overlapped a window, this run + int64_t longest_turn_us{0}; + rate_estimate turn_rate_clean; // turns that overlapped no window, across runs + std::array turn_slowdown_pct{}; // turns that overlapped a window, this run, by % slower + int64_t turns_overlapped{0}; + int64_t turn_slowdown_ppm_max{0}; int64_t turns_clean{0}; rate_estimate rate_no_window; // guarded by mu rate_estimate rate_in_window; // guarded by mu From 65b3ada329bc527aba75907f32258a7354f42a99 Mon Sep 17 00:00:00 2001 From: Joel Teply Date: Tue, 6 Oct 2026 08:40:45 -0500 Subject: [PATCH 3/4] server /train: prune window spans against the oldest turn still in flight (Cormac on #39) Pruning against the longest FINISHED turn dropped spans that a longer-than-ever turn still in flight had overlapped; that turn then landed in the clean baseline and drifted it slow. The server loop knows which slots are working: send_final_response passes the earliest start of any other turn in flight, and the trainer drops only windows that ended before it (on every finished turn, measured or not). A bucketed percentile is clamped to the exact max, which 1% rounding could exceed. Smoke on the 0.8B (Metal): overlapped turns p95 35% (engine) vs 38% (client per-request). Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01LoTjvf5j3Ez13g6k8mRkFo --- tools/server/server-context.cpp | 9 ++++++++- tools/server/server-train.cpp | 17 +++++++++-------- tools/server/server-train.h | 9 +++++---- 3 files changed, 22 insertions(+), 13 deletions(-) diff --git a/tools/server/server-context.cpp b/tools/server/server-context.cpp index 04f384c4e1b..cb99e47171c 100644 --- a/tools/server/server-context.cpp +++ b/tools/server/server-context.cpp @@ -2024,7 +2024,14 @@ struct server_context_impl { // t_prompt_last / t_gen_last are ggml_time_us TIMESTAMPS (server-common.h): generation // began at the prompt's last batch and ended at the last token if (trainer && slot.stats.t_gen_last > 0) { - trainer->on_turn(slot.stats.t_prompt_last, slot.stats.t_gen_last, (int64_t) slot.stats.n_gen_steps()); + // the earliest start of a turn still in flight: the trainer drops windows before it + int64_t oldest_open_us = slot.stats.t_gen_last; + for (const auto & other : slots) { + if (&other != &slot && other.is_processing() && other.stats.t_start > 0) { + oldest_open_us = std::min(oldest_open_us, other.stats.t_start); + } + } + trainer->on_turn(slot.stats.t_prompt_last, slot.stats.t_gen_last, (int64_t) slot.stats.n_gen_steps(), oldest_open_us); } auto res = std::make_unique(); diff --git a/tools/server/server-train.cpp b/tools/server/server-train.cpp index f3869c7773e..007890bb921 100644 --- a/tools/server/server-train.cpp +++ b/tools/server/server-train.cpp @@ -497,16 +497,16 @@ int64_t server_trainer::yield_for_slowdown(int64_t d_ms, int64_t max_slowdown_pp return y > 0 ? (int64_t) std::ceil(y) : 0; } -void server_trainer::on_turn(int64_t gen_start_us, int64_t gen_end_us, int64_t steps) { - if (steps < TURN_MIN_STEPS || gen_end_us <= gen_start_us) { - return; - } +void server_trainer::on_turn(int64_t gen_start_us, int64_t gen_end_us, int64_t steps, int64_t oldest_open_us) { std::lock_guard lock(mu); - longest_turn_us = std::max(longest_turn_us, gen_end_us - gen_start_us); - const int64_t horizon = gen_end_us - longest_turn_us; // no unfinished turn began before this + // pruned on every finished turn, measured or not: the turns still in flight are the only + // ones a span can still be asked about, and none of them began before oldest_open_us window_spans.erase(std::remove_if(window_spans.begin(), window_spans.end(), - [horizon](const std::pair & w) { return w.second != 0 && w.second < horizon; }), + [oldest_open_us](const std::pair & w) { return w.second != 0 && w.second < oldest_open_us; }), window_spans.end()); + if (steps < TURN_MIN_STEPS || gen_end_us <= gen_start_us) { + return; + } int64_t overlap_us = 0; for (const auto & [start, end] : window_spans) { const int64_t e = end == 0 ? gen_end_us : end; // the window still running @@ -736,7 +736,8 @@ json server_trainer::status() const { for (size_t b = 0; b < turn_slowdown_pct.size(); ++b) { seen += turn_slowdown_pct[b]; if (seen > rank) { - return (int64_t) b * 10000; + // a bucket rounds to the nearest 1%, which can sit above the exact max + return std::min((int64_t) b * 10000, turn_slowdown_ppm_max); } } return (int64_t) 1000000; diff --git a/tools/server/server-train.h b/tools/server/server-train.h index 8e1fc877c1e..f05062c6a42 100644 --- a/tools/server/server-train.h +++ b/tools/server/server-train.h @@ -108,7 +108,9 @@ class server_trainer { // gen_end_us (ggml_time_us) over `steps` decode steps. A turn that overlapped no training // window is her per-turn baseline; one that overlapped a window is a per-turn slowdown // sample against it (the bound is on the cycle's average, Joel's bar is per turn: Cormac). - void on_turn(int64_t gen_start_us, int64_t gen_end_us, int64_t steps); + // oldest_open_us: the earliest start of any turn still in flight (gen_end_us when none); no + // window that ended before it can overlap a turn still to finish. + void on_turn(int64_t gen_start_us, int64_t gen_end_us, int64_t steps, int64_t oldest_open_us); // Stops the running job at its next training window (no adapter is written); ok=false when // nothing is running. common_json cancel(); @@ -169,10 +171,9 @@ class server_trainer { // one runs), every turn's overlap measured against them; a turn's rate is per slot, so its // baseline is other turns' rate, never the lane's total rate. Guarded by mu. static constexpr int64_t TURN_MIN_STEPS = 16; // fewer decode steps time the scheduler, not her - // Bounded: spans that ended before (now - the longest turn seen) can overlap no turn still - // to finish, and are dropped; samples live in a fixed histogram of 1% buckets. + // Bounded: spans that ended before the oldest turn still in flight began can overlap no + // turn still to finish, and are dropped; samples live in a fixed histogram of 1% buckets. std::vector> window_spans; - int64_t longest_turn_us{0}; rate_estimate turn_rate_clean; // turns that overlapped no window, across runs std::array turn_slowdown_pct{}; // turns that overlapped a window, this run, by % slower int64_t turns_overlapped{0}; From bbdc2b33f3d2887123ebadbd70194ce98756538c Mon Sep 17 00:00:00 2001 From: Joel Teply Date: Tue, 6 Oct 2026 12:09:05 -0500 Subject: [PATCH 4/4] server /train: a turn's overlap is measured BEFORE the windows are pruned (Cormac on #39) on_turn pruned the window spans against oldest_open_us before summing this turn's overlap, and with no other turn in flight the server passes this turn's own END as oldest_open_us, so every window that ended inside the turn was erased before it was counted: a turn a finished window slowed filed as clean, and the baseline drifted slow. The decision is now one pure function, server_train_turn_overlap_then_prune (new server-train-spans.h): it sums the overlap first, then drops windows that ended before the oldest turn still in flight (after this turn is measured, a window that ended before every open turn began can overlap no turn still to finish, and new turns start after now). It runs on every finished turn, a short one too. tests/test-server-train-spans.cpp pins it, starting with Cormac's case (span [100, 200], turn [50, 300], nothing else open: overlap 100, then the span goes), plus a running window, a window kept for another open turn, and an empty turn that still prunes. Known-positive: with the prune moved back in front of the sum, the first case FAILS. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01LoTjvf5j3Ez13g6k8mRkFo --- tests/CMakeLists.txt | 4 +++ tests/test-server-train-spans.cpp | 41 +++++++++++++++++++++++++++++++ tools/server/server-train-spans.h | 32 ++++++++++++++++++++++++ tools/server/server-train.cpp | 14 +++-------- 4 files changed, 81 insertions(+), 10 deletions(-) create mode 100644 tests/test-server-train-spans.cpp create mode 100644 tools/server/server-train-spans.h diff --git a/tests/CMakeLists.txt b/tests/CMakeLists.txt index c3305283012..5fc7d54312d 100644 --- a/tests/CMakeLists.txt +++ b/tests/CMakeLists.txt @@ -150,6 +150,10 @@ endif () llama_build(test-recurrent-state-rollback.cpp) +# the server trainer's per-turn receipt, pure (tools/server/server-train-spans.h) +llama_build_and_test(test-server-train-spans.cpp) +target_include_directories(test-server-train-spans PRIVATE ${PROJECT_SOURCE_DIR}/tools/server) + if (NOT WIN32 OR NOT BUILD_SHARED_LIBS) # these tests are disabled on Windows because they use internal functions not exported with LLAMA_API (when building with shared libraries) llama_build_and_test(test-unicode.cpp) diff --git a/tests/test-server-train-spans.cpp b/tests/test-server-train-spans.cpp new file mode 100644 index 00000000000..a3e2cdbd312 --- /dev/null +++ b/tests/test-server-train-spans.cpp @@ -0,0 +1,41 @@ +// what this catches (Cormac on #39): the per-turn receipt pruned the training windows BEFORE +// summing a turn's overlap, against a horizon at the turn's own end, so with no other turn in +// flight every window that ended inside the turn was erased first and the slowed turn filed +// as clean. Measured, then pruned. + +#include "server-train-spans.h" + +#include +#include + +#define CHECK(cond) do { if (!(cond)) { fprintf(stderr, "FAILED %s:%d: %s\n", __FILE__, __LINE__, #cond); exit(1); } } while (0) + +int main() { + // Cormac's case: a span [100, 200], one turn [50, 300], no other turn open (the server + // passes this turn's end as the oldest open start). Overlap 100 us, then the span goes. + { + std::vector> spans = {{100, 200}}; + CHECK(server_train_turn_overlap_then_prune(spans, 50, 300, 300) == 100); + CHECK(spans.empty()); + } + // a window still running counts to the turn's end and is never pruned + { + std::vector> spans = {{250, 0}}; + CHECK(server_train_turn_overlap_then_prune(spans, 50, 300, 300) == 50); + CHECK(spans.size() == 1); + } + // another turn in flight since 150 keeps a window that ended at 200: it overlapped her + { + std::vector> spans = {{100, 200}}; + CHECK(server_train_turn_overlap_then_prune(spans, 220, 300, 150) == 0); + CHECK(spans.size() == 1); + } + // an empty or inverted turn measures nothing but still prunes + { + std::vector> spans = {{100, 200}}; + CHECK(server_train_turn_overlap_then_prune(spans, 300, 300, 300) == 0); + CHECK(spans.empty()); + } + printf("test-server-train-spans: OK\n"); + return 0; +} diff --git a/tools/server/server-train-spans.h b/tools/server/server-train-spans.h new file mode 100644 index 00000000000..989a9653a95 --- /dev/null +++ b/tools/server/server-train-spans.h @@ -0,0 +1,32 @@ +#pragma once + +// The per-turn receipt's one decision, pure (see server_trainer::on_turn): how much of a +// finished turn a training window overlapped, and which windows can be forgotten after it. + +#include +#include +#include +#include + +// spans: this run's training windows as [start, end) in ggml_time_us, end 0 while one runs. +// Returns the overlap of [gen_start_us, gen_end_us) with every window, a running one counted +// to the turn's end. THEN drops the windows that ended before oldest_open_us, the earliest +// start of a turn still in flight (none of them can overlap a window that ended before it, +// and turns not yet begun start after now). Measured BEFORE pruned: pruned first, a window +// that ended inside this very turn was gone before it was counted, so a slowed turn filed as +// clean (Cormac on #39). +inline int64_t server_train_turn_overlap_then_prune(std::vector> & spans, + int64_t gen_start_us, int64_t gen_end_us, + int64_t oldest_open_us) { + int64_t overlap_us = 0; + if (gen_end_us > gen_start_us) { + for (const auto & [start, end] : spans) { + const int64_t e = end == 0 ? gen_end_us : end; // the window still running + overlap_us += std::max(0, std::min(e, gen_end_us) - std::max(start, gen_start_us)); + } + } + spans.erase(std::remove_if(spans.begin(), spans.end(), + [oldest_open_us](const std::pair & w) { return w.second != 0 && w.second < oldest_open_us; }), + spans.end()); + return overlap_us; +} diff --git a/tools/server/server-train.cpp b/tools/server/server-train.cpp index 007890bb921..2b531cfe30a 100644 --- a/tools/server/server-train.cpp +++ b/tools/server/server-train.cpp @@ -1,4 +1,5 @@ #include "server-train.h" +#include "server-train-spans.h" #include "common.h" #include "log.h" @@ -499,19 +500,12 @@ int64_t server_trainer::yield_for_slowdown(int64_t d_ms, int64_t max_slowdown_pp void server_trainer::on_turn(int64_t gen_start_us, int64_t gen_end_us, int64_t steps, int64_t oldest_open_us) { std::lock_guard lock(mu); - // pruned on every finished turn, measured or not: the turns still in flight are the only - // ones a span can still be asked about, and none of them began before oldest_open_us - window_spans.erase(std::remove_if(window_spans.begin(), window_spans.end(), - [oldest_open_us](const std::pair & w) { return w.second != 0 && w.second < oldest_open_us; }), - window_spans.end()); + // measured, then pruned, on every finished turn (a short one too: it still ends windows' + // relevance); see server-train-spans.h + const int64_t overlap_us = server_train_turn_overlap_then_prune(window_spans, gen_start_us, gen_end_us, oldest_open_us); if (steps < TURN_MIN_STEPS || gen_end_us <= gen_start_us) { return; } - int64_t overlap_us = 0; - for (const auto & [start, end] : window_spans) { - const int64_t e = end == 0 ? gen_end_us : end; // the window still running - overlap_us += std::max(0, std::min(e, gen_end_us) - std::max(start, gen_start_us)); - } const double ms = (double) (gen_end_us - gen_start_us) / 1000.0; if (overlap_us == 0) { turn_rate_clean.add(ms, (double) steps);