diff --git a/CLAUDE.md b/CLAUDE.md index 672d893d0..c4a1a1118 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -108,6 +108,16 @@ 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 +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 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..4062847cb 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,55 @@ 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 ( + 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") + ): + 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 +968,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 +1040,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 +1253,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..7290436b3 --- /dev/null +++ b/src/skulk/master/tests/test_recurrent_memory_admission.py @@ -0,0 +1,292 @@ +"""Hybrid-model admission must reserve actual slots, rollback state and KV together.""" + +import pytest + +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, 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 ( + 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 + ) + + +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()}, + ) 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..743594a6f 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, @@ -53,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 @@ -72,7 +74,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 @@ -197,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 @@ -1903,6 +1927,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 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``, ``TextToImage``, ``ImageToImage``, ``TextToSpeech``, ``SpeechToText``, @@ -2049,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 @@ -2141,6 +2205,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 +2386,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 +2415,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 +2458,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 +3047,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 +3162,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 +3177,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 +3186,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 +3210,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/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_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/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/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..824d05f45 100644 --- a/website/docs/api-guide.md +++ b/website/docs/api-guide.md @@ -561,6 +561,24 @@ 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. +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. +`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..c6b01cd24 100644 --- a/website/docs/architecture-reference.md +++ b/website/docs/architecture-reference.md @@ -12,6 +12,19 @@ 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. + `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 - **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..f88df61cc 100644 --- a/website/docs/architecture.md +++ b/website/docs/architecture.md @@ -293,6 +293,25 @@ 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. 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. Placement is capability-aware so each model runs only where it can.