diff --git a/dag/extdeps/linux/cgroup_v2_memory.dag b/dag/extdeps/linux/cgroup_v2_memory.dag index ef874c98223..174c4270027 100644 --- a/dag/extdeps/linux/cgroup_v2_memory.dag +++ b/dag/extdeps/linux/cgroup_v2_memory.dag @@ -1,5 +1,7 @@ module extdeps.linux.cgroup_v2_memory +import v2.std.algebra { filter } +import v2.std.collection { list_at_optional } import std.algebra { trim } import std.checked_arithmetic { nat_magnitude } import extdeps.external_authority { ExternalAuthority, ExternalModelScope, ExternalSubjectRef } @@ -7,6 +9,7 @@ import std.decl_ref { DeclarationRef, WholeDeclaration } import extdeps.uri { Https, Uri } import std.measure { ByteSize, byte_size, byte_size_count } import std.types { Int, NonEmptyStr, String, List } +import std.nat { Nat } import v2.std.optional { Present, Absent } data extdeps_external_authority_anchor: ExternalAuthority = ExternalAuthority { @@ -192,3 +195,68 @@ fn cgroup_v2_bounding_dirs(mount_point: NonEmptyStr, relative_path: NonEmptyStr) CgroupBoundingDirsFold { prefix: next, dirs: acc.dirs |> list_push(next as NonEmptyStr) } }).dirs } + +// ── memory.stat: THE FIVE KEYS A RESIDENCY READING NEEDS ─────────────────────────────────────── +// +// cgroup-v2 ("memory.stat") documents the file as flat "key value" lines. Its keys have two +// different kinds. anon, file and file_mapped are Instantaneous byte counts. workingset_refault_file +// and pgmajfault are CumulativeCounter event counts, so only a before/after pair of them measures an +// interval. The kinds are named per field here because this one file mixes them; that is also why +// memory.stat is not a CgroupMemoryInterfaceFile variant with a single kind. +type CgroupMemoryStat { + anon: ByteSize + file: ByteSize + file_mapped: ByteSize + workingset_refault_file: Nat + pgmajfault: Nat +} + +type CgroupMemoryStatReading + = CgroupMemoryStatRead { stat: CgroupMemoryStat } + | CgroupMemoryStatIncomplete { missing: List } + +fn cgroup_memory_stat_value(body: String, key: String) -> Nat? { + list_at_optional(xs: flat_map(split(s: body, delimiter: "\n"), line => { + let words = filter(split(s: trim(s: line), delimiter: " "), w => w != "") + match list_at_optional(xs: words, index: 0) { + Absent => [] as List + Present { value: k } => + if k != key { [] as List } else { + match list_at_optional(xs: words, index: 1) { + Absent => [] as List + Present { value: v } => match parse_int(s: v) { + Absent => [] as List + Present { value: n } => if n < 0 { [] as List } else { [nat_magnitude(a: n)] } + } + } + } + } + }), index: 0) +} + +fn parse_cgroup_memory_stat(body: String) -> CgroupMemoryStatReading { + let keys = ["anon", "file", "file_mapped", "workingset_refault_file", "pgmajfault"] + let missing = filter(keys, k => match cgroup_memory_stat_value(body: body, key: k) { Present { value: _ } => false Absent => true }) + match cgroup_memory_stat_value(body: body, key: "anon") { + Absent => CgroupMemoryStatIncomplete { missing: missing } + Present { value: anon } => match cgroup_memory_stat_value(body: body, key: "file") { + Absent => CgroupMemoryStatIncomplete { missing: missing } + Present { value: file } => match cgroup_memory_stat_value(body: body, key: "file_mapped") { + Absent => CgroupMemoryStatIncomplete { missing: missing } + Present { value: mapped } => match cgroup_memory_stat_value(body: body, key: "workingset_refault_file") { + Absent => CgroupMemoryStatIncomplete { missing: missing } + Present { value: refault } => match cgroup_memory_stat_value(body: body, key: "pgmajfault") { + Absent => CgroupMemoryStatIncomplete { missing: missing } + Present { value: majfault } => CgroupMemoryStatRead { stat: CgroupMemoryStat { + anon: byte_size(count: anon), + file: byte_size(count: file), + file_mapped: byte_size(count: mapped), + workingset_refault_file: refault, + pgmajfault: majfault, + } } + } + } + } + } + } +} diff --git a/dag/extdeps/linux/diskstats.dag b/dag/extdeps/linux/diskstats.dag new file mode 100644 index 00000000000..5aad6b70d62 --- /dev/null +++ b/dag/extdeps/linux/diskstats.dag @@ -0,0 +1,144 @@ +module extdeps.linux.diskstats + +import v2.std.algebra { filter } +import v2.std.collection { list_at_optional } +import std.types { Bool, List, NonEmptyStr, String } +import std.algebra { trim } +import std.nat { Nat } +import std.checked_arithmetic { nat_magnitude } +import std.measure { ByteSize, byte_size, byte_size_count, Microsecond, microsecond, Millisecond, millisecond, millisecond_count } +import std.decl_ref { DeclarationRef, WholeDeclaration } +import extdeps.external_authority { ExternalAuthority, ExternalModelScope, ExternalSubjectRef } +import extdeps.uri { Https, Uri } +import extdeps.linux.kernel { LinuxKernelRelease } +import v2.std.optional { Present, Absent } + +// /proc/diskstats, Documentation/admin-guide/iostats.rst. The 14-field-after-name layout this +// module reads has been stable since 2.6; later kernels APPEND discard (4.18) and flush (5.5) +// fields, which is why fields are read by position from the front and never by line length. +data extdeps_external_authority_anchor: ExternalAuthority = ExternalAuthority { + uri: Uri { + scheme: Https + locator: "docs.kernel.org/admin-guide/iostats.html" + } +} + +data extdeps_model_scope: ExternalModelScope = ExternalModelScope { + subject: ExternalSubjectRef { + declaration: DeclarationRef { + module_path: "extdeps.linux.diskstats", + decl_name: "DiskstatsRow", + field: WholeDeclaration + } + }, + first_citation: extdeps_external_authority_anchor, + further_citations: [] +} + +data diskstats_positional_layout_since_release: LinuxKernelRelease = "2.6" as LinuxKernelRelease + +// iostats.rst: "sectors" in this file are ALWAYS 512 bytes, whatever the device's logical block +// size is. Multiplying by the device's own sector size would be a wrong reading on 4Kn drives. +data diskstats_sector_bytes: ByteSize = byte_size(count: 512) + +// The read-side fields, by their iostats.rst numbers (field 1 = first after the device name). +// Write-side fields are not read here because the measurement that consumes this is read-only +// against the Engram files; a later consumer adds them as fields, not as a second row type. +type DiskstatsRow { + device: NonEmptyStr + reads_completed: Nat + sectors_read: Nat + read_ticks: Millisecond + in_flight: Nat +} + +type DiskstatsReading + = DiskstatsDeviceRead { row: DiskstatsRow } + | DiskstatsDeviceAbsent { device: NonEmptyStr } + | DiskstatsUnparseable { device: NonEmptyStr, line: String } + +fn diskstats_nat_at(words: List, index: Nat) -> Nat? { + match list_at_optional(xs: words, index: index) { + Absent => none + Present { value: w } => match parse_int(s: w) { + Absent => none + Present { value: n } => if n < 0 { none } else { Present { value: nat_magnitude(a: n) } } + } + } +} + +// Word 2 is the device name; words 3, 5, 6 and 11 are iostats fields 1, 3, 4 and 9. +fn diskstats_row_of_line(device: NonEmptyStr, line: String) -> DiskstatsReading? { + let words = filter(split(s: trim(s: line), delimiter: " "), w => w != "") + match list_at_optional(xs: words, index: 2) { + Absent => none + Present { value: name } => + if name != (device as String) { + none + } else { + match diskstats_nat_at(words: words, index: 3) { + Absent => Present { value: DiskstatsUnparseable { device: device, line: line } } + Present { value: reads } => match diskstats_nat_at(words: words, index: 5) { + Absent => Present { value: DiskstatsUnparseable { device: device, line: line } } + Present { value: sectors } => match diskstats_nat_at(words: words, index: 6) { + Absent => Present { value: DiskstatsUnparseable { device: device, line: line } } + Present { value: ticks } => match diskstats_nat_at(words: words, index: 11) { + Absent => Present { value: DiskstatsUnparseable { device: device, line: line } } + Present { value: inflight } => Present { value: DiskstatsDeviceRead { row: DiskstatsRow { + device: device, + reads_completed: reads, + sectors_read: sectors, + read_ticks: millisecond(count: ticks), + in_flight: inflight, + } } } + } + } + } + } + } + } +} + +fn parse_diskstats(body: String, device: NonEmptyStr) -> DiskstatsReading { + match list_at_optional(xs: flat_map(split(s: body, delimiter: "\n"), l => match diskstats_row_of_line(device: device, line: l) { + Present { value: r } => [r] + Absent => [] as List + }), index: 0) { + Present { value: r } => r + Absent => DiskstatsDeviceAbsent { device: device } + } +} + +// A READ INTERVAL IS A DELTA OF TWO CUMULATIVE ROWS. iostats.rst warns the counters may wrap; a +// field that went backwards is a wrap or a device re-registration, and the pair measures nothing. +// The mean await is the kernel's own definition (read ticks spent per completed read), floored to +// whole microseconds; with no completed reads there is no await, which is a distinct arm from zero. +type DiskReadInterval + = DiskReadMeasured { read_bytes: ByteSize, reads: Nat, mean_read_await: Microsecond? } + | DiskReadUnmeasurable { cause: NonEmptyStr } + +fn disk_read_interval(before: DiskstatsReading, after: DiskstatsReading) -> DiskReadInterval { + match before { + DiskstatsDeviceAbsent { device: d } => DiskReadUnmeasurable { cause: "the device is absent from the before reading" as NonEmptyStr } + DiskstatsUnparseable { device: _, line: _ } => DiskReadUnmeasurable { cause: "the before line is unparseable" as NonEmptyStr } + DiskstatsDeviceRead { row: b } => match after { + DiskstatsDeviceAbsent { device: d } => DiskReadUnmeasurable { cause: "the device is absent from the after reading" as NonEmptyStr } + DiskstatsUnparseable { device: _, line: _ } => DiskReadUnmeasurable { cause: "the after line is unparseable" as NonEmptyStr } + DiskstatsDeviceRead { row: a } => + if b.device != a.device { + DiskReadUnmeasurable { cause: "the two readings name different devices" as NonEmptyStr } + } else if a.reads_completed < b.reads_completed || a.sectors_read < b.sectors_read + || millisecond_count(m: a.read_ticks) < millisecond_count(m: b.read_ticks) { + DiskReadUnmeasurable { cause: "a cumulative field went backwards (wrap or re-registration)" as NonEmptyStr } + } else { + let reads = a.reads_completed - b.reads_completed + let ticks = millisecond_count(m: a.read_ticks) - millisecond_count(m: b.read_ticks) + DiskReadMeasured { + read_bytes: byte_size(count: (a.sectors_read - b.sectors_read) * byte_size_count(b: diskstats_sector_bytes)), + reads: reads, + mean_read_await: if reads == 0 { none } else { Present { value: microsecond(count: (ticks * 1000) / reads) } }, + } + } + } + } +} diff --git a/dag/extdeps/linux/mincore.dag b/dag/extdeps/linux/mincore.dag new file mode 100644 index 00000000000..17305bddfd1 --- /dev/null +++ b/dag/extdeps/linux/mincore.dag @@ -0,0 +1,66 @@ +module extdeps.linux.mincore + +import std.types { Bool, NonEmptyStr, String } +import std.nat { Nat } +import std.measure { ByteSize, byte_size, byte_size_count } +import std.decl_ref { DeclarationRef, WholeDeclaration } +import extdeps.external_authority { ExternalAuthority, ExternalModelScope, ExternalSubjectRef } +import extdeps.uri { Https, Uri } + +// mincore(2): "determine whether pages are resident in memory". For a mapping of a file it returns +// one byte per page of the range, whose least significant bit is set when the page is resident in +// the page cache. It answers RESIDENCY at the instant of the call, and nothing about which pages any +// process touched: a page read by readahead, by a checksum pass, or by another process is resident +// exactly like a page a lookup fetched. That is the whole reason this reading is never a touch count. +data extdeps_external_authority_anchor: ExternalAuthority = ExternalAuthority { + uri: Uri { + scheme: Https + locator: "man7.org/linux/man-pages/man2/mincore.2.html" + } +} + +data extdeps_model_scope: ExternalModelScope = ExternalModelScope { + subject: ExternalSubjectRef { + declaration: DeclarationRef { + module_path: "extdeps.linux.mincore", + decl_name: "MincoreFileResidency", + field: WholeDeclaration + } + }, + first_citation: extdeps_external_authority_anchor, + further_citations: [] +} + +// The vector covers ceil(length / page_size) pages; the resident count is the number of bytes in it +// with bit 0 set. Resident BYTES are pages times page size, so the file's partial last page counts +// whole -- that is what the page cache holds, and it can exceed the file length by less than a page. +type MincoreFileResidency sole_constructor { + file: NonEmptyStr + file_bytes: ByteSize + page_bytes: ByteSize + resident_pages: Nat +} + +type MincoreResidencyAdmission + = MincoreResidencyAdmitted { residency: MincoreFileResidency } + | MincoreResidencyRefused { file: NonEmptyStr, cause: NonEmptyStr } + +fn mincore_page_count(file_bytes: ByteSize, page_bytes: ByteSize) -> Nat { + (byte_size_count(b: file_bytes) + byte_size_count(b: page_bytes) - 1) / byte_size_count(b: page_bytes) +} + +fn mincore_file_residency(file: NonEmptyStr, file_bytes: ByteSize, page_bytes: ByteSize, resident_pages: Nat) -> MincoreResidencyAdmission { + if byte_size_count(b: page_bytes) == 0 { + MincoreResidencyRefused { file: file, cause: "a zero page size is not a page size" as NonEmptyStr } + } else if resident_pages > mincore_page_count(file_bytes: file_bytes, page_bytes: page_bytes) { + MincoreResidencyRefused { file: file, cause: "more resident pages than the file has pages; the vector covered a different range" as NonEmptyStr } + } else { + MincoreResidencyAdmitted { residency: MincoreFileResidency { + file: file, file_bytes: file_bytes, page_bytes: page_bytes, resident_pages: resident_pages, + } } + } +} + +fn mincore_resident_bytes(r: MincoreFileResidency) -> ByteSize { + byte_size(count: r.resident_pages * byte_size_count(b: r.page_bytes)) +} diff --git a/dag/extdeps/linux/psi.dag b/dag/extdeps/linux/psi.dag new file mode 100644 index 00000000000..7c28988eb92 --- /dev/null +++ b/dag/extdeps/linux/psi.dag @@ -0,0 +1,122 @@ +module extdeps.linux.psi + +import v2.std.algebra { filter } +import v2.std.collection { list_at_optional } +import std.types { Bool, List, NonEmptyStr, String } +import std.algebra { trim } +import std.nat { Nat } +import std.checked_arithmetic { nat_magnitude } +import std.measure { Microsecond, microsecond, microsecond_count } +import std.decl_ref { DeclarationRef, WholeDeclaration } +import extdeps.external_authority { ExternalAuthority, ExternalModelScope, ExternalSubjectRef } +import extdeps.uri { Https, Uri } +import extdeps.linux.kernel { LinuxKernelRelease } +import v2.std.optional { Present, Absent } + +// Pressure Stall Information, Documentation/accounting/psi.rst. The interface first shipped in +// Linux 4.20; the per-cgroup files (memory.pressure, io.pressure, cpu.pressure) in cgroup-v2 share +// this exact format (the same document, "Cgroup2 interface"). +data extdeps_external_authority_anchor: ExternalAuthority = ExternalAuthority { + uri: Uri { + scheme: Https + locator: "docs.kernel.org/accounting/psi.html" + } +} + +data extdeps_model_scope: ExternalModelScope = ExternalModelScope { + subject: ExternalSubjectRef { + declaration: DeclarationRef { + module_path: "extdeps.linux.psi", + decl_name: "PsiReading", + field: WholeDeclaration + } + }, + first_citation: extdeps_external_authority_anchor, + further_citations: [] +} + +// The first kernel release that carries the interface, as psi.rst's version history states it. +data psi_first_kernel_release: LinuxKernelRelease = "4.20" as LinuxKernelRelease + +// THE FILE IS TWO LINES, `some` AND `full`, each "avg10=A avg60=B avg300=C total=T". The averages +// are the kernel's own exponentially-decayed ratios over windows the kernel chose; `total` is the +// cumulative stall time in MICROSECONDS since boot (or cgroup creation). Only `total` is carried: +// a stall fraction over an interval the measurement itself chose is the delta of two totals over +// that interval's length, and the kernel's decayed averages would be a second, differently-windowed +// answer to the same question. `full` is absent from the system-wide cpu file on kernels before +// 5.13 and is carried as optional for that reason, never defaulted to zero. +type PsiLineKind = PsiSome | PsiFull + +type PsiReading + = PsiRead { some_total: Microsecond, full_total: Microsecond? } + | PsiUnparseable { body: String, cause: NonEmptyStr } + +fn psi_line_keyword(k: PsiLineKind) -> String { + match k { + PsiSome => "some" + PsiFull => "full" + } +} + +fn psi_total_of_line(line: String, kind: PsiLineKind) -> Microsecond? { + let words = filter(split(s: trim(s: line), delimiter: " "), w => w != "") + match list_at_optional(xs: words, index: 0) { + Absent => none + Present { value: head } => + if head != psi_line_keyword(k: kind) { + none + } else { + match list_at_optional(xs: filter(words, w => starts_with(s: w, prefix: "total=")), index: 0) { + Absent => none + Present { value: field } => + match parse_int(s: substring(s: field, start: 6, end: string_length(s: field))) { + Absent => none + Present { value: n } => if n < 0 { none } else { Present { value: microsecond(count: nat_magnitude(a: n)) } } + } + } + } + } +} + +fn psi_total_in_body(body: String, kind: PsiLineKind) -> Microsecond? { + list_at_optional(xs: flat_map(split(s: body, delimiter: "\n"), l => match psi_total_of_line(line: l, kind: kind) { + Present { value: t } => [t] + Absent => [] as List + }), index: 0) +} + +fn parse_psi(body: String) -> PsiReading { + match psi_total_in_body(body: body, kind: PsiSome) { + Absent => PsiUnparseable { body: body, cause: "no `some ... total=N` line" as NonEmptyStr } + Present { value: s } => PsiRead { some_total: s, full_total: psi_total_in_body(body: body, kind: PsiFull) } + } +} + +// THE STALL OVER AN INTERVAL IS A DELTA OF TWO CUMULATIVE TOTALS, and a total that went BACKWARDS +// is not a small interval: it is a reset (a new cgroup, a reboot) and the pair measures nothing. +type PsiStallInterval + = PsiStallMeasured { stalled: Microsecond, window: Microsecond } + | PsiStallUnmeasurable { cause: NonEmptyStr } + +fn psi_full_stall_interval(before: PsiReading, after: PsiReading, window: Microsecond) -> PsiStallInterval { + match before { + PsiUnparseable { body: _, cause: c } => PsiStallUnmeasurable { cause: c } + PsiRead { some_total: _, full_total: fb } => match after { + PsiUnparseable { body: _, cause: c } => PsiStallUnmeasurable { cause: c } + PsiRead { some_total: _, full_total: fa } => match fb { + Absent => PsiStallUnmeasurable { cause: "the `full` line is absent from the before reading" as NonEmptyStr } + Present { value: b } => match fa { + Absent => PsiStallUnmeasurable { cause: "the `full` line is absent from the after reading" as NonEmptyStr } + Present { value: a } => + if microsecond_count(m: a) < microsecond_count(m: b) { + PsiStallUnmeasurable { cause: "the cumulative `full` total went backwards; the counter was reset between readings" as NonEmptyStr } + } else if microsecond_count(m: window) == 0 { + PsiStallUnmeasurable { cause: "a zero-length window measures no stall fraction" as NonEmptyStr } + } else { + PsiStallMeasured { stalled: microsecond(count: microsecond_count(m: a) - microsecond_count(m: b)), window: window } + } + } + } + } + } +} diff --git a/dag/gunbc/spark/v41_capacity_measurement.dag b/dag/gunbc/spark/v41_capacity_measurement.dag new file mode 100644 index 00000000000..2643a812892 --- /dev/null +++ b/dag/gunbc/spark/v41_capacity_measurement.dag @@ -0,0 +1,852 @@ +module gunbc.spark.v41_capacity_measurement + +import v2.std.algebra { filter } +import v2.std.collection { list_at_optional } +import std.types { Bool, Int, List, NonEmptyStr, String } +import std.nat { Nat, nat_min } +import std.measure { + ByteSize, byte_size, byte_size_count, TokenCount, token_count, token_count_value, + Millisecond, millisecond, millisecond_count, Microsecond, microsecond, microsecond_count, + BasisPoint, basis_point, basis_point_count, concurrent_request_hundredths_count, + Gibibyte, gibibyte, gibibyte_to_byte_size, +} +import std.content_hash { ContentHash } +import std.dissolution { DissolutionCondition, unbound_dissolution } +import v2.std.optional { Present, Absent } +import extdeps.external_authority { CitedFigureStanding, TranscribedUncited } +import extdeps.vllm.server { VllmSourceRevision } +import extdeps.vllm.kv_cache { VllmObservedKvPool, VllmResolvedKvFormat, ResolvedFp8DsMla, ResolvedUnobserved } +import extdeps.vllm.kv_layout { VllmReportedCapacity, observed_revisions_agree } +import extdeps.linux.psi { PsiReading, PsiStallInterval, PsiStallMeasured, PsiStallUnmeasurable, psi_full_stall_interval } +import extdeps.linux.diskstats { DiskstatsReading, DiskReadInterval, DiskReadMeasured, DiskReadUnmeasurable, disk_read_interval } +import extdeps.linux.mincore { MincoreFileResidency, mincore_resident_bytes } +import extdeps.linux.cgroup_v2_memory { CgroupMemoryStat } +import extdeps.container.oci.digest { OciContentDigest, render_oci_content_digest_wire } +import gunbc.floor_memory_demand { CgroupMemoryEvents } +import gunbc.spark.v41_runtime_realized { V41OciImageDigest } +import gunbc.spark.serving_load_probe { ClientBenchLatency, ClientBenchLatencyUnread, ClientBenchLatencyReported } + +// ══ V4.1 CAPACITY MEASUREMENT ═══════════════════════════════════════════════════════════════════ +// +// THE FIRST SERVING ACTION FOR V4.1 FLASH on Group A (srv5-8, one TP4 replica, fp8_ds_mla on SM12, +// Engram file-backed) is a MEASUREMENT, not production traffic (operator approval, 2026-09-25). +// This module holds its plan, its receipt carriers, and the three capacities it reports. It +// computes every answer from readings and refuses, typed and naming the missing reading, whenever +// a reading is absent. It never interpolates: a capacity is one of the concurrencies the staircase +// or the probe actually ran (extdeps.vllm.kv_cache VllmCapPoolObservation carries the same rule for +// the context cap / pool relation). +// +// THE SERVE THAT MEASURES IS ITS OWN ADMITTED REALIZATION. The production gate at the bottom refuses +// V4.1 production admission until a capacity report is admitted; it does not gate the measurement +// serve, which must be up to be measured. The two are distinct variants so the gate can tell them +// apart (ruling from proud-deer-538, 2026-09-25). + +// ── THE PLAN ──────────────────────────────────────────────────────────────────────────────────── + +// The TP4 ranks, by their rank index. +data v41_measurement_ranks: List = [0, 1, 2, 3] + +data v41_measurement_rank_count: Nat = length(v41_measurement_ranks) + +data v41_measurement_max_model_len: TokenCount = token_count(count: 262144) + +// The staircase concurrencies, in order. A capacity is always one of these, never a value between two. +data v41_staircase_concurrencies: List = [1, 2, 4, 6, 8, 12, 16] + +type V41ProbeCell { context: TokenCount, concurrency: Nat } + +data v41_probe_cells: List = [ + V41ProbeCell { context: token_count(count: 32768), concurrency: 1 }, + V41ProbeCell { context: token_count(count: 262144), concurrency: 1 }, + V41ProbeCell { context: token_count(count: 262144), concurrency: 2 }, + V41ProbeCell { context: token_count(count: 262144), concurrency: 4 }, + V41ProbeCell { context: token_count(count: 262144), concurrency: 6 }, +] + +// THE ENGRAM STARTUP ARMS ARE IMAGES, NOT A RUNTIME FLAG. Arm A keeps the page cache after the +// checksum pass: it is the pre-release overlay on main. Arm B releases it: it is the image built +// after the fetch_rows / page-release change lands. There is deliberately no switch inside one +// image that keeps the cache (the release refuses if pages stay resident). So the arm is read off +// the image's patch digest by a binding, and an image no binding names has no arm and refuses. +type V41EngramStartupArm = EngramCacheRetained | EngramCacheReleased + +type V41EngramArmBinding { arm: V41EngramStartupArm, patch_digest: ContentHash } + +// Empty until the first srv8 image exists: its patch digest is read from that image, never typed +// in from a prefix relayed in chat. +data v41_engram_arm_bindings: List = [] + +data v41_engram_arm_binding_frontier: DissolutionCondition = unbound_dissolution(description: "TRIGGER: the first srv8 V4.1 runtime image is built (after #12294) and its source-patch digest is read from the image (gunbc.spark.v41_runtime_realized); that digest is bound to EngramCacheRetained. The image built after the Engram page-release change lands is bound to EngramCacheReleased. SUFFICIENT FOR: every residency timeline cell names the arm by image content, not by a flag or a label." as NonEmptyStr) + +fn v41_engram_arm_of(patch_digest: ContentHash, bindings: List) -> V41EngramStartupArm? { + match list_at_optional(xs: filter(bindings, b => b.patch_digest == patch_digest), index: 0) { + Present { value: b } => Present { value: b.arm } + Absent => none + } +} + +fn v41_arm_wire(a: V41EngramStartupArm) -> String { + match a { + EngramCacheRetained => "A-cache-retained" + EngramCacheReleased => "B-cache-released" + } +} + +// ── BETS TO FALSIFY ───────────────────────────────────────────────────────────────────────────── +// +// These figures come from a side-chat derivation. They are carried as bets (DESIGN 4d): each has +// its standing and the reading that would settle it, and each gets a verdict against the reading. +// No capacity is ever computed from a bet. +type V41ByteRangeBet { + key: NonEmptyStr + low: ByteSize + high: ByteSize + standing: CitedFigureStanding +} + +type V41EngramResidentBet { + concurrency: Nat + resident: Gibibyte + standing: CitedFigureStanding +} + +data v41_kv_bytes_per_token_bet: V41ByteRangeBet = V41ByteRangeBet { + key: "kv bytes per logical token per rank (MLA replicated under TP4, DCP=1): 1,790 B content = (584 main + 132 indexer) x 2.5 states from compress ratios [2,2,2,1] on kv_source_layers [2,8,14,20], plus block alignment" as NonEmptyStr, + low: byte_size(count: 1800), + high: byte_size(count: 1940), + standing: TranscribedUncited { read_obligation: "the min rank's available KV bytes divided by its GPU KV cache size tokens, from the startup receipt (v41_min_rank_bytes_per_token)" as NonEmptyStr }, +} + +data v41_engram_resident_bets: List = [ + V41EngramResidentBet { concurrency: 2, resident: gibibyte(count: 19), standing: TranscribedUncited { read_obligation: "mincore resident bytes over the rank's Engram files after requests at 262k x2, disjoint prompts" as NonEmptyStr } }, + V41EngramResidentBet { concurrency: 4, resident: gibibyte(count: 30), standing: TranscribedUncited { read_obligation: "mincore resident bytes over the rank's Engram files after requests at 262k x4, disjoint prompts" as NonEmptyStr } }, + V41EngramResidentBet { concurrency: 6, resident: gibibyte(count: 37), standing: TranscribedUncited { read_obligation: "mincore resident bytes over the rank's Engram files after requests at 262k x6, disjoint prompts" as NonEmptyStr } }, + V41EngramResidentBet { concurrency: 16, resident: gibibyte(count: 46), standing: TranscribedUncited { read_obligation: "mincore resident bytes over the rank's Engram files after requests at 262k x16, disjoint prompts" as NonEmptyStr } }, +] + +type V41BetVerdict + = BetConfirmed { key: NonEmptyStr, observed: ByteSize } + | BetFalsified { key: NonEmptyStr, observed: ByteSize } + | BetUnreadable { key: NonEmptyStr, cause: NonEmptyStr } + +fn v41_judge_byte_range_bet(bet: V41ByteRangeBet, observed: ByteSize?) -> V41BetVerdict { + match observed { + Absent => BetUnreadable { key: bet.key, cause: "the reading the bet names was not taken" as NonEmptyStr } + Present { value: o } => + if byte_size_count(b: o) >= byte_size_count(b: bet.low) && byte_size_count(b: o) <= byte_size_count(b: bet.high) { + BetConfirmed { key: bet.key, observed: o } + } else { + BetFalsified { key: bet.key, observed: o } + } + } +} + +// THE ENGRAM FIGURES ARE POINT BETS ("about 19 GiB") WITH NO STATED TOLERANCE, so no band exists to +// confirm or falsify them against, and inventing one here would be a heuristic. Their verdict pairs +// the observation with the bet, both in bytes, and a reader decides with both in view. A reading +// that is missing is Unreadable, exactly as for the ranged bet. +type V41PointBetVerdict + = PointBetObserved { concurrency: Nat, bet: ByteSize, observed: ByteSize } + | PointBetUnreadable { concurrency: Nat, cause: NonEmptyStr } + +fn v41_judge_engram_resident_bet(bet: V41EngramResidentBet, observed: ByteSize?) -> V41PointBetVerdict { + match observed { + Absent => PointBetUnreadable { concurrency: bet.concurrency, cause: "no after-requests mincore residency was read at this concurrency" as NonEmptyStr } + Present { value: o } => PointBetObserved { concurrency: bet.concurrency, bet: gibibyte_to_byte_size(g: bet.resident), observed: o } + } +} + +// The observed figure the Engram bets name: a rank's mincore residency after requests, from the +// admitted timeline. It is read per rank; the bet is per rank. +fn v41_after_requests_resident(t: V41ResidencyTimeline, rank: Nat) -> ByteSize? { + match t { + ResidencyTimelineRefused { causes: _ } => none + ResidencyTimelineAdmitted { run_id: _, arm: _, cells: cs } => + match list_at_optional(xs: filter(cs, c => c.rank == rank && v41_phase_wire(p: c.phase) == v41_phase_wire(p: AfterRequests)), index: 0) { + Present { value: c } => Present { value: v41_cell_engram_resident_bytes(c: c) } + Absent => none + } + } +} + +// ── MEASUREMENT 1: THE STARTUP MEMORY RECEIPT, FROM EVERY RANK ───────────────────────────────── +// +// A READING carries every field as optional, because the wet reader takes them from log lines and +// endpoints that may not print. ADMISSION requires every field and names each missing one. A rank +// that is missing anything refuses the whole receipt: the replica is judged by its MINIMUM rank, and +// a missing rank could be that minimum. +type V41RankParallelism { tp: Nat, dcp: Nat, pp: Nat, dp: Nat } + +type V41RankMemoryBreakdown { + free: ByteSize + total: ByteSize + weights: ByteSize + non_torch: ByteSize + activation_peak: ByteSize + cuda_graph: ByteSize + available_kv: ByteSize +} + +type V41RankStartupReading { + run_id: NonEmptyStr + rank: Nat + host: NonEmptyStr + image: V41OciImageDigest? + vllm_revision: VllmSourceRevision? + parallelism: V41RankParallelism? + kv_format: VllmResolvedKvFormat + block_size: Nat? + max_model_len: TokenCount? + max_num_seqs: Nat? + max_num_batched_tokens: Nat? + spec_decode_config: NonEmptyStr? + memory: V41RankMemoryBreakdown? + reported: VllmReportedCapacity? + num_gpu_blocks: Nat? + kv_cache_groups: Nat? +} + +type V41RankStartupReceipt sole_constructor { + run_id: NonEmptyStr + rank: Nat + host: NonEmptyStr + image: V41OciImageDigest + vllm_revision: VllmSourceRevision + parallelism: V41RankParallelism + block_size: Nat + max_model_len: TokenCount + max_num_seqs: Nat + max_num_batched_tokens: Nat + spec_decode_config: NonEmptyStr + memory: V41RankMemoryBreakdown + reported: VllmReportedCapacity + num_gpu_blocks: Nat + kv_cache_groups: Nat +} + +fn v41_missing_if(v: T?, name: String) -> List { + match v { + Present { value: _ } => [] as List + Absent => [name] + } +} + +fn v41_rank_reading_missing(r: V41RankStartupReading) -> List { + concat(concat(concat(concat(concat(concat(concat(concat(concat(concat(concat(concat( + v41_missing_if(v: r.image, name: "image digest"), + v41_missing_if(v: r.vllm_revision, name: "vLLM commit")), + v41_missing_if(v: r.parallelism, name: "TP/DCP/PP/DP")), + match r.kv_format { ResolvedFp8DsMla => [] as List ResolvedUnobserved => ["resolved kv dtype"] }), + v41_missing_if(v: r.block_size, name: "resolved block size")), + v41_missing_if(v: r.max_model_len, name: "max_model_len")), + v41_missing_if(v: r.max_num_seqs, name: "max_num_seqs")), + v41_missing_if(v: r.max_num_batched_tokens, name: "max_num_batched_tokens")), + v41_missing_if(v: r.spec_decode_config, name: "spec-decode config")), + v41_missing_if(v: r.memory, name: "memory breakdown (free/total, weights, non-torch, activation peak, CUDA graphs, available KV)")), + v41_missing_if(v: r.reported, name: "GPU KV cache size tokens and max concurrency")), + v41_missing_if(v: r.num_gpu_blocks, name: "num_gpu_blocks")), + v41_missing_if(v: r.kv_cache_groups, name: "KV cache groups")) +} + +type V41RankStartupAdmission + = RankStartupAdmitted { receipt: V41RankStartupReceipt } + | RankStartupRefused { rank: Nat, missing: List } + +fn v41_admit_rank_startup(r: V41RankStartupReading) -> V41RankStartupAdmission { + match r.image { Absent => RankStartupRefused { rank: r.rank, missing: v41_rank_reading_missing(r: r) } + Present { value: image } => match r.vllm_revision { Absent => RankStartupRefused { rank: r.rank, missing: v41_rank_reading_missing(r: r) } + Present { value: rev } => match r.parallelism { Absent => RankStartupRefused { rank: r.rank, missing: v41_rank_reading_missing(r: r) } + Present { value: par } => match r.block_size { Absent => RankStartupRefused { rank: r.rank, missing: v41_rank_reading_missing(r: r) } + Present { value: block } => match r.max_model_len { Absent => RankStartupRefused { rank: r.rank, missing: v41_rank_reading_missing(r: r) } + Present { value: mml } => match r.max_num_seqs { Absent => RankStartupRefused { rank: r.rank, missing: v41_rank_reading_missing(r: r) } + Present { value: seqs } => match r.max_num_batched_tokens { Absent => RankStartupRefused { rank: r.rank, missing: v41_rank_reading_missing(r: r) } + Present { value: batched } => match r.spec_decode_config { Absent => RankStartupRefused { rank: r.rank, missing: v41_rank_reading_missing(r: r) } + Present { value: spec } => match r.memory { Absent => RankStartupRefused { rank: r.rank, missing: v41_rank_reading_missing(r: r) } + Present { value: mem } => match r.reported { Absent => RankStartupRefused { rank: r.rank, missing: v41_rank_reading_missing(r: r) } + Present { value: rep } => match r.num_gpu_blocks { Absent => RankStartupRefused { rank: r.rank, missing: v41_rank_reading_missing(r: r) } + Present { value: blocks } => match r.kv_cache_groups { Absent => RankStartupRefused { rank: r.rank, missing: v41_rank_reading_missing(r: r) } + Present { value: groups } => match r.kv_format { + ResolvedUnobserved => RankStartupRefused { rank: r.rank, missing: v41_rank_reading_missing(r: r) } + ResolvedFp8DsMla => RankStartupAdmitted { receipt: V41RankStartupReceipt { + run_id: r.run_id, rank: r.rank, host: r.host, image: image, vllm_revision: rev, + parallelism: par, block_size: block, max_model_len: mml, max_num_seqs: seqs, + max_num_batched_tokens: batched, spec_decode_config: spec, memory: mem, reported: rep, + num_gpu_blocks: blocks, kv_cache_groups: groups, + } } + } } } } } } } } } } } } } +} + +// EFFECTIVE BYTES PER TOKEN IS AVAILABLE KV BYTES OVER THE ENGINE'S OWN TOKEN FIGURE. Blocks times +// block size is not a substitute: under MLA with the DeepSeek fp8 layout it is a different quantity +// (extdeps.vllm.kv_cache records that reconstruction as a refusal subject). +fn v41_rank_bytes_per_token(r: V41RankStartupReceipt) -> ByteSize? { + let tokens = token_count_value(t: r.reported.pool.tokens) + if tokens == 0 { none } else { Present { value: byte_size(count: byte_size_count(b: r.memory.available_kv) / tokens) } } +} + +type V41StartupReceipt sole_constructor { + run_id: NonEmptyStr + ranks: List + min_rank: V41RankStartupReceipt +} + +type V41StartupAdmission + = StartupAdmitted { receipt: V41StartupReceipt } + | StartupRefused { causes: List } + +fn v41_image_wire(i: V41OciImageDigest) -> String { + join([render_oci_content_digest_wire(d: i.manifest) as String, "/", render_oci_content_digest_wire(d: i.config) as String], "") +} + +// The min rank is the one with the fewest GPU KV cache tokens: that pool bounds the replica. +fn v41_min_rank(first: V41RankStartupReceipt, rest: List) -> V41RankStartupReceipt { + fold(rest, init: first, f: (acc, r) => + if token_count_value(t: r.reported.pool.tokens) < token_count_value(t: acc.reported.pool.tokens) { r } else { acc }) +} + +fn v41_admit_startup(readings: List) -> V41StartupAdmission { + let admissions = map(readings, r => v41_admit_rank_startup(r: r)) + let refusals = flat_map(admissions, a => match a { + RankStartupAdmitted { receipt: _ } => [] as List + RankStartupRefused { rank: k, missing: m } => [join(["rank ", to_string(k), " missing: ", join(m, ", ")], "")] + }) + let receipts = flat_map(admissions, a => match a { + RankStartupAdmitted { receipt: x } => [x] + RankStartupRefused { rank: _, missing: _ } => [] as List + }) + let absent_ranks = flat_map(v41_measurement_ranks, k => + if length(filter(readings, r => r.rank == k)) == 0 { [join(["rank ", to_string(k), " has no reading"], "")] } else { [] as List }) + let duplicate_ranks = flat_map(v41_measurement_ranks, k => + if length(filter(readings, r => r.rank == k)) > 1 { [join(["rank ", to_string(k), " has more than one reading"], "")] } else { [] as List }) + let foreign_ranks = flat_map(readings, r => + if r.rank >= v41_measurement_rank_count { [join(["rank ", to_string(r.rank), " is outside TP4"], "")] } else { [] as List }) + let shape = concat(concat(concat(refusals, absent_ranks), duplicate_ranks), foreign_ranks) + if length(shape) != 0 { + StartupRefused { causes: shape } + } else { + match list_at_optional(xs: receipts, index: 0) { + Absent => StartupRefused { causes: ["no rank reading was supplied"] } + Present { value: first } => { + let disagreements = flat_map(receipts, r => concat(concat( + if r.run_id != first.run_id { [join(["rank ", to_string(r.rank), " is from another run"], "")] } else { [] as List }, + if v41_image_wire(i: r.image) != v41_image_wire(i: first.image) { [join(["rank ", to_string(r.rank), " ran another image"], "")] } else { [] as List }), + if !observed_revisions_agree(a: r.vllm_revision, b: first.vllm_revision) { [join(["rank ", to_string(r.rank), " ran another vLLM commit"], "")] } else { [] as List })) + if length(disagreements) != 0 { + StartupRefused { causes: disagreements } + } else { + StartupAdmitted { receipt: V41StartupReceipt { + run_id: first.run_id, + ranks: receipts, + min_rank: v41_min_rank(first: first, rest: receipts), + } } + } + } + } + } +} + +fn v41_min_rank_bytes_per_token(s: V41StartupAdmission) -> ByteSize? { + match s { + StartupAdmitted { receipt: r } => v41_rank_bytes_per_token(r: r.min_rank) + StartupRefused { causes: _ } => none + } +} + +// THE TOKEN THRESHOLDS ARE CONCURRENCY TIMES THE CONTEXT CAP, not three authored numbers: x6 at 262k +// is 1,572,864 tokens and x16 is 4,194,304. The production x16 row adds a headroom margin, which is +// operator policy (V41CapacityPolicy.production_pool_margin) and not derived here. +fn v41_pool_tokens_needed(concurrency: Nat) -> Nat { + concurrency * token_count_value(t: v41_measurement_max_model_len) +} + +// ── MEASUREMENT 2: THE ENGRAM RESIDENCY TIMELINE ──────────────────────────────────────────────── +type V41ResidencyPhase = BeforeStart | AfterChecksum | AfterKvAllocation | AfterRequests + +data v41_residency_phases: List = [BeforeStart, AfterChecksum, AfterKvAllocation, AfterRequests] + +fn v41_phase_wire(p: V41ResidencyPhase) -> String { + match p { + BeforeStart => "before-start" + AfterChecksum => "after-checksum" + AfterKvAllocation => "after-kv-allocation" + AfterRequests => "after-requests" + } +} + +type V41TorchMemory { free: ByteSize, reserved: ByteSize, allocated: ByteSize } + +// ONE CELL: one rank at one phase under one image arm. mincore residency is RESIDENCY, never touches. +type V41ResidencyCell { + run_id: NonEmptyStr + arm: V41EngramStartupArm + phase: V41ResidencyPhase + rank: Nat + engram_files: List + memory_current: ByteSize + memory_stat: CgroupMemoryStat + memory_events: CgroupMemoryEvents + memory_pressure: PsiReading + io_pressure: PsiReading + nvme: DiskstatsReading + torch: V41TorchMemory? +} + +type V41ResidencyTimeline + = ResidencyTimelineAdmitted { run_id: NonEmptyStr, arm: V41EngramStartupArm, cells: List } + | ResidencyTimelineRefused { causes: List } + +// BeforeStart has no process, so it has no torch reading; every later phase owes one. +fn v41_cell_owes_torch(p: V41ResidencyPhase) -> Bool { + match p { + BeforeStart => false + AfterChecksum => true + AfterKvAllocation => true + AfterRequests => true + } +} + +// THE TIMELINE IS AN IDENTITY JOIN over (phase x rank). A missing cell is named, a duplicate is +// named, and a count that happens to match is not accepted in place of the join. +fn v41_admit_residency_timeline(run_id: NonEmptyStr, arm: V41EngramStartupArm, cells: List) -> V41ResidencyTimeline { + let owed = flat_map(v41_residency_phases, p => map(v41_measurement_ranks, k => V41ResidencyKey { phase: p, rank: k })) + let join_causes = flat_map(owed, key => { + let matching = filter(cells, c => v41_phase_wire(p: c.phase) == v41_phase_wire(p: key.phase) && c.rank == key.rank) + if length(matching) == 0 { + [join(["no cell for ", v41_phase_wire(p: key.phase), " rank ", to_string(key.rank)], "")] + } else if length(matching) > 1 { + [join(["more than one cell for ", v41_phase_wire(p: key.phase), " rank ", to_string(key.rank)], "")] + } else { [] as List } + }) + let cell_causes = flat_map(cells, c => concat(concat(concat( + if c.run_id != run_id { [join(["a ", v41_phase_wire(p: c.phase), " cell is from another run"], "")] } else { [] as List }, + if v41_arm_wire(a: c.arm) != v41_arm_wire(a: arm) { [join(["a ", v41_phase_wire(p: c.phase), " cell is from the other image arm"], "")] } else { [] as List }), + if c.rank >= v41_measurement_rank_count { [join(["a cell names rank ", to_string(c.rank), ", outside TP4"], "")] } else { [] as List }), + if v41_cell_owes_torch(p: c.phase) && v41_torch_absent(t: c.torch) { [join(["no torch memory reading at ", v41_phase_wire(p: c.phase), " rank ", to_string(c.rank)], "")] } else { [] as List })) + let causes = concat(join_causes, cell_causes) + if length(causes) != 0 { ResidencyTimelineRefused { causes: causes } } else { ResidencyTimelineAdmitted { run_id: run_id, arm: arm, cells: cells } } +} + +type V41ResidencyKey { phase: V41ResidencyPhase, rank: Nat } + +fn v41_torch_absent(t: V41TorchMemory?) -> Bool { + match t { + Present { value: _ } => false + Absent => true + } +} + +fn v41_cell_engram_resident_bytes(c: V41ResidencyCell) -> ByteSize { + byte_size(count: fold(c.engram_files, init: 0, f: (acc, f) => acc + byte_size_count(b: mincore_resident_bytes(r: f)))) +} + +// ── MEASUREMENT 3: THE COLD/WARM ENGRAM PROBE AND THE STAIRCASE ──────────────────────────────── +// +// DISTINCT TOUCHES ARE NOT RESIDENCY. The per-rank distinct-page counter and the fetch_rows timing +// have to come from inside the fetch_rows path. That path is being rewritten by the Engram +// page-release lane, and the counter is a follow-up lane dispatched once that change lands (ruling, +// proud-deer-538, 2026-09-25). Until the counter exists, every touch reading is +// TouchCounterUnavailable, and the Engram-resident capacity refuses on it. mincore residency is +// still recorded, as its own reading. +type V41EngramLookupObservable + = TouchCounterUnavailable { frontier: DissolutionCondition } + | EngramLookupObserved { distinct_pages_touched: Nat, fetch_rows_calls: Nat, fetch_rows_time: Microsecond } + +data v41_touch_counter_frontier: DissolutionCondition = unbound_dissolution(description: "TRIGGER: the follow-up lane after the Engram page-release / batched fetch_rows change exposes, per rank, a distinct-page count (bitmap or sketch, no row contents) and fetch_rows call count and time, readable at an epoch boundary. SUFFICIENT FOR: each probe cell carries EngramLookupObserved, and the Engram-resident capacity is computed from touches rather than refused." as NonEmptyStr) + +type V41ProbeTemperature = ProbeCold | ProbeWarm + +type V41ProbeRankReading { + rank: Nat + lookup: V41EngramLookupObservable + residency_before: List + residency_after: List + stat_before: CgroupMemoryStat + stat_after: CgroupMemoryStat + nvme_before: DiskstatsReading + nvme_after: DiskstatsReading +} + +type V41ProbeReading { + run_id: NonEmptyStr + cell: V41ProbeCell + temperature: V41ProbeTemperature + prompts_disjoint: Bool + prefix_caching_defeated: Bool + latency: ClientBenchLatency + ranks: List +} + +// A WARM PASS WITH NO MAJOR FAULTS ON ANY RANK means the cell's Engram working set was resident. A +// cold pass is the contrast, not the criterion. +type V41ProbeVerdict + = ProbeWorkingSetResident { cell: V41ProbeCell } + | ProbeWorkingSetFaulted { cell: V41ProbeCell, major_faults: Nat } + | ProbeRefused { cell: V41ProbeCell, causes: List } + +fn v41_probe_rank_major_faults(r: V41ProbeRankReading) -> Nat? { + if r.stat_after.pgmajfault < r.stat_before.pgmajfault { none } else { Present { value: r.stat_after.pgmajfault - r.stat_before.pgmajfault } } +} + +fn v41_judge_probe(p: V41ProbeReading) -> V41ProbeVerdict { + let causes = concat(concat(concat(concat( + if !p.prompts_disjoint { ["prompts were not disjoint"] } else { [] as List }, + if !p.prefix_caching_defeated { ["prefix caching was not defeated"] } else { [] as List }), + match p.temperature { ProbeWarm => [] as List ProbeCold => ["a cold pass is the contrast; the verdict is read on the warm pass"] }), + flat_map(v41_measurement_ranks, k => if length(filter(p.ranks, r => r.rank == k)) == 1 { [] as List } else { [join(["rank ", to_string(k), " has no single probe reading"], "")] })), + flat_map(p.ranks, r => match r.lookup { + TouchCounterUnavailable { frontier: _ } => [join(["rank ", to_string(r.rank), ": the distinct-touch counter is unavailable"], "")] + EngramLookupObserved { distinct_pages_touched: _, fetch_rows_calls: _, fetch_rows_time: _ } => [] as List + })) + if length(causes) != 0 { + ProbeRefused { cell: p.cell, causes: causes } + } else { + let faults = map(p.ranks, r => v41_probe_rank_major_faults(r: r)) + if length(filter(faults, f => match f { Present { value: _ } => false Absent => true })) != 0 { + ProbeRefused { cell: p.cell, causes: ["a pgmajfault counter went backwards between readings"] } + } else { + let total = fold(faults, init: 0, f: (acc, f) => match f { Present { value: n } => acc + n Absent => acc }) + if total == 0 { ProbeWorkingSetResident { cell: p.cell } } else { ProbeWorkingSetFaulted { cell: p.cell, major_faults: total } } + } + } +} + +// ── THE STAIRCASE AND ITS STOP RULES ──────────────────────────────────────────────────────────── +// +// THE THRESHOLDS ARE OPERATOR POLICY. No SLO or sustained-window threshold exists in the corpus: TTFT +// is only a selection axis, and the load probe measures TTFT and TPOT without bounds. So the policy +// row is carried with its standing. While it is Proposed, the staircase refuses; it does not run +// against a proposal (ruling, proud-deer-538, 2026-09-25). +type V41CapacityPolicy { + ttft_ceiling: Millisecond + tpot_ceiling: Millisecond + psi_full_ceiling: BasisPoint + psi_window: Microsecond + nvme_read_await_ceiling: Microsecond + refault_growth_ceiling: BasisPoint + production_pool_margin: BasisPoint +} + +type V41PolicyStanding + = PolicyProposed { proposal: V41CapacityPolicy, reasons: NonEmptyStr } + | PolicyOperatorSigned { policy: V41CapacityPolicy, ruling: NonEmptyStr } + +data v41_capacity_policy: V41PolicyStanding = PolicyProposed { + proposal: V41CapacityPolicy { + ttft_ceiling: millisecond(count: 30000), + tpot_ceiling: millisecond(count: 100), + psi_full_ceiling: basis_point(count: 1000), + psi_window: microsecond(count: 60000000), + nvme_read_await_ceiling: microsecond(count: 5000), + refault_growth_ceiling: basis_point(count: 15000), + production_pool_margin: basis_point(count: 1000), + }, + reasons: "PROPOSED by clever-gull-48, awaiting operator sign-off; not policy. ttft 30 s: a cold 262k prefill is a few tens of seconds on one TP4 Spark replica. tpot 100 ms: 10 tok/s per stream is the floor of interactive use. PSI full 10% over a 60 s window: a tenth of wall time with every task stalled on memory is thrash, not a transient. NVMe read await 5 ms: NVMe reads are sub-millisecond when healthy, so 5 ms is queueing. Refault growth 1.5x per concurrency unit between adjacent steps: this is the superlinear test, since linear growth keeps the per-request ratio near 1.0. Production pool margin 10%: this takes x16 at 262k from 4,194,304 to ~4.6M tokens, which matches the brief's production figure." as NonEmptyStr, +} + +type V41StaircaseStepReading { + run_id: NonEmptyStr + concurrency: Nat + preemptions: Nat? + refault_file_before: Nat + refault_file_after: Nat + memory_pressure_before: PsiReading + memory_pressure_after: PsiReading + nvme_before: DiskstatsReading + nvme_after: DiskstatsReading + latency: ClientBenchLatency +} + +type V41StaircaseStop + = StopKvPreemption { concurrency: Nat, preemptions: Nat } + | StopSuperlinearRefault { concurrency: Nat, per_request: Nat, previous_per_request: Nat } + | StopPsiFullSustained { concurrency: Nat, stalled: Microsecond, window: Microsecond } + | StopNvmeLatency { concurrency: Nat, mean_read_await: Microsecond } + | StopSloBreach { concurrency: Nat, which: NonEmptyStr, observed: Millisecond } + | StopReadingMissing { concurrency: Nat, cause: NonEmptyStr } + +type V41StaircaseOutcome + = StaircaseStopped { last_sustained: Nat?, stop: V41StaircaseStop } + | StaircaseCompleted { last_sustained: Nat } + | StaircaseRefused { cause: NonEmptyStr } + +type V41StairScan { outcome: V41StaircaseOutcome?, last: Nat?, previous_per_request: Nat? } + +fn v41_step_stop(policy: V41CapacityPolicy, step: V41StaircaseStepReading, previous_per_request: Nat?) -> V41StaircaseStop? { + let c = step.concurrency + match step.preemptions { + Absent => Present { value: StopReadingMissing { concurrency: c, cause: "the preemption counter was not read" as NonEmptyStr } } + Present { value: pre } => + if pre > 0 { Present { value: StopKvPreemption { concurrency: c, preemptions: pre } } } + else if step.refault_file_after < step.refault_file_before { + Present { value: StopReadingMissing { concurrency: c, cause: "workingset_refault_file went backwards" as NonEmptyStr } } + } else { + let per_request = (step.refault_file_after - step.refault_file_before) / c + let superlinear = match previous_per_request { + Absent => false + Present { value: prev } => prev > 0 && per_request * 10000 > prev * basis_point_count(bp: policy.refault_growth_ceiling) + } + if superlinear { + Present { value: StopSuperlinearRefault { concurrency: c, per_request: per_request, previous_per_request: match previous_per_request { Present { value: p } => p Absent => 0 } } } + } else { + match psi_full_stall_interval(before: step.memory_pressure_before, after: step.memory_pressure_after, window: policy.psi_window) { + PsiStallUnmeasurable { cause: why } => Present { value: StopReadingMissing { concurrency: c, cause: why } } + PsiStallMeasured { stalled: s, window: w } => + if microsecond_count(m: s) * 10000 > microsecond_count(m: w) * basis_point_count(bp: policy.psi_full_ceiling) { + Present { value: StopPsiFullSustained { concurrency: c, stalled: s, window: w } } + } else { + match disk_read_interval(before: step.nvme_before, after: step.nvme_after) { + DiskReadUnmeasurable { cause: why } => Present { value: StopReadingMissing { concurrency: c, cause: why } } + DiskReadMeasured { read_bytes: _, reads: _, mean_read_await: aw } => + match aw { + Present { value: a } => + if microsecond_count(m: a) > microsecond_count(m: policy.nvme_read_await_ceiling) { + Present { value: StopNvmeLatency { concurrency: c, mean_read_await: a } } + } else { v41_slo_stop(policy: policy, c: c, latency: step.latency) } + Absent => v41_slo_stop(policy: policy, c: c, latency: step.latency) + } + } + } + } + } + } + } +} + +fn v41_slo_stop(policy: V41CapacityPolicy, c: Nat, latency: ClientBenchLatency) -> V41StaircaseStop? { + match latency { + ClientBenchLatencyUnread { cause: why } => Present { value: StopReadingMissing { concurrency: c, cause: why } } + ClientBenchLatencyReported { output_throughput: _, tpot_milli: tpot, itl_milli: _, ttft_milli: ttft, window: _ } => + if millisecond_count(m: ttft) > millisecond_count(m: policy.ttft_ceiling) { + Present { value: StopSloBreach { concurrency: c, which: "ttft" as NonEmptyStr, observed: ttft } } + } else if millisecond_count(m: tpot) > millisecond_count(m: policy.tpot_ceiling) { + Present { value: StopSloBreach { concurrency: c, which: "tpot" as NonEmptyStr, observed: tpot } } + } else { none } + } +} + +// THE STAIRCASE RUNS IN PLAN ORDER AND STOPS AT THE FIRST STOP RULE. A step reading that is missing +// for a planned concurrency is a StopReadingMissing, never a skipped step: a staircase with a hole +// cannot claim anything above the hole. +fn v41_run_staircase(policy: V41PolicyStanding, run_id: NonEmptyStr, steps: List) -> V41StaircaseOutcome { + match policy { + PolicyProposed { proposal: _, reasons: _ } => + StaircaseRefused { cause: "the staircase thresholds are an unsigned proposal; operator sign-off is required before any stop rule is applied" as NonEmptyStr } + PolicyOperatorSigned { policy: p, ruling: _ } => { + let scan = fold(v41_staircase_concurrencies, init: V41StairScan { outcome: none, last: none, previous_per_request: none }, f: (acc, c) => + match acc.outcome { + Present { value: _ } => acc + Absent => { + let at = filter(steps, s => s.concurrency == c && s.run_id == run_id) + if length(at) != 1 { + V41StairScan { outcome: Present { value: StaircaseStopped { last_sustained: acc.last, stop: StopReadingMissing { concurrency: c, cause: "no single step reading for this run at this concurrency" as NonEmptyStr } } }, last: acc.last, previous_per_request: acc.previous_per_request } + } else { + match list_at_optional(xs: at, index: 0) { + Absent => acc + Present { value: step } => match v41_step_stop(policy: p, step: step, previous_per_request: acc.previous_per_request) { + Present { value: stop } => V41StairScan { outcome: Present { value: StaircaseStopped { last_sustained: acc.last, stop: stop } }, last: acc.last, previous_per_request: acc.previous_per_request } + Absent => V41StairScan { outcome: none, last: Present { value: c }, previous_per_request: Present { value: (step.refault_file_after - step.refault_file_before) / c } } + } + } + } + } + }) + match scan.outcome { + Present { value: o } => o + Absent => match scan.last { + Present { value: l } => StaircaseCompleted { last_sustained: l } + Absent => StaircaseRefused { cause: "the staircase plan is empty" as NonEmptyStr } + } + } + } + } +} + +// ── THE THREE CAPACITIES, REPORTED SEPARATELY ────────────────────────────────────────────────── +type V41Capacity + = CapacityAt { concurrency: Nat } + | CapacityBelowFirstStep + | CapacityRefused { cause: NonEmptyStr } + +type V41CapacityReport { + run_id: NonEmptyStr + kv_admission: V41Capacity + engram_resident: V41Capacity + sla_sustainable: V41Capacity +} + +// KV ADMISSION: the largest planned concurrency whose max_model_len requests fit the MIN rank's pool. +fn v41_kv_admission_capacity(s: V41StartupAdmission) -> V41Capacity { + match s { + StartupRefused { causes: cs } => CapacityRefused { cause: join(["the startup receipt refused: ", join(cs, "; ")], "") as NonEmptyStr } + StartupAdmitted { receipt: r } => { + let pool = token_count_value(t: r.min_rank.reported.pool.tokens) + let fits = filter(v41_staircase_concurrencies, c => v41_pool_tokens_needed(concurrency: c) <= pool) + match list_at_optional(xs: reverse(fits), index: 0) { + Present { value: c } => CapacityAt { concurrency: c } + Absent => CapacityBelowFirstStep + } + } + } +} + +// ENGRAM RESIDENT: the largest 262k probe concurrency whose warm pass faulted nothing, over a prefix +// of the probe plan with no hole. It refuses on any refused probe, which today means on the absent +// touch counter. +fn v41_engram_resident_capacity(verdicts: List) -> V41Capacity { + let refused = flat_map(verdicts, v => match v { + ProbeRefused { cell: c, causes: cs } => [join(["probe ", to_string(token_count_value(t: c.context)), " x", to_string(c.concurrency), ": ", join(cs, "; ")], "")] + ProbeWorkingSetResident { cell: _ } => [] as List + ProbeWorkingSetFaulted { cell: _, major_faults: _ } => [] as List + }) + let missing = flat_map(v41_probe_cells, cell => + if length(filter(verdicts, v => v41_probe_verdict_cell_is(v: v, cell: cell))) == 1 { [] as List } else { [join(["probe ", to_string(token_count_value(t: cell.context)), " x", to_string(cell.concurrency), " has no single verdict"], "")] }) + let causes = concat(refused, missing) + if length(causes) != 0 { + CapacityRefused { cause: join(causes, "; ") as NonEmptyStr } + } else { + let long = filter(v41_probe_cells, cell => token_count_value(t: cell.context) == token_count_value(t: v41_measurement_max_model_len)) + let scan = fold(long, init: V41ResidentScan { broken: false, best: none }, f: (acc, cell) => + if acc.broken { acc } else { + if length(filter(verdicts, v => match v { ProbeWorkingSetResident { cell: x } => v41_same_cell(a: x, b: cell) ProbeWorkingSetFaulted { cell: _, major_faults: _ } => false ProbeRefused { cell: _, causes: _ } => false })) == 1 { + V41ResidentScan { broken: false, best: Present { value: cell.concurrency } } + } else { V41ResidentScan { broken: true, best: acc.best } } + }) + match scan.best { + Present { value: c } => CapacityAt { concurrency: c } + Absent => CapacityBelowFirstStep + } + } +} + +type V41ResidentScan { broken: Bool, best: Nat? } + +fn v41_same_cell(a: V41ProbeCell, b: V41ProbeCell) -> Bool { + token_count_value(t: a.context) == token_count_value(t: b.context) && a.concurrency == b.concurrency +} + +fn v41_probe_verdict_cell_is(v: V41ProbeVerdict, cell: V41ProbeCell) -> Bool { + match v { + ProbeWorkingSetResident { cell: c } => v41_same_cell(a: c, b: cell) + ProbeWorkingSetFaulted { cell: c, major_faults: _ } => v41_same_cell(a: c, b: cell) + ProbeRefused { cell: c, causes: _ } => v41_same_cell(a: c, b: cell) + } +} + +// SLA SUSTAINABLE: the last staircase step before the first stop. +fn v41_sla_capacity(o: V41StaircaseOutcome) -> V41Capacity { + match o { + StaircaseRefused { cause: c } => CapacityRefused { cause: c } + StaircaseCompleted { last_sustained: l } => CapacityAt { concurrency: l } + StaircaseStopped { last_sustained: l, stop: s } => match s { + StopReadingMissing { concurrency: c, cause: why } => CapacityRefused { cause: join(["a reading is missing at x", to_string(c), ": ", why as String], "") as NonEmptyStr } + StopKvPreemption { concurrency: _, preemptions: _ } => v41_last_or_below(l: l) + StopSuperlinearRefault { concurrency: _, per_request: _, previous_per_request: _ } => v41_last_or_below(l: l) + StopPsiFullSustained { concurrency: _, stalled: _, window: _ } => v41_last_or_below(l: l) + StopNvmeLatency { concurrency: _, mean_read_await: _ } => v41_last_or_below(l: l) + StopSloBreach { concurrency: _, which: _, observed: _ } => v41_last_or_below(l: l) + } + } +} + +fn v41_last_or_below(l: Nat?) -> V41Capacity { + match l { + Present { value: c } => CapacityAt { concurrency: c } + Absent => CapacityBelowFirstStep + } +} + +fn v41_capacity_report( + run_id: NonEmptyStr, + startup: V41StartupAdmission, + probes: List, + staircase: V41StaircaseOutcome, +) -> V41CapacityReport { + V41CapacityReport { + run_id: run_id, + kv_admission: v41_kv_admission_capacity(s: startup), + engram_resident: v41_engram_resident_capacity(verdicts: map(probes, p => v41_judge_probe(p: p))), + sla_sustainable: v41_sla_capacity(o: staircase), + } +} + +fn v41_capacity_wire(c: V41Capacity) -> String { + match c { + CapacityAt { concurrency: n } => join(["x", to_string(n)], "") + CapacityBelowFirstStep => "below-x1" + CapacityRefused { cause: why } => join(["refused: ", why as String], "") + } +} + +// ── THE PRODUCTION GATE, WHICH DOES NOT GATE THE MEASUREMENT SERVE ───────────────────────────── +// +// The measurement serve is admitted with a no-production-traffic contract: it is what produces the +// report, so gating it on the report would be circular. Production admission requires all three +// capacities to be CapacityAt, and it admits the MINIMUM of the three. The larger two are headroom +// on axes that are not the binding one. +type V41ServingRealization + = V41MeasurementServe { run_id: NonEmptyStr } + | V41ProductionServe + +type V41ServingAdmission + = MeasurementServeAdmitted { run_id: NonEmptyStr } + | ProductionAdmitted { run_id: NonEmptyStr, concurrency: Nat } + | ProductionRefused { cause: NonEmptyStr } + +fn v41_capacity_at(c: V41Capacity) -> Nat? { + match c { + CapacityAt { concurrency: n } => Present { value: n } + CapacityBelowFirstStep => none + CapacityRefused { cause: _ } => none + } +} + +fn v41_serving_admission(realization: V41ServingRealization, report: V41CapacityReport?) -> V41ServingAdmission { + match realization { + V41MeasurementServe { run_id: r } => MeasurementServeAdmitted { run_id: r } + V41ProductionServe => match report { + Absent => ProductionRefused { cause: "no V4.1 capacity report has been admitted; the first serving action is the measurement" as NonEmptyStr } + Present { value: rep } => match v41_capacity_at(c: rep.kv_admission) { + Absent => ProductionRefused { cause: join(["KV admission: ", v41_capacity_wire(c: rep.kv_admission)], "") as NonEmptyStr } + Present { value: kv } => match v41_capacity_at(c: rep.engram_resident) { + Absent => ProductionRefused { cause: join(["Engram resident: ", v41_capacity_wire(c: rep.engram_resident)], "") as NonEmptyStr } + Present { value: en } => match v41_capacity_at(c: rep.sla_sustainable) { + Absent => ProductionRefused { cause: join(["SLA sustainable: ", v41_capacity_wire(c: rep.sla_sustainable)], "") as NonEmptyStr } + Present { value: sla } => ProductionAdmitted { run_id: rep.run_id, concurrency: nat_min(a: nat_min(a: kv, b: en), b: sla) } + } + } + } + } + } +} + +// ── THE DRY RUN ──────────────────────────────────────────────────────────────────────────────── +// +// The plan the wet runner will execute, rendered without touching a host: the startup receipt from +// every rank, the residency timeline for each bound image arm, the probe cells, the staircase, and +// the standing of every precondition that is not met yet. The wet runner is a fleet mode that lands +// after the srv8 runtime image exists (v41_wet_runner_frontier), and it consumes this same plan. +data v41_wet_runner_frontier: DissolutionCondition = unbound_dissolution(description: "TRIGGER: the srv8 V4.1 runtime image exists (after #12294) and a Group A fleet mode (a gunbc.ci.ci_spec GunbcRunStepTarget) brings up the measurement serve and reads the startup log lines, cgroup files, PSI, diskstats, mincore residency, torch memory, and vLLM metrics on every rank, producing the readings this module admits. SUFFICIENT FOR: v41_capacity_report is computed from readings taken on srv5-8 under one run id." as NonEmptyStr) + +fn v41_capacity_measurement_plan_text() -> String { + let arms = if length(v41_engram_arm_bindings) == 0 { + "engram arms: unbound (no image patch digest bound yet)" + } else { + join(["engram arms: ", join(map(v41_engram_arm_bindings, b => v41_arm_wire(a: b.arm)), ", ")], "") + } + let policy = match v41_capacity_policy { + PolicyProposed { proposal: _, reasons: _ } => "staircase policy: PROPOSED, awaiting operator sign-off; the staircase refuses" + PolicyOperatorSigned { policy: _, ruling: r } => join(["staircase policy: signed (", r as String, ")"], "") + } + join([ + join(["V4.1 capacity measurement: Group A TP4 x1, ", to_string(v41_measurement_rank_count), " ranks, max_model_len ", to_string(token_count_value(t: v41_measurement_max_model_len))], ""), + "measurement 1: startup memory receipt from every rank, judged by the min rank", + join([" thresholds: x6 needs ", to_string(v41_pool_tokens_needed(concurrency: 6)), " tokens, x16 needs ", to_string(v41_pool_tokens_needed(concurrency: 16))], ""), + join(["measurement 2: residency timeline at ", join(map(v41_residency_phases, p => v41_phase_wire(p: p)), ", ")], ""), + join([" ", arms], ""), + join(["measurement 3: probe cells ", join(map(v41_probe_cells, c => join([to_string(token_count_value(t: c.context)), "x", to_string(c.concurrency)], "")), ", ")], ""), + " distinct touches: TouchCounterUnavailable until the follow-up counter lane lands; Engram-resident refuses", + join([" staircase: ", join(map(v41_staircase_concurrencies, c => join(["x", to_string(c)], "")), " -> ")], ""), + join([" ", policy], ""), + "production admission: refused until a capacity report is admitted; the measurement serve is admitted on its own", + ], "\n") +} diff --git a/dag/test/claim/spark/v41_capacity_measurement_witness_test.dag b/dag/test/claim/spark/v41_capacity_measurement_witness_test.dag new file mode 100644 index 00000000000..9305894df4d --- /dev/null +++ b/dag/test/claim/spark/v41_capacity_measurement_witness_test.dag @@ -0,0 +1,446 @@ +module test.claim.spark.v41_capacity_measurement_witness + +import v2.std.algebra { filter } +import v2.std.collection { list_at_optional } +import std.types { Bool, Int, List, NonEmptyStr, String } +import std.nat { Nat } +import std.measure { + ByteSize, byte_size, byte_size_count, token_count, millisecond, microsecond, microsecond_count, + milli_tokens_per_second, concurrent_request_hundredths, basis_point, +} +import std.content_hash { Sha256Digest, Sha256DigestHex } +import v2.std.optional { Present, Absent } +import v2.std.live_tree { LiveTreeDisposition, SubstrateInputsOnly } +import extdeps.container.oci.digest { OciSha256Digest } +import extdeps.vllm.server { VllmSourceRevision, vllm_source_revision } +import extdeps.vllm.log_lines { vllm_startup_reading, GpuKvCacheSize, MaximumConcurrency } +import extdeps.vllm.kv_cache { VllmObservedKvPool, KvPoolAdmitted, KvPoolRefused, vllm_observed_kv_pool_from_reading, ResolvedFp8DsMla, ResolvedUnobserved } +import extdeps.vllm.kv_layout { VllmReportedCapacity, ReportedCapacityAdmitted, ReportedCapacityRefused, vllm_reported_capacity } +import extdeps.linux.psi { PsiReading, PsiRead, PsiUnparseable, parse_psi, PsiStallMeasured, PsiStallUnmeasurable, psi_full_stall_interval } +import extdeps.linux.diskstats { parse_diskstats, DiskstatsDeviceRead, DiskstatsDeviceAbsent, DiskReadMeasured, DiskReadUnmeasurable, disk_read_interval } +import extdeps.linux.mincore { MincoreFileResidency, MincoreResidencyAdmitted, MincoreResidencyRefused, mincore_file_residency } +import extdeps.linux.cgroup_v2_memory { CgroupMemoryStat, CgroupMemoryStatRead, CgroupMemoryStatIncomplete, parse_cgroup_memory_stat } +import gunbc.floor_memory_demand { CgroupMemoryEvents } +import gunbc.spark.v41_runtime_realized { V41OciImageDigest, v41_oci_image_digest } +import gunbc.spark.serving_load_probe { ClientBenchLatency, ClientBenchLatencyReported, ForeignTrafficNotObserved, DrainObservedIdle } +import gunbc.spark.v41_capacity_measurement { + V41RankStartupReading, V41RankParallelism, V41RankMemoryBreakdown, + V41StartupAdmission, StartupAdmitted, StartupRefused, v41_admit_startup, v41_min_rank_bytes_per_token, + v41_kv_bytes_per_token_bet, V41BetVerdict, BetConfirmed, BetFalsified, BetUnreadable, v41_judge_byte_range_bet, + V41ResidencyCell, V41ResidencyTimeline, ResidencyTimelineAdmitted, ResidencyTimelineRefused, v41_admit_residency_timeline, + EngramCacheRetained, EngramCacheReleased, BeforeStart, AfterChecksum, AfterKvAllocation, AfterRequests, V41TorchMemory, + V41ProbeCell, V41ProbeReading, V41ProbeRankReading, ProbeWarm, TouchCounterUnavailable, EngramLookupObserved, + v41_touch_counter_frontier, v41_judge_probe, ProbeRefused, ProbeWorkingSetResident, ProbeWorkingSetFaulted, + V41CapacityPolicy, PolicyOperatorSigned, v41_capacity_policy, V41StaircaseStepReading, v41_run_staircase, + V41PolicyStanding, StaircaseStopped, StaircaseCompleted, StaircaseRefused, StopKvPreemption, StopReadingMissing, + StopSuperlinearRefault, StopPsiFullSustained, StopNvmeLatency, StopSloBreach, V41ResidencyPhase, + V41Capacity, CapacityAt, CapacityBelowFirstStep, CapacityRefused, v41_capacity_wire, + v41_kv_admission_capacity, v41_engram_resident_capacity, v41_sla_capacity, v41_capacity_report, + V41CapacityReport, V41MeasurementServe, V41ProductionServe, v41_serving_admission, + MeasurementServeAdmitted, ProductionAdmitted, ProductionRefused, + v41_capacity_measurement_plan_text, v41_measurement_ranks, + v41_engram_resident_bets, v41_judge_engram_resident_bet, v41_after_requests_resident, PointBetObserved, PointBetUnreadable, +} + +data live_tree_disposition: LiveTreeDisposition = SubstrateInputsOnly + +// THE FOLDS AT THEIR OWN INTERFACE, OVER SUPPLIED READINGS. The producer of real readings is the wet +// fleet mode (gunbc.spark.v41_capacity_measurement v41_wet_runner_frontier), which cannot run until +// the srv8 runtime image exists; its inhabitance claim lands with it. + +fn w_rev() -> VllmSourceRevision { + vllm_source_revision(repository: "github.com/vllm-project/vllm" as NonEmptyStr, full_revision: "1111111111111111111111111111111111111111" as NonEmptyStr) +} + +fn w_rev_other() -> VllmSourceRevision { + vllm_source_revision(repository: "github.com/vllm-project/vllm" as NonEmptyStr, full_revision: "2222222222222222222222222222222222222222" as NonEmptyStr) +} + +fn w_reported(tokens: Int, rev: VllmSourceRevision) -> VllmReportedCapacity? { + match vllm_observed_kv_pool_from_reading(tokens: token_count(count: tokens), reading: vllm_startup_reading(line: GpuKvCacheSize, captured_text: "GPU KV cache size: n tokens" as NonEmptyStr, observed_at_revision: rev)) { + KvPoolRefused { cause: _ } => none + KvPoolAdmitted { pool: p } => match vllm_reported_capacity(pool: p, max_concurrency: concurrent_request_hundredths(count: (tokens * 100) / 262144), reading: vllm_startup_reading(line: MaximumConcurrency, captured_text: "Maximum concurrency for 262,144 tokens per request: n" as NonEmptyStr, observed_at_revision: rev)) { + ReportedCapacityAdmitted { capacity: c } => Present { value: c } + ReportedCapacityRefused { cause: _ } => none + } + } +} + +fn w_image() -> V41OciImageDigest { + v41_oci_image_digest( + manifest: OciSha256Digest(Sha256Digest { hex: "3333333333333333333333333333333333333333333333333333333333333333" as Sha256DigestHex }), + config: OciSha256Digest(Sha256Digest { hex: "4444444444444444444444444444444444444444444444444444444444444444" as Sha256DigestHex }), + ) +} + +// 1,850 bytes per token over the pool, so the bet (1,800-1,940) is confirmed on the min rank. +fn w_rank(rank: Nat, tokens: Int, rev: VllmSourceRevision) -> V41RankStartupReading { + V41RankStartupReading { + run_id: "run-1" as NonEmptyStr, + rank: rank, + host: "srv5" as NonEmptyStr, + image: Present { value: w_image() }, + vllm_revision: Present { value: rev }, + parallelism: Present { value: V41RankParallelism { tp: 4, dcp: 1, pp: 1, dp: 1 } }, + kv_format: ResolvedFp8DsMla, + block_size: Present { value: 64 }, + max_model_len: Present { value: token_count(count: 262144) }, + max_num_seqs: Present { value: 16 }, + max_num_batched_tokens: Present { value: 8192 }, + spec_decode_config: Present { value: "off" as NonEmptyStr }, + memory: Present { value: V41RankMemoryBreakdown { + free: byte_size(count: 100000000000), total: byte_size(count: 130663231488), + weights: byte_size(count: 40000000000), non_torch: byte_size(count: 2000000000), + activation_peak: byte_size(count: 3000000000), cuda_graph: byte_size(count: 1000000000), + available_kv: byte_size(count: tokens * 1850), + } }, + reported: w_reported(tokens: tokens, rev: rev), + num_gpu_blocks: Present { value: 1000 }, + kv_cache_groups: Present { value: 2 }, + } +} + +// The same reading with the memory breakdown missing and the kv dtype unresolved. +fn w_rank_broken(rank: Nat) -> V41RankStartupReading { + let r = w_rank(rank: rank, tokens: 2000000, rev: w_rev()) + V41RankStartupReading { + run_id: r.run_id, rank: r.rank, host: r.host, image: r.image, vllm_revision: r.vllm_revision, + parallelism: r.parallelism, kv_format: ResolvedUnobserved, block_size: r.block_size, + max_model_len: r.max_model_len, max_num_seqs: r.max_num_seqs, max_num_batched_tokens: r.max_num_batched_tokens, + spec_decode_config: r.spec_decode_config, memory: none, reported: r.reported, + num_gpu_blocks: r.num_gpu_blocks, kv_cache_groups: r.kv_cache_groups, + } +} + +// On the MIN rank (1,500,000) x4 fits and x6 (1,572,864) does not; the other ranks alone would admit x6. +fn w_four_ranks() -> List { + [w_rank(rank: 0, tokens: 2000000, rev: w_rev()), w_rank(rank: 1, tokens: 1500000, rev: w_rev()), + w_rank(rank: 2, tokens: 2000000, rev: w_rev()), w_rank(rank: 3, tokens: 2000000, rev: w_rev())] +} + +test fn the_startup_receipt_is_judged_by_the_min_rank() -> Bool { + v41_capacity_wire(c: v41_kv_admission_capacity(s: v41_admit_startup(readings: w_four_ranks()))) == "x4" +} + +test fn effective_bytes_per_token_confirms_the_bet_in_range() -> Bool { + match v41_judge_byte_range_bet(bet: v41_kv_bytes_per_token_bet, observed: v41_min_rank_bytes_per_token(s: v41_admit_startup(readings: w_four_ranks()))) { + BetConfirmed { key: _, observed: o } => byte_size_count(b: o) == 1850 + BetFalsified { key: _, observed: _ } => false + BetUnreadable { key: _, cause: _ } => false + } +} + +test fn a_reading_outside_the_bet_falsifies_it_and_a_missing_one_is_unreadable() -> Bool { + let falsified = match v41_judge_byte_range_bet(bet: v41_kv_bytes_per_token_bet, observed: Present { value: byte_size(count: 2400) }) { + BetFalsified { key: _, observed: _ } => true + BetConfirmed { key: _, observed: _ } => false + BetUnreadable { key: _, cause: _ } => false + } + let unreadable = match v41_judge_byte_range_bet(bet: v41_kv_bytes_per_token_bet, observed: Absent) { + BetUnreadable { key: _, cause: _ } => true + BetConfirmed { key: _, observed: _ } => false + BetFalsified { key: _, observed: _ } => false + } + falsified && unreadable +} + +// THE REDS: a missing rank, a missing field, an unresolved kv dtype, and a rank on another commit each +// refuse the whole receipt, and nothing about KV admission is computed. +fn w_refused(s: V41StartupAdmission, needle: String) -> Bool { + match s { + StartupRefused { causes: cs } => length(filter(cs, c => contains(c, needle))) > 0 + StartupAdmitted { receipt: _ } => false + } +} + +test fn a_missing_rank_refuses_the_startup_receipt() -> Bool { + let three = filter(w_four_ranks(), r => r.rank != 3) + w_refused(s: v41_admit_startup(readings: three), needle: "rank 3 has no reading") + && contains(v41_capacity_wire(c: v41_kv_admission_capacity(s: v41_admit_startup(readings: three))), "refused") +} + +test fn a_missing_field_is_named_in_the_refusal() -> Bool { + let broken = w_rank_broken(rank: 2) + let readings = [w_rank(rank: 0, tokens: 2000000, rev: w_rev()), w_rank(rank: 1, tokens: 2000000, rev: w_rev()), broken, w_rank(rank: 3, tokens: 2000000, rev: w_rev())] + w_refused(s: v41_admit_startup(readings: readings), needle: "memory breakdown") + && w_refused(s: v41_admit_startup(readings: readings), needle: "resolved kv dtype") +} + +test fn a_rank_on_another_vllm_commit_refuses() -> Bool { + let readings = [w_rank(rank: 0, tokens: 2000000, rev: w_rev()), w_rank(rank: 1, tokens: 2000000, rev: w_rev_other()), w_rank(rank: 2, tokens: 2000000, rev: w_rev()), w_rank(rank: 3, tokens: 2000000, rev: w_rev())] + w_refused(s: v41_admit_startup(readings: readings), needle: "another vLLM commit") +} + +// ── the kernel readers ── + +data w_psi_before: String = "some avg10=0.00 avg60=0.00 avg300=0.00 total=1000\nfull avg10=0.00 avg60=0.00 avg300=0.00 total=500\n" +data w_psi_after: String = "some avg10=9.00 avg60=3.00 avg300=1.00 total=9000000\nfull avg10=8.00 avg60=2.00 avg300=1.00 total=7000500\n" + +test fn a_psi_full_stall_is_the_delta_of_two_totals() -> Bool { + match psi_full_stall_interval(before: parse_psi(body: w_psi_before), after: parse_psi(body: w_psi_after), window: microsecond(count: 60000000)) { + PsiStallMeasured { stalled: s, window: _ } => microsecond_count(m: s) == 7000000 + PsiStallUnmeasurable { cause: _ } => false + } +} + +test fn a_psi_total_that_went_backwards_is_unmeasurable() -> Bool { + match psi_full_stall_interval(before: parse_psi(body: w_psi_after), after: parse_psi(body: w_psi_before), window: microsecond(count: 60000000)) { + PsiStallUnmeasurable { cause: _ } => true + PsiStallMeasured { stalled: _, window: _ } => false + } +} + +data w_disk_before: String = " 8 0 sda 1 0 8 1 0 0 0 0 0 0 0\n 259 0 nvme0n1 1000 0 80000 500 10 0 80 5 0 400 505 0 0 0 0\n" +data w_disk_after: String = " 259 0 nvme0n1 3000 0 4080000 2500 10 0 80 5 4 1400 2505 0 0 0 0\n" + +test fn a_disk_read_interval_counts_512_byte_sectors_and_mean_await() -> Bool { + match disk_read_interval(before: parse_diskstats(body: w_disk_before, device: "nvme0n1" as NonEmptyStr), after: parse_diskstats(body: w_disk_after, device: "nvme0n1" as NonEmptyStr)) { + DiskReadMeasured { read_bytes: b, reads: n, mean_read_await: aw } => + byte_size_count(b: b) == 4000000 * 512 && n == 2000 + && match aw { Present { value: a } => microsecond_count(m: a) == 1000 Absent => false } + DiskReadUnmeasurable { cause: _ } => false + } +} + +test fn an_absent_device_is_not_a_zero_read() -> Bool { + match disk_read_interval(before: parse_diskstats(body: w_disk_before, device: "nvme1n1" as NonEmptyStr), after: parse_diskstats(body: w_disk_after, device: "nvme1n1" as NonEmptyStr)) { + DiskReadUnmeasurable { cause: _ } => true + DiskReadMeasured { read_bytes: _, reads: _, mean_read_await: _ } => false + } +} + +test fn mincore_refuses_more_resident_pages_than_the_file_has() -> Bool { + let ok = match mincore_file_residency(file: "rank0.engram" as NonEmptyStr, file_bytes: byte_size(count: 8193), page_bytes: byte_size(count: 4096), resident_pages: 3) { + MincoreResidencyAdmitted { residency: _ } => true + MincoreResidencyRefused { file: _, cause: _ } => false + } + let red = match mincore_file_residency(file: "rank0.engram" as NonEmptyStr, file_bytes: byte_size(count: 8193), page_bytes: byte_size(count: 4096), resident_pages: 4) { + MincoreResidencyRefused { file: _, cause: _ } => true + MincoreResidencyAdmitted { residency: _ } => false + } + ok && red +} + +test fn memory_stat_reads_the_five_keys_and_names_a_missing_one() -> Bool { + let full = "anon 100\nfile 200\nkernel 5\nfile_mapped 50\nworkingset_refault_file 7\npgmajfault 9\n" + let read = match parse_cgroup_memory_stat(body: full) { + CgroupMemoryStatRead { stat: s } => byte_size_count(b: s.file) == 200 && s.pgmajfault == 9 && s.workingset_refault_file == 7 + CgroupMemoryStatIncomplete { missing: _ } => false + } + let red = match parse_cgroup_memory_stat(body: "anon 100\nfile 200\nfile_mapped 50\npgmajfault 9\n") { + CgroupMemoryStatIncomplete { missing: m } => m == ["workingset_refault_file"] + CgroupMemoryStatRead { stat: _ } => false + } + read && red +} + +// ── measurement 2 ── + +fn w_stat(refault: Nat, majfault: Nat) -> CgroupMemoryStat { + CgroupMemoryStat { anon: byte_size(count: 1), file: byte_size(count: 2), file_mapped: byte_size(count: 1), workingset_refault_file: refault, pgmajfault: majfault } +} + +fn w_cell(phase: V41ResidencyPhase, rank: Nat) -> V41ResidencyCell { + V41ResidencyCell { + run_id: "run-1" as NonEmptyStr, + arm: EngramCacheRetained, + phase: phase, + rank: rank, + engram_files: [] as List, + memory_current: byte_size(count: 1), + memory_stat: w_stat(refault: 0, majfault: 0), + memory_events: CgroupMemoryEvents { high_events: 0, max_events: 0, oom_kills: 0 }, + memory_pressure: parse_psi(body: w_psi_before), + io_pressure: parse_psi(body: w_psi_before), + nvme: parse_diskstats(body: w_disk_before, device: "nvme0n1" as NonEmptyStr), + torch: Present { value: V41TorchMemory { free: byte_size(count: 1), reserved: byte_size(count: 1), allocated: byte_size(count: 1) } }, + } +} + +fn w_all_cells() -> List { + flat_map([BeforeStart, AfterChecksum, AfterKvAllocation, AfterRequests], p => map(v41_measurement_ranks, k => w_cell(phase: p, rank: k))) +} + +test fn a_complete_residency_timeline_is_admitted_and_a_missing_cell_is_named() -> Bool { + let ok = match v41_admit_residency_timeline(run_id: "run-1" as NonEmptyStr, arm: EngramCacheRetained, cells: w_all_cells()) { + ResidencyTimelineAdmitted { run_id: _, arm: _, cells: _ } => true + ResidencyTimelineRefused { causes: _ } => false + } + let holed = filter(w_all_cells(), c => !(c.rank == 2 && match c.phase { AfterChecksum => true BeforeStart => false AfterKvAllocation => false AfterRequests => false })) + let red = match v41_admit_residency_timeline(run_id: "run-1" as NonEmptyStr, arm: EngramCacheRetained, cells: holed) { + ResidencyTimelineRefused { causes: cs } => cs == ["no cell for after-checksum rank 2"] + ResidencyTimelineAdmitted { run_id: _, arm: _, cells: _ } => false + } + ok && red +} + +test fn a_cell_from_the_other_image_arm_refuses_the_timeline() -> Bool { + match v41_admit_residency_timeline(run_id: "run-1" as NonEmptyStr, arm: EngramCacheReleased, cells: w_all_cells()) { + ResidencyTimelineRefused { causes: cs } => length(cs) == 16 + ResidencyTimelineAdmitted { run_id: _, arm: _, cells: _ } => false + } +} + +// ── measurement 3 ── + +fn w_probe_rank(rank: Nat, observed: Bool, majfault_after: Nat) -> V41ProbeRankReading { + V41ProbeRankReading { + rank: rank, + lookup: if observed { EngramLookupObserved { distinct_pages_touched: 10, fetch_rows_calls: 3, fetch_rows_time: microsecond(count: 9) } } else { TouchCounterUnavailable { frontier: v41_touch_counter_frontier } }, + residency_before: [] as List, + residency_after: [] as List, + stat_before: w_stat(refault: 0, majfault: 5), + stat_after: w_stat(refault: 0, majfault: majfault_after), + nvme_before: parse_diskstats(body: w_disk_before, device: "nvme0n1" as NonEmptyStr), + nvme_after: parse_diskstats(body: w_disk_after, device: "nvme0n1" as NonEmptyStr), + } +} + +fn w_latency(ttft: Nat) -> ClientBenchLatency { + ClientBenchLatencyReported { output_throughput: milli_tokens_per_second(count: 1000), tpot_milli: millisecond(count: 50), itl_milli: none, ttft_milli: millisecond(count: ttft), window: ForeignTrafficNotObserved { drain: DrainObservedIdle } } +} + +fn w_probe(context: Nat, concurrency: Nat, observed: Bool, majfault_after: Nat) -> V41ProbeReading { + V41ProbeReading { + run_id: "run-1" as NonEmptyStr, + cell: V41ProbeCell { context: token_count(count: context), concurrency: concurrency }, + temperature: ProbeWarm, + prompts_disjoint: true, + prefix_caching_defeated: true, + latency: w_latency(ttft: 1000), + ranks: map(v41_measurement_ranks, k => w_probe_rank(rank: k, observed: observed, majfault_after: majfault_after)), + } +} + +// TODAY'S STATE: no touch counter, so every probe refuses and the Engram-resident capacity refuses, +// even though mincore residency and fault counters are present. +test fn engram_resident_refuses_while_the_touch_counter_is_unavailable() -> Bool { + let probes = [w_probe(context: 32768, concurrency: 1, observed: false, majfault_after: 5), w_probe(context: 262144, concurrency: 1, observed: false, majfault_after: 5), + w_probe(context: 262144, concurrency: 2, observed: false, majfault_after: 5), w_probe(context: 262144, concurrency: 4, observed: false, majfault_after: 5), + w_probe(context: 262144, concurrency: 6, observed: false, majfault_after: 5)] + let c = v41_engram_resident_capacity(verdicts: map(probes, p => v41_judge_probe(p: p))) + contains(v41_capacity_wire(c: c), "distinct-touch counter is unavailable") +} + +// ONCE THE COUNTER EXISTS: the largest 262k concurrency whose warm pass faulted nothing, stopping at +// the first faulting cell. +test fn engram_resident_is_the_last_warm_cell_without_major_faults() -> Bool { + let probes = [w_probe(context: 32768, concurrency: 1, observed: true, majfault_after: 5), w_probe(context: 262144, concurrency: 1, observed: true, majfault_after: 5), + w_probe(context: 262144, concurrency: 2, observed: true, majfault_after: 5), w_probe(context: 262144, concurrency: 4, observed: true, majfault_after: 90), + w_probe(context: 262144, concurrency: 6, observed: true, majfault_after: 5)] + v41_capacity_wire(c: v41_engram_resident_capacity(verdicts: map(probes, p => v41_judge_probe(p: p)))) == "x2" +} + +// ── the staircase ── + +fn w_signed() -> V41PolicyStanding { + PolicyOperatorSigned { + policy: V41CapacityPolicy { + ttft_ceiling: millisecond(count: 30000), tpot_ceiling: millisecond(count: 100), + psi_full_ceiling: basis_point(count: 1000), psi_window: microsecond(count: 60000000), + nvme_read_await_ceiling: microsecond(count: 5000), refault_growth_ceiling: basis_point(count: 15000), + production_pool_margin: basis_point(count: 1000), + }, + ruling: "witness fixture, not the operator's policy" as NonEmptyStr, + } +} + +fn w_step(c: Nat, preemptions: Nat) -> V41StaircaseStepReading { + V41StaircaseStepReading { + run_id: "run-1" as NonEmptyStr, + concurrency: c, + preemptions: Present { value: preemptions }, + refault_file_before: 0, + refault_file_after: c * 100, + memory_pressure_before: parse_psi(body: w_psi_before), + memory_pressure_after: parse_psi(body: "some avg10=0.00 avg60=0.00 avg300=0.00 total=2000\nfull avg10=0.00 avg60=0.00 avg300=0.00 total=600\n"), + nvme_before: parse_diskstats(body: w_disk_before, device: "nvme0n1" as NonEmptyStr), + nvme_after: parse_diskstats(body: w_disk_after, device: "nvme0n1" as NonEmptyStr), + latency: w_latency(ttft: 1000), + } +} + +test fn the_staircase_refuses_while_the_policy_is_an_unsigned_proposal() -> Bool { + match v41_run_staircase(policy: v41_capacity_policy, run_id: "run-1" as NonEmptyStr, steps: map([1, 2, 4, 6, 8, 12, 16], c => w_step(c: c, preemptions: 0))) { + StaircaseRefused { cause: _ } => true + StaircaseCompleted { last_sustained: _ } => false + StaircaseStopped { last_sustained: _, stop: _ } => false + } +} + +test fn the_staircase_stops_at_the_first_kv_preemption() -> Bool { + let steps = [w_step(c: 1, preemptions: 0), w_step(c: 2, preemptions: 0), w_step(c: 4, preemptions: 0), w_step(c: 6, preemptions: 0), w_step(c: 8, preemptions: 3), w_step(c: 12, preemptions: 9)] + let o = v41_run_staircase(policy: w_signed(), run_id: "run-1" as NonEmptyStr, steps: steps) + let stopped_at_8 = match o { + StaircaseStopped { last_sustained: _, stop: s } => match s { StopKvPreemption { concurrency: c, preemptions: p } => c == 8 && p == 3 StopReadingMissing { concurrency: _, cause: _ } => false StopSuperlinearRefault { concurrency: _, per_request: _, previous_per_request: _ } => false StopPsiFullSustained { concurrency: _, stalled: _, window: _ } => false StopNvmeLatency { concurrency: _, mean_read_await: _ } => false StopSloBreach { concurrency: _, which: _, observed: _ } => false } + StaircaseCompleted { last_sustained: _ } => false + StaircaseRefused { cause: _ } => false + } + stopped_at_8 && v41_capacity_wire(c: v41_sla_capacity(o: o)) == "x6" +} + +test fn a_hole_in_the_staircase_refuses_the_sla_capacity_above_it() -> Bool { + let steps = [w_step(c: 1, preemptions: 0), w_step(c: 4, preemptions: 0)] + contains(v41_capacity_wire(c: v41_sla_capacity(o: v41_run_staircase(policy: w_signed(), run_id: "run-1" as NonEmptyStr, steps: steps))), "a reading is missing at x2") +} + +// ── the production gate ── + +test fn the_measurement_serve_is_admitted_and_production_refuses_without_a_report() -> Bool { + let m = match v41_serving_admission(realization: V41MeasurementServe { run_id: "run-1" as NonEmptyStr }, report: none) { + MeasurementServeAdmitted { run_id: _ } => true + ProductionAdmitted { run_id: _, concurrency: _ } => false + ProductionRefused { cause: _ } => false + } + let p = match v41_serving_admission(realization: V41ProductionServe, report: none) { + ProductionRefused { cause: _ } => true + MeasurementServeAdmitted { run_id: _ } => false + ProductionAdmitted { run_id: _, concurrency: _ } => false + } + m && p +} + +test fn production_admits_the_minimum_of_three_capacities_and_refuses_on_any_refusal() -> Bool { + let all = V41CapacityReport { run_id: "run-1" as NonEmptyStr, kv_admission: CapacityAt { concurrency: 6 }, engram_resident: CapacityAt { concurrency: 4 }, sla_sustainable: CapacityAt { concurrency: 8 } } + let admitted = match v41_serving_admission(realization: V41ProductionServe, report: Present { value: all }) { + ProductionAdmitted { run_id: _, concurrency: c } => c == 4 + ProductionRefused { cause: _ } => false + MeasurementServeAdmitted { run_id: _ } => false + } + let today = V41CapacityReport { run_id: "run-1" as NonEmptyStr, kv_admission: CapacityAt { concurrency: 6 }, sla_sustainable: CapacityAt { concurrency: 8 }, engram_resident: CapacityRefused { cause: "the distinct-touch counter is unavailable" as NonEmptyStr } } + let refused = match v41_serving_admission(realization: V41ProductionServe, report: Present { value: today }) { + ProductionRefused { cause: c } => contains(c as String, "Engram resident") + ProductionAdmitted { run_id: _, concurrency: _ } => false + MeasurementServeAdmitted { run_id: _ } => false + } + admitted && refused +} + +test fn the_dry_run_plan_renders_the_derived_thresholds_and_the_open_preconditions() -> Bool { + let t = v41_capacity_measurement_plan_text() + contains(t, "x6 needs 1572864 tokens, x16 needs 4194304") + && contains(t, "engram arms: unbound") + && contains(t, "TouchCounterUnavailable") + && contains(t, "PROPOSED") +} + +// THE ENGRAM POINT BETS ARE JUDGED IN BYTES AGAINST THE AFTER-REQUESTS RESIDENCY, and a refused +// timeline leaves them unreadable rather than zero. +test fn an_engram_bet_is_paired_with_the_after_requests_residency_in_bytes() -> Bool { + let timeline = v41_admit_residency_timeline(run_id: "run-1" as NonEmptyStr, arm: EngramCacheRetained, cells: w_all_cells()) + match list_at_optional(xs: v41_engram_resident_bets, index: 0) { + Absent => false + Present { value: bet } => { + let observed = match v41_judge_engram_resident_bet(bet: bet, observed: v41_after_requests_resident(t: timeline, rank: 0)) { + PointBetObserved { concurrency: c, bet: b, observed: o } => c == 2 && byte_size_count(b: b) == 19 * 1073741824 && byte_size_count(b: o) == 0 + PointBetUnreadable { concurrency: _, cause: _ } => false + } + let refused = v41_admit_residency_timeline(run_id: "run-1" as NonEmptyStr, arm: EngramCacheReleased, cells: w_all_cells()) + let unreadable = match v41_judge_engram_resident_bet(bet: bet, observed: v41_after_requests_resident(t: refused, rank: 0)) { + PointBetUnreadable { concurrency: _, cause: _ } => true + PointBetObserved { concurrency: _, bet: _, observed: _ } => false + } + observed && unreadable + } + } +} diff --git a/src/v2/workflow/floor_unimported_bare_provider_debt_roster.dag b/src/v2/workflow/floor_unimported_bare_provider_debt_roster.dag index e77a16ca580..210e7c2b3bc 100644 --- a/src/v2/workflow/floor_unimported_bare_provider_debt_roster.dag +++ b/src/v2/workflow/floor_unimported_bare_provider_debt_roster.dag @@ -224,7 +224,7 @@ data unimported_bare_provider_dispositions: List