Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
46 commits
Select commit Hold shift + click to select a range
cd5bd14
Wire vllm bench serve into serving-load carriers that refuse a rate w…
Sep 18, 2026
7c29847
Import Int for the Poisson millirate wire.
Sep 18, 2026
74a4542
Carry workload admission as unread instead of dividing the printed KV…
Sep 18, 2026
41625b9
Refuse a rate when client tokens are unread or disagree, and when occ…
Sep 18, 2026
40d2b26
Declare nccl-tests as shared-window traffic and keep the 49 ms/token …
Sep 18, 2026
a54eb01
Carry served alias and tokenizer as two protocol facts, and name the …
Sep 18, 2026
7fd5445
Default this arm to shared-with-operator traffic and prove wiring on …
Sep 18, 2026
ada7d7b
Derive the probe rate from the client ledger under declared contention.
Sep 18, 2026
c574acf
Carry the contended 256/128 sweep as shared-window client latency, no…
Sep 18, 2026
3411e20
Fold the shared-window 256/128 smoke as client output rate, ITL tail,…
Sep 18, 2026
4b45d0e
Fix heads-only parse of the serving-load witness and name the live im…
Sep 18, 2026
1b1d3a8
Stop prometheus label braces from opening string interpolation in the…
Sep 18, 2026
358922e
Fix floor type errors on the serving-load probe.
Sep 18, 2026
3f6cbed
Cross Int to Nat via checked_int_magnitude, not a claim-time cast.
Sep 18, 2026
c86463d
Stop treating client bucket peak and ITL tail as engine contention.
Sep 18, 2026
9035df0
Merge origin/main so ObservedRankIncarnation from #11560 is available.
Sep 18, 2026
c0b9bb3
Scope generation_tokens_total to one ObservedRankIncarnation.
Sep 18, 2026
9a9d100
Qualify the client rate from subject, window, and foreign-work standing.
Sep 18, 2026
0e72440
Name the ObservedRankIncarnation on a scoped generation delta.
Sep 18, 2026
f18bbb3
Include backend, model, tokenizer, and arrival in load-protocol ident…
Sep 18, 2026
a255729
Close the SOURCE HOLD on mixing, arrival policy, window standing, occ…
Sep 18, 2026
6f6ced8
Merge remote-tracking branch 'origin/main' into session/deep-otter-201
Sep 18, 2026
e26bf8f
Bind the live serving-load subject to GunbcGlm53DerivedGb10.
Sep 18, 2026
42ce5a4
Gate bench-serve client facts on the remote execution.
Sep 18, 2026
15273a9
Land serving-load as telemetry capture, not qualification.
Sep 18, 2026
bfb364f
serving-load: typed metrics scrape, Millisecond latencies, counts are…
Sep 19, 2026
421df04
HistogramDeltaRead sum stays Nat: the fold also carries booked-token …
Sep 19, 2026
0cbec26
serving-load: histogram sums keep their metric's dimension; rates are…
Sep 19, 2026
db8c13c
serving-load: subtract histogram sums exactly before projecting a unit
Sep 19, 2026
63c8333
serving-load: refuse negative cumulative sums; exact_decimal_sub is c…
Sep 19, 2026
0515d2b
serving-load: delete constant-valued predicates; print the parsed ITL…
Sep 19, 2026
9892356
serving-load: the foreign-actor roster is an annotation, not a String…
Sep 19, 2026
ac17446
Regenerate the std_measure bootstrap mirror for MilliTokensPerSecond …
Sep 19, 2026
29c0529
serving-load: delete the constant-valued qualification and capacity f…
Sep 19, 2026
c082914
serving-load: remove every non-discriminating witness and the declara…
Sep 19, 2026
0f6205f
serving-load: replay the retained stdout against its execution at the…
Sep 19, 2026
d775e07
serving-load: milli projections are checked and exact; ArrivalSkew re…
Sep 19, 2026
1353060
serving-load: delete the dead arrival comparison gate; name the produ…
Sep 19, 2026
7e9854b
serving-load: provenance is supplied by the caller that ran the opera…
Sep 19, 2026
749bad6
serving-load: token counts are TokenCount and the derived rate is Tok…
Sep 19, 2026
765b959
serving-load: bench URL from the profile's port; name the thousandths…
Sep 19, 2026
98e4421
serving-load: the metrics route comes from extdeps.vllm.server; the f…
Sep 19, 2026
4654dd6
prometheus: carry the full classic # TYPE vocabulary and state which …
Sep 19, 2026
a546ec6
serving-load: the live subject reads max_num_seqs from its profile an…
Sep 19, 2026
d5e2bcc
serving-load: drop the unenforced sole_constructor on a coproduct; on…
Sep 19, 2026
1073814
Merge remote-tracking branch 'origin/main' into session/deep-otter-201
Sep 19, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
21 changes: 21 additions & 0 deletions dag/extdeps/prometheus/client.dag
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,27 @@ data extdeps_model_scope: ExternalModelScope = ExternalModelScope {
further_citations: []
}

// THE METRIC FAMILY TYPE of the classic text exposition format's `# TYPE` line
// (prometheus.io/docs/instrumenting/exposition_formats): counter, gauge, histogram, summary, and
// untyped, which is also the default when no TYPE line is given. Both histogram and summary families
// export `_sum` and `_count` series; counters, gauges and untyped families do not.
type PrometheusMetricType
= PrometheusCounter
| PrometheusGauge
| PrometheusHistogram
| PrometheusSummary
| PrometheusUntyped

fn prometheus_metric_type_word(t: PrometheusMetricType) -> String {
match t {
PrometheusCounter => "counter"
PrometheusGauge => "gauge"
PrometheusHistogram => "histogram"
PrometheusSummary => "summary"
PrometheusUntyped => "untyped"
}
}

// THE PROCESS COLLECTOR IS THE CLIENT LIBRARY'S, NOT THE SERVING ENGINE'S.
//
// Anything exporting through prometheus_client on Linux also exports a small population of series
Expand Down
11 changes: 11 additions & 0 deletions dag/extdeps/tools/curl.dag
Original file line number Diff line number Diff line change
Expand Up @@ -57,6 +57,17 @@ fn curl_localhost_bounded_get_argv_prefix() -> String {
// load-bearing are recorded once at extdeps.exec.command program_spelling_is_an_observation_claim.
data curl_program: NonEmptyStr = "curl"

// A bounded GET as argv words: -f makes an HTTP error a non-zero exit instead of a body, -sS keeps
// the error text on stderr while suppressing the progress meter.
fn curl_bounded_get_argv(url: String, connect_timeout: Second, max_time: Second) -> List<String> {
[
curl_program as String, "-fsS",
"--connect-timeout", to_string(second_count(s: connect_timeout)),
"--max-time", to_string(second_count(s: max_time)),
url,
]
}

fn curl_download_command(url: String, destination: String) -> ArgvCommand {
argv_command(program: curl_program, arguments: ["-fsL", "-o", destination, url])
}
Expand Down
173 changes: 173 additions & 0 deletions dag/extdeps/vllm/bench_serve.dag
Original file line number Diff line number Diff line change
@@ -0,0 +1,173 @@
module extdeps.vllm.bench_serve

import std.types { Bool, List, NonEmptyStr, String, Int }
import std.nat { Nat }
import v2.std.optional { Present, Absent }
import std.measure {
TokenCount, token_count, token_count_value,
Measure, Dimensionless, Milli, measure_count,
MilliRequestsPerSecond, milli_requests_per_second_count,
}
import std.decl_ref { DeclarationRef, WholeDeclaration }
import std.decimal { fixed_point_wire }
import extdeps.external_authority { ExternalAuthority, ExternalModelScope, ExternalSubjectRef, CitedFigureStanding, TranscribedUncited }
import extdeps.uri { Uri, Https }
import extdeps.vllm.server { VllmRoute, vllm_route_chat_completions, vllm_route_completions, vllm_route_path }

data extdeps_external_authority_anchor: ExternalAuthority = ExternalAuthority {
uri: Uri {
scheme: Https
locator: "github.com/vllm-project/vllm/blob/main/vllm/benchmarks/serve.py"
}
}

data extdeps_model_scope: ExternalModelScope = ExternalModelScope {
subject: ExternalSubjectRef {
declaration: DeclarationRef {
module_path: "extdeps.vllm.bench_serve",
decl_name: "VllmBenchServeRequest",
field: WholeDeclaration
}
},
first_citation: extdeps_external_authority_anchor,
further_citations: []
}

// ── UPSTREAM'S LOAD DRIVER, NOT OURS ──────────────────────────────────────────────────────────
//
// vLLM ships `vllm bench serve` with Poisson arrivals (`--request-rate` finite and
// `--burstiness` 1), request-rate control, max-concurrency, random prefixes, and trace replay.
// A second load generator in this repository would be a fork of that driver. This module names
// the request the CLI actually takes and renders the argv; it does not invent a client.

data vllm_bench_serve_source_standing: CitedFigureStanding = TranscribedUncited {
read_obligation: "bind this argv model to the serve.py file digest inside the running image; citing github.com/vllm-project/vllm/blob/main/vllm/benchmarks/serve.py is not the installed revision" as NonEmptyStr
}

type VllmBenchBackend
= OpenaiCompletions
| OpenaiChat

fn vllm_bench_backend_wire(b: VllmBenchBackend) -> NonEmptyStr {
match b {
OpenaiCompletions => "openai" as NonEmptyStr
OpenaiChat => "openai-chat" as NonEmptyStr
}
}

fn vllm_bench_backend_route(b: VllmBenchBackend) -> VllmRoute {
match b {
OpenaiCompletions => vllm_route_completions
OpenaiChat => vllm_route_chat_completions
}
}

// FINITE RATE IS POISSON WHEN BURSTINESS IS 1; INF SENDS EVERY REQUEST AT TIME ZERO.
// That last arm is the all-at-once launch whose wall clock the bash-loop instrument confused
// with engine time. It remains representable because the upstream CLI offers it; a holistic
// session simulation must not select it.
type BenchArrival
= SendAllImmediately
| PoissonProcess { request_rate: MilliRequestsPerSecond, burstiness: Measure<Dimensionless, Milli, Nat> }

type RandomDatasetSpec {
input_len: TokenCount
output_len: TokenCount
prefix_len: TokenCount
}

type VllmBenchServeRequest {
backend: VllmBenchBackend
base_url: NonEmptyStr
model: NonEmptyStr
tokenizer: NonEmptyStr
dataset: RandomDatasetSpec
num_prompts: Nat
arrival: BenchArrival
max_concurrency: Nat?
ignore_eos: Bool
seed: Nat
}

fn vllm_bench_serve_request(
backend: VllmBenchBackend,
base_url: NonEmptyStr,
model: NonEmptyStr,
tokenizer: NonEmptyStr,
dataset: RandomDatasetSpec,
num_prompts: Nat,
arrival: BenchArrival,
max_concurrency: Nat?,
ignore_eos: Bool,
seed: Nat,
) -> VllmBenchServeRequest {
VllmBenchServeRequest {
backend: backend,
base_url: base_url,
model: model,
tokenizer: tokenizer,
dataset: dataset,
num_prompts: num_prompts,
arrival: arrival,
max_concurrency: max_concurrency,
ignore_eos: ignore_eos,
seed: seed,
}
}

fn vllm_bench_arrival_rate_wire(a: BenchArrival) -> String {
match a {
SendAllImmediately => "inf"
PoissonProcess { request_rate: n, burstiness: _ } =>
fixed_point_wire(value: milli_requests_per_second_count(r: n) as Int, decimals: 3)
}
}

fn vllm_bench_burstiness_words(a: BenchArrival) -> List<String> {
match a {
SendAllImmediately => []
PoissonProcess { request_rate: _, burstiness: b } =>
["--burstiness", fixed_point_wire(value: measure_count(b) as Int, decimals: 3)]
}
}

fn vllm_bench_concurrency_words(n: Nat?) -> List<String> {
match n {
Absent => []
Present { value: c } => ["--max-concurrency", to_string(c)]
}
}

fn vllm_bench_eos_words(ignore: Bool) -> List<String> {
if ignore { ["--ignore-eos"] } else { [] }
}

// THE ARGV IS DERIVED FROM THE REQUEST, so a protocol that asked for Poisson arrivals cannot
// silently render `inf`, and a protocol that named a shared prefix cannot drop `--random-prefix-len`.
// SERVED NAME AND TOKENIZER IDENTITY ARE DIFFERENT FACTS. --model is what the OpenAI route answers;
// --tokenizer is the tree the client actually loads. A protocol that carries only the served alias
// tries HuggingFace and cannot be replayed against this engine.
fn vllm_bench_serve_arguments(req: VllmBenchServeRequest) -> List<String> {
concat(
[
"bench", "serve",
"--backend", vllm_bench_backend_wire(b: req.backend) as String,
"--base-url", req.base_url as String,
"--model", req.model as String,
"--tokenizer", req.tokenizer as String,
"--endpoint", vllm_route_path(r: vllm_bench_backend_route(b: req.backend)) as String,
"--dataset-name", "random",
"--random-input-len", to_string(token_count_value(t: req.dataset.input_len)),
"--random-output-len", to_string(token_count_value(t: req.dataset.output_len)),
"--random-prefix-len", to_string(token_count_value(t: req.dataset.prefix_len)),
"--num-prompts", to_string(req.num_prompts),
"--seed", to_string(req.seed),
"--request-rate", vllm_bench_arrival_rate_wire(a: req.arrival),
],
concat(
vllm_bench_burstiness_words(a: req.arrival),
concat(vllm_bench_concurrency_words(n: req.max_concurrency), vllm_bench_eos_words(ignore: req.ignore_eos)),
),
)
}

44 changes: 44 additions & 0 deletions dag/extdeps/vllm/metrics.dag
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import std.types { Bool, String, NonEmptyStr, List }
import std.decl_ref { DeclarationRef, WholeDeclaration }
import extdeps.external_authority { ExternalAuthority, ExternalModelScope, ExternalSubjectRef }
import extdeps.uri { Uri, Https }
import extdeps.prometheus.client { PrometheusMetricType, PrometheusCounter, PrometheusGauge, PrometheusHistogram }
import extdeps.vllm.server { VllmSourceRevision, vllm_source_revision_read }

data extdeps_external_authority_anchor: ExternalAuthority = ExternalAuthority {
Expand Down Expand Up @@ -444,3 +445,46 @@ fn interconnect_counter_note(c: InterconnectCounter) -> NonEmptyStr {
"port_xmit_data and port_rcv_data, in 4-byte words, are the RoCE traffic; measured 2026-09-07 at ~7.9 Gb/s per rank on a 50 Gb/s link during a cold prefill. That is a BULK-TRANSFER reading and bounds bandwidth saturation only: a low average is equally consistent with latency-bound collectives, a straggler rank, or host gaps between kernels, so it does not establish what share of the critical path was communication" as NonEmptyStr
}
}

// ── HISTOGRAM SUM AND COUNT ARE THE ENGINE CLOCK FOR A RATE, NOT A SECOND METRIC NAME ─────────
//
// A RATE DERIVED FROM SHELL WALL CLOCK is the instrument that reported a 4-to-8 stream cliff that
// did not exist (launcher arrival stagger landing in the wall) and understated aggregate at 16
// streams by about 40%. The engine's own histogram exposition carries `_sum` and `_count` on the
// families that already live in VllmServingMetric. Those suffixes are Prometheus histogram
// statistics, not new serving questions, so they are projections of a family rather than new
// metric variants.
//
// InterTokenLatencySeconds: count is completed decode intervals, sum is their seconds. That ratio
// is a latency distribution, not aggregate tokens/s. IterationTokensTotal: count is recorded
// iterations, sum is tokens booked in them, not occupancy and not steps. Aggregate throughput is
// vllm:generation_tokens_total on one exporter across a scrape interval.
type HistogramStatistic
= HistogramSum
| HistogramCount

fn vllm_histogram_series_name(m: VllmServingMetric, s: HistogramStatistic) -> NonEmptyStr {
match s {
HistogramSum => join([vllm_metric_name(m: m) as String, "_sum"], "") as NonEmptyStr
HistogramCount => join([vllm_metric_name(m: m) as String, "_count"], "") as NonEmptyStr
}
}

// Each family's Prometheus type as vLLM's metrics module REGISTERS it (Histogram, Counter or Gauge).
// A projection from the catalog, not a reading of the running exporter's `# TYPE` line.
fn vllm_metric_type(m: VllmServingMetric) -> PrometheusMetricType {
match m {
InterTokenLatencySeconds => PrometheusHistogram
IterationTokensTotal => PrometheusHistogram
KvCacheUsagePerc => PrometheusGauge
NumRequestsRunning => PrometheusGauge
NumRequestsWaiting => PrometheusGauge
NumRequestsWaitingByReason => PrometheusGauge
PrefixCacheQueries => PrometheusCounter
PrefixCacheHits => PrometheusCounter
PromptTokensTotal => PrometheusCounter
GenerationTokensTotal => PrometheusCounter
NumPreemptionsTotal => PrometheusCounter
EstimatedFlopsPerGpuTotal => PrometheusCounter
}
}
7 changes: 7 additions & 0 deletions dag/extdeps/vllm/server.dag
Original file line number Diff line number Diff line change
Expand Up @@ -91,11 +91,17 @@ fn vllm_route(path: NonEmptyStr, population: VllmRoutePopulation) -> VllmRoute {
VllmRoute { path: path, population: population, established_against: vllm_source_revision_read }
}

fn vllm_route_path(r: VllmRoute) -> NonEmptyStr {
r.path
}

data vllm_route_models: VllmRoute = vllm_route(path: "/v1/models" as NonEmptyStr, population: OrdinarilyRegistered)
data vllm_route_version: VllmRoute = vllm_route(path: "/version" as NonEmptyStr, population: OrdinarilyRegistered)
data vllm_route_health: VllmRoute = vllm_route(path: "/health" as NonEmptyStr, population: OrdinarilyRegistered)
data vllm_route_ping: VllmRoute = vllm_route(path: "/ping" as NonEmptyStr, population: OrdinarilyRegistered)
data vllm_route_metrics: VllmRoute = vllm_route(path: "/metrics" as NonEmptyStr, population: OrdinarilyRegistered)
data vllm_route_completions: VllmRoute = vllm_route(path: "/v1/completions" as NonEmptyStr, population: OrdinarilyRegistered)
data vllm_route_chat_completions: VllmRoute = vllm_route(path: "/v1/chat/completions" as NonEmptyStr, population: OrdinarilyRegistered)

data vllm_route_server_info: VllmRoute = vllm_route(path: "/server_info" as NonEmptyStr, population: RegisteredOnlyUnderDevelopmentMode)
data vllm_route_get_world_size: VllmRoute = vllm_route(path: "/get_world_size" as NonEmptyStr, population: RegisteredOnlyUnderDevelopmentMode)
Expand Down Expand Up @@ -190,6 +196,7 @@ data vllm_http_parallel_evidence_routes: List<VllmParallelEvidenceRoute> = [

data vllm_routes: List<VllmRoute> = [
vllm_route_models, vllm_route_version, vllm_route_health, vllm_route_ping, vllm_route_metrics,
vllm_route_completions, vllm_route_chat_completions,
vllm_route_server_info, vllm_route_get_world_size, vllm_route_collective_rpc,
]

Expand Down
69 changes: 68 additions & 1 deletion dag/gunbc/harness/harness_backend.dag
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import std.types { String, Bool, Int, List, NonEmptyStr, Timestamp, EpochSecs }
import std.nat { Nat }
import std.checked_arithmetic { CheckedNatReady, CheckedNatOverflow, checked_int_magnitude }
import std.algebra { trim }
import v2.std.algebra { filter }
import v2.std.optional { Present, Absent }
import product.fabric.demand { ObservationReceiptRef }
import std.measure { Second, second, second_count }
Expand All @@ -14,7 +15,8 @@ import gunbc.serving.serving_availability {
}
import extdeps.vllm.metrics {
VllmServingMetric, NumRequestsWaiting, NumRequestsRunning, NumRequestsWaitingByReason, GenerationTokensTotal,
vllm_metric_name, vllm_waiting_reason_label, vllm_waiting_reason_capacity,
HistogramStatistic, vllm_metric_name, vllm_histogram_series_name,
vllm_waiting_reason_label, vllm_waiting_reason_capacity,
}
import extdeps.prometheus.client { ProcessStartTimeSeconds, prometheus_process_series_name }
import gunbc.spark.fabric_switch_observed { FabricGroup, fabric_group_wire }
Expand Down Expand Up @@ -361,6 +363,14 @@ fn harness_gauge_value(body: String, metric: VllmServingMetric) -> String? {
harness_metric_line_value(body: body, prefix: join([vllm_metric_name(m: metric) as String, "{"], ""), selector: SeriesByName)
}

fn harness_histogram_statistic_value(body: String, metric: VllmServingMetric, statistic: HistogramStatistic) -> String? {
harness_metric_line_value(
body: body,
prefix: join([vllm_histogram_series_name(m: metric, s: statistic) as String, "{"], ""),
selector: SeriesByName,
)
}

// THE LABEL PAIR IS BUILT ONCE, HERE, from the label key and value rows the upstream module carries,
// so no consumer spells a Prometheus label rendering of its own.
fn harness_labelled_gauge_value(body: String, metric: VllmServingMetric, label: NonEmptyStr, value: NonEmptyStr) -> String? {
Expand Down Expand Up @@ -546,3 +556,60 @@ fn harness_observe_engine_at(endpoint: ServingRouteEndpoint, at: EpochSecs) -> H
}
}
}

fn line_is_sudo_password_prompt(line: String) -> Bool {
match first(split(s: trim(s: line), delimiter: "[sudo] password for")) {
Absent => false
Present { value: prefix } => prefix == ""
}
}

fn text_without_sudo_prompt(body: String) -> String {
join(filter(xs: split(s: body, delimiter: "\n"), predicate: fn(line) { !line_is_sudo_password_prompt(line: line) }), "\n")
}

fn text_labeled_field(body: String, label: String) -> String? {
let needle = join([label, ":"], "")
let lines = split(s: body, delimiter: "\n")
let matching = filter(lines, l => starts_with(s: trim(s: l), prefix: needle))
if length(matching) != 1 {
none
} else {
match first(matching) {
Absent => none
Present { value: line } => {
let parts = split(s: line, delimiter: ":")
if length(parts) < 2 {
none
} else {
last(parts)
}
}
}
}
}

fn text_field_last_word(body: String, label: String) -> String? {
match text_labeled_field(body: body, label: label) {
Absent => none
Present { value: t } => last(split(s: trim(s: t), delimiter: " "))
}
}

fn text_labeled_count(body: String, label: String) -> Nat? {
match text_field_last_word(body: body, label: label) {
Absent => none
Present { value: last_word } => harness_gauge_count(text: trim(s: last_word))
}
}

fn text_labeled_head_count(body: String, label: String) -> Nat? {
match text_field_last_word(body: body, label: label) {
Absent => none
Present { value: last_word } =>
match first(split(s: trim(s: last_word), delimiter: ".")) {
Absent => none
Present { value: whole } => harness_gauge_count(text: whole)
}
}
}
Loading