From 59e1b3006ed63c854b697cc6d538de2d6d3e8455 Mon Sep 17 00:00:00 2001 From: Tom Tupper <11300465+ttupper92618@users.noreply.github.com> Date: Tue, 8 Sep 2026 02:56:30 -0500 Subject: [PATCH 1/3] feat: account for hybrid GGUF recurrent cache allocations --- CLAUDE.md | 7 + .../api/tests/test_data_plane_dispatch.py | 30 ++- src/skulk/master/main.py | 14 +- src/skulk/master/placement.py | 69 ++++- src/skulk/master/placement_utils.py | 36 ++- .../tests/test_recurrent_memory_admission.py | 243 ++++++++++++++++++ src/skulk/shared/models/gguf_memory.py | 125 +++++++++ .../shared/models/llama_server_settings.py | 55 ++++ src/skulk/shared/models/memory_estimate.py | 140 +++++++++- src/skulk/shared/models/model_cards.py | 97 ++++--- .../shared/models/tests/test_gguf_cards.py | 40 +++ .../shared/models/tests/test_gguf_memory.py | 95 +++++++ src/skulk/shared/types/profiling.py | 16 ++ src/skulk/shared/types/worker/shards.py | 5 + src/skulk/worker/main.py | 2 + .../worker/runner/llama_server/runner.py | 34 ++- website/docs/api-guide.md | 10 + website/docs/architecture-reference.md | 8 + website/docs/architecture.md | 11 + 19 files changed, 958 insertions(+), 79 deletions(-) create mode 100644 src/skulk/master/tests/test_recurrent_memory_admission.py create mode 100644 src/skulk/shared/models/gguf_memory.py create mode 100644 src/skulk/shared/models/llama_server_settings.py create mode 100644 src/skulk/shared/models/tests/test_gguf_memory.py diff --git a/CLAUDE.md b/CLAUDE.md index 672d893d0..39efa2d25 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -108,6 +108,13 @@ If `nix fmt` changes any files, stage them before committing. The CI runs `nix f ### Node Composition +GGUF hybrid cache admission uses `ModelCard.gguf_cache_geometry` from the exact +artifact header. Qwen3.5 scalar metadata separates attention, recurrent and NextN +layers; recurrent FP32 buffers scale with slots and speculative rollback rows. +NodeResources advertises `llama_server_settings`, placement stamps them on shards, +and the runner refuses changed local settings before spawning. Registry geometry +must enter through a new signed card revision, never a silent runtime overlay. + Exact non-RPC placements resolve omitted backends from advertised compatible engines before memory admission; restored unstamped GPU-host shards reserve conservatively. Exact and quick-launch master refusals retain correlated instance diff --git a/src/skulk/api/tests/test_data_plane_dispatch.py b/src/skulk/api/tests/test_data_plane_dispatch.py index 13d50b869..dad3c1bce 100644 --- a/src/skulk/api/tests/test_data_plane_dispatch.py +++ b/src/skulk/api/tests/test_data_plane_dispatch.py @@ -107,23 +107,25 @@ async def test_state_surfaces_split_data_transport_health() -> None: payload = await api.get_cluster_state() assert payload["nodeResources"] == { - "remote-management-node": { - "backends": [], - "engineBuilds": {}, - "hardwareClasses": [], - "participation": "management", - "apiAvailable": True, - "dataTransport": "zenoh", + "remote-management-node": { + "backends": [], + "engineBuilds": {}, + "llamaServerSettings": None, + "hardwareClasses": [], + "participation": "management", + "apiAvailable": True, + "dataTransport": "zenoh", "zenohConnectedPeers": None, "capabilityConflicts": [], }, - "worker-node": { - "backends": ["mlx"], - "engineBuilds": {}, - "hardwareClasses": [], - "participation": "full", - "apiAvailable": True, - "dataTransport": "gossipsub", + "worker-node": { + "backends": ["mlx"], + "engineBuilds": {}, + "llamaServerSettings": None, + "hardwareClasses": [], + "participation": "full", + "apiAvailable": True, + "dataTransport": "gossipsub", "zenohConnectedPeers": None, "capabilityConflicts": [], }, diff --git a/src/skulk/master/main.py b/src/skulk/master/main.py index 2e9e96026..5e6491c8b 100644 --- a/src/skulk/master/main.py +++ b/src/skulk/master/main.py @@ -1127,7 +1127,12 @@ def _record_freed_instance(self, instance: Instance) -> None: fraction = shard_fraction_of_model(shard) if fraction is None or fraction <= 0.0: continue - footprint = estimate_shard_footprint(shard.model_card, fraction) + footprint = estimate_shard_footprint( + shard.model_card, + fraction, + resolved_backend=shard.resolved_backend, + llama_server_settings=shard.llama_server_settings, + ) self._recently_freed_bytes.setdefault(node_id, []).append( (footprint.in_bytes, deadline) ) @@ -3468,7 +3473,12 @@ def _steward_replacement_memory_inputs( fraction = shard_fraction_of_model(shard) if fraction is None or fraction <= 0.0: continue - footprint = estimate_shard_footprint(shard.model_card, fraction) + footprint = estimate_shard_footprint( + shard.model_card, + fraction, + resolved_backend=shard.resolved_backend, + llama_server_settings=shard.llama_server_settings, + ) credit[node_id] = credit.get(node_id, 0) + footprint.in_bytes memory = { node_id: ( diff --git a/src/skulk/master/placement.py b/src/skulk/master/placement.py index b2d5bc496..544eb57b1 100644 --- a/src/skulk/master/placement.py +++ b/src/skulk/master/placement.py @@ -85,7 +85,11 @@ instance_meta_of, ) from skulk.shared.types.worker.runners import ShardAssignments -from skulk.shared.types.worker.shards import Sharding, TensorShardMetadata +from skulk.shared.types.worker.shards import ( + RpcDonorShardMetadata, + Sharding, + TensorShardMetadata, +) # Ring/coordinator/donor listener ports are drawn from a band BELOW every # OS's default ephemeral range: macOS assigns outgoing-connection local ports @@ -233,6 +237,7 @@ def add_instance_to_placements( assignments = assignments.model_copy( update={"runner_to_shard": resolved_shards} ) + assignments = _stamp_llama_server_settings(assignments, node_resources or {}) fixed_memory_by_node: dict[NodeId, Memory] = {} first_shard = next(iter(assignments.runner_to_shard.values()), None) if first_shard is not None: @@ -290,12 +295,44 @@ def add_instance_to_placements( context_budget=ceiling if ceiling is not None else KV_CONTEXT_BUDGET_TOKENS, + resolved_backend=shard.resolved_backend, + llama_server_settings=shard.llama_server_settings, ) if ceiling == 0 or footprint > available: raise PlacementError("Insufficient GPU memory for the exact placement") return {**current_instances, instance.instance_id: instance} +def _stamp_llama_server_settings( + assignments: ShardAssignments, resources_by_node: Mapping[NodeId, NodeResources] +) -> ShardAssignments: + shards = dict(assignments.runner_to_shard) + has_rpc_donors = any( + isinstance(shard, RpcDonorShardMetadata) for shard in shards.values() + ) + for node_id, runner_id in assignments.node_to_runner.items(): + shard = shards[runner_id] + if isinstance(shard, RpcDonorShardMetadata): + # Donors provide memory to the driver's one context; they do not + # launch a served instance with their own parallel-slot controls. + continue + if not has_rpc_donors and ( + shard.resolved_backend is None + or not shard.resolved_backend.startswith("llama_server") + ): + continue + resources = resources_by_node.get(node_id) + settings = resources.llama_server_settings if resources is not None else None + if settings is None and shard.model_card.gguf_cache_geometry is not None: + raise PlacementError( + "Serving settings are required for recurrent memory admission" + ) + # Exact requests cannot substitute client-supplied slot settings for the + # node's observation. The worker verifies this stamp before spawning. + shards[runner_id] = shard.model_copy(update={"llama_server_settings": settings}) + return assignments.model_copy(update={"runner_to_shard": shards}) + + def _get_node_download_fraction( node_id: NodeId, model_id: ModelId, @@ -920,6 +957,20 @@ def place_instance( command.sharding, node_vram=node_vram, fixed_memory_by_node=standard_fixed, + resolved_backends={ + node_id: resolve_node_backend( + _card_platform_backends(command.model_card, resources), + command.model_card.placement.backend_preference, + resources.backends, + ) + for node_id in standard_cycle.node_ids + if (resources := resolved_node_resources.get(node_id)) is not None + }, + llama_server_settings={ + node_id: resources.llama_server_settings + for node_id in standard_cycle.node_ids + if (resources := resolved_node_resources.get(node_id)) is not None + }, ) cycles_with_sufficient_memory.extend(standard_fit) memory_diagnostics.pending_info_node_ids.extend( @@ -978,6 +1029,18 @@ def place_instance( node_vram=rpc_vram_map, exact_pipeline_layers=False, fixed_memory_by_node=rpc_fixed, + resolved_backends={ + node_id: "llama_server" for node_id in rpc_cycle.node_ids + }, + llama_server_settings={ + node_id: ( + resources.llama_server_settings + if (resources := resolved_node_resources.get(driver)) + is not None + else None + ) + for node_id in rpc_cycle.node_ids + }, ) cycles_with_sufficient_memory.extend(rpc_fit) if rpc_fit: @@ -1179,6 +1242,10 @@ def place_instance( node_to_runner=shard_assignments.node_to_runner, ) + shard_assignments = _stamp_llama_server_settings( + shard_assignments, resolved_node_resources + ) + # Stamp the context-admission ceiling into the placement decision (#279 # slice 2). Computed once here from the hosting nodes' static ram_total, so # every rank reads the identical value off replicated state rather than diff --git a/src/skulk/master/placement_utils.py b/src/skulk/master/placement_utils.py index 2d78da349..0b66b98ef 100644 --- a/src/skulk/master/placement_utils.py +++ b/src/skulk/master/placement_utils.py @@ -5,13 +5,16 @@ from loguru import logger from pydantic import Field +from skulk.shared.models.llama_server_settings import LlamaServerSettings from skulk.shared.models.memory_estimate import ( GPU_VRAM_WORKING_SET_FRACTION, GPU_WORKING_SET_FRACTION, UMA_GPU_OS_HEADROOM, backend_offloads_to_vram, + estimate_recurrent_cache_bytes, estimate_shard_footprint, memory_overhead_factor, + per_token_kv_bytes, shard_fraction_of_model, ) from skulk.shared.models.memory_estimate import ( @@ -247,6 +250,8 @@ def reserve_instance_vram( footprint = estimate_shard_footprint( shard.model_card, fraction, + resolved_backend=shard.resolved_backend, + llama_server_settings=shard.llama_server_settings, context_budget=( instance.context_token_limit if instance.context_token_limit is not None @@ -458,6 +463,8 @@ def filter_cycles_by_memory( *, exact_pipeline_layers: bool = True, fixed_memory_by_node: Mapping[NodeId, Memory] | None = None, + resolved_backends: Mapping[NodeId, str | None] | None = None, + llama_server_settings: Mapping[NodeId, LlamaServerSettings | None] | None = None, ) -> tuple[list[Cycle], CycleMemoryDiagnostics]: """Keep cycles whose every node can hold its shard with runtime headroom. @@ -559,10 +566,37 @@ def filter_cycles_by_memory( overhead_factor = memory_overhead_factor(model_card) overloaded: list[str] = [] for node_id, share in node_shares.items(): - kv_share = share * kv_ratio + backend = (resolved_backends or {}).get(node_id) + settings = (llama_server_settings or {}).get(node_id) + node_kv_ratio = kv_ratio + recurrent_share = Memory() + if model_card.gguf_cache_geometry is not None: + if ( + backend is not None + and backend.startswith("llama_server") + and settings is None + ): + overloaded.append( + f"node {node_id} has no observed llama-server serving settings" + ) + continue + node_kv_ratio = ( + per_token_kv_bytes( + model_card, + resolved_backend=backend, + llama_server_settings=settings, + ) + * context_budget + / required_memory.in_bytes + ) + recurrent_share = estimate_recurrent_cache_bytes( + model_card, resolved_backend=backend, llama_server_settings=settings + ) * (share.in_bytes / required_memory.in_bytes) + kv_share = share * node_kv_ratio required_with_overhead = ( share * overhead_factor + kv_share + + recurrent_share + PLACEMENT_MEMORY_OVERHEAD_FLOOR + fixed_memory_by_node.get(node_id, Memory()) ) diff --git a/src/skulk/master/tests/test_recurrent_memory_admission.py b/src/skulk/master/tests/test_recurrent_memory_admission.py new file mode 100644 index 000000000..81a2552e9 --- /dev/null +++ b/src/skulk/master/tests/test_recurrent_memory_admission.py @@ -0,0 +1,243 @@ +"""Hybrid-model admission must reserve actual slots, rollback state and KV together.""" + +import pytest + +from skulk.master.placement import PlacementError, add_instance_to_placements +from skulk.master.placement_utils import ( + filter_cycles_by_memory, + get_shard_assignments_for_llama_rpc, + usable_vram_by_node, +) +from skulk.master.tests.conftest import create_node_memory +from skulk.shared.models.gguf_memory import qwen35_cache_geometry +from skulk.shared.models.llama_server_settings import LlamaServerSettings +from skulk.shared.models.memory_estimate import ( + estimate_recurrent_cache_bytes, + estimate_shard_footprint, + instance_context_token_limit, + per_token_kv_bytes, +) +from skulk.shared.models.model_cards import ( + ModelCard, + ModelId, + ModelTask, + PlacementCardConfig, + RuntimeCapabilityCardConfig, +) +from skulk.shared.models.tests.test_gguf_memory import metadata +from skulk.shared.topology import Topology +from skulk.shared.types.commands import CreateInstance +from skulk.shared.types.common import NodeId +from skulk.shared.types.memory import Memory +from skulk.shared.types.profiling import ( + AcceleratorMetrics, + NodeResources, + SystemPerformanceProfile, +) +from skulk.shared.types.topology import Cycle +from skulk.shared.types.worker.instances import ( + InstanceId, + LlamaRpcInstance, + MlxRingInstance, +) +from skulk.shared.types.worker.runners import RunnerId, ShardAssignments +from skulk.shared.types.worker.shards import PipelineShardMetadata + + +def hybrid_instance(node: NodeId, slots: int = 16) -> MlxRingInstance: + """Build a synthetic admission input with source-derived tensor dimensions.""" + card = ModelCard( + model_id=ModelId("test/hybrid"), + storage_size=Memory.from_bytes(5868826976), + n_layers=33, + hidden_size=4096, + supports_tensor=False, + num_key_value_heads=4, + context_length=262144, + tasks=[ModelTask.TextGeneration], + gguf_file="model.gguf", + gguf_cache_geometry=qwen35_cache_geometry("qwen35", metadata()), + runtime=RuntimeCapabilityCardConfig( + served_spec_type="draft_mtp", served_spec_n_max=3 + ), + placement=PlacementCardConfig( + compatible_backends=frozenset({"llama_server-cuda"}) + ), + ) + runner = RunnerId() + return MlxRingInstance( + instance_id=InstanceId(), + shard_assignments=ShardAssignments( + model_id=card.model_id, + node_to_runner={node: runner}, + runner_to_shard={ + runner: PipelineShardMetadata( + model_card=card, + device_rank=0, + world_size=1, + start_layer=0, + end_layer=33, + n_layers=33, + resolved_backend="llama_server-cuda", + llama_server_settings=LlamaServerSettings(parallel_slots=slots), + ) + }, + ), + hosts_by_node={}, + ephemeral_port=50000, + context_token_limit=8192, + ) + + +@pytest.mark.parametrize("slots", [1, 16, 64]) +def test_context_ceiling_and_footprint_agree_at_the_memory_boundary(slots: int) -> None: + """Context can grow only after fixed per-slot state has been reserved.""" + node = NodeId() + instance = hybrid_instance(node, slots) + shard = next(iter(instance.shard_assignments.runner_to_shard.values())) + available = Memory.from_gb(24) + limit = instance_context_token_limit( + instance.shard_assignments, + {node: Memory.from_gb(64)}, + node_vram={node: available}, + ) + assert limit is not None and 8192 < limit <= 262144 + footprint = estimate_shard_footprint( + shard.model_card, + 1.0, + limit, + resolved_backend=shard.resolved_backend, + llama_server_settings=shard.llama_server_settings, + ) + assert footprint <= available + if limit < 262144: + assert ( + footprint.in_bytes + per_token_kv_bytes(shard.model_card) + > available.in_bytes + ) + + +def test_actual_slot_count_is_reserved_before_load_telemetry() -> None: + """A pending 64-slot placement cannot expose the same free GPU budget as 16.""" + node = NodeId() + system = { + node: SystemPerformanceProfile( + accelerator=AcceleratorMetrics( + vendor="nvidia", + vram_total_bytes=Memory.from_gb(48).in_bytes, + vram_used_bytes=0, + ) + ) + } + ordinary, wide = hybrid_instance(node, 16), hybrid_instance(node, 64) + ordinary_free = usable_vram_by_node( + system, current_instances={ordinary.instance_id: ordinary} + )[node] + wide_free = usable_vram_by_node(system, current_instances={wide.instance_id: wide})[ + node + ] + assert ordinary_free.in_bytes - wide_free.in_bytes == 3 * 3372220416 + + +def test_in_process_and_plain_served_caches_keep_their_distinct_costs() -> None: + """Disabling speculation removes rollback rows, but does not remove slot state.""" + card = next( + iter(hybrid_instance(NodeId()).shard_assignments.runner_to_shard.values()) + ).model_card + assert ( + estimate_recurrent_cache_bytes(card, resolved_backend="llama_cpp-cuda").in_bytes + == 52690944 + ) + assert per_token_kv_bytes(card, resolved_backend="llama_cpp-cuda") == 32768 + assert ( + estimate_recurrent_cache_bytes( + card, + resolved_backend="llama_server-cuda", + llama_server_settings=LlamaServerSettings(speculation_enabled=False), + ).in_bytes + == 843055104 + ) + + +def test_exact_placement_uses_observed_settings_instead_of_caller_settings() -> None: + """A client cannot submit a cheap one-slot estimate for a 64-slot worker.""" + node = NodeId() + instance = hybrid_instance(node, 1) + resources = { + node: NodeResources( + backends=frozenset({"llama_server-cuda"}), + llama_server_settings=LlamaServerSettings(parallel_slots=64), + ) + } + with pytest.raises(PlacementError, match="Insufficient GPU memory"): + add_instance_to_placements( + CreateInstance(instance=instance), + Topology(), + {}, + {node: create_node_memory(Memory.from_gb(64).in_bytes)}, + node_vram={node: Memory.from_gb(12)}, + node_resources=resources, + ) + + +def test_cycle_filter_uses_the_same_fixed_slot_cost() -> None: + """Auto placement rejects a wide instance before selecting an impossible cycle.""" + node = NodeId() + card = next( + iter(hybrid_instance(node).shard_assignments.runner_to_shard.values()) + ).model_card + memory = {node: create_node_memory(Memory.from_gb(64).in_bytes)} + for slots, accepted in [(16, True), (64, False)]: + cycles, _ = filter_cycles_by_memory( + [Cycle(node_ids=[node])], + memory, + card, + node_vram={node: Memory.from_gb(12)}, + resolved_backends={node: "llama_server-cuda"}, + llama_server_settings={node: LlamaServerSettings(parallel_slots=slots)}, + ) + assert bool(cycles) is accepted + + +def test_rpc_captures_only_the_driver_serving_settings() -> None: + """RPC donors expose memory, not independent speculative server contexts.""" + driver, donor = NodeId(), NodeId() + card = next( + iter(hybrid_instance(driver).shard_assignments.runner_to_shard.values()) + ).model_card + assignments = get_shard_assignments_for_llama_rpc( + card, Cycle(node_ids=[driver, donor]), driver + ) + instance = LlamaRpcInstance( + instance_id=InstanceId(), + shard_assignments=assignments, + driver_node=driver, + donor_endpoints={donor: "localhost:50000"}, + context_token_limit=8192, + ) + settings = LlamaServerSettings(parallel_slots=64) + result = add_instance_to_placements( + CreateInstance(instance=instance), + Topology(), + {}, + { + node: create_node_memory(Memory.from_gb(64).in_bytes) + for node in [driver, donor] + }, + node_vram={node: Memory.from_gb(24) for node in [driver, donor]}, + node_resources={ + driver: NodeResources( + backends=frozenset({"llama_server-cuda"}), + llama_server_settings=settings, + ), + donor: NodeResources(backends=frozenset({"llama_server-cuda"})), + }, + )[instance.instance_id].shard_assignments + assert ( + result.runner_to_shard[result.node_to_runner[driver]].llama_server_settings + == settings + ) + assert ( + result.runner_to_shard[result.node_to_runner[donor]].llama_server_settings + is None + ) diff --git a/src/skulk/shared/models/gguf_memory.py b/src/skulk/shared/models/gguf_memory.py new file mode 100644 index 000000000..02452a049 --- /dev/null +++ b/src/skulk/shared/models/gguf_memory.py @@ -0,0 +1,125 @@ +"""Artifact-derived cache geometry for llama.cpp memory admission. + +Geometry describes tensors, not a measured total or a model-name heuristic. +Engine configuration supplies slot count and recurrent rollback depth separately. +The Qwen3.5 mapping follows llama.cpp b10753's qwen35 loader, hparams and +recurrent-memory implementation; other recurrent architectures need their own +verified mapping before they can use this estimate. +""" + +from collections.abc import Mapping +from typing import Self, final + +from pydantic import Field, model_validator + +from skulk.utils.pydantic_ext import FrozenModel + + +@final +class GgufCacheGeometry(FrozenModel): + """Fixed recurrent state and per-token attention dimensions of one artifact.""" + + attention_layers: int = Field(ge=0, description="Target full-attention layers.") + recurrent_layers: int = Field(ge=0, description="Target recurrent layers.") + nextn_layers: int = Field(ge=0, description="Embedded MTP attention layers.") + key_width: int = Field(gt=0, description="Key elements per token per layer.") + value_width: int = Field(gt=0, description="Value elements per token per layer.") + convolution_width: int = Field( + ge=0, description="Convolution-state elements per recurrent layer and row." + ) + recurrent_width: int = Field( + ge=0, description="Recurrent-state elements per recurrent layer and row." + ) + + @model_validator(mode="after") + def validate_recurrent_state(self) -> Self: + """Reject incomplete geometry instead of treating an unknown cost as zero.""" + if self.attention_layers + self.recurrent_layers <= 0: + raise ValueError("cache geometry requires target layers") + if self.recurrent_layers > 0 and self.recurrent_width <= 0: + raise ValueError("recurrent layers require their state width") + if self.recurrent_layers == 0 and ( + self.convolution_width or self.recurrent_width + ): + raise ValueError("recurrent state requires recurrent layers") + return self + + def recurrent_bytes(self, *, parallel_slots: int, rollback_depth: int) -> int: + """Return FP32 recurrent buffers, including each slot's rollback rows. + + llama.cpp allocates both state tensors as FP32, independent of weight + quantization or attention-cache dtype. Rollback copies multiply rows; + they do not duplicate the model's weights or all attention caches. + """ + if parallel_slots < 1 or rollback_depth < 0: + raise ValueError("invalid recurrent cache configuration") + return ( + (self.convolution_width + self.recurrent_width) + * self.recurrent_layers + * 4 + * parallel_slots + * (1 + rollback_depth) + ) + + def attention_bytes_per_token(self, *, embedded_mtp: bool) -> int: + """Return FP16 K/V bytes for target and, when enabled, embedded MTP. + + The MTP context shares target weights but owns only the NextN layers' + attention cache. Its full context window is charged once, not per slot. + """ + layers = self.attention_layers + (self.nextn_layers if embedded_mtp else 0) + return layers * (self.key_width + self.value_width) * 2 + + +def qwen35_cache_geometry( + architecture: str, + metadata: Mapping[str, int], + *, + has_recurrent_layer_override: bool = False, +) -> GgufCacheGeometry | None: + """Resolve scalar Qwen3.5 GGUF dimensions, or return unknown for other layouts. + + Keys are architecture-relative GGUF names. An explicit recurrent-layer + vector supersedes the interval in llama.cpp; until that vector is decoded, + its presence must never be mistaken for the regular interval layout. + Missing dimensions remain unknown and malformed complete dimensions raise. + """ + if architecture != "qwen35" or has_recurrent_layer_override: + return None + required = ( + "block_count", + "attention.head_count_kv", + "attention.key_length", + "attention.value_length", + "ssm.conv_kernel", + "ssm.inner_size", + "ssm.state_size", + "ssm.group_count", + ) + if any(key not in metadata for key in required): + return None + if any(metadata[key] <= 0 for key in required): + raise ValueError("GGUF cache dimensions must be positive") + nextn = metadata.get("nextn_predict_layers", 0) + layers = metadata["block_count"] - nextn + interval = metadata.get("full_attention_interval", 4) + if nextn < 0 or layers <= 0 or interval <= 0: + raise ValueError("invalid GGUF attention layer layout") + attention_layers = layers // interval + recurrent_layers = layers - attention_layers + inner = metadata["ssm.inner_size"] + state = metadata["ssm.state_size"] + groups = metadata["ssm.group_count"] + heads = metadata["attention.head_count_kv"] + return GgufCacheGeometry( + attention_layers=attention_layers, + recurrent_layers=recurrent_layers, + nextn_layers=nextn, + key_width=heads * metadata["attention.key_length"], + value_width=heads * metadata["attention.value_length"], + convolution_width=(metadata["ssm.conv_kernel"] - 1) + * (inner + 2 * groups * state) + if recurrent_layers + else 0, + recurrent_width=state * inner if recurrent_layers else 0, + ) diff --git a/src/skulk/shared/models/llama_server_settings.py b/src/skulk/shared/models/llama_server_settings.py new file mode 100644 index 000000000..9df876430 --- /dev/null +++ b/src/skulk/shared/models/llama_server_settings.py @@ -0,0 +1,55 @@ +"""Serving settings shared by node advertisement, admission and process launch.""" + +from collections.abc import Mapping +from typing import Final, final + +from pydantic import Field + +from skulk.utils.pydantic_ext import FrozenModel + +LLAMA_SERVER_DEFAULT_PARALLEL: Final = 16 +LLAMA_SERVER_DEFAULT_DRAFT_DEPTH: Final = 3 + + +@final +class LlamaServerSettings(FrozenModel): + """Node settings that change llama-server's persistent cache allocation.""" + + parallel_slots: int = Field( + default=LLAMA_SERVER_DEFAULT_PARALLEL, + gt=0, + description="Operator-selected concurrent slots before model-specific limits.", + ) + speculation_enabled: bool = Field( + default=True, + description="Whether the node permits the model card's speculative mode.", + ) + + def effective_slots(self, *, speculative_vision: bool) -> int: + """Apply the existing serial limit only to speculative vision serving.""" + return ( + 1 + if speculative_vision and self.speculation_enabled + else self.parallel_slots + ) + + +def resolve_llama_server_settings( + environment: Mapping[str, str], +) -> LlamaServerSettings: + """Resolve existing environment controls without changing their defaults. + + Invalid slot declarations retain the runner's historical default. The runner + remains responsible for warning the operator about such declarations. + """ + try: + slots = int(environment.get("SKULK_LLAMA_SERVER_PARALLEL", "").strip()) + except ValueError: + slots = LLAMA_SERVER_DEFAULT_PARALLEL + if slots < 1: + slots = LLAMA_SERVER_DEFAULT_PARALLEL + disabled = environment.get("SKULK_LLAMA_SERVER_FORCE_NO_SPEC", "").strip().lower() + return LlamaServerSettings( + parallel_slots=slots, + speculation_enabled=disabled not in ("1", "true", "yes", "on"), + ) diff --git a/src/skulk/shared/models/memory_estimate.py b/src/skulk/shared/models/memory_estimate.py index d77594b33..d5525f5d1 100644 --- a/src/skulk/shared/models/memory_estimate.py +++ b/src/skulk/shared/models/memory_estimate.py @@ -19,6 +19,10 @@ from collections.abc import Set as AbstractSet from typing import Final +from skulk.shared.models.llama_server_settings import ( + LLAMA_SERVER_DEFAULT_DRAFT_DEPTH, + LlamaServerSettings, +) from skulk.shared.models.model_cards import ModelCard from skulk.shared.types.common import NodeId from skulk.shared.types.memory import Memory @@ -127,6 +131,10 @@ def estimate_kv_cache_bytes( ) -> Memory: """Estimate KV-cache bytes for ``n_layers`` layers at ``context_tokens``. + A GGUF cache geometry uses its actual attention layers and widths, scaled + by the requested layer fraction. It does not fold recurrent state into KV. + Without that geometry the legacy approximation below remains in use. + The cache holds a key and a value vector per token, per layer, each sized ``num_key_value_heads * head_dim``:: @@ -137,11 +145,13 @@ def estimate_kv_cache_bytes( non-positive — the weight-overhead factor must absorb the slack then. ``head_dim`` falls back to ``KV_HEAD_DIM_FALLBACK`` (cards omit it). """ - kv_heads = model_card.num_key_value_heads - if kv_heads is None or context_tokens <= 0 or n_layers <= 0: + if context_tokens <= 0 or n_layers <= 0: return Memory() kv_bytes = ( - 2 * n_layers * context_tokens * kv_heads * KV_HEAD_DIM_FALLBACK * KV_DTYPE_BYTES + per_token_kv_bytes(model_card) + * n_layers + * context_tokens + // model_card.n_layers ) return Memory.from_bytes(kv_bytes) @@ -150,10 +160,13 @@ def estimate_shard_footprint( model_card: ModelCard, shard_fraction: float, context_budget: int = KV_CONTEXT_BUDGET_TOKENS, + *, + resolved_backend: str | None = None, + llama_server_settings: LlamaServerSettings | None = None, ) -> Memory: """Estimate resident memory for a shard holding ``shard_fraction`` of a model. - ``weights_share * MEMORY_OVERHEAD_FACTOR + kv_share + MEMORY_OVERHEAD_FLOOR`` + ``weights_share * overhead + kv_share + recurrent_share + overhead_floor`` where weights and KV both scale by ``shard_fraction``. That single fraction works for every sharding because both quantities are linear in it: @@ -163,15 +176,31 @@ def estimate_shard_footprint( ``1/world_size`` of each weight matrix and of the KV heads). ``shard_fraction == 1.0`` gives the whole-model footprint (single node). + ``resolved_backend`` and ``llama_server_settings`` select the actual engine's + slot and speculation costs; omitted settings use the shipped server defaults + for advisory planning. Persisted placements supply their stamped settings. """ if shard_fraction <= 0.0: return Memory() weights_share = model_card.storage_size * shard_fraction - full_kv = estimate_kv_cache_bytes(model_card, model_card.n_layers, context_budget) + full_kv = Memory.from_bytes( + per_token_kv_bytes( + model_card, + resolved_backend=resolved_backend, + llama_server_settings=llama_server_settings, + ) + * max(0, context_budget) + ) kv_share = full_kv * shard_fraction footprint = ( weights_share * memory_overhead_factor(model_card) + kv_share + + estimate_recurrent_cache_bytes( + model_card, + resolved_backend=resolved_backend, + llama_server_settings=llama_server_settings, + ) + * shard_fraction + MEMORY_OVERHEAD_FLOOR ) # A single-node GGUF vision runner owns the complete projector in addition @@ -265,14 +294,86 @@ def backend_offloads_to_vram(resolved_backend: str | None) -> bool: ) and not resolved_backend.endswith("-cpu") -def per_token_kv_bytes(model_card: ModelCard) -> int: +def _served_speculative_mode( + model_card: ModelCard, + resolved_backend: str | None, + settings: LlamaServerSettings, +) -> str | None: + if ( + not settings.speculation_enabled + or ( + resolved_backend is not None + and not resolved_backend.startswith("llama_server") + ) + or model_card.runtime is None + ): + return None + return model_card.runtime.served_spec_type + + +def estimate_recurrent_cache_bytes( + model_card: ModelCard, + *, + resolved_backend: str | None = None, + llama_server_settings: LlamaServerSettings | None = None, +) -> Memory: + """Return fixed FP32 recurrent state for the configured llama.cpp instance. + + Unknown geometry retains the legacy approximation; it must not be described + as a proven zero-cost recurrent model. Non-llama engines use their own memory + contract. Unresolved planning uses shipped llama-server settings. + """ + geometry = model_card.gguf_cache_geometry + if geometry is None or ( + resolved_backend is not None + and not resolved_backend.startswith(("llama_cpp", "llama_server")) + ): + return Memory() + settings = llama_server_settings or LlamaServerSettings() + mode = _served_speculative_mode(model_card, resolved_backend, settings) + depth = ( + (model_card.runtime.served_spec_n_max or LLAMA_SERVER_DEFAULT_DRAFT_DEPTH) + if mode in ("draft_mtp", "draft_eagle3", "draft_dflash") + and model_card.runtime is not None + else 0 + ) + slots = ( + 1 + if resolved_backend is not None and resolved_backend.startswith("llama_cpp") + else settings.effective_slots( + speculative_vision=model_card.vision is not None + and mode not in (None, "none") + ) + ) + return Memory.from_bytes( + geometry.recurrent_bytes(parallel_slots=slots, rollback_depth=depth) + ) + + +def per_token_kv_bytes( + model_card: ModelCard, + *, + resolved_backend: str | None = None, + llama_server_settings: LlamaServerSettings | None = None, +) -> int: """Whole-model KV-cache bytes consumed by ONE token of context. - Covers all layers at fp16 (``KV_DTYPE_BYTES``); a node holding + Uses artifact attention geometry when present, including embedded MTP only + for a speculative served instance. Otherwise covers all layers at fp16 + (``KV_DTYPE_BYTES``); a node holding ``shard_fraction`` of the model pays ``per_token_kv_bytes * shard_fraction`` per token. Returns 0 when the card lacks ``num_key_value_heads`` — callers must treat 0 as "KV cost unknown, cannot enforce a memory ceiling". """ + geometry = model_card.gguf_cache_geometry + if geometry is not None and ( + resolved_backend is None + or resolved_backend.startswith(("llama_cpp", "llama_server")) + ): + mode = _served_speculative_mode( + model_card, resolved_backend, llama_server_settings or LlamaServerSettings() + ) + return geometry.attention_bytes_per_token(embedded_mtp=mode == "draft_mtp") kv_heads = model_card.num_key_value_heads if kv_heads is None or model_card.n_layers <= 0: return 0 @@ -361,21 +462,38 @@ def instance_context_token_limit( for shard in shard_assignments.runner_to_shard.values() ) - whole_model_token_bytes = per_token_kv_bytes(model_card) - if whole_model_token_bytes > 0: + if per_token_kv_bytes(model_card) > 0: node_to_runner = shard_assignments.node_to_runner for node_id, runner_id in node_to_runner.items(): shard = shard_assignments.runner_to_shard[runner_id] + whole_model_token_bytes = per_token_kv_bytes( + model_card, + resolved_backend=shard.resolved_backend, + llama_server_settings=shard.llama_server_settings, + ) fraction = shard_fraction_of_model(shard) ram_total = node_ram_totals.get(node_id) - if fraction is None or fraction <= 0.0 or ram_total is None: + if ( + fraction is None + or fraction <= 0.0 + or ram_total is None + or whole_model_token_bytes <= 0 + ): memory_limit = None break working_set = node_vram.get(node_id) or gpu_working_set_ceiling(ram_total) kv_budget = ( working_set - - model_card.storage_size * fraction * memory_overhead_factor(model_card) + - model_card.storage_size + * fraction + * memory_overhead_factor(model_card) - MEMORY_OVERHEAD_FLOOR + - estimate_recurrent_cache_bytes( + model_card, + resolved_backend=shard.resolved_backend, + llama_server_settings=shard.llama_server_settings, + ) + * fraction - fixed_memory_by_node.get(node_id, Memory()) ) node_tokens = max( diff --git a/src/skulk/shared/models/model_cards.py b/src/skulk/shared/models/model_cards.py index 98f198c0a..ecf1e68d2 100644 --- a/src/skulk/shared/models/model_cards.py +++ b/src/skulk/shared/models/model_cards.py @@ -45,6 +45,7 @@ SKULK_MODELS_DIRS, SKULK_OFFLINE, ) +from skulk.shared.models.gguf_memory import GgufCacheGeometry, qwen35_cache_geometry from skulk.shared.models.registry import ( EMBEDDED_REGISTRY_ROOT, RegistryAdvisory, @@ -72,7 +73,7 @@ # whatever that generator got wrong (the fresh-fleet audit found exactly that: # pre-#652-fix cards forcing serial in-process llama_cpp over the served # engine). -CARD_GENERATOR_REVISION: Final[int] = 2 +CARD_GENERATOR_REVISION: Final[int] = 3 # kinda ugly... # TODO: load search path from config.toml @@ -1903,6 +1904,14 @@ class ModelCard(CamelCaseModel): num_key_value_heads: PositiveInt | None = None """KV-head count for grouped-query attention, used in KV-cache sizing. ``None`` when unknown/not applicable.""" + gguf_cache_geometry: GgufCacheGeometry | None = None + """Artifact-derived attention and recurrent cache dimensions for GGUF admission. + + This is intrinsic model metadata. Slot count, speculation and runtime buffer + overhead are separate engine inputs. An absent value means unknown, never + zero recurrent cost; registry cards require a newly signed metadata revision + to add it without changing their accepted identity silently. + """ tasks: list[ModelTask] """The task types this model serves (``TextGeneration``, ``TextEmbedding``, ``TextToImage``, ``ImageToImage``, ``TextToSpeech``, ``SpeechToText``, @@ -2141,6 +2150,20 @@ def _validate_pipeline_split_limit(self) -> "ModelCard": ) return self + @model_validator(mode="after") + def _validate_gguf_cache_geometry(self) -> "ModelCard": + """Keep artifact-derived cache dimensions attached to the matching layer set.""" + geometry = self.gguf_cache_geometry + if geometry is not None and ( + self.gguf_file is None + or geometry.attention_layers + + geometry.recurrent_layers + + geometry.nextn_layers + != self.n_layers + ): + raise ValueError("GGUF cache geometry must match the artifact layer count") + return self + @field_validator("tasks", mode="before") @classmethod def _validate_tasks(cls, v: list[str | ModelTask]) -> list[ModelTask]: @@ -2308,11 +2331,9 @@ async def _fetch_gguf_from_hf( Sizes the weights from the GGUF file the runner would actually load (`select_gguf_file` picks the first sorted file + its shard group). - Structural fields come from `config.json` when the repo ships one (most - community GGUF repos do); a bare repo with no usable config.json has its - metadata read straight from the selected GGUF file's binary header via a - ranged read of the file start (#327), so neither path fabricates the - layer/hidden sizes placement's memory and KV-budget math depend on. + Structural fields come from the selected GGUF binary header through + ranged reads. Repository config remains a source of vision metadata, + but cannot override the artifact's layer count or cache dimensions. Stamps the llama.cpp backend tags so placement routes the model only to nodes with a llama.cpp engine and prefers a GPU backend. """ @@ -2339,40 +2360,21 @@ async def _fetch_gguf_from_hf( except (FileNotFoundError, ValidationError): config_data = None - if ( - config_data is not None - and config_data.layer_count - and (config_data.hidden_size) - ): - n_layers = config_data.layer_count - hidden_size = config_data.hidden_size - num_key_value_heads = config_data.num_key_value_heads - context_length = config_data.max_position_embeddings or 0 - else: - # No usable config.json: read the structural fields from the GGUF - # binary header. ``selected`` is the first shard, which carries the - # metadata block, so a ranged read of its start is enough. - reason = "absent" if config_data is None else "missing layer/hidden sizes" - logger.info( - f"GGUF repo {model_id} config.json {reason}; reading model " - f"metadata from the GGUF header of {selected}" + # A repository config can describe the base checkpoint rather than the + # selected GGUF (notably omitting embedded NextN layers). Read artifact + # dimensions even when that config exists; retain it for vision metadata. + from skulk.download.download_utils import range_read + + async def _fetch(offset: int, length: int) -> bytes: + return await range_read( + model_id, source_revision or "main", selected, offset, length ) - from skulk.download.download_utils import range_read - - async def _fetch(offset: int, length: int) -> bytes: - return await range_read( - model_id, - source_revision or "main", - selected, - offset, - length, - ) - fields = await read_gguf_structural_fields(_fetch) - n_layers = fields.n_layers - hidden_size = fields.hidden_size - num_key_value_heads = fields.num_key_value_heads - context_length = fields.context_length + fields = await read_gguf_structural_fields(_fetch) + n_layers = fields.n_layers + hidden_size = fields.hidden_size + num_key_value_heads = fields.num_key_value_heads + context_length = fields.context_length # Mark the card vision-capable so the llama.cpp runner loads the mmproj # projector (#128, #346). Prefer the rich config.json vision_config when @@ -2401,6 +2403,7 @@ async def _fetch(offset: int, length: int) -> bytes: # llama.cpp runs single-node in Skulk (no tensor parallelism). supports_tensor=False, num_key_value_heads=num_key_value_heads, + gguf_cache_geometry=fields.cache_geometry, context_length=context_length, tasks=[ModelTask.TextGeneration], quantization=_gguf_quant_label(selected), @@ -2989,6 +2992,7 @@ class GgufStructuralFields(NamedTuple): hidden_size: int num_key_value_heads: int | None context_length: int + cache_geometry: GgufCacheGeometry | None = None class _GgufHeaderReader: @@ -3103,6 +3107,7 @@ async def parse_structural_fields(self) -> GgufStructuralFields: architecture: str | None = None collected: dict[str, int] = {} + recurrent_layer_override = False def _has(suffixes: "tuple[str, ...]") -> bool: return architecture is not None and all( @@ -3117,6 +3122,8 @@ def _has(suffixes: "tuple[str, ...]") -> bool: # tokenizer array for best-effort fields that are simply absent. if key.startswith("tokenizer.") and _has(_GGUF_REQUIRED_SUFFIXES): break + if key.endswith(".attention.recurrent_layers"): + recurrent_layer_override = True value = await self.read_value(await self._u32()) if key == "general.architecture" and isinstance(value, str): architecture = value @@ -3124,7 +3131,10 @@ def _has(suffixes: "tuple[str, ...]") -> bool: collected[key] = value # Fast path: every wanted key (including best-effort ones) is in, # before any tokenizer array. - if _has(_GGUF_WANTED_SUFFIXES): + # Hybrid metadata follows the basic dimensions. Stopping at the + # first four scalars loses its fixed recurrent-state allocation and + # can also miss an explicit layer vector overriding the interval. + if architecture != "qwen35" and _has(_GGUF_WANTED_SUFFIXES): break if architecture is None: @@ -3145,6 +3155,15 @@ def field(suffix: str) -> int | None: hidden_size=hidden_size, num_key_value_heads=field(_GGUF_KEY_HEAD_COUNT_KV) or None, context_length=field(_GGUF_KEY_CONTEXT_LENGTH) or 0, + cache_geometry=qwen35_cache_geometry( + architecture, + { + key.removeprefix(architecture + "."): value + for key, value in collected.items() + if key.startswith(architecture + ".") + }, + has_recurrent_layer_override=recurrent_layer_override, + ), ) diff --git a/src/skulk/shared/models/tests/test_gguf_cards.py b/src/skulk/shared/models/tests/test_gguf_cards.py index 114e1f9b3..0e7805de1 100644 --- a/src/skulk/shared/models/tests/test_gguf_cards.py +++ b/src/skulk/shared/models/tests/test_gguf_cards.py @@ -95,6 +95,28 @@ def _factory( return _factory +def _mock_dense_header(monkeypatch: pytest.MonkeyPatch, layers: int = 32) -> None: + """Serve an exact GGUF header independently of the repository config.""" + from skulk.download import download_utils + + blob = _build_gguf( + [ + _kv_string("general.architecture", "llama"), + _kv_u32("llama.block_count", layers), + _kv_u32("llama.embedding_length", 4096), + _kv_u32("llama.attention.head_count_kv", 8), + _kv_u32("llama.context_length", 8192), + ] + ) + + async def read_range( + _model_id: object, _revision: str, _path: str, start: int, length: int + ) -> bytes: + return blob[start : start + length] + + monkeypatch.setattr(download_utils, "range_read", read_range) + + def test_gguf_weight_siblings_filters_gguf_and_mmproj( monkeypatch: pytest.MonkeyPatch, ) -> None: @@ -144,6 +166,7 @@ def test_shard_base_detection() -> None: async def test_fetch_gguf_card_stamps_both_llama_engines( monkeypatch: pytest.MonkeyPatch, ) -> None: + _mock_dense_header(monkeypatch) monkeypatch.setattr(model_cards, "model_info", _fake_model_info(["model-q4.gguf"])) async def _fake_config( @@ -186,6 +209,7 @@ async def _fake_config( async def test_fetch_gguf_card_honors_requested_file( monkeypatch: pytest.MonkeyPatch, ) -> None: + _mock_dense_header(monkeypatch) requested = "model-IQ3_XXS.gguf" monkeypatch.setattr( model_cards, @@ -482,6 +506,7 @@ def test_select_preferred_gguf_sharded_group() -> None: async def test_gguf_card_pins_selected_quant(monkeypatch: pytest.MonkeyPatch) -> None: + _mock_dense_header(monkeypatch, layers=16) monkeypatch.setattr( model_cards, "model_info", @@ -531,3 +556,18 @@ def test_default_gguf_selection_never_picks_companion_artifacts() -> None: assert "dspark" not in select_preferred_gguf(files) # A drafter-only repo (a published draft companion) still resolves. assert select_preferred_gguf([("gemma-mtp-draft-Q8_0.gguf", 1)]) + + +async def test_header_keeps_hybrid_fields_after_basic_structural_dimensions() -> None: + """The real header ordering must not trigger the old four-field early exit.""" + from skulk.shared.models.tests.test_gguf_memory import metadata + + blob = _build_gguf( + [_kv_string("general.architecture", "qwen35")] + + [_kv_u32("qwen35." + key, value) for key, value in metadata().items()] + + [_kv_string("tokenizer.ggml.model", "unused")] + ) + fields = await model_cards.read_gguf_structural_fields(_mem_fetch(blob)) + assert fields.cache_geometry == model_cards.qwen35_cache_geometry( + "qwen35", metadata() + ) diff --git a/src/skulk/shared/models/tests/test_gguf_memory.py b/src/skulk/shared/models/tests/test_gguf_memory.py new file mode 100644 index 000000000..60442d0ee --- /dev/null +++ b/src/skulk/shared/models/tests/test_gguf_memory.py @@ -0,0 +1,95 @@ +"""Cache accounting from artifact geometry and configured engine dimensions.""" + +import pytest +from pydantic import ValidationError + +from skulk.shared.models.gguf_memory import qwen35_cache_geometry + + +def metadata() -> dict[str, int]: + """Return the intrinsic dimensions of the pinned Qwen3.5 9B artifact.""" + return { + "block_count": 33, + "embedding_length": 4096, + "attention.head_count_kv": 4, + "attention.key_length": 256, + "attention.value_length": 256, + "context_length": 262144, + "ssm.conv_kernel": 4, + "ssm.inner_size": 4096, + "ssm.state_size": 128, + "ssm.group_count": 16, + "full_attention_interval": 4, + "nextn_predict_layers": 1, + } + + +def test_recurrent_allocation_matches_engine_tensor_dimensions() -> None: + """Reserve both FP32 state buffers and all 64 retained sequence rows.""" + geometry = qwen35_cache_geometry("qwen35", metadata()) + assert geometry is not None + assert geometry.attention_layers == 8 + assert geometry.recurrent_layers == 24 + assert geometry.nextn_layers == 1 + assert geometry.convolution_width == 24576 + assert geometry.recurrent_width == 524288 + assert geometry.recurrent_bytes(parallel_slots=16, rollback_depth=3) == 3372220416 + assert geometry.recurrent_bytes(parallel_slots=1, rollback_depth=0) == 52690944 + assert geometry.attention_bytes_per_token(embedded_mtp=False) == 32768 + assert geometry.attention_bytes_per_token(embedded_mtp=True) == 36864 + + +def test_geometry_uses_dimensions_without_model_size_or_quantization_guess() -> None: + """Different layer/width/slot/depth choices change their own tensor factors.""" + values = metadata() | { + "block_count": 17, + "ssm.inner_size": 2048, + "ssm.group_count": 8, + "attention.value_length": 128, + } + geometry = qwen35_cache_geometry("qwen35", values) + assert geometry is not None + assert geometry.recurrent_bytes(parallel_slots=8, rollback_depth=2) == 316145664 + assert geometry.attention_bytes_per_token(embedded_mtp=True) == 15360 + + +@pytest.mark.parametrize("architecture", ["qwen35moe", "qwen3next", "mamba", "llama"]) +def test_unmapped_architectures_do_not_inherit_qwen35_formula( + architecture: str, +) -> None: + """Similar names and partial fields cannot prove the same engine allocation.""" + assert qwen35_cache_geometry(architecture, metadata()) is None + + +def test_incomplete_or_overridden_layout_stays_unknown() -> None: + """Unknown dimensions and explicit layer vectors cannot become zero-cost state.""" + values = metadata() + del values["ssm.state_size"] + assert qwen35_cache_geometry("qwen35", values) is None + assert ( + qwen35_cache_geometry("qwen35", metadata(), has_recurrent_layer_override=True) + is None + ) + + +@pytest.mark.parametrize( + "field,value", + [ + ("nextn_predict_layers", 33), + ("nextn_predict_layers", -1), + ("ssm.conv_kernel", 0), + ("full_attention_interval", 0), + ], +) +def test_malformed_complete_geometry_is_rejected(field: str, value: int) -> None: + """Complete but invalid geometry is not accepted as a usable memory bound.""" + with pytest.raises(ValueError): + qwen35_cache_geometry("qwen35", metadata() | {field: value}) + + +def test_geometry_requires_strict_numeric_fields() -> None: + """Serialized intrinsic geometry cannot coerce strings into dimensions.""" + geometry = qwen35_cache_geometry("qwen35", metadata()) + assert geometry is not None + with pytest.raises(ValidationError): + type(geometry).model_validate(geometry.model_dump() | {"key_width": "1024"}) diff --git a/src/skulk/shared/types/profiling.py b/src/skulk/shared/types/profiling.py index 39e2097bd..966af44a8 100644 --- a/src/skulk/shared/types/profiling.py +++ b/src/skulk/shared/types/profiling.py @@ -9,6 +9,10 @@ import psutil from pydantic import UUID4, BaseModel, Field, field_serializer, field_validator +from skulk.shared.models.llama_server_settings import ( + LlamaServerSettings, + resolve_llama_server_settings, +) from skulk.shared.types.memory import Memory from skulk.shared.types.node_facts import CapabilityConflict from skulk.shared.types.thunderbolt import ThunderboltIdentifier @@ -313,6 +317,10 @@ class NodeResources(CamelCaseModel): default_factory=dict, description="Exact installed build identities keyed by engine and backend tag.", ) + llama_server_settings: LlamaServerSettings | None = Field( + default=None, + description="Observed serving controls needed to budget per-slot recurrent state.", + ) hardware_classes: frozenset[str] = Field( default_factory=frozenset, description="Open observed hardware identifiers for support constraints.", @@ -425,6 +433,14 @@ async def gather( return cls( backends=derivation.backends, engine_builds=engine_builds, + llama_server_settings=( + resolve_llama_server_settings(os.environ) + if any( + backend.startswith("llama_server") + for backend in derivation.backends + ) + else None + ), hardware_classes=hardware_class_inventory(facts), participation=participation, api_available=api_available, diff --git a/src/skulk/shared/types/worker/shards.py b/src/skulk/shared/types/worker/shards.py index 4ae401d2c..2141d2e5c 100644 --- a/src/skulk/shared/types/worker/shards.py +++ b/src/skulk/shared/types/worker/shards.py @@ -3,6 +3,7 @@ from pydantic import Field +from skulk.shared.models.llama_server_settings import LlamaServerSettings from skulk.shared.models.model_cards import ModelCard from skulk.utils.pydantic_ext import TaggedModel @@ -31,6 +32,10 @@ class BaseShardMetadata(TaggedModel): # lacked the node's resources at placement); the worker then falls back to # its local backend probe. See #330. resolved_backend: str | None = None + llama_server_settings: LlamaServerSettings | None = Field( + default=None, + description="Node serving settings captured by placement for memory admission.", + ) # Error handling; equivalent to monkey-patch, but we can't monkey-patch runner.py # This is kinda annoying because it allocates memory in the ShardMetadata object. Can be rethought after Shanghai. diff --git a/src/skulk/worker/main.py b/src/skulk/worker/main.py index c10ab09fd..c8dcb52aa 100644 --- a/src/skulk/worker/main.py +++ b/src/skulk/worker/main.py @@ -920,6 +920,8 @@ def _local_shard_fit_error( shard.model_card, self._shard_memory_fraction(shard), context_budget=kv_context, + resolved_backend=shard.resolved_backend, + llama_server_settings=shard.llama_server_settings, ) # On a discrete-GPU node the engine allocates from VRAM, not system RAM, # so size the guard against local usable VRAM or it would falsely refuse diff --git a/src/skulk/worker/runner/llama_server/runner.py b/src/skulk/worker/runner/llama_server/runner.py index 7596f2fbc..7d3e915cc 100644 --- a/src/skulk/worker/runner/llama_server/runner.py +++ b/src/skulk/worker/runner/llama_server/runner.py @@ -44,6 +44,11 @@ from skulk.shared.backends import LLAMA_SERVER_BIN_ENV from skulk.shared.constants import MAX_OUTPUT_TOKENS from skulk.shared.models.capabilities import resolve_model_capability_profile +from skulk.shared.models.llama_server_settings import ( + LLAMA_SERVER_DEFAULT_DRAFT_DEPTH, + LLAMA_SERVER_DEFAULT_PARALLEL, + resolve_llama_server_settings, +) from skulk.shared.models.model_cards import ModelCard, OutputParserType from skulk.shared.types.chunks import ErrorChunk, TokenChunk, ToolCallChunk from skulk.shared.types.common import CommandId, ModelId @@ -150,12 +155,7 @@ def _force_no_spec() -> bool: MTP on-vs-off throughput comparison and for debugging a misbehaving spec pairing; unset in normal operation. """ - return os.environ.get("SKULK_LLAMA_SERVER_FORCE_NO_SPEC", "").strip().lower() in ( - "1", - "true", - "yes", - "on", - ) + return not resolve_llama_server_settings(os.environ).speculation_enabled _LLAMA_SERVER_PARALLEL_ENV: Final = "SKULK_LLAMA_SERVER_PARALLEL" @@ -164,7 +164,7 @@ def _force_no_spec() -> bool: # Sixteen is the exercised fleet setting. A unified KV buffer keeps every slot's # advertised context window truthful without allocating N private caches; users # dominated by near-window prompts can still opt back to serial explicitly. -_DEFAULT_LLAMA_SERVER_PARALLEL: Final = 16 +_DEFAULT_LLAMA_SERVER_PARALLEL: Final = LLAMA_SERVER_DEFAULT_PARALLEL # An omitted OpenAI ``max_tokens`` value otherwise lets llama-server consume the # remainder of the shared KV pool. Bound it to the same normal-generation width # as Skulk's MLX path so aggregate admission has a finite reservation and one @@ -218,7 +218,7 @@ def _llama_server_parallel() -> int: f"using {_DEFAULT_LLAMA_SERVER_PARALLEL}" ) return _DEFAULT_LLAMA_SERVER_PARALLEL - return value + return resolve_llama_server_settings(os.environ).parallel_slots def _request_context_reservation( @@ -653,6 +653,13 @@ def __init__( # _spawn_server) keeps every slot's context at the full stamped window, # while the weighted gate below prevents their aggregate reservations # from exhausting the one shared pool. + settings = self.shard_metadata.llama_server_settings + if settings is not None and settings != resolve_llama_server_settings( + os.environ + ): + raise ValueError( + "llama-server settings changed since memory admission; re-place the instance" + ) effective_parallel = _effective_server_parallel(self.shard_metadata.model_card) self._init_concurrent_dispatch(effective_parallel, "llama-gen") if ( @@ -935,9 +942,14 @@ def _spawn_server( ) if draft_args is not None: cmd += ["--spec-type", flag] - n_max = getattr(runtime, "served_spec_n_max", None) - if n_max is not None: - cmd += ["--spec-draft-n-max", str(n_max)] + # Admission reserves rollback rows at this depth. Pass the + # shared default explicitly so an engine upgrade cannot + # silently change the allocation after placement. + n_max = ( + getattr(runtime, "served_spec_n_max", None) + or LLAMA_SERVER_DEFAULT_DRAFT_DEPTH + ) + cmd += ["--spec-draft-n-max", str(n_max)] cmd += draft_args self.server_log_path = ( diff --git a/website/docs/api-guide.md b/website/docs/api-guide.md index 645ffafc1..9673ca122 100644 --- a/website/docs/api-guide.md +++ b/website/docs/api-guide.md @@ -561,6 +561,16 @@ The mounted model must also declare `TextGeneration`; speech-only cards return ### Context-length limits +GGUF cards may include `gguf_cache_geometry`, derived from the selected artifact's +header, to distinguish per-token attention cache from fixed recurrent state. +For supported hybrid layouts, memory requirements include FP32 recurrent buffers +for each serving slot and speculative rollback row. Embedded MTP adds its own +attention layers without charging a second copy of the target weights. +`NodeResources.llama_server_settings` reports the slot count and speculation +override; placement captures those settings on each served shard. A settings +change requires a new placement before the runner can start. Missing geometry +retains the legacy estimate and is not proof that recurrent state costs zero. + Ordinary text requests prefer ready or running placements over placements still loading. Among equally ready placements, ordinary model instances take precedence over the resident steward, then active text-task counts balance requests. Retained diff --git a/website/docs/architecture-reference.md b/website/docs/architecture-reference.md index 833610019..6d6117572 100644 --- a/website/docs/architecture-reference.md +++ b/website/docs/architecture-reference.md @@ -12,6 +12,14 @@ This file is intentionally dense. If you find a stale fact, fix it inline rather ## Components +- **GGUF cache admission:** `shared/models/gguf_memory.py` resolves Qwen3.5 scalar + header geometry into fixed FP32 recurrent state and per-token FP16 attention + widths. `ModelCard.gguf_cache_geometry` is artifact metadata; unknown layouts + remain unknown. `NodeResources.llama_server_settings` carries configured slots + and speculation enablement. Placement copies it to served shard metadata; + `memory_estimate.py` charges rollback rows separately from weights/KV, and the + llama-server runner rejects local settings that differ from the stamp. + ### Master - **Role:** elects + acts as cluster coordinator; indexes events; plans instance placements; publishes snapshots diff --git a/website/docs/architecture.md b/website/docs/architecture.md index 5eabe6f49..d076aa3f5 100644 --- a/website/docs/architecture.md +++ b/website/docs/architecture.md @@ -293,6 +293,17 @@ A snapshot-bootstrap rollout has one operational rule: once a master starts comp ### Heterogeneous nodes and capability-aware placement +GGUF memory admission separates artifact geometry from node serving settings. +The selected header supplies attention and recurrent dimensions through +`GgufCacheGeometry`; generated cards use those artifact dimensions even when a +repository config describes a different base-layer count. For the supported +Qwen3.5 scalar layout, admission charges FP32 recurrent state across configured +slots and rollback rows, plus the target and embedded-MTP attention caches. +`NodeResources.llama_server_settings` advertises the existing environment +controls, and placement stamps them into shard metadata. The runner rejects a +changed stamp before launch. Geometry remains signed card content for registry +models; loading a card does not silently enrich or replace its accepted identity. + A cluster can mix node types: Apple Silicon nodes serving MLX models and non-Mac (for example AMD/Linux) nodes serving GGUF models through llama.cpp. Placement is capability-aware so each model runs only where it can. From bd673c826b0e46993c7c938c8a1fbbd192752bf2 Mon Sep 17 00:00:00 2001 From: tt92618 <11300465+ttupper92618@users.noreply.github.com> Date: Tue, 8 Sep 2026 03:46:25 -0500 Subject: [PATCH 2/3] feat: verify signed GGUF header projections for memory admission --- CLAUDE.md | 5 +- src/skulk/shared/models/model_cards.py | 59 +++++- src/skulk/shared/models/registry.py | 95 +++++++++- .../shared/models/registry_gguf_metadata.py | 86 +++++++++ .../shared/models/tests/test_registry.py | 179 +++++++++++++++++- website/docs/api-guide.md | 8 + website/docs/architecture-reference.md | 5 + website/docs/architecture.md | 12 +- 8 files changed, 441 insertions(+), 8 deletions(-) create mode 100644 src/skulk/shared/models/registry_gguf_metadata.py diff --git a/CLAUDE.md b/CLAUDE.md index 39efa2d25..c4a1a1118 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -113,7 +113,10 @@ artifact header. Qwen3.5 scalar metadata separates attention, recurrent and Next layers; recurrent FP32 buffers scale with slots and speculative rollback rows. NodeResources advertises `llama_server_settings`, placement stamps them on shards, and the runner refuses changed local settings before spawning. Registry geometry -must enter through a new signed card revision, never a silent runtime overlay. +is explicitly derived from the separately signed `v1/gguf-metadata.json` target, +bound to the catalog snapshot, signed role version and exact artifact. Canonical +cards remain unchanged for older clients. Runtime cards retain the header evidence +in `registry_gguf_metadata`; it participates in the full-card approval digest. Exact non-RPC placements resolve omitted backends from advertised compatible engines before memory admission; restored unstamped GPU-host shards reserve diff --git a/src/skulk/shared/models/model_cards.py b/src/skulk/shared/models/model_cards.py index ecf1e68d2..743594a6f 100644 --- a/src/skulk/shared/models/model_cards.py +++ b/src/skulk/shared/models/model_cards.py @@ -54,6 +54,7 @@ RegistryEngineSupportClaim, TufRegistryClient, ) +from skulk.shared.models.registry_gguf_metadata import RegistryGgufArtifactMetadata from skulk.shared.types.common import CommandId, ModelId from skulk.shared.types.memory import Memory from skulk.shared.types.text_generation import ReasoningEffort @@ -198,6 +199,28 @@ def registry_model_cards(catalog: RegistryCatalog) -> list["ModelCard"]: "is_custom": False, } ) + header_evidence = ( + catalog.gguf_metadata.artifacts.get(envelope.card_id) + if catalog.gguf_metadata is not None + else None + ) + payload["registry_gguf_metadata"] = header_evidence + if header_evidence is not None: + header = header_evidence.header + geometry = qwen35_cache_geometry( + header.architecture, + header.scalars, + has_recurrent_layer_override=header.has_recurrent_layer_override, + ) + payload["gguf_cache_geometry"] = geometry + if geometry is not None: + # Repository config can omit an embedded NextN block. Exact + # signed header facts govern this runtime projection, while + # canonical card bytes and their content-derived ID stay intact. + payload["n_layers"] = header.scalars["block_count"] + if "embedding_length" in header.scalars: + payload["hidden_size"] = header.scalars["embedding_length"] + payload["num_key_value_heads"] = header.scalars["attention.head_count_kv"] card = ModelCard.model_validate(payload) envelope_bundle = envelope.artifact.bundle card_bundle = card.artifact_bundle @@ -1909,8 +1932,8 @@ class ModelCard(CamelCaseModel): This is intrinsic model metadata. Slot count, speculation and runtime buffer overhead are separate engine inputs. An absent value means unknown, never - zero recurrent cost; registry cards require a newly signed metadata revision - to add it without changing their accepted identity silently. + zero recurrent cost. Registry geometry is an explicit projection of the + separately signed header target retained in ``registry_gguf_metadata``. """ tasks: list[ModelTask] """The task types this model serves (``TextGeneration``, ``TextEmbedding``, @@ -2058,6 +2081,38 @@ def _require_bundle_revision_and_selected_file(self) -> "ModelCard": """Exact signed artifact format used for support-matrix joins.""" registry_capability_claims: tuple[RegistryCapabilityClaim, ...] = () """Open signed model/artifact capability claims, independent of engines.""" + registry_gguf_metadata: RegistryGgufArtifactMetadata | None = None + """Exact signed header evidence used for this runtime geometry projection. + + It participates in the full-card authorization digest. A changed projection + cannot silently reuse an approval for a different memory contract. + """ + + @model_validator(mode="after") + def _validate_registry_gguf_metadata(self) -> "ModelCard": + """Keep retained signed header evidence and derived runtime geometry bound.""" + evidence = self.registry_gguf_metadata + if evidence is None: + return self + if ( + self.registry_card_id is None + or evidence.repository != str(self.source_repository or self.model_id) + or evidence.revision != self.source_revision + or evidence.selected_file != self.gguf_file + ): + raise ValueError( + "registry GGUF header evidence does not match runtime artifact" + ) + geometry = qwen35_cache_geometry( + evidence.header.architecture, + evidence.header.scalars, + has_recurrent_layer_override=evidence.header.has_recurrent_layer_override, + ) + if geometry != self.gguf_cache_geometry: + raise ValueError( + "runtime GGUF geometry disagrees with signed header evidence" + ) + return self @field_validator("registry_capability_claims", mode="before") @classmethod diff --git a/src/skulk/shared/models/registry.py b/src/skulk/shared/models/registry.py index b814abe6e..d578fba0f 100644 --- a/src/skulk/shared/models/registry.py +++ b/src/skulk/shared/models/registry.py @@ -8,10 +8,19 @@ from typing import Any, Literal, Self, cast from filelock import FileLock -from pydantic import BaseModel, ConfigDict, Field, field_validator, model_validator +from pydantic import ( + BaseModel, + ConfigDict, + Field, + PrivateAttr, + field_validator, + model_validator, +) from tuf.ngclient.updater import Updater from tuf.ngclient.urllib3_fetcher import Urllib3Fetcher +from skulk.shared.models.registry_gguf_metadata import RegistryGgufMetadata + EMBEDDED_REGISTRY_ROOT = Path(__file__).with_name("model_registry_root.json") """Public TUF trust root shipped inside the Skulk Python package.""" @@ -291,6 +300,31 @@ class RegistryCatalog(BaseModel): note: str cards: tuple[RegistryCard, ...] card_metadata: dict[str, RegistryCardMetadata] + _gguf_metadata: RegistryGgufMetadata | None = PrivateAttr(default=None) + + @property + def gguf_metadata(self) -> RegistryGgufMetadata | None: + """Return separately verified header facts, never an inline catalog field.""" + return self._gguf_metadata + + def with_gguf_metadata(self, metadata: RegistryGgufMetadata | None) -> Self: + """Bind auxiliary evidence without modifying canonical catalog bytes.""" + if metadata is not None: + if metadata.snapshot_id != self.snapshot_id: + raise ValueError("GGUF metadata catalog snapshot mismatch") + cards = {card.card_id: card for card in self.cards} + for card_id, entry in metadata.artifacts.items(): + card = cards.get(card_id) + if card is None or ( + card.artifact.format != "gguf" + or entry.repository != card.artifact.repository + or entry.revision != card.artifact.revision + or entry.selected_file != card.artifact.selected_file + ): + raise ValueError("GGUF metadata artifact identity mismatch") + copied = self.model_copy() + copied._gguf_metadata = metadata + return copied class RegistryAdvisory(BaseModel): @@ -559,6 +593,8 @@ def __init__( self._targets_dir = cache_dir / "targets" self._last_known_good_path = cache_dir / "last-known-good-catalog.json" self._cache_record_path = cache_dir / "last-known-good.json" + self._gguf_metadata_path = cache_dir / "last-known-good-gguf-metadata.json" + self._gguf_cache_record_path = cache_dir / "last-known-good-gguf-record.json" self._last_known_good_advisories_path = ( cache_dir / "last-known-good-advisories.json" ) @@ -739,11 +775,67 @@ def _refresh( downloaded = Path(updater.download_target(target)) payload = downloaded.read_bytes() catalog = RegistryCatalog.model_validate_json(payload, strict=False) + gguf_target = updater.get_targetinfo("v1/gguf-metadata.json") + gguf_metadata: RegistryGgufMetadata | None = None + if gguf_target is not None: + gguf_metadata = RegistryGgufMetadata.model_validate_json( + Path(updater.download_target(gguf_target)).read_bytes() + ) + trusted_targets = _TrustedTargetsMetadataVersion.model_validate_json( + (self._metadata_dir / "targets.json").read_bytes(), strict=False + ) + if gguf_metadata.target_version != trusted_targets.signed.version: + raise ValueError("GGUF metadata version does not match signed targets") + catalog = catalog.with_gguf_metadata(gguf_metadata) if catalog_validator is not None: catalog_validator(catalog) + self._write_verified_gguf_cache(catalog) self._write_verified_cache(payload, catalog) return catalog + def _write_verified_gguf_cache(self, catalog: RegistryCatalog) -> None: + """Retain a hash-bound auxiliary projection without changing old cache schemas.""" + payload = ( + catalog.gguf_metadata.model_dump_json().encode() + if catalog.gguf_metadata is not None + else b"null" + ) + record = _VerifiedCacheRecord( + sha256=hashlib.sha256(payload).hexdigest(), + verified_at=datetime.now(UTC), + snapshot_id=catalog.snapshot_id, + ) + self._atomic_write(self._gguf_metadata_path, payload) + self._atomic_write( + self._gguf_cache_record_path, record.model_dump_json().encode() + ) + + def _attach_cached_gguf_metadata(self, catalog: RegistryCatalog) -> RegistryCatalog: + """Recover only auxiliary facts tied to this exact cached catalog.""" + if ( + not self._gguf_cache_record_path.exists() + and not self._gguf_metadata_path.exists() + ): + # A cache created by an older reader has no auxiliary evidence. + return catalog + record = _VerifiedCacheRecord.model_validate_json( + self._gguf_cache_record_path.read_bytes(), strict=False + ) + if record.snapshot_id != catalog.snapshot_id: + raise ValueError("cached GGUF metadata snapshot mismatch") + if datetime.now(UTC) - record.verified_at.astimezone(UTC) > timedelta( + days=self._max_stale_days + ): + raise ValueError("cached GGUF metadata is too old") + payload = self._gguf_metadata_path.read_bytes() + if hashlib.sha256(payload).hexdigest() != record.sha256: + raise ValueError("cached GGUF metadata hash mismatch") + return catalog.with_gguf_metadata( + None + if payload == b"null" + else RegistryGgufMetadata.model_validate_json(payload) + ) + def _write_verified_cache(self, payload: bytes, catalog: RegistryCatalog) -> None: """Atomically retain bytes that the updater just verified.""" record = _VerifiedCacheRecord( @@ -805,6 +897,7 @@ def _load_last_known_good( catalog = RegistryCatalog.model_validate_json(payload, strict=False) if catalog.snapshot_id != record.snapshot_id: raise ValueError("last-known-good snapshot identity mismatch") + catalog = self._attach_cached_gguf_metadata(catalog) if catalog_validator is not None: catalog_validator(catalog) return catalog diff --git a/src/skulk/shared/models/registry_gguf_metadata.py b/src/skulk/shared/models/registry_gguf_metadata.py new file mode 100644 index 000000000..355f6f366 --- /dev/null +++ b/src/skulk/shared/models/registry_gguf_metadata.py @@ -0,0 +1,86 @@ +"""Strict wire contract for independently signed GGUF header metadata.""" + +from datetime import datetime +from typing import Annotated, Final, Literal, final + +from pydantic import BaseModel, ConfigDict, Field, field_validator + +GGUF_SCALAR_FIELDS: Final = frozenset( + { + "block_count", + "embedding_length", + "attention.head_count", + "attention.head_count_kv", + "attention.key_length", + "attention.value_length", + "full_attention_interval", + "nextn_predict_layers", + "ssm.conv_kernel", + "ssm.inner_size", + "ssm.state_size", + "ssm.group_count", + "ssm.time_step_rank", + } +) + + +@final +class RegistryGgufHeaderMetadata(BaseModel): + """ + Bounded facts read from one exact artifact's complete GGUF metadata area. + + The digest covers the file prefix through the last metadata value, including + the GGUF preamble. It is evidence identity, not a full-artifact checksum. + Consumers own architecture-specific interpretation and engine allocation. + """ + + model_config = ConfigDict(frozen=True, strict=True, extra="forbid") + + architecture: str = Field(min_length=1, max_length=128) + scalars: dict[str, int] = Field( + max_length=len(GGUF_SCALAR_FIELDS), + description="Architecture-relative integer GGUF fields, without defaults.", + ) + has_recurrent_layer_override: bool = Field( + description="Whether an explicit attention.recurrent_layers field exists.", + ) + metadata_sha256: str = Field(pattern=r"^[0-9a-f]{64}$") + metadata_bytes: int = Field(ge=24, le=16 * 1024 * 1024) + + @field_validator("scalars") + @classmethod + def validate_scalars(cls, values: dict[str, int]) -> dict[str, int]: + """Keep only the bounded, unsigned integer vocabulary this version reads.""" + if not set(values).issubset(GGUF_SCALAR_FIELDS): + raise ValueError("unknown GGUF scalar field") + if any(value < 0 or value > 2**64 - 1 for value in values.values()): + raise ValueError("GGUF scalar lies outside unsigned 64-bit range") + return values + + +@final +class RegistryGgufArtifactMetadata(BaseModel): + """Header evidence bound to one immutable repository file.""" + + model_config = ConfigDict(frozen=True, strict=True, extra="forbid") + + repository: str = Field(min_length=3, max_length=512) + revision: str = Field(pattern=r"^[0-9a-f]{40}$") + selected_file: str = Field(min_length=1, max_length=4096) + header: RegistryGgufHeaderMetadata + + +@final +class RegistryGgufMetadata(BaseModel): + """One verified auxiliary target bound to its catalog and signed role version.""" + + model_config = ConfigDict(frozen=True, strict=True, extra="forbid") + + schema_version: Literal[1] + snapshot_id: str = Field(min_length=1, max_length=120) + target_version: int = Field(ge=1) + generated_at: datetime + artifacts: dict[ + Annotated[str, Field(pattern=r"^card_[a-z2-7]{52}$")], + RegistryGgufArtifactMetadata, + ] diff --git a/src/skulk/shared/models/tests/test_registry.py b/src/skulk/shared/models/tests/test_registry.py index d3a64b5ba..68db84060 100644 --- a/src/skulk/shared/models/tests/test_registry.py +++ b/src/skulk/shared/models/tests/test_registry.py @@ -29,6 +29,7 @@ RegistryUnavailableError, TufRegistryClient, ) +from skulk.shared.models.registry_gguf_metadata import RegistryGgufMetadata from skulk.shared.types.common import ModelId from skulk.shared.types.memory import Memory from skulk.shared.types.worker.shards import PipelineShardMetadata @@ -832,8 +833,8 @@ def __init__(self, **kwargs: object) -> None: def refresh(self) -> None: pass - def get_targetinfo(self, _: str) -> object: - return object() + def get_targetinfo(self, target_path: str) -> object | None: + return object() if target_path == "v1/catalog.json" else None def download_target(self, _: object) -> str: return str(payload_path) @@ -1439,3 +1440,177 @@ async def ignore_progress(*_: object) -> None: assert observed == [ModelId("org/multi")] assert path.name == "org--multi@q4" + + +def _gguf_metadata_payload() -> bytes: + return json.dumps( + { + "schema_version": 1, + "snapshot_id": "snapshot_1_test", + "target_version": 7, + "generated_at": "2026-09-08T00:00:00Z", + "artifacts": { + f"card_{'a' * 52}": { + "repository": "org/multi-gguf", + "revision": "b" * 40, + "selected_file": "model-Q4_K_M.gguf", + "header": { + "architecture": "qwen35", + "scalars": { + "block_count": 33, + "embedding_length": 4096, + "attention.head_count_kv": 4, + "attention.key_length": 256, + "attention.value_length": 256, + "ssm.conv_kernel": 4, + "ssm.inner_size": 4096, + "ssm.state_size": 128, + "ssm.group_count": 16, + "nextn_predict_layers": 1, + }, + "has_recurrent_layer_override": False, + "metadata_sha256": "c" * 64, + "metadata_bytes": 10943341, + }, + } + }, + } + ).encode() + + +def test_signed_header_projects_geometry_without_changing_canonical_card() -> None: + """Auxiliary artifact facts correct runtime dimensions and bind approvals.""" + original = RegistryCatalog.model_validate_json(_catalog_payload(), strict=False) + metadata = RegistryGgufMetadata.model_validate_json(_gguf_metadata_payload()) + catalog = original.with_gguf_metadata(metadata) + assert original.gguf_metadata is None + assert catalog.model_dump_json() == original.model_dump_json() + card = registry_model_cards(catalog)[0] + assert card.registry_card_id == original.cards[0].card_id + assert original.cards[0].card["n_layers"] == 4 + assert card.n_layers == 33 + assert card.hidden_size == 4096 + assert card.gguf_cache_geometry is not None + assert ( + card.gguf_cache_geometry.recurrent_bytes(parallel_slots=16, rollback_depth=3) + == 3372220416 + ) + assert card.registry_gguf_metadata is not None + restored = ModelCard.model_validate_json(card.model_dump_json()) + assert restored == card + assert model_cards_module.authorized_model_card_digest( + card + ) != model_cards_module.authorized_model_card_digest( + registry_model_cards(original)[0] + ) + + +@pytest.mark.parametrize("field", ["repository", "revision", "selected_file"]) +def test_header_metadata_rejects_cross_artifact_binding(field: str) -> None: + """A same-name artifact or copied card ID cannot supply another file's geometry.""" + catalog = RegistryCatalog.model_validate_json(_catalog_payload(), strict=False) + metadata = RegistryGgufMetadata.model_validate_json(_gguf_metadata_payload()) + key = catalog.cards[0].card_id + changed = metadata.artifacts[key].model_copy( + update={field: "d" * 40 if field == "revision" else "other/file"} + ) + with pytest.raises(ValueError, match="artifact identity mismatch"): + catalog.with_gguf_metadata( + metadata.model_copy(update={"artifacts": {key: changed}}) + ) + + +def test_header_metadata_rejects_cross_snapshot_and_numeric_coercion() -> None: + """Even signed metadata must match the exact snapshot and strict scalar types.""" + catalog = RegistryCatalog.model_validate_json(_catalog_payload(), strict=False) + metadata = RegistryGgufMetadata.model_validate_json(_gguf_metadata_payload()) + with pytest.raises(ValueError, match="snapshot mismatch"): + catalog.with_gguf_metadata(metadata.model_copy(update={"snapshot_id": "other"})) + payload = _gguf_metadata_payload().replace( + b'"block_count": 33', b'"block_count": true' + ) + with pytest.raises(ValueError): + RegistryGgufMetadata.model_validate_json(payload) + + +def test_client_recovers_snapshot_bound_header_cache_and_rejects_tampering( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + """Network loss retains verified dimensions, while bad auxiliary bytes fail closed.""" + catalog_path = tmp_path / "catalog.json" + catalog_path.write_bytes(_catalog_payload()) + header_path = tmp_path / "gguf-metadata.json" + header_path.write_bytes(_gguf_metadata_payload()) + root = tmp_path / "root.json" + root.write_text("{}") + offline = False + + class MetadataUpdater: + def __init__(self, **kwargs: object) -> None: + self.metadata_dir = Path(cast("str", kwargs["metadata_dir"])) + + def refresh(self) -> None: + if offline: + raise OSError("offline") + (self.metadata_dir / "targets.json").write_text('{"signed":{"version":7}}') + + def get_targetinfo(self, path: str) -> str: + return str(header_path if path == "v1/gguf-metadata.json" else catalog_path) + + def download_target(self, target: str) -> str: + return target + + monkeypatch.setattr(registry_module, "Updater", MetadataUpdater) + client = TufRegistryClient( + base_url="https://registry.example/", + cache_dir=tmp_path / "cache", + embedded_root=root, + timeout_seconds=1, + max_stale_days=30, + ) + first = client.load_catalog(registry_model_cards) + assert first.gguf_metadata is not None + # A rejected newer target cannot displace previously verified facts. + header_path.write_bytes( + _gguf_metadata_payload().replace(b'"target_version": 7', b'"target_version": 8') + ) + assert ( + client.load_catalog(registry_model_cards).gguf_metadata == first.gguf_metadata + ) + offline = True + recovered = client.load_catalog(registry_model_cards) + assert registry_model_cards(recovered)[0] == registry_model_cards(first)[0] + (tmp_path / "cache/last-known-good-gguf-metadata.json").write_text("tampered") + with pytest.raises(RegistryUnavailableError): + client.load_catalog(registry_model_cards) + + +def test_installed_same_card_keeps_current_signed_geometry(tmp_path: Path) -> None: + """An older installed sidecar cannot hide newly verified same-artifact metadata.""" + original = RegistryCatalog.model_validate_json(_catalog_payload(), strict=False) + metadata = RegistryGgufMetadata.model_validate_json(_gguf_metadata_payload()) + old_card = registry_model_cards(original)[0] + current = registry_model_cards(original.with_gguf_metadata(metadata))[0] + (tmp_path / "model-Q4_K_M.gguf").write_bytes(b"weights") + (tmp_path / ".skulk-source-revision").write_text(f"{old_card.source_revision}\n") + record = build_installed_card_record(tmp_path, old_card) + assert record.verification == "registry_verified" + prior_cache = dict(model_cards_module._card_cache) + prior_current = dict(model_cards_module._registry_current_cards) + prior_installed = dict(model_cards_module._installed_card_cache) + prior_current_ids = dict(model_cards_module._installed_current_registry_ids) + try: + model_cards_module._card_cache[current.model_id] = current + model_cards_module._registry_current_cards[current.model_id] = current + model_cards_module._apply_installed_card_snapshot([record], scan_version=10**9) + assert model_cards_module._card_cache[current.model_id] == current + finally: + model_cards_module._card_cache.clear() + model_cards_module._card_cache.update(prior_cache) + model_cards_module._registry_current_cards.clear() + model_cards_module._registry_current_cards.update(prior_current) + model_cards_module._installed_card_cache.clear() + model_cards_module._installed_card_cache.update(prior_installed) + model_cards_module._installed_current_registry_ids.clear() + model_cards_module._installed_current_registry_ids.update(prior_current_ids) diff --git a/website/docs/api-guide.md b/website/docs/api-guide.md index 9673ca122..824d05f45 100644 --- a/website/docs/api-guide.md +++ b/website/docs/api-guide.md @@ -563,6 +563,14 @@ The mounted model must also declare `TextGeneration`; speech-only cards return GGUF cards may include `gguf_cache_geometry`, derived from the selected artifact's header, to distinguish per-token attention cache from fixed recurrent state. +Registry-backed runtime cards also expose `registry_gguf_metadata`, the separately +signed exact-file header evidence used for this projection. It names the repository, +immutable revision, selected file, architecture, scalar dimensions, inspected-prefix +length and digest, and any recurrent-layer override. This is not a whole-weight +checksum. Skulk verifies its catalog snapshot and signed target version before use; +canonical card IDs remain unchanged, but a changed runtime projection changes the +full-card approval digest. Older installed metadata for the same canonical card +does not override the current verified projection. For supported hybrid layouts, memory requirements include FP32 recurrent buffers for each serving slot and speculative rollback row. Embedded MTP adds its own attention layers without charging a second copy of the target weights. diff --git a/website/docs/architecture-reference.md b/website/docs/architecture-reference.md index 6d6117572..c6b01cd24 100644 --- a/website/docs/architecture-reference.md +++ b/website/docs/architecture-reference.md @@ -19,6 +19,11 @@ This file is intentionally dense. If you find a stale fact, fix it inline rather and speculation enablement. Placement copies it to served shard metadata; `memory_estimate.py` charges rollback rows separately from weights/KV, and the llama-server runner rejects local settings that differ from the stamp. + `registry_gguf_metadata.py` defines the separate signed header target. + `TufRegistryClient` binds it to the same catalog snapshot and targets-role + version. `registry_model_cards` retains the exact-file evidence as + `ModelCard.registry_gguf_metadata` and projects supported structural dimensions; + canonical cards are unchanged and full-card approvals bind the projection. ### Master diff --git a/website/docs/architecture.md b/website/docs/architecture.md index d076aa3f5..f88df61cc 100644 --- a/website/docs/architecture.md +++ b/website/docs/architecture.md @@ -301,8 +301,16 @@ Qwen3.5 scalar layout, admission charges FP32 recurrent state across configured slots and rollback rows, plus the target and embedded-MTP attention caches. `NodeResources.llama_server_settings` advertises the existing environment controls, and placement stamps them into shard metadata. The runner rejects a -changed stamp before launch. Geometry remains signed card content for registry -models; loading a card does not silently enrich or replace its accepted identity. +changed stamp before launch. For registry models, the TUF client reads the separate +`v1/gguf-metadata.json` target from the same verified metadata refresh. Its target +version, catalog snapshot, card ID and exact artifact must match before header +facts become runtime geometry. The canonical catalog and card bytes are unchanged, +so older readers can ignore the auxiliary target. Runtime cards retain the exact +header evidence in `registry_gguf_metadata`; supported header dimensions correct +base-config counts that omit NextN blocks. This projection participates in the +full-card authorization digest, so changed geometry cannot reuse a different +approved memory contract. Bounded hash- and snapshot-bound cached evidence survives +temporary registry outages; mismatched or corrupted cache pairs are rejected. A cluster can mix node types: Apple Silicon nodes serving MLX models and non-Mac (for example AMD/Linux) nodes serving GGUF models through llama.cpp. From e78f978b7443932d98c1a53ce4ccf88be8985d25 Mon Sep 17 00:00:00 2001 From: tt92618 <11300465+ttupper92618@users.noreply.github.com> Date: Tue, 8 Sep 2026 05:09:40 -0500 Subject: [PATCH 3/3] fix: reject unresolved hybrid placement before serving admission --- src/skulk/master/placement.py | 11 ++++ .../tests/test_recurrent_memory_admission.py | 53 ++++++++++++++++++- 2 files changed, 62 insertions(+), 2 deletions(-) diff --git a/src/skulk/master/placement.py b/src/skulk/master/placement.py index 544eb57b1..4062847cb 100644 --- a/src/skulk/master/placement.py +++ b/src/skulk/master/placement.py @@ -316,6 +316,17 @@ def _stamp_llama_server_settings( # Donors provide memory to the driver's one context; they do not # launch a served instance with their own parallel-slot controls. continue + if ( + shard.resolved_backend is None + and shard.model_card.gguf_cache_geometry is not None + and not has_rpc_donors + ): + # A worker may resolve this shard to a served engine after telemetry + # catches up. Default or caller-supplied slots cannot authorize that + # later process's recurrent allocation. + raise PlacementError( + "Backend telemetry is required for recurrent memory admission" + ) if not has_rpc_donors and ( shard.resolved_backend is None or not shard.resolved_backend.startswith("llama_server") diff --git a/src/skulk/master/tests/test_recurrent_memory_admission.py b/src/skulk/master/tests/test_recurrent_memory_admission.py index 81a2552e9..7290436b3 100644 --- a/src/skulk/master/tests/test_recurrent_memory_admission.py +++ b/src/skulk/master/tests/test_recurrent_memory_admission.py @@ -2,13 +2,18 @@ import pytest -from skulk.master.placement import PlacementError, add_instance_to_placements +from skulk.master.placement import ( + PlacementError, + add_instance_to_placements, + place_instance, +) from skulk.master.placement_utils import ( filter_cycles_by_memory, get_shard_assignments_for_llama_rpc, usable_vram_by_node, ) -from skulk.master.tests.conftest import create_node_memory +from skulk.master.tests.conftest import create_node_memory, create_node_network +from skulk.master.tests.test_placement import place_instance_command from skulk.shared.models.gguf_memory import qwen35_cache_geometry from skulk.shared.models.llama_server_settings import LlamaServerSettings from skulk.shared.models.memory_estimate import ( @@ -241,3 +246,47 @@ def test_rpc_captures_only_the_driver_serving_settings() -> None: result.runner_to_shard[result.node_to_runner[donor]].llama_server_settings is None ) + + +def test_unresolved_exact_hybrid_placement_requires_backend_observation() -> None: + """Legacy RAM-only admission cannot approve a hybrid with unknown serving slots.""" + node = NodeId() + instance = hybrid_instance(node, 1) + assignments = instance.shard_assignments + instance = instance.model_copy( + update={ + "shard_assignments": assignments.model_copy( + update={ + "runner_to_shard": { + runner: shard.model_copy(update={"resolved_backend": None}) + for runner, shard in assignments.runner_to_shard.items() + } + } + ) + } + ) + with pytest.raises(PlacementError, match="Backend telemetry.*recurrent"): + add_instance_to_placements( + CreateInstance(instance=instance), + Topology(), + {}, + {node: create_node_memory(Memory.from_gb(64).in_bytes)}, + ) + + +def test_auto_hybrid_placement_waits_for_backend_observation() -> None: + """A known topology and RAM reading cannot substitute for serving controls.""" + node = NodeId() + topology = Topology() + topology.add_node(node) + card = next( + iter(hybrid_instance(node).shard_assignments.runner_to_shard.values()) + ).model_card + with pytest.raises(PlacementError, match="Backend telemetry.*recurrent"): + place_instance( + place_instance_command(card), + topology, + {}, + {node: create_node_memory(Memory.from_gb(64).in_bytes)}, + {node: create_node_network()}, + )