Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
5 changes: 2 additions & 3 deletions dag/extdeps/huggingface/hub.dag
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@ module extdeps.huggingface.hub

import std.types { String, Bool, List, NonEmptyStr }
import std.nat { Nat }
import std.measure { ByteSize, byte_size_count, Gibibyte, gibibyte, gibibyte_scale_factor_bytes }
import std.measure { ByteSize, Gibibyte, gibibyte_ceiling }
import extdeps.external_authority { ExternalAuthority }
import extdeps.uri { Uri, Https, uri_wire }

Expand Down Expand Up @@ -275,6 +275,5 @@ fn hub_checkpoint_manifest(
// test would admit, never admit one it would refuse. Feasibility arithmetic that must be exact reads
// `safetensors_bytes` and does not round at all.
fn hub_checkpoint_manifest_gibibytes(m: HubCheckpointManifest) -> Gibibyte {
let scale = gibibyte_scale_factor_bytes()
gibibyte(count: (byte_size_count(b: m.safetensors_bytes) + (scale - 1)) / scale)
gibibyte_ceiling(b: m.safetensors_bytes)
}
35 changes: 35 additions & 0 deletions dag/extdeps/vllm/engine_args.dag
Original file line number Diff line number Diff line change
Expand Up @@ -155,6 +155,41 @@ fn vllm_prefix_caching_argv(r: VllmPrefixCachingRequest) -> List<String> {
}
}

// ── THE ENGRAM EMBEDDING STORAGE REQUEST ─────────────────────────────────────────────────────────
//
// vLLM's EngramConfig as upstream ships it, read at github.com/vllm-project/vllm revision
// d2d649e674c75425d2d6975c87eb89fd4d55fff8 (vllm/config/engram.py): three booleans, cpu_offload
// (pinned host memory read over UVA; its default is the VLLM_PLE_CPU_OFFLOAD environment variable,
// and an explicit value takes precedence), embedding_across_dp and dp_shared_memory (which requires
// cpu_offload). vllm/engine/arg_utils.py registers the group as --engram-config, a JSON value.
//
// THE WIRE IS THE DOTTED FORM, NOT JSON, AND THAT IS UPSTREAM'S OWN SYNTAX. vllm/utils/argparse_utils.py
// at the same revision rewrites `--engram-config.cpu_offload true` into the nested dict before
// argparse sees it, json-decoding each value. A JSON object would need quotes, which a systemd Exec
// line cannot carry unmodelled; the dotted words are plain.
//
// ONLY THE TWO FIELDS A LAUNCH HAS ASKED FOR ARE MODELLED. The file-backed Engram overlay
// (gunbc.spark.v41_runtime_candidate v41_engram_file_backed_overlay_path) requires cpu_offload true and
// dp_shared_memory false and refuses at model init otherwise; embedding_across_dp has no consumer and is
// left to upstream's default by saying nothing. Omitted renders no flag at all.
type VllmEngramConfigRequest
= EngramConfigOmitted
| EngramConfigExplicit { cpu_offload: Bool, dp_shared_memory: Bool }

fn vllm_engram_bool_wire(b: Bool) -> String {
if b { "true" } else { "false" }
}

fn vllm_engram_config_argv(r: VllmEngramConfigRequest) -> List<String> {
match r {
EngramConfigOmitted => []
EngramConfigExplicit { cpu_offload: c, dp_shared_memory: d } => [
"--engram-config.cpu_offload", vllm_engram_bool_wire(b: c),
"--engram-config.dp_shared_memory", vllm_engram_bool_wire(b: d),
]
}
}

// THE KV CACHE DTYPE AS REQUESTED, WHICH IS NOT THE DTYPE AS RESOLVED. Passing `fp8` for a
// DeepSeek MLA checkpoint resolves to `fp8_ds_mla`, and the engine says so on its own startup line.
// A realization that records the REQUESTED spelling and calls it the served format has recorded the
Expand Down
35 changes: 23 additions & 12 deletions dag/gunbc/spark/glm_canary_converge.dag
Original file line number Diff line number Diff line change
Expand Up @@ -292,20 +292,31 @@ fn glm_canary_reasoning_projection(arm: GlmCanaryArm) -> ReasoningResponseProjec
// THE MOUNT IS THE PLACEMENT'S HOST PATH, so the CLI writes to and verifies exactly the tree the arm
// declares, and the container path is where that tree appears to the CLI. Both come from the one
// placement value rather than being spelled beside it.
// A RUNTIME WITH NO HUB CLI READ INSIDE IT GETS THE SAME EXECUTOR AS AN UNMODELLED ONE: no path was
// read, so none is spelled, and the sentinel the executor runs refuses on the absent program exactly
// as it does for a runtime with no row.
fn glm_canary_unmodelled_acquisition_executor() -> SnapshotAcquisitionExecutor {
ContainerizedCli {
image: "(unmodelled runtime)" as NonEmptyStr,
cli_path: "/nonexistent" as NonEmptyStr,
host_path: "/nonexistent" as NonEmptyStr,
container_path: "/out" as NonEmptyStr,
}
}

fn glm_canary_acquisition_executor(arm: GlmCanaryArm) -> SnapshotAcquisitionExecutor {
match serving_runtime_row(r: arm.runtime) {
Absent => ContainerizedCli {
image: "(unmodelled runtime)" as NonEmptyStr,
cli_path: "/nonexistent" as NonEmptyStr,
host_path: "/nonexistent" as NonEmptyStr,
container_path: "/out" as NonEmptyStr,
}
Present { value: rrow } => ContainerizedCli {
image: runtime_artifact_wire(a: rrow.artifact),
cli_path: rrow.hub_cli_path,
host_path: glm_canary_host_path(arm: arm),
container_path: "/out" as NonEmptyStr,
}
Absent => glm_canary_unmodelled_acquisition_executor()
Present { value: rrow } =>
match rrow.hub_cli_path {
Absent => glm_canary_unmodelled_acquisition_executor()
Present { value: cli } => ContainerizedCli {
image: runtime_artifact_wire(a: rrow.artifact),
cli_path: cli,
host_path: glm_canary_host_path(arm: arm),
container_path: "/out" as NonEmptyStr,
}
}
}
}

Expand Down
3 changes: 2 additions & 1 deletion dag/gunbc/spark/glm_group_b_launch.dag
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,7 @@ import gunbc.spark.fabric_switch_observed { FabricGroup }
import gunbc.spark.serving_arm {
ArmIntent,
}
import gunbc.spark.serving_arm_launch { ArmLaunchPlan, plan_arm_launch_for }
import gunbc.spark.serving_arm_launch { ArmLaunchPlan, plan_arm_launch_for, glm_native_arm_container }

// ── THE ONE CALL SITE THAT SAYS WHAT GROUP B LAUNCHES ─────────────────────────────────────────
//
Expand Down Expand Up @@ -136,5 +136,6 @@ fn glm_group_b_launch_plan(intent: ArmIntent) -> ArmLaunchPlan {
intent: intent,
profile: a.profile,
deployment_env: glm_group_b_deployment_env(),
container: glm_native_arm_container,
)
}
12 changes: 12 additions & 0 deletions dag/gunbc/spark/model_snapshot_acquisition.dag
Original file line number Diff line number Diff line change
Expand Up @@ -395,6 +395,18 @@ data spark_hub_cache_root: NonEmptyStr = "/home/briansrls/.cache/huggingface/hub
// than a prefix re-spelled per row.
data spark_models_root: NonEmptyStr = "/home/briansrls/models" as NonEmptyStr

// WHERE A PINNED CHECKPOINT MATERIALIZED UNDER THE MODELS ROOT LIVES: the Hub's own repo folder name
// under the root, and the pinned revision beneath it, so a revision bump lands in a new directory rather
// than on top of the old one. The materializer that installs it and the launch that mounts it both
// read this, so the install path and the mount cannot be two derivations of one directory.
fn spark_local_checkpoint_repo_directory(repo: HubRepoId) -> NonEmptyStr {
join([spark_models_root as String, "/", hub_repo_cache_directory_name(repo: repo) as String], "") as NonEmptyStr
}

fn spark_local_checkpoint_directory(repo: HubRepoId, revision: NonEmptyStr) -> NonEmptyStr {
join([spark_local_checkpoint_repo_directory(repo: repo) as String, "/", revision as String], "") as NonEmptyStr
}

fn model_snapshot_install_path(snapshot: ModelSnapshotDesired) -> NonEmptyStr? {
match model_snapshot_subject(snapshot: snapshot, cache_root: spark_hub_cache_root) {
Absent => none
Expand Down
71 changes: 46 additions & 25 deletions dag/gunbc/spark/native_serving_apply.dag
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
module gunbc.spark.native_serving_apply

import std.types { String, Bool, List, NonEmptyStr, Secret, FilePath }
import std.types { String, Bool, List, NonEmptyStr, Secret, FilePath, Port }
import std.algebra { trim }
import v2.std.algebra { zip_map, list_head, list_tail, HeadFound, HeadAbsent, TailFound, TailAbsent }
import v2.std.optional { Present, Absent }
Expand Down Expand Up @@ -47,10 +47,10 @@ import gunbc.spark.pair_serving_apply {
}
import gunbc.spark.native_serving_realization {
GlmNativeUnitRender, GlmNativeUnitRendered, GlmNativeUnitRenderRefused,
glm_native_all_unit_renders, glm_native_ranked_hosts, glm_native_head_unit_name,
glm_native_head_log_path, glm_native_rank_log_path,
glm_native_all_unit_renders, glm_native_ranked_hosts,
ArmUnitIdentity, glm_native_arm_units, arm_unit_log_path_for,
}
import gunbc.spark.serving_arm_launch { glm_native_container_name, ArmNodeLaunch, arm_rail_gid_selection }
import gunbc.spark.serving_arm_launch { ArmNodeLaunch, arm_rail_gid_selection }
import gunbc.spark.serving_rank_observe {
ObservedRankSubject, observed_rank_subject,
RankInspectReading, RankInspected, RankInspectRefused, observe_rank_container,
Expand All @@ -66,7 +66,7 @@ import gunbc.spark.serving_relaunch_transaction {
}
import extdeps.nvidia.nccl { NcclTransportReading, NcclTransportIbVerbs, NcclTransportSocket }
import gunbc.spark.model_snapshot_materialize { SparkMaterializeContext, HostCli }
import gunbc.spark.serving_arm { glm_native_concurrent_profile }
import gunbc.spark.serving_arm { glm_native_concurrent_profile, ServingLaunchProfile }
import gunbc.spark.native_serving_realization { glm_native_launch_lines, glm_native_master_endpoint }
import std.nat { Nat }

Expand Down Expand Up @@ -138,8 +138,8 @@ data spark_native_apply_attempt_raw: String = "spark-native-apply"
// routes on how a remote step is expressed, which is the DESIGN section 3 fork this PR spent its
// length removing from the launch. That is the reason this change declares the marker rather than
// rewriting the bodies: the rewrite belongs with the pair route, not beside it.
fn spark_native_host_apply_script(unit_name: NonEmptyStr, log_path: NonEmptyStr, unit_text: String) -> String {
let ctr = glm_native_container_name as String
fn spark_native_host_apply_script(unit_name: NonEmptyStr, log_path: NonEmptyStr, unit_text: String, container_name: NonEmptyStr) -> String {
let ctr = container_name as String
join([
"set -e\n",
"install -d -m 755 /var/log/gunbc\n",
Expand Down Expand Up @@ -593,8 +593,8 @@ fn spark_native_preflight_script(unit_name: NonEmptyStr, unit_text: String, spec
// THIS STAGE DOES NOT MUTATE THE SERVING ARM. The active unit and the running container are untouched;
// everything written is this transaction's own. That is what lets a failure here refuse rather than
// roll back, which is the difference between reporting a problem and reloading 328 GB of weights.
fn spark_native_preserve_script(unit_name: NonEmptyStr, transaction_id: String) -> String {
let ctr = glm_native_container_name as String
fn spark_native_preserve_script(unit_name: NonEmptyStr, transaction_id: String, container_name: NonEmptyStr) -> String {
let ctr = container_name as String
let u = spark_native_unit_path(unit_name: unit_name)
join([
"set -e\n",
Expand Down Expand Up @@ -706,8 +706,8 @@ fn spark_native_commit_script(unit_name: NonEmptyStr, transaction_id: String) ->
//
// THE REMOVE BRANCH MAY CLEAR, because its evidence is stronger and is an ABSENCE: the unit is gone and
// the daemon reports no container by that name. There is no unverified running thing left to adopt.
fn spark_native_rollback_script(unit_name: NonEmptyStr, transaction_id: String) -> String {
let ctr = glm_native_container_name as String
fn spark_native_rollback_script(unit_name: NonEmptyStr, transaction_id: String, container_name: NonEmptyStr) -> String {
let ctr = container_name as String
let u = spark_native_unit_path(unit_name: unit_name)
join([
"set -e\n",
Expand All @@ -733,8 +733,11 @@ fn spark_native_rollback_script(unit_name: NonEmptyStr, transaction_id: String)

// ── THE PER-STAGE PLANS, ONE PER HOST, IN THE ARM'S OWN RANK ORDER ────────────────────────────

fn spark_native_is_head_unit(name: NonEmptyStr) -> Bool {
(name as String) == (glm_native_head_unit_name as String)
// THE HEAD IS RANK ZERO OF THE PLAN, read off the step's own launch rather than by comparing a unit
// name against one deployment's head-unit constant -- which made every other deployment's head
// unrecognisable to this module.
fn spark_native_step_is_head(st: SparkNativeHostStep) -> Bool {
st.launch.node.node_rank == 0
}

// ONE ROW PER HOST CARRYING EVERYTHING EVERY STAGE NEEDS, so a stage is a projection of the arm rather
Expand All @@ -743,22 +746,28 @@ type SparkNativeHostStep sole_constructor {
host: HostIdentity
unit_name: NonEmptyStr
unit_text: String
log_path: NonEmptyStr
launch: ArmNodeLaunch
serve_port: Port
served_alias: NonEmptyStr
}

type SparkNativeArmSteps
= NativeArmSteps { steps: List<SparkNativeHostStep>, master_endpoint: NonEmptyStr }
| NativeArmStepsRefused { cause: NonEmptyStr }

fn spark_native_step_of(launch: ArmNodeLaunch, render: GlmNativeUnitRender) -> SparkNativeHostStep? {
fn spark_native_step_of(launch: ArmNodeLaunch, render: GlmNativeUnitRender, units: ArmUnitIdentity, profile: ServingLaunchProfile) -> SparkNativeHostStep? {
match render {
GlmNativeUnitRenderRefused { unit_name: _, cause: _ } => none
GlmNativeUnitRendered { unit_name: name, text: text } =>
Present { value: SparkNativeHostStep {
host: launch.node.host,
unit_name: name,
unit_text: text,
log_path: arm_unit_log_path_for(units: units, launch: launch),
launch: launch,
serve_port: profile.serve_port,
served_alias: profile.served_alias,
} }
}
}
Expand All @@ -767,12 +776,23 @@ fn spark_native_step_of(launch: ArmNodeLaunch, render: GlmNativeUnitRender) -> S
// tensor-parallel degree is a property of the arm -- so one unrendered unit refuses the transaction
// before any credential is materialized, let alone any host touched.
fn spark_native_arm_steps() -> SparkNativeArmSteps {
let launches = glm_native_launch_lines()
let renders = glm_native_all_unit_renders()
match glm_native_master_endpoint() {
spark_native_arm_steps_of(
launches: glm_native_launch_lines(),
renders: glm_native_all_unit_renders(),
master_endpoint: glm_native_master_endpoint(),
units: glm_native_arm_units,
profile: glm_native_concurrent_profile,
)
}

// THE SAME WHOLE-OR-NOTHING PLANNING FOR ANY ARM, given that arm's launch lines, its rendered units and
// the names and profile it deploys under. Group B's arm is one call; the V4.1 arm on Group A is another
// (gunbc.spark.v41_group_a_launch), and the stage scripts below are shared by both.
fn spark_native_arm_steps_of(launches: List<ArmNodeLaunch>, renders: List<GlmNativeUnitRender>, master_endpoint: NonEmptyStr?, units: ArmUnitIdentity, profile: ServingLaunchProfile) -> SparkNativeArmSteps {
match master_endpoint {
Absent => NativeArmStepsRefused { cause: "the arm's rendezvous endpoint is unknown, so nothing could be verified after a restart" as NonEmptyStr }
Present { value: endpoint } => {
let steps = flat_map(zip_map(a: launches, b: renders, f: fn(l, r) { spark_native_step_of(launch: l, render: r) }),
let steps = flat_map(zip_map(a: launches, b: renders, f: fn(l, r) { spark_native_step_of(launch: l, render: r, units: units, profile: profile) }),
o => match o { Absent => [] as List<SparkNativeHostStep> Present { value: v } => [v] })
if count(steps) == count(launches) && count(steps) > 0 {
NativeArmSteps { steps: steps, master_endpoint: endpoint }
Expand All @@ -794,14 +814,15 @@ fn spark_native_preflight_plans(steps: List<SparkNativeHostStep>) -> List<SparkP
}

fn spark_native_preserve_plans(steps: List<SparkNativeHostStep>, transaction_id: String) -> List<SparkPairHostApplyPlan> {
map(steps, st => spark_native_plan_of(step: st, script: spark_native_preserve_script(unit_name: st.unit_name, transaction_id: transaction_id)))
map(steps, st => spark_native_plan_of(step: st, script: spark_native_preserve_script(unit_name: st.unit_name, transaction_id: transaction_id, container_name: st.launch.container.name)))
}

fn spark_native_apply_stage_plans(steps: List<SparkNativeHostStep>) -> List<SparkPairHostApplyPlan> {
map(steps, st => spark_native_plan_of(step: st, script: spark_native_host_apply_script(
unit_name: st.unit_name,
log_path: if spark_native_is_head_unit(name: st.unit_name) { glm_native_head_log_path } else { glm_native_rank_log_path },
unit_text: st.unit_text)))
log_path: st.log_path,
unit_text: st.unit_text,
container_name: st.launch.container.name)))
}

// THE ONLY SHELL LEG LEFT IN THE READBACK IS THE HEAD'S FRONT DOOR, and it runs on the head alone
Expand All @@ -810,10 +831,10 @@ fn spark_native_apply_stage_plans(steps: List<SparkNativeHostStep>) -> List<Spar
// below and is not a script at all.
fn spark_native_front_door_plans(steps: List<SparkNativeHostStep>) -> List<SparkPairHostApplyPlan> {
flat_map(steps, st =>
if spark_native_is_head_unit(name: st.unit_name) {
if spark_native_step_is_head(st: st) {
[spark_native_plan_of(step: st, script: spark_native_front_door_script(
port_text: to_string(glm_native_concurrent_profile.serve_port),
served_alias: glm_native_concurrent_profile.served_alias))]
port_text: to_string(st.serve_port),
served_alias: st.served_alias))]
} else {
[] as List<SparkPairHostApplyPlan>
})
Expand All @@ -828,7 +849,7 @@ fn spark_native_commit_plans(steps: List<SparkNativeHostStep>, transaction_id: S
}

fn spark_native_rollback_plans(steps: List<SparkNativeHostStep>, transaction_id: String) -> List<SparkPairHostApplyPlan> {
map(steps, st => spark_native_plan_of(step: st, script: spark_native_rollback_script(unit_name: st.unit_name, transaction_id: transaction_id)))
map(steps, st => spark_native_plan_of(step: st, script: spark_native_rollback_script(unit_name: st.unit_name, transaction_id: transaction_id, container_name: st.launch.container.name)))
}

// ── RUNNING ONE STAGE ─────────────────────────────────────────────────────────────────────────
Expand Down
Loading