Files
3cloud-backend/app/scheduling_filters.py
T
JamesBhattarai 35106a2ed7 Feat: GPU VM & GPU-Dynamic Passthrough
API to register GPU (worker does at initialization), list all GPUs (to fetch during creation menu)
Schedular to Match with relevant hosts filter
A Table to store GPU models
 A table to map gpu to pcie to  workload and host
Dynamic binding unbinding of vfio, nvidia driver at runtime for passthrough
2026-05-16 03:59:00 +00:00

694 lines
25 KiB
Python

from abc import ABC, abstractmethod
from typing import Any, Dict, List, Optional
from app.models.models import WorkloadHost
from logger import logger
import ipaddress
import random
class SchedulingFilter(ABC):
@abstractmethod
def apply(self, hosts, requirements=None):
pass
class ExcludeAllOfflineHosts(SchedulingFilter):
def apply(self, hosts, requirements=None):
host_ids = [str(host.id) for host in hosts]
logger.debug(
f"Starting ExcludeAllOfflineHosts filter. Input hosts ({len(hosts)}): {host_ids}"
)
filtered_hosts = [host for host in hosts if host.status == "online"]
filtered_ids = [str(host.id) for host in filtered_hosts]
logger.debug(
f"Finished ExcludeAllOfflineHosts. "
f"Input: {len(hosts)}, Filtered out: {len(hosts) - len(filtered_hosts)}, Remaining: {len(filtered_hosts)}\n"
f"Removed hosts: {list(set(host_ids) - set(filtered_ids))}\n"
f"Remaining hosts: {filtered_ids}"
)
return filtered_hosts
class ExcludeAllDisabledHosts(SchedulingFilter):
def apply(self, hosts, requirements=None):
host_ids = [str(host.id) for host in hosts]
logger.debug(
f"Starting ExcludeAllDisabledHosts filter. Input hosts ({len(hosts)}): {host_ids}"
)
filtered_hosts = [host for host in hosts if host.available_for_scheduling]
filtered_ids = [str(host.id) for host in filtered_hosts]
logger.debug(
f"Finished ExcludeAllDisabledHosts. "
f"Input: {len(hosts)}, Filtered out: {len(hosts) - len(filtered_hosts)}, Remaining: {len(filtered_hosts)}\n"
f"Removed hosts: {list(set(host_ids) - set(filtered_ids))}\n"
f"Remaining hosts: {filtered_ids}"
)
return filtered_hosts
class CapableHostsFilter(SchedulingFilter):
def __init__(self, capabilities, requires_all=True):
if isinstance(capabilities, str):
self.capabilities = [capabilities]
else:
self.capabilities = capabilities
self.requires_all = requires_all
def apply(self, hosts, requirements=None):
host_ids = [str(host.id) for host in hosts]
logger.debug(
f"Starting CapableHostsFilter. Capabilities: {self.capabilities}, "
f"requires_all: {self.requires_all}\n"
f"Input hosts ({len(hosts)}): {host_ids}"
)
filtered_hosts = []
for host in hosts:
host_labels = {label.label_key: label.label_value.lower()
for label in host.getLabels()}
matches = [
key in host_labels and host_labels[key] in ['true', 'yes', '1']
for key in self.capabilities
]
if (self.requires_all and all(matches)) or (not self.requires_all and any(matches)):
filtered_hosts.append(host)
filtered_ids = [str(host.id) for host in filtered_hosts]
logger.debug(
f"Finished CapableHostsFilter. "
f"Input: {len(hosts)}, Filtered out: {len(hosts) - len(filtered_hosts)}, Remaining: {len(filtered_hosts)}\n"
f"Removed hosts: {list(set(host_ids) - set(filtered_ids))}\n"
f"Remaining hosts: {filtered_ids}"
)
return filtered_hosts
class LibvirtCapableHosts(CapableHostsFilter):
def __init__(self):
super().__init__('libvirt-capable')
class DockerCapableHosts(CapableHostsFilter):
def __init__(self):
super().__init__('docker-capable')
class GPUCapableHosts(SchedulingFilter):
"""
Filter hosts that have sufficient available GPUs, sourced from the
GPU_HOST table (populated by the worker's PCI device scanner on startup).
Accepts either:
- ``required_gpu_count`` (int): total GPUs needed of any model.
- ``gpu_model_requests`` (list of {id, count}): specific models required.
When provided, each model must have enough available slots on the host.
"""
def __init__(self, required_gpu_count: int = 1, gpu_model_requests: list = None):
self.required_gpu_count = required_gpu_count
self.gpu_model_requests = gpu_model_requests # e.g. [{"id": "<uuid>", "count": 2}]
def apply(self, hosts, requirements=None):
from app.compute.models.gpu import GPU_HOST
from app import db
host_ids = [str(host.id) for host in hosts]
logger.debug(
f"Starting GPUCapableHosts filter. Required GPUs: {self.required_gpu_count}, "
f"model_requests: {self.gpu_model_requests}\n"
f"Input hosts ({len(hosts)}): {host_ids}"
)
filtered_hosts = []
for host in hosts:
host_id_str = str(host.id)
if self.gpu_model_requests:
# Model-based: every requested model must have enough available on this host
eligible = True
for req in self.gpu_model_requests:
available = (
db.session.query(GPU_HOST)
.filter_by(
workload_host_id=host_id_str,
gpu_model_id=req['id'],
gpu_status='available',
deleted=False,
)
.count()
)
needed = int(req.get('count', 1))
if available < needed:
logger.debug(
f"Host {host.id} ({host.hostname}) has {available} of model "
f"{req['id']} available, need {needed} - FAIL"
)
eligible = False
break
if eligible:
filtered_hosts.append(host)
logger.debug(f"Host {host.id} ({host.hostname}) satisfies all GPU model requests - PASS")
else:
# Integer format: just check total available count on this host
available = (
db.session.query(GPU_HOST)
.filter_by(
workload_host_id=host_id_str,
gpu_status='available',
deleted=False,
)
.count()
)
if available >= self.required_gpu_count:
filtered_hosts.append(host)
logger.debug(
f"Host {host.id} ({host.hostname}) has {available} GPU(s) available "
f"(required: {self.required_gpu_count}) - PASS"
)
else:
logger.debug(
f"Host {host.id} ({host.hostname}) has {available} GPU(s) available "
f"(required: {self.required_gpu_count}) - FAIL"
)
filtered_ids = [str(host.id) for host in filtered_hosts]
logger.debug(
f"Finished GPUCapableHosts. "
f"Input: {len(hosts)}, Filtered out: {len(hosts) - len(filtered_hosts)}, Remaining: {len(filtered_hosts)}\n"
f"Removed hosts: {list(set(host_ids) - set(filtered_ids))}\n"
f"Remaining hosts: {filtered_ids}"
)
return filtered_hosts
class MostAvailableCapacity(SchedulingFilter):
def apply(self, hosts, requirements=None):
host_ids = [str(host.id) for host in hosts]
logger.debug(
f"Starting MostAvailableCapacity filter. Input hosts ({len(hosts)}): {host_ids}"
)
filtered_hosts = sorted(
hosts,
key=lambda host: host.placement_priority,
reverse=True
)
sorted_ids = [str(host.id) for host in filtered_hosts]
logger.debug(
f"Finished MostAvailableCapacity. "
f"Input: {len(hosts)}, Filtered out: 0, Remaining: {len(filtered_hosts)}\n"
f"Priority order: {sorted_ids}"
)
return filtered_hosts
class HasSufficientResources(SchedulingFilter):
"""
Scheduling filter that determines whether each WorkloadHost in a candidate
list can satisfy this workload's resource requirements after applying
per-resource overcommit ratios.
This class is fully generic. It does not assume that only CPU and RAM are
relevant. You can schedule on any resource type that both:
1. appears in `required_resources`
2. is reported by WorkloadHost.get_resource_utilization()
Core concepts
-------------
required_resources:
Dict of {resource_type: required_amount}.
For example:
{
"cpu": 10.0,
"ram": 16384,
"bandwidth_mbps": 5000,
"monitor_ports": 1,
}
Interpretation: the workload being scheduled needs at least the given
amount of each listed resource on whatever host we pick.
overcommit_ratios:
Dict of {resource_type: ratio}.
For example:
{
"cpu": 10.0,
"ram": 2.0,
"bandwidth_mbps": 1.0
}
The overcommit ratio describes how aggressively we are allowed to
oversubscribe that particular resource. The scheduler treats
physical_free * ratio as "effective available".
If a resource_type is not present in overcommit_ratios,
a default of 1.0 (no overcommit) is assumed.
Host utilization model
----------------------
Each WorkloadHost is expected to implement get_resource_utilization() and
return a dictionary shaped like:
{
"resources": {
"<resource_type>": {
"total_quantity": <float|int>,
"quantity_in_use": <float|int>,
},
...
}
}
For each resource_type in required_resources:
physical_available = total_quantity - quantity_in_use
effective_available = physical_available * overcommit_ratio
host passes if effective_available >= required_amount
Any failure rejects that host.
Logging
-------
- Debug logs capture:
- The full host utilization snapshot.
- Per-resource pass/fail evaluation.
- Final accept/reject list.
- This is important for audit and scheduler trace replay.
Attributes
----------
required_resources : Dict[str, float | int]
Required amounts per resource type for the workload.
overcommit_ratios : Dict[str, float]
Overcommit multiplier per resource type.
"""
def __init__(
self,
required_resources: Dict[str, float | int],
overcommit_ratios: Optional[Dict[str, float]] = None,
):
"""
Initialize the HasSufficientResources filter with a generic resource
requirement set and optional overcommit ratios.
Parameters
----------
required_resources : Dict[str, float | int]
Mapping of resource_type -> required amount for this workload.
Example:
{
"cpu": 10.0,
"ram": 16384,
"bandwidth_mbps": 5000,
}
overcommit_ratios : Dict[str, float], optional
Mapping of resource_type -> oversubscription multiplier.
If a resource_type is missing here, the filter will assume 1.0
for that resource_type (which means no overcommit allowed).
Example:
{
"cpu": 10.0,
"ram": 2.0,
"bandwidth_mbps": 1.0,
}
"""
self.required_resources = required_resources or {}
self.overcommit_ratios = overcommit_ratios or {}
def apply(
self,
hosts: List[WorkloadHost],
requirements: Optional[Dict] = None,
) -> List[WorkloadHost]:
"""
Filter a list of WorkloadHost objects and return only the hosts that are
capable of satisfying this workload's resource requirements.
Parameters
----------
hosts : List[WorkloadHost]
Candidate hosts under consideration for scheduling.
requirements : Dict, optional
Unused in this refactored implementation, but kept in the signature
for compatibility with the SchedulingFilter interface. The required
resources are provided in self.required_resources at construction
time instead.
Returns
-------
List[WorkloadHost]
The subset of input hosts that satisfy all required resource
constraints after applying per-resource overcommit rules.
"""
if requirements is None:
requirements = {}
host_ids = [str(host.id) for host in hosts]
logger.debug(
"HasSufficientResources.apply() start. "
"Candidate hosts (%d): %s. "
"Workload required_resources=%s, overcommit_ratios=%s",
len(hosts),
host_ids,
self.required_resources,
self.overcommit_ratios,
)
filtered_hosts = [
host for host in hosts if self._host_satisfies_resources(host)
]
filtered_ids = [str(host.id) for host in filtered_hosts]
removed_ids = list(set(host_ids) - set(filtered_ids))
logger.debug(
"HasSufficientResources.apply() complete. "
"Input hosts: %d, Filtered out: %d, Remaining: %d.\n"
"Removed hosts: %s\n"
"Remaining hosts: %s",
len(hosts),
len(removed_ids),
len(filtered_hosts),
removed_ids,
filtered_ids,
)
return filtered_hosts
def _host_satisfies_resources(self, host: WorkloadHost) -> bool:
"""
Determine whether the provided WorkloadHost can satisfy ALL required
resources for this workload, after considering per-resource overcommit
ratios.
The evaluation is generic and applies to any arbitrary resource type.
Process per resource_type:
1. Look up host resource snapshot for that resource_type.
2. Calculate physical_available:
total_quantity - quantity_in_use
3. Look up overcommit_ratio for that resource_type, default 1.0.
4. Calculate effective_available:
physical_available * overcommit_ratio
5. If effective_available < required_amount, reject the host.
The first failing resource rejects the host immediately.
Parameters
----------
host : WorkloadHost
The WorkloadHost instance to evaluate.
Returns
-------
bool
True if the host satisfies all resource requirements (after applying
overcommit policy for each resource). False otherwise. When False,
this logs a specific rejection reason for audit and traceability.
"""
# Retrieve live utilization snapshot for this host
host_utilization = host.get_resource_utilization()
if not host_utilization or "resources" not in host_utilization:
logger.debug(
"Host %s rejected by HasSufficientResources: "
"missing utilization data",
host.id,
)
return False
# Log full snapshot for deep trace / post-mortem replay
logger.debug(
"Host %s utilization snapshot for scheduling evaluation: %s",
host.id,
host_utilization,
)
# Iterate through every resource we require
for resource_type, required_amount in self.required_resources.items():
# Ignore non-positive requirements; they do not constrain placement
if required_amount is None or required_amount <= 0:
logger.debug(
"Host %s resource '%s' requirement is non-positive (%s); "
"skipping check.",
host.id,
resource_type,
required_amount,
)
continue
# Pull what the host reports for this resource
resource_snapshot = host_utilization["resources"].get(resource_type)
# Resource type missing entirely on this host
if resource_snapshot is None:
logger.debug(
"Host %s rejected by HasSufficientResources: "
"resource '%s' not reported by host.",
host.id,
resource_type,
)
return False
total_quantity = resource_snapshot.get("total_quantity")
quantity_in_use = resource_snapshot.get("quantity_in_use", 0.0)
# Bail if host capacity model is unusable
if total_quantity is None or total_quantity <= 0:
logger.debug(
"Host %s rejected by HasSufficientResources: resource '%s' "
"has missing or zero total_quantity.",
host.id,
resource_type,
)
return False
total_quantity = resource_snapshot.get("total_quantity")
quantity_in_use = resource_snapshot.get("quantity_in_use", 0.0)
if total_quantity is None or total_quantity <= 0:
logger.debug(
"Host %s rejected by HasSufficientResources: resource '%s' "
"has missing or zero total_quantity.",
host.id,
resource_type,
)
return False
overcommit_ratio = self.overcommit_ratios.get(resource_type, 1.0) or 1.0
# total_effective_capacity represents the max schedulable capacity for this resource
# after applying oversubscription policy. For example:
# - If there are 2 physical CPU cores and overcommit_ratio is 10.0,
# then total_effective_capacity for 'cpu' is 20 logical cores.
total_effective_capacity = total_quantity * overcommit_ratio
# effective_available is how much schedulable capacity we have left
# after subtracting what is already allocated on this host.
effective_available = total_effective_capacity - quantity_in_use
# Check sufficiency
if effective_available < required_amount:
logger.debug(
"Host %s rejected by HasSufficientResources (%s): "
"required=%s, "
"overcommit_ratio=%s, effective_available=%s",
host.id,
resource_type,
required_amount,
overcommit_ratio,
effective_available,
)
return False
# Resource passes, log details
logger.debug(
"Host %s passes %s check: required=%s, "
" overcommit_ratio=%s, "
"effective_available=%s",
host.id,
resource_type,
required_amount,
overcommit_ratio,
effective_available,
)
# All required resources passed checks
return True
class ExcludeHostsWithoutNorthSouthIP(SchedulingFilter):
"""
Filter that removes every WorkloadHost lacking a valid north-south IP address.
A host is **kept** only if:
• `ip_address_northsouth` is not None / empty, **and**
• its value parses as a valid IPv4 or IPv6 address.
All others are filtered out.
"""
def apply(self, hosts, requirements=None):
"""
Iterate over *hosts* and return only those with a valid
`ip_address_northsouth`.
Parameters
----------
hosts : list[WorkloadHost]
Candidate hosts.
Returns
-------
list[WorkloadHost]
Hosts that passed the IP-validation check.
"""
host_ids = [str(h.id) for h in hosts]
logger.debug(
f"Starting ExcludeHostsWithoutNorthSouthIP filter. "
f"Input hosts ({len(hosts)}): {host_ids}"
)
filtered_hosts, removed_ids = [], []
for host in hosts:
ip_addr = getattr(host, "ip_address_northsouth", None)
try:
# ipaddress throws ValueError on bad input
if ip_addr and ipaddress.ip_address(ip_addr):
filtered_hosts.append(host)
else:
removed_ids.append(str(host.id))
except ValueError:
removed_ids.append(str(host.id))
filtered_ids = [str(h.id) for h in filtered_hosts]
logger.debug(
f"Finished ExcludeHostsWithoutNorthSouthIP. "
f"Input: {len(hosts)}, Filtered out: {len(hosts) - len(filtered_hosts)}, Remaining: {len(filtered_hosts)}\n"
f"Removed hosts: {removed_ids}\n"
f"Remaining hosts: {filtered_ids}"
)
return filtered_hosts
class RandomizeHostsOrder(SchedulingFilter):
"""
Randomizes the order of the remaining candidate hosts.
This can help with spreading load more evenly over time
when multiple hosts are equally eligible after filtering.
"""
def apply(self, hosts, requirements=None):
"""
Shuffle the list of hosts randomly.
Parameters
----------
hosts : list[WorkloadHost]
Candidate hosts.
Returns
-------
list[WorkloadHost]
Shuffled list of hosts.
"""
host_ids = [str(host.id) for host in hosts]
logger.debug(
f"Starting RandomizeHostsOrder filter. "
f"Input hosts ({len(hosts)}): {host_ids}"
)
shuffled_hosts = hosts[:]
random.shuffle(shuffled_hosts)
shuffled_ids = [str(host.id) for host in shuffled_hosts]
logger.debug(
f"Finished RandomizeHostsOrder. "
f"Input: {len(hosts)}, Filtered out: 0, Remaining: {len(shuffled_hosts)}\n"
f"Original order: {host_ids}\n"
f"Shuffled order: {shuffled_ids}"
)
return shuffled_hosts
class SortByPlacementPriority(SchedulingFilter):
"""
Sorts workload hosts in descending order of placement priority.
Hosts with higher `placement_priority` values are preferred first.
"""
def apply(self, hosts, requirements=None):
logger.debug(self)
logger.debug(hosts)
logger.debug(requirements)
host_ids = [str(host.id) for host in hosts]
logger.debug(
f"Starting SortByPlacementPriority filter. "
f"Input hosts ({len(hosts)}): {host_ids}"
)
sorted_hosts = sorted(
hosts,
key=lambda host: host.placement_priority,
reverse=True # highest value first
)
sorted_ids = [str(host.id) for host in sorted_hosts]
logger.debug(
f"Finished SortByPlacementPriority. "
f"Input: {len(hosts)}, Filtered out: 0, Remaining: {len(sorted_hosts)}\n"
f"Sorted host IDs by descending placement priority: {sorted_ids}"
)
return sorted_hosts
class AntiAffinityFilter(SchedulingFilter):
def __init__(self, label_key: str, customer_id: str, all_workloads: list):
self.label_key = label_key
self.customer_id = customer_id
self.all_workloads = all_workloads
def apply(self, hosts, requirements=None):
conflicting_workloads = [
w for w in self.all_workloads
if w.owner_id == self.customer_id and self.label_key in (w.labels or {})
]
excluded_host_ids = {w.workload_host_id for w in conflicting_workloads if w.workload_host_id}
filtered_hosts = [host for host in hosts if host.id not in excluded_host_ids]
logger.debug(
f"AntiAffinityFilter: Excluding host IDs with label '{self.label_key}': {excluded_host_ids}"
)
return filtered_hosts
class AffinityFilter(SchedulingFilter):
def __init__(self, label_key: str, customer_id: str, all_workloads: list):
self.label_key = label_key
self.customer_id = customer_id
self.all_workloads = all_workloads
def apply(self, hosts, requirements=None):
matching_workloads = [
w for w in self.all_workloads
if w.owner_id == self.customer_id and self.label_key in (w.labels or {})
]
preferred_host_ids = {w.workload_host_id for w in matching_workloads if w.workload_host_id}
filtered_hosts = [host for host in hosts if host.id in preferred_host_ids]
logger.debug(
f"AffinityFilter: Preferred host IDs for label '{self.label_key}': {preferred_host_ids}"
)
return filtered_hosts