From bb2bd94beb9302557d91994e224eb5ff8189d3e8 Mon Sep 17 00:00:00 2001 From: Brian Searls Date: Sun, 28 Jun 2026 18:30:40 +0000 Subject: [PATCH] Revert "Re-model compute_fabric as a namespace need<->opportunity connector; shape DERIVED from the .dag program, no heavy/workload/locality (replaces discarded #5866) (#5889)" This reverts commit 989ba32a3ae72be2135eae78ef8daa929a1954b3. --- dsl/product/compute_fabric.dag | 1399 ++++++++++++++++- .../compute_fabric_resource_witness_test.dag | 24 + .../claim/compute_fabric_witness_test.dag | 20 - ...ector_compute_fabric_import_repro_test.dag | 4 +- 4 files changed, 1338 insertions(+), 109 deletions(-) create mode 100644 dsl/test/claim/compute_fabric_resource_witness_test.dag delete mode 100644 dsl/test/claim/compute_fabric_witness_test.dag diff --git a/dsl/product/compute_fabric.dag b/dsl/product/compute_fabric.dag index 75c0a842107..80a4809a6b1 100644 --- a/dsl/product/compute_fabric.dag +++ b/dsl/product/compute_fabric.dag @@ -1,134 +1,1359 @@ module product.compute_fabric -import std.types { NonEmptyStr, Bool, List } -import std.measure { HardwareThreadCount, hardware_thread_count, hardware_thread_count_value, measure_le } +import std.nat { Nat } +import std.types { NonEmptyStr, ContentHash, Bool, Int, LogicalTime, List } +import extdeps.vendor { Vendor } +import extdeps.hardware { Hardware } +import extdeps.vendor.nvidia { nvidia } +import extdeps.vendor.amd { amd } +import extdeps.vendor.ampere { ampere } +import std.os.types { + FileSystemSemantics, + Linux, + OperatingSystemSurface, + Posix, + PosixProcess, + ProcessSemantics, +} +import std.float { Float } +import std.measure { + Measure, + Time, + Currency, + Micro, + One, + Nano, + ByteSize, + MoneyAmount, + HardwareThreadCount, + Count, + byte_size, + byte_size_count, + measure_count, + hardware_thread_count, + hardware_thread_count_value, + measure_scale_fraction_ceil, + hertz, + watt, + Bandwidth, +} +import std.currency { CurrencyCode } +import std.disposition { Disposition, Scaffold, SingleAuthority } +import std.decl_ref { DeclarationRef, WholeDeclaration } +import std.verification { NanosecondDuration, nanosecond_duration, nanosecond_duration_count } +import std.realization_schedule { CostAccount } +import std.realization_measurement { cost_account_measured_from_time, cost_account_measured_from_time_and_space } +import std.content_hash { + content_hash_atom, + content_hash_combine, + content_hash_tagged, +} +import std.cache_identity { CacheInterfaceId, ArtifactKindId, ExecutionReceiptDigest } +import extdeps.cpu { + CpuFacts, + cpu_facts_architecture, + cpu_facts_sustained_per_thread_hz, + cpu_facts_threads, +} +import extdeps.cpu.types { + CpuDeploymentFacts, + CpuModelCatalogRow, + CpuSocket, + SocketPackage, +} +import extdeps.storage.types { NvmeDeviceCatalogRow, PcieLink } +import extdeps.ctrl.jobserver { CtrlJobserverConfig } +import extdeps.toolchain.types { + Architecture, + Aarch64, + X86_64, + Armv7, + Riscv64, + Wasm32, +} +import product.placement_supply { HostIdentity, PlacementSupplyRow } +import std.baseboard { ServerBaseboard } +import std.realization_width { bounded_host_spawn_width, memory_aware_spawn_width, int_min } + +type ComputeWitness = Node +type Outcome = Node + +type UpsertInputRef = Node +type ChangeSet = Node +type ArtifactRef = Node +type ArtifactSpec = Node +type WorkAction = Node +type WorkGraph = Node +type Partitioner = Node +type Reducer = Node +type EffectBoundary = Node +type SymbolicCost = Node +type Associativity = Node +type Commutativity = Node +type Idempotency = Node +type IdentityElement = Node +type InstructionSet = Node + +type NetworkLocality = + Loopback + | Lan + | Internet + | WslNat + | WslBridged + | DockerBridge + | TailscaleOverlay +type NetworkAddressability = Unicast | Multicast | Broadcast +type NetworkEgressClass = + None + | CratesIo + | InternalMirror + | Unrestricted +type IsolationBoundary = SharedHostHome | PerJobFilesystem | ContainerHermetic | Vm +type PersistenceKind = Ephemeral | Persistent | CachedMirror +type MemoryKind = Dram | Hbm | UnifiedShared +type StorageMedium = Nvme | Ssd | Hdd | NetworkAttached | EphemeralFs +type LatencyClass = UltraLow | Low | Medium | High +type GpuRuntime = Cuda | Rocm | Metal | OpenCl +type ContainerRuntime = Docker | Podman | Gvisor | Ubicloud | GcloudRun +type CostClass = OwnedMarginalZero | PerSecondBilled | Metered | DeveloperAttention +type OomSignalClass = KillProcess | Throttle | ReportOnly + +type Duration = Measure + +type MoneyMicros = MoneyAmount + +fn money_micros(count: Nat) -> MoneyMicros { + MoneyMicros { count: count } +} + +fn money_micros_count(m: MoneyMicros) -> Nat { + measure_count(m: m) +} + +type ProviderIdentity = NonEmptyStr where brand("ProviderIdentity") + +type ComputeArtifactLocality + = InProcess + | PerRunnerColocation + | PerHostColocation + | CrossHostFetch + +type StorageMount { + mount_path: NonEmptyStr + device: StorageDevice + lifecycle: PersistenceKind +} + +type MemoryFacts { capacity: ByteSize, memory_kind: MemoryKind } + +type GpuComputeCapability { + vendor: Vendor + major: Int + minor: Int +} + +type UnsubstantiatedGpuCapabilityLabel { vendor: Vendor } -type Fabric = NonEmptyStr where brand("Fabric") +type GpuCapabilityLabelOutcome + = GpuCapabilityLabelOk { label: String } + | GpuCapabilityLabelRejected { reason: UnsubstantiatedGpuCapabilityLabel } + +fn gpu_compute_capability_sm_label(cap: GpuComputeCapability) -> GpuCapabilityLabelOutcome { + if cap.vendor == nvidia { + GpuCapabilityLabelOk { + label: concat("sm_", concat(to_string(cap.major), concat("_", to_string(cap.minor)))) + } + } else { + if cap.vendor == amd { + GpuCapabilityLabelOk { + label: concat("gfx", concat(to_string(cap.major), to_string(cap.minor))) + } + } else { + GpuCapabilityLabelRejected { + reason: UnsubstantiatedGpuCapabilityLabel { vendor: cap.vendor } + } + } + } +} -type HardRequirements { threads: HardwareThreadCount } +type GpuFacts { + vendor: Vendor + model: NonEmptyStr + compute_capability: GpuComputeCapability? + memory: MemoryFacts + supported_runtimes: List +} -type Shape { hard: HardRequirements } +type AcceleratorFacts { vendor: Vendor, model: NonEmptyStr } -fn shape_covers(offered: Shape, needed: Shape) -> Bool { - measure_le(a: needed.hard.threads, b: offered.hard.threads) +type MemoryDevice { + capacity: ByteSize + memory_kind: MemoryKind + bandwidth: Bandwidth? } -type Program { - id: NonEmptyStr - hard: HardRequirements +type StorageDevice + = NvmeDevice { + catalog: NvmeDeviceCatalogRow + serial: NonEmptyStr? + persistence: PersistenceKind + } + | BareDevice { + capacity: ByteSize + medium: StorageMedium + serial: NonEmptyStr? + read_bandwidth: Bandwidth? + write_bandwidth: Bandwidth? + persistence: PersistenceKind + } + +fn nvme_storage_device(catalog: NvmeDeviceCatalogRow, persistence: PersistenceKind) -> StorageDevice { + NvmeDevice { + catalog: catalog, + serial: none, + persistence: persistence, + } } -fn derive_shape(p: Program) -> Shape { - Shape { hard: p.hard } +fn storage_device_capacity(d: StorageDevice) -> ByteSize { + match d { + NvmeDevice { catalog: c, serial: _, persistence: _ } => c.capacity_bytes + BareDevice { capacity: cap, medium: _, serial: _, read_bandwidth: _, write_bandwidth: _, persistence: _ } => cap + } } -type ComputeNeed { - fabric: Fabric - program: Program +fn storage_device_medium(d: StorageDevice) -> StorageMedium { + match d { + NvmeDevice { catalog: _, serial: _, persistence: _ } => Nvme + BareDevice { capacity: _, medium: m, serial: _, read_bandwidth: _, write_bandwidth: _, persistence: _ } => m + } } -type Opportunity { - fabric: Fabric - offers: Shape - tag: NonEmptyStr +fn storage_device_read_bandwidth(d: StorageDevice) -> Bandwidth? { + match d { + NvmeDevice { catalog: c, serial: _, persistence: _ } => Present { value: c.sequential_read } + BareDevice { capacity: _, medium: _, serial: _, read_bandwidth: bw, write_bandwidth: _, persistence: _ } => bw + } } -type Connection - = Bound { opportunity: Opportunity } - | Pending - | Unmet { reason: NonEmptyStr } +fn storage_device_write_bandwidth(d: StorageDevice) -> Bandwidth? { + match d { + NvmeDevice { catalog: c, serial: _, persistence: _ } => Present { value: c.sequential_write } + BareDevice { capacity: _, medium: _, serial: _, read_bandwidth: _, write_bandwidth: bw, persistence: _ } => bw + } +} -fn connect(need: ComputeNeed, opportunities: List) -> Connection { - let shape = derive_shape(p: need.program) +fn storage_device_pcie_link(d: StorageDevice) -> PcieLink? { + match d { + NvmeDevice { catalog: c, serial: _, persistence: _ } => Present { value: c.pcie_link } + BareDevice { capacity: _, medium: _, serial: _, read_bandwidth: _, write_bandwidth: _, persistence: _ } => none + } +} + +fn storage_device_serial(d: StorageDevice) -> NonEmptyStr? { + match d { + NvmeDevice { catalog: _, serial: s, persistence: _ } => s + BareDevice { capacity: _, medium: _, serial: s, read_bandwidth: _, write_bandwidth: _, persistence: _ } => s + } +} + +fn storage_device_persistence(d: StorageDevice) -> PersistenceKind { + match d { + NvmeDevice { catalog: _, serial: _, persistence: p } => p + BareDevice { capacity: _, medium: _, serial: _, read_bandwidth: _, write_bandwidth: _, persistence: p } => p + } +} + +type NetworkInterface { + addressability: NetworkAddressability + bandwidth: Bandwidth? + latency_class: LatencyClass? + locality: NetworkLocality +} + +type ExecutionSurface { + os: OperatingSystemSurface + isolation: IsolationBoundary + container_runtime: ContainerRuntime? + toolchains: List + mounted_storage: List + network: List +} + +type ProcessorKind + = CpuProcessor { cpu: CpuFacts } + | GpuProcessor { gpu: GpuFacts } + | AcceleratorProcessor { accelerator: AcceleratorFacts } + +type ComputeHost { + identity: HostIdentity + baseboard: ServerBaseboard? + processors: List + memory: List + storage: List + network_interfaces: List +} + +type ProcessorCostModel { placeholder: Bool } +type ProviderCostModel { placeholder: Bool } + +type CostModel { + cost_class: CostClass + rates: ProviderCostModel +} + +type ComputeSupplyFacts { + physical: ComputeHost + execution: ExecutionSurface + cost: CostModel? + observed_performance: List +} + +type ProviderConstraint + = HostJobserverFifo { fifo_path: NonEmptyStr, config_ref: CtrlJobserverConfig } + | SharedHomeRoot { path: NonEmptyStr } + | MaxConcurrentRunners { cap: Int } + | OomBehavior { signal: OomSignalClass } + +type AvailabilityWindow { open: LogicalTime, close: LogicalTime? } +type CostEstimate { model: CostModel, unit_price: MoneyMicros, currency: CurrencyCode } + +fn per_second_billed_cost_quote(unit_price: MoneyMicros, currency: CurrencyCode) -> CostEstimate { + CostEstimate { + model: CostModel { + cost_class: PerSecondBilled, + rates: ProviderCostModel { placeholder: true }, + }, + unit_price: unit_price, + currency: currency, + } +} + +fn cost_estimate_unit_price_micros(quote: CostEstimate) -> Nat { + money_micros_count(m: quote.unit_price) +} + +type ComputeOffer { + provider: ProviderIdentity + supply: ComputeSupplyFacts + available_window: AvailabilityWindow + cost_quote: CostEstimate? + constraints: List +} + +type MachineView { label: NonEmptyStr, host: HostIdentity } + +fn machine_view(host: ComputeHost, surface: ExecutionSurface) -> MachineView { + MachineView { + label: machine_view_label(host: host, surface: surface), + host: host.identity, + } +} + +fn machine_view_label(host: ComputeHost, surface: ExecutionSurface) -> NonEmptyStr { + match surface.os.distro_or_product { + Present { value: distro } => distro + Absent => host.identity + } +} + +fn placement_supply_row(identity: HostIdentity, cpu: CpuFacts, ram_bytes_total: ByteSize) -> PlacementSupplyRow { + PlacementSupplyRow { + identity: identity, + hardware_threads: cpu_facts_threads(cpu), + clock_hz: cpu_facts_sustained_per_thread_hz(cpu), + ram_bytes: ram_bytes_total, + } +} + +fn compute_host_primary_cpu(host: ComputeHost) -> CpuFacts? { + match host.processors.first() { + Absent => none + Present { value: kind } => + match kind { + CpuProcessor { cpu: c } => Present { value: c } + GpuProcessor { gpu: _ } => none + AcceleratorProcessor { accelerator: _ } => none + } + } +} + +fn compute_host_ram_bytes_total(host: ComputeHost) -> ByteSize { + byte_size( + host.memory + |> fold( + init: 0, + f: (acc, dev) => acc + byte_size_count(dev.capacity), + ), + ) +} + +fn offer_placement_supply_row(offer: ComputeOffer) -> PlacementSupplyRow? { + match compute_host_primary_cpu(host: offer.supply.physical) { + Absent => none + Present { value: cpu } => + Present { + value: placement_supply_row( + identity: offer.supply.physical.identity, + cpu: cpu, + ram_bytes_total: compute_host_ram_bytes_total(host: offer.supply.physical), + ), + } + } +} + +fn hardware_threads_honor_demand(supply: HardwareThreadCount, demand: HardwareThreadCount) -> Bool { + hardware_thread_count_value(supply) >= hardware_thread_count_value(demand) +} + +fn cpu_threads_satisfies(supply: PlacementSupplyRow, demand: CpuRequirement?) -> Bool { + match demand { + Absent => true + Present { value: req } => + hardware_threads_honor_demand(supply: supply.hardware_threads, demand: req.min_threads) + } +} + +fn cpu_architecture_satisfies(host: ComputeHost, demand: CpuRequirement?) -> Bool { + match demand { + Absent => true + Present { value: req } => + match req.architecture { + Absent => true + Present { value: arch } => + match compute_host_primary_cpu(host: host) { + Absent => false + Present { value: cpu } => cpu_facts_architecture(cpu: cpu) == arch + } + } + } +} + +fn memory_bytes_honor_demand(supply: ByteSize, demand: ByteSize) -> Bool { + byte_size_count(supply) >= byte_size_count(demand) +} + +fn memory_requirement_satisfies(supply: PlacementSupplyRow, demand: MemoryRequirement?) -> Bool { + match demand { + Absent => true + Present { value: req } => + memory_bytes_honor_demand(supply: supply.ram_bytes, demand: req.min_bytes) + } +} + +fn requires_placement_facts(demand: WorkDemand) -> Bool { + match demand.resources.cpu { + Present { value: _ } => true + Absent => + match demand.resources.memory { + Present { value: _ } => true + Absent => false + } + } +} + +fn satisfied_dimensions_for(demand: WorkDemand) -> List { + let parallelism_dims = match demand.parallelism { + IndependentShards { shard_count: _ } => [DemandParallelism] + SingleWorkItem => [] + DependencyGraphParallel { graph: _ } => [] + PartitionedReduce { + partitioner: _, + map: _, + reduce: _, + laws: _, + } => [] + } + concat( + [DemandIsolation], + concat( + parallelism_dims, + concat( + match demand.resources.cpu { + Absent => [] + Present { value: _ } => [DemandCpu] + }, + match demand.resources.memory { + Absent => [] + Present { value: _ } => [DemandMemory] + }, + ), + ), + ) +} + +fn work_demand_shard_count(demand: WorkDemand) -> Int { + match demand.parallelism { + SingleWorkItem => 1 + IndependentShards { shard_count } => shard_count + DependencyGraphParallel { graph: _ } => 1 + PartitionedReduce { + partitioner: _, + map: _, + reduce: _, + laws: _, + } => 1 + } +} + +type ResourceEnvelope { + cpu: CpuRequirement? + gpu: GpuRequirement? + memory: MemoryRequirement? + storage: StorageRequirement? + network: NetworkRequirement? +} + +type CpuRequirement { min_threads: HardwareThreadCount, architecture: Architecture? } +type GpuRequirement { min_vram: ByteSize, runtimes: List } +type MemoryRequirement { min_bytes: ByteSize } +type StorageRequirement { min_bytes: ByteSize, persistence: PersistenceKind } +type NetworkRequirement { egress: NetworkEgressClass?, ambient_allowed: Bool } +type OperatingSystemRequirement { surface: OperatingSystemSurface } +type IsolationRequirement { boundary: IsolationBoundary } +type ToolchainRequirement { name: NonEmptyStr, version: NonEmptyStr } + +type ArtifactLocalityRequirement { + artifact_kind: ArtifactKindId + locality: ComputeArtifactLocality +} + +type ToolchainEnvIsolation + = SharedHomeAcrossJobs + | PerJobCargoHome { cargo_home_var: NonEmptyStr, rustup_home_var: NonEmptyStr } + | HermeticContainer + +type BuildWrapperKind = CtrlBuild | BakedInImage | NoBuildWrapper + +type ToolchainCapability { + name: NonEmptyStr + version: NonEmptyStr + env_isolation: ToolchainEnvIsolation + build_wrapper: BuildWrapperKind? + linked_cache_interface: CacheInterfaceId? +} + +type InputSizeAxis + = WitnessCount + | SourceNodeCount + | CorpusNodeCount + +type InputBound { + axis: InputSizeAxis + max: Measure +} + +type InputEnvelope + = BoundedInput { bounds: List } + | EnvelopeUnknown + +type AdmissionVerdict + = Admitted + | RefusedOverEnvelope { axis: InputSizeAxis } + | RefusedUndeclared + +fn find_input_bound_for_axis(declared: List, axis: InputSizeAxis) -> InputBound? { fold( - opportunities, - init: Unmet { reason: "no_matching_opportunity" }, - f: (acc, opp) => + declared, + init: none, + f: (acc, b) => match acc { - Bound { opportunity: _ } => acc - Pending => acc - Unmet { reason: _ } => - if opp.fabric == need.fabric && shape_covers(offered: opp.offers, needed: shape) { - Bound { opportunity: opp } - } else { - acc + Present { value: _ } => acc + Absent => if b.axis == axis { Present { value: b } } else { none } + } + ) +} + +fn input_admitted(envelope: InputEnvelope, actual: List) -> AdmissionVerdict { + match envelope { + EnvelopeUnknown => RefusedUndeclared + BoundedInput { bounds: declared } => + fold( + actual, + init: Admitted, + f: (verdict, actual_bound) => + match verdict { + RefusedOverEnvelope { axis: _ } => verdict + RefusedUndeclared => verdict + Admitted => + match find_input_bound_for_axis(declared: declared, axis: actual_bound.axis) { + Absent => RefusedOverEnvelope { axis: actual_bound.axis } + Present { value: decl_bound } => + if actual_bound.max.count > decl_bound.max.count { + RefusedOverEnvelope { axis: actual_bound.axis } + } else { + Admitted + } + } } + ) + } +} + +type WorkDemand { + resources: ResourceEnvelope + os: OperatingSystemRequirement? + isolation: IsolationRequirement + toolchains: List + parallelism: ParallelismShape + data_locality: List + effects: List + input_envelope: InputEnvelope +} + +fn projected_concurrent_peak(width: HardwareThreadCount, per_shard_peak: ByteSize) -> ByteSize { + measure_scale_fraction_ceil(m: per_shard_peak, num: hardware_thread_count_value(width), den: 1) +} + +fn parallel_run_demand_envelope(width: HardwareThreadCount, per_shard_peak: ByteSize) -> ResourceEnvelope { + ResourceEnvelope { + cpu: Present { value: CpuRequirement { min_threads: width, architecture: none } }, + gpu: none, + memory: Present { + value: MemoryRequirement { + min_bytes: projected_concurrent_peak(width: width, per_shard_peak: per_shard_peak) } + }, + storage: none, + network: none, + } +} + +fn parallel_run_demand_envelope_from_budget( + shard_count: Int, + hardware_threads: HardwareThreadCount, + memory_budget: ByteSize, + per_shard_peak: ByteSize, + conservative_fallback_width: Int +) -> ResourceEnvelope { + parallel_run_demand_envelope( + width: memory_aware_spawn_width( + shard_count: shard_count, + hardware_threads: hardware_threads, + memory_budget: memory_budget, + per_shard_peak: per_shard_peak, + conservative_fallback_width: conservative_fallback_width + ), + per_shard_peak: per_shard_peak ) } -data example_fabric: Fabric = "ci-fleet" +fn demand_envelope_memory(envelope: ResourceEnvelope) -> ByteSize { + match envelope.memory { + Absent => byte_size(0) + Present { value: req } => req.min_bytes + } +} -data example_small_program: Program = Program { - id: "small-prog", - hard: HardRequirements { threads: hardware_thread_count(4) }, +fn demand_envelope_thread_count(envelope: ResourceEnvelope) -> Nat { + match envelope.cpu { + Absent => 0 + Present { value: req } => hardware_thread_count_value(req.min_threads) + } } -data example_large_program: Program = Program { - id: "large-prog", - hard: HardRequirements { threads: hardware_thread_count(64) }, +fn demand_envelope_fits_budget(envelope: ResourceEnvelope, budget: ByteSize) -> Bool { + memory_bytes_honor_demand(supply: budget, demand: demand_envelope_memory(envelope: envelope)) } -data example_opportunity: Opportunity = Opportunity { - fabric: example_fabric, - offers: Shape { hard: HardRequirements { threads: hardware_thread_count(128) } }, - tag: "srv1", +type ParallelismShape + = SingleWorkItem + | IndependentShards { shard_count: Int } + | DependencyGraphParallel { graph: WorkGraph } + | PartitionedReduce { + partitioner: Partitioner + map: WorkAction + reduce: Reducer + laws: ReducerLaws + } + +type ReducerLaws { + associative: ComputeWitness + commutative: ComputeWitness? + identity: IdentityElement? + idempotent: ComputeWitness? } -data example_small_opportunity: Opportunity = Opportunity { - fabric: example_fabric, - offers: Shape { hard: HardRequirements { threads: hardware_thread_count(8) } }, - tag: "srv1", +type ExecutionBudget { + expected_cost: SymbolicCost + provider_model: ProviderCostModel + watchdog: WatchdogLimit +} + +type WatchdogLimit { max_wall: Duration } + +type DemandDimension + = DemandCpu + | DemandGpu + | DemandMemory + | DemandStorage + | DemandNetwork + | DemandOs + | DemandIsolation + | DemandToolchains + | DemandParallelism + | DemandDataLocality + | DemandEffects + +type MissingDemandFact { dimension: DemandDimension, required: NonEmptyStr } + +type ComputeLeaseWitness { + provider: ProviderIdentity + satisfied: List } -fn witness_need_is_only_fabric_plus_program() -> Bool { - let need = ComputeNeed { fabric: example_fabric, program: example_small_program } - hardware_thread_count_value(t: need.program.hard.threads) == 4 +type ComputeLeaseEligibility + = Eligible { witness: ComputeLeaseWitness } + | Rejected { reason: MissingDemandFact } + +type AllocationReceipt { + provider: ProviderIdentity + scope: NonEmptyStr + acquired_at: LogicalTime } -fn witness_shape_is_derived_from_program() -> Bool { - let small_shape = derive_shape(p: example_small_program) - let large_shape = derive_shape(p: example_large_program) - small_shape.hard.threads != large_shape.hard.threads +type ComputeLease { + offer: ComputeOffer + demand: WorkDemand + allocation: AllocationReceipt + eligibility: ComputeWitness } -fn witness_fabric_binds_need_to_opportunity() -> Bool { - let need = ComputeNeed { fabric: example_fabric, program: example_small_program } - match connect(need: need, opportunities: [example_opportunity]) { - Bound { opportunity: _ } => true - Pending => false - Unmet { reason: _ } => false +fn isolation_boundary_ordinal(boundary: IsolationBoundary) -> Int { + match boundary { + SharedHostHome => 0 + PerJobFilesystem => 1 + ContainerHermetic => 2 + Vm => 3 } } -fn witness_exec_and_ci_are_just_programs() -> Bool { - let exec_prog = Program { - id: "exec-cmd", - hard: HardRequirements { threads: hardware_thread_count(4) }, +fn isolation_boundary_label(boundary: IsolationBoundary) -> NonEmptyStr { + match boundary { + SharedHostHome => "Vm/ContainerHermetic/PerJobFilesystem/SharedHostHome" + PerJobFilesystem => "Vm/ContainerHermetic/PerJobFilesystem" + ContainerHermetic => "Vm/ContainerHermetic" + Vm => "Vm" } - let ci_job = Program { - id: "ci-job", - hard: HardRequirements { threads: hardware_thread_count(4) }, +} + +fn isolation_satisfies(supply: IsolationBoundary, demand: IsolationBoundary) -> Bool { + isolation_boundary_ordinal(boundary: supply) >= isolation_boundary_ordinal(boundary: demand) +} + +fn satisfies(offer: ComputeOffer, demand: WorkDemand) -> ComputeLeaseEligibility { + if !isolation_satisfies(supply: offer.supply.execution.isolation, demand: demand.isolation.boundary) { + Rejected { + reason: MissingDemandFact { + dimension: DemandIsolation, + required: isolation_boundary_label(boundary: demand.isolation.boundary), + }, + } + } else if !requires_placement_facts(demand: demand) { + Eligible { + witness: ComputeLeaseWitness { + provider: offer.provider, + satisfied: [DemandIsolation], + }, + } + } else { + match offer_placement_supply_row(offer: offer) { + Absent => + Rejected { + reason: MissingDemandFact { + dimension: DemandCpu, + required: "primary_cpu", + }, + } + Present { value: row } => + if !cpu_threads_satisfies(supply: row, demand: demand.resources.cpu) { + Rejected { + reason: MissingDemandFact { + dimension: DemandCpu, + required: "min_threads", + }, + } + } else if !cpu_architecture_satisfies(host: offer.supply.physical, demand: demand.resources.cpu) { + Rejected { + reason: MissingDemandFact { + dimension: DemandCpu, + required: "cpu_architecture", + }, + } + } else if !memory_requirement_satisfies(supply: row, demand: demand.resources.memory) { + Rejected { + reason: MissingDemandFact { + dimension: DemandMemory, + required: "min_bytes", + }, + } + } else { + Eligible { + witness: ComputeLeaseWitness { + provider: offer.provider, + satisfied: satisfied_dimensions_for(demand: demand), + }, + } + } + } } - let need_exec = ComputeNeed { fabric: example_fabric, program: exec_prog } - let need_ci = ComputeNeed { fabric: example_fabric, program: ci_job } - match connect(need: need_exec, opportunities: [example_opportunity]) { - Bound { opportunity: _ } => - match connect(need: need_ci, opportunities: [example_opportunity]) { - Bound { opportunity: _ } => true - Pending => false - Unmet { reason: _ } => false - } - Pending => false - Unmet { reason: _ } => false +} + +type WorkUnitId = NonEmptyStr where brand("WorkUnitId") + +type WorkUnit { + id: WorkUnitId + inputs: List> + action: WorkAction + output: ArtifactSpec + requirements: WorkDemand +} + +type MeasurementConfidence + = SingleSample + | Range { low: Float, high: Float } + | DistributionSummary { mean: Float, variance: Float } + +type PerformanceReceipt { + provider: ProviderIdentity + work_shape: WorkUnitId + cost: CostAccount + cpu_seconds: Float? + sample_count: Int + measurement_context: ExecutionSurface + confidence: MeasurementConfidence + cache_state_summary: NonEmptyStr? +} + +fn performance_receipts_total_time(receipts: List) -> NanosecondDuration { + nanosecond_duration( + count: fold( + receipts, + init: 0, + f: (acc, receipt) => acc + measure_count(m: receipt.cost.time) + ) + ) +} + +fn cost_account_from_performance_receipts(receipts: List) -> CostAccount { + cost_account_measured_from_time_and_space( + time: performance_receipts_total_time(receipts: receipts), + space: byte_size(count: fold( + receipts, + init: 0, + f: (acc, receipt) => acc + byte_size_count(receipt.cost.space) + )) + ) +} + +data host_v1_interpreter_provider: ProviderIdentity = "host:v1-interpreter" + +data host_v1_interpreter_execution_surface: ExecutionSurface = ExecutionSurface { + os: OperatingSystemSurface { + kernel: Linux, + distro_or_product: none, + version: none, + filesystem_semantics: Posix, + process_semantics: PosixProcess, + }, + isolation: SharedHostHome, + container_runtime: none, + toolchains: [], + mounted_storage: [], + network: [], +} + +fn performance_receipt_host_single_sample( + work_shape: WorkUnitId, + wall_time: NanosecondDuration, + space: ByteSize, +) -> PerformanceReceipt { + PerformanceReceipt { + provider: host_v1_interpreter_provider, + work_shape: work_shape, + cost: cost_account_measured_from_time_and_space(time: wall_time, space: space), + cpu_seconds: none, + sample_count: 1, + measurement_context: host_v1_interpreter_execution_surface, + confidence: SingleSample, + cache_state_summary: none, + } +} + +type PricingSource + = VendorDocument { citation: NonEmptyStr } + | ObservedBill { receipt_id: NonEmptyStr } + | NegotiatedRate { contract_ref: NonEmptyStr } + +type AmortizationScope + = PerLease + | PerCalendarMonth + | PerWorkShape + +type CostReceipt { + provider: ProviderIdentity + cost_class: CostClass + amount: MoneyMicros? + pricing_source: PricingSource + amortization_scope: AmortizationScope? +} + +type ExecutionReceipt { + work: WorkUnit + lease: ComputeLease + output: Outcome> + performance: PerformanceReceipt + cost: CostReceipt + started_at: LogicalTime + finished_at: LogicalTime +} + +data example_host_identity: HostIdentity = "example-host" +data example_shared_host_provider: ProviderIdentity = "example-shared-host" +data example_vm_provider: ProviderIdentity = "example-vm" +data example_availability_open: LogicalTime = "example:always-open" + +data example_linux_os_surface: OperatingSystemSurface = OperatingSystemSurface { + kernel: Linux, + distro_or_product: none, + version: none, + filesystem_semantics: Posix, + process_semantics: PosixProcess, +} + +data example_shared_host_execution_surface: ExecutionSurface = ExecutionSurface { + os: example_linux_os_surface, + isolation: SharedHostHome, + container_runtime: none, + toolchains: [], + mounted_storage: [], + network: [], +} + +data example_vm_execution_surface: ExecutionSurface = ExecutionSurface { + os: example_linux_os_surface, + isolation: Vm, + container_runtime: none, + toolchains: [], + mounted_storage: [], + network: [], +} + +data example_compute_host: ComputeHost = ComputeHost { + identity: example_host_identity, + baseboard: none, + processors: [], + memory: [], + storage: [], + network_interfaces: [], +} + +data example_owned_cost: CostModel = CostModel { + cost_class: OwnedMarginalZero, + rates: ProviderCostModel { placeholder: true }, +} + +data example_shared_host_offer: ComputeOffer = ComputeOffer { + provider: example_shared_host_provider, + supply: ComputeSupplyFacts { + physical: example_compute_host, + execution: example_shared_host_execution_surface, + cost: Present { value: example_owned_cost }, + observed_performance: [], + }, + available_window: AvailabilityWindow { open: example_availability_open, close: none }, + cost_quote: none, + constraints: [], +} + +data example_vm_offer: ComputeOffer = ComputeOffer { + provider: example_vm_provider, + supply: ComputeSupplyFacts { + physical: example_compute_host, + execution: example_vm_execution_surface, + cost: Present { value: example_owned_cost }, + observed_performance: [], + }, + available_window: AvailabilityWindow { open: example_availability_open, close: none }, + cost_quote: none, + constraints: [], +} + +data vm_isolation_work_demand: WorkDemand = WorkDemand { + resources: ResourceEnvelope { + cpu: none, + gpu: none, + memory: none, + storage: none, + network: none, + }, + os: none, + isolation: IsolationRequirement { boundary: Vm }, + toolchains: [], + parallelism: SingleWorkItem, + data_locality: [], + effects: [], + input_envelope: EnvelopeUnknown, +} + +fn witness_shared_host_home_rejects_vm() -> Bool { + isolation_satisfies(supply: SharedHostHome, demand: Vm) == false +} + +fn witness_shared_host_home_admits_shared_host_home() -> Bool { + isolation_satisfies(supply: SharedHostHome, demand: SharedHostHome) +} + +fn witness_vm_admits_shared_host_home() -> Bool { + isolation_satisfies(supply: Vm, demand: SharedHostHome) +} + +fn witness_shared_host_offer_rejects_vm_demand() -> Bool { + match satisfies(offer: example_shared_host_offer, demand: vm_isolation_work_demand) { + Rejected { reason: _ } => true + Eligible { witness: _ } => false + } +} + +fn witness_vm_offer_admits_vm_demand() -> Bool { + match satisfies(offer: example_vm_offer, demand: vm_isolation_work_demand) { + Eligible { witness: _ } => true + Rejected { reason: _ } => false + } +} + +data example_resource_cpu_catalog: CpuModelCatalogRow = CpuModelCatalogRow { + model_number: "GENERIC-ARM-64C", + oem_listing_id: none, + vendor: ampere, + architecture: Aarch64, + cores: 64, + threads: hardware_thread_count(64), + nominal_sustained_per_thread_hz: hertz(3000000000), + socket: CpuSocket { package: LandGridArray, contact_count: 4926 }, + tdp_watts: watt(250), +} + +data example_resource_cpu: CpuFacts = CpuFacts { + catalog: example_resource_cpu_catalog, + deployment: CpuDeploymentFacts { sustained_per_thread_hz: hertz(3000000000) }, +} + +data example_resource_host: ComputeHost = ComputeHost { + identity: "example-resource-host", + baseboard: none, + processors: [CpuProcessor { cpu: example_resource_cpu }], + memory: [MemoryDevice { capacity: byte_size(68719476736), memory_kind: Dram, bandwidth: none }], + storage: [], + network_interfaces: [], +} + +data example_resource_offer: ComputeOffer = ComputeOffer { + provider: "example-resource-offer", + supply: ComputeSupplyFacts { + physical: example_resource_host, + execution: example_shared_host_execution_surface, + cost: Present { value: example_owned_cost }, + observed_performance: [], + }, + available_window: AvailabilityWindow { open: example_availability_open, close: none }, + cost_quote: none, + constraints: [], +} + +data example_thread_work_demand: WorkDemand = WorkDemand { + resources: ResourceEnvelope { + cpu: Present { + value: CpuRequirement { + min_threads: hardware_thread_count(128), + architecture: Present { value: Aarch64 }, + }, + }, + gpu: none, + memory: none, + storage: none, + network: none, + }, + os: none, + isolation: IsolationRequirement { boundary: SharedHostHome }, + toolchains: [], + parallelism: SingleWorkItem, + data_locality: [], + effects: [], + input_envelope: EnvelopeUnknown, +} + +data example_memory_work_demand: WorkDemand = WorkDemand { + resources: ResourceEnvelope { + cpu: none, + gpu: none, + memory: Present { value: MemoryRequirement { min_bytes: byte_size(137438953472) } }, + storage: none, + network: none, + }, + os: none, + isolation: IsolationRequirement { boundary: SharedHostHome }, + toolchains: [], + parallelism: SingleWorkItem, + data_locality: [], + effects: [], + input_envelope: EnvelopeUnknown, +} + +data example_feasible_resource_work_demand: WorkDemand = WorkDemand { + resources: ResourceEnvelope { + cpu: Present { + value: CpuRequirement { + min_threads: hardware_thread_count(2), + architecture: Present { value: Aarch64 }, + }, + }, + gpu: none, + memory: Present { value: MemoryRequirement { min_bytes: byte_size(1073741824) } }, + storage: none, + network: none, + }, + os: none, + isolation: IsolationRequirement { boundary: SharedHostHome }, + toolchains: [], + parallelism: SingleWorkItem, + data_locality: [], + effects: [], + input_envelope: EnvelopeUnknown, +} + +fn witness_resource_offer_rejects_insufficient_threads() -> Bool { + match satisfies(offer: example_resource_offer, demand: example_thread_work_demand) { + Rejected { reason: r } => r.dimension == DemandCpu + Eligible { witness: _ } => false + } +} + +fn witness_resource_offer_rejects_insufficient_memory() -> Bool { + match satisfies(offer: example_resource_offer, demand: example_memory_work_demand) { + Rejected { reason: r } => r.dimension == DemandMemory + Eligible { witness: _ } => false + } +} + +fn witness_resource_offer_admits_feasible_demand() -> Bool { + match satisfies(offer: example_resource_offer, demand: example_feasible_resource_work_demand) { + Eligible { witness: w } => + w.satisfied.length() == 3 + && w.satisfied.first() == DemandIsolation + && w.satisfied.skip(n: 1).first() == DemandCpu + && w.satisfied.skip(n: 2).first() == DemandMemory + Rejected { reason: _ } => false + } +} + +fn witness_run_demand_envelope_fits_budget() -> Bool { + let envelope = parallel_run_demand_envelope_from_budget( + shard_count: 7, + hardware_threads: hardware_thread_count(128), + memory_budget: byte_size(67108864000), + per_shard_peak: byte_size(2100028416), + conservative_fallback_width: 4 + ) + demand_envelope_thread_count(envelope: envelope) == 7 + && demand_envelope_fits_budget(envelope: envelope, budget: byte_size(67108864000)) +} + +fn witness_high_demand_shrinks_envelope_and_fits() -> Bool { + let envelope = parallel_run_demand_envelope_from_budget( + shard_count: 7, + hardware_threads: hardware_thread_count(128), + memory_budget: byte_size(67108864000), + per_shard_peak: byte_size(20000000000), + conservative_fallback_width: 4 + ) + demand_envelope_thread_count(envelope: envelope) == 2 + && demand_envelope_fits_budget(envelope: envelope, budget: byte_size(67108864000)) +} + +fn witness_fail_open_memory_blind_demand_exceeds_budget() -> Bool { + let blind_width = bounded_host_spawn_width( + shard_count: 7, + hardware_threads: hardware_thread_count(128) + ) + let blind_envelope = parallel_run_demand_envelope( + width: blind_width, + per_shard_peak: byte_size(20000000000) + ) + demand_envelope_fits_budget(envelope: blind_envelope, budget: byte_size(67108864000)) == false +} + +fn witness_unreadable_budget_uses_conservative_demand() -> Bool { + let envelope = parallel_run_demand_envelope_from_budget( + shard_count: 7, + hardware_threads: hardware_thread_count(128), + memory_budget: byte_size(0), + per_shard_peak: byte_size(2100028416), + conservative_fallback_width: 4 + ) + demand_envelope_thread_count(envelope: envelope) == 4 +} + +data execution_receipt_digest_schema: NonEmptyStr = "execution-receipt-canonical-v1" + +fn isolation_boundary_digest(boundary: IsolationBoundary) -> ContentHash { + content_hash_atom(value: isolation_boundary_label(boundary: boundary)) +} + +fn cpu_architecture_digest(arch: Architecture) -> ContentHash { + match arch { + X86_64 => content_hash_atom(value: "arch:X86_64") + X86 => content_hash_atom(value: "arch:X86") + Aarch64 => content_hash_atom(value: "arch:Aarch64") + Arm => content_hash_atom(value: "arch:Arm") + Armv7 => content_hash_atom(value: "arch:Armv7") + Mips => content_hash_atom(value: "arch:Mips") + Mipsel => content_hash_atom(value: "arch:Mipsel") + Mips64 => content_hash_atom(value: "arch:Mips64") + Mips64el => content_hash_atom(value: "arch:Mips64el") + Riscv64 => content_hash_atom(value: "arch:Riscv64") + Wasm32 => content_hash_atom(value: "arch:Wasm32") + } +} + +fn optional_cpu_requirement_digest(req: CpuRequirement?) -> ContentHash { + match req { + Absent => content_hash_atom(value: "cpu:none") + Present { value: r } => + content_hash_combine( + content_hash_atom(value: "cpu"), + content_hash_combine( + content_hash_atom(value: to_string(hardware_thread_count_value(r.min_threads))), + match r.architecture { + Absent => content_hash_atom(value: "arch:any") + Present { value: arch } => cpu_architecture_digest(arch: arch) + }, + ), + ) + } +} + +fn optional_memory_requirement_digest(req: MemoryRequirement?) -> ContentHash { + match req { + Absent => content_hash_atom(value: "memory:none") + Present { value: r } => + content_hash_combine( + content_hash_atom(value: "memory"), + content_hash_atom(value: to_string(byte_size_count(r.min_bytes))), + ) } } -fn witness_hard_requirement_unmet_is_unmet() -> Bool { - let big_need = ComputeNeed { fabric: example_fabric, program: example_large_program } - match connect(need: big_need, opportunities: [example_small_opportunity]) { - Unmet { reason: _ } => true - Bound { opportunity: _ } => false - Pending => false +fn resource_envelope_digest(envelope: ResourceEnvelope) -> ContentHash { + content_hash_combine( + optional_cpu_requirement_digest(req: envelope.cpu), + content_hash_combine( + optional_memory_requirement_digest(req: envelope.memory), + content_hash_atom(value: "resource-envelope-v1"), + ), + ) +} + +fn parallelism_shape_digest(shape: ParallelismShape) -> ContentHash { + match shape { + SingleWorkItem => content_hash_atom(value: "parallelism:SingleWorkItem") + IndependentShards { shard_count: n } => + content_hash_combine( + content_hash_atom(value: "parallelism:IndependentShards"), + content_hash_atom(value: to_string(n)), + ) + DependencyGraphParallel { graph: graph } => + content_hash_combine( + content_hash_atom(value: "parallelism:DependencyGraphParallel"), + content_hash_atom(value: "parallelism-payload-pending:WorkGraph"), + ) + PartitionedReduce { partitioner: partitioner, map: map, reduce: reduce, laws: laws } => + content_hash_combine( + content_hash_atom(value: "parallelism:PartitionedReduce"), + content_hash_combine( + content_hash_atom(value: "parallelism-payload-pending:Partitioner"), + content_hash_combine( + content_hash_atom(value: "parallelism-payload-pending:WorkAction-map"), + content_hash_combine( + content_hash_atom(value: "parallelism-payload-pending:Reducer"), + content_hash_atom(value: "parallelism-payload-pending:ReducerLaws"), + ), + ), + ), + ) + } +} + +fn work_demand_content_digest(demand: WorkDemand) -> ContentHash { + content_hash_combine( + resource_envelope_digest(envelope: demand.resources), + content_hash_combine( + isolation_boundary_digest(boundary: demand.isolation.boundary), + content_hash_combine( + parallelism_shape_digest(shape: demand.parallelism), + content_hash_atom(value: "work-demand-v1"), + ), + ), + ) +} + +fn execution_receipt_subject_digest( + work_id: WorkUnitId, + demand: WorkDemand, +) -> ExecutionReceiptDigest { + content_hash_tagged( + tag: execution_receipt_digest_schema, + payload: content_hash_combine( + content_hash_atom(value: work_id), + work_demand_content_digest(demand: demand), + ), + ) +} + +fn execution_receipt_digest(receipt: ExecutionReceipt) -> ExecutionReceiptDigest { + execution_receipt_subject_digest( + work_id: receipt.work.id, + demand: receipt.work.requirements, + ) +} + +fn witness_execution_receipt_digest_not_raw_work_id( + work_id: WorkUnitId, + demand: WorkDemand, +) -> Bool { + execution_receipt_subject_digest(work_id: work_id, demand: demand) != work_id +} + +fn witness_execution_receipt_digest_differs_by_demand( + work_id: WorkUnitId, + left: WorkDemand, + right: WorkDemand, +) -> Bool { + execution_receipt_subject_digest(work_id: work_id, demand: left) != execution_receipt_subject_digest(work_id: work_id, demand: right) +} + +data gunbc_ci_corpus_envelope: InputEnvelope = BoundedInput { + bounds: [ + InputBound { + axis: WitnessCount, + max: Measure { count: 2000 } + }, + InputBound { + axis: SourceNodeCount, + max: Measure { count: 500000 } + }, + InputBound { + axis: CorpusNodeCount, + max: Measure { count: 5000000 } + }, + ] +} + +data gunbc_ci_corpus_envelope_ceiling_scaffold: Disposition = Scaffold { + dissolves_to: SingleAuthority, + bind: DeclarationRef { + module_path: "gunbc.plans.input_envelope_roadmap", + decl_name: "input_envelope_roadmap_plan", + field: WholeDeclaration } } diff --git a/dsl/test/claim/compute_fabric_resource_witness_test.dag b/dsl/test/claim/compute_fabric_resource_witness_test.dag new file mode 100644 index 00000000000..ca14963d86e --- /dev/null +++ b/dsl/test/claim/compute_fabric_resource_witness_test.dag @@ -0,0 +1,24 @@ +module test.claim.compute_fabric_resource_witness + +import product.compute_fabric { + witness_resource_offer_rejects_insufficient_threads, + witness_resource_offer_rejects_insufficient_memory, + witness_resource_offer_admits_feasible_demand, + witness_run_demand_envelope_fits_budget, + witness_high_demand_shrinks_envelope_and_fits, + witness_fail_open_memory_blind_demand_exceeds_budget, + witness_unreadable_budget_uses_conservative_demand, +} + +test fn compute_fabric_resource_feasibility_holds() -> Bool { + witness_resource_offer_rejects_insufficient_threads() + && witness_resource_offer_rejects_insufficient_memory() + && witness_resource_offer_admits_feasible_demand() +} + +test fn ci_floor_run_demand_envelope_model_holds() -> Bool { + witness_run_demand_envelope_fits_budget() + && witness_high_demand_shrinks_envelope_and_fits() + && witness_fail_open_memory_blind_demand_exceeds_budget() + && witness_unreadable_budget_uses_conservative_demand() +} diff --git a/dsl/test/claim/compute_fabric_witness_test.dag b/dsl/test/claim/compute_fabric_witness_test.dag deleted file mode 100644 index 21c1a201259..00000000000 --- a/dsl/test/claim/compute_fabric_witness_test.dag +++ /dev/null @@ -1,20 +0,0 @@ -module test.claim.compute_fabric_witness - -import product.compute_fabric { - witness_need_is_only_fabric_plus_program, - witness_shape_is_derived_from_program, - witness_fabric_binds_need_to_opportunity, - witness_exec_and_ci_are_just_programs, - witness_hard_requirement_unmet_is_unmet, -} - -test fn compute_fabric_broker_shape_and_need_holds() -> Bool { - witness_need_is_only_fabric_plus_program() - && witness_shape_is_derived_from_program() -} - -test fn compute_fabric_broker_connect_holds() -> Bool { - witness_fabric_binds_need_to_opportunity() - && witness_exec_and_ci_are_just_programs() - && witness_hard_requirement_unmet_is_unmet() -} diff --git a/dsl/test/claim/probe_selector_compute_fabric_import_repro_test.dag b/dsl/test/claim/probe_selector_compute_fabric_import_repro_test.dag index 57f2f1a412f..5c32375d951 100644 --- a/dsl/test/claim/probe_selector_compute_fabric_import_repro_test.dag +++ b/dsl/test/claim/probe_selector_compute_fabric_import_repro_test.dag @@ -1,7 +1,7 @@ module test.claim.probe_selector_compute_fabric_import_repro -import product.compute_fabric { example_fabric } +import product.compute_fabric { example_shared_host_offer } test fn probe_selector_compute_fabric_import_repro_holds() -> Bool { - example_fabric == example_fabric + example_shared_host_offer.provider == example_shared_host_offer.provider }