Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions frontend/src/pages/Offers/List/index.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -67,6 +67,7 @@ const getRequestParams = ({
profile: { name: 'default', default: false },
ssh_key_pub: '(dummy)',
},
full_offers: true,
};
};

Expand Down
1 change: 1 addition & 0 deletions frontend/src/types/gpu.d.ts
Original file line number Diff line number Diff line change
Expand Up @@ -88,6 +88,7 @@ declare type TGpusListQueryParams = {
profile?: { name: string; default?: boolean };
ssh_key_pub: string;
};
full_offers: boolean;
};

declare type TGpusListQueryResponse = {
Expand Down
7 changes: 7 additions & 0 deletions src/dstack/_internal/cli/commands/offer.py
Original file line number Diff line number Diff line change
Expand Up @@ -59,6 +59,11 @@ def _register(self):
type=int,
default=50,
)
self._parser.add_argument(
"--full-offers",
action="store_true",
help="Show full offers not adjusted by requirements",
)
resources_group = self._parser.add_argument_group("Resources")
register_resources_args(resources_group)
# TODO: register only relevant options
Expand All @@ -79,6 +84,7 @@ def _list_offers(self, args: argparse.Namespace) -> None:
project_name=self.api.project,
run_spec=run_spec,
max_offers=args.max_offers,
full_offers=args.full_offers,
)
job_plan = run_plan.job_plans[0]
if args.format == "plain":
Expand All @@ -103,6 +109,7 @@ def _list_gpus(self, args: argparse.Namespace, group_by: list[str]) -> None:
project_name=self.api.project,
run_spec=run_spec,
group_by=[g for g in group_by if g != "gpu"],
full_offers=args.full_offers,
)
if args.format == "plain":
print_gpu_table(gpus, run_spec, group_by, self.api.project)
Expand Down
7 changes: 7 additions & 0 deletions src/dstack/_internal/cli/services/configurators/run.py
Original file line number Diff line number Diff line change
Expand Up @@ -143,6 +143,8 @@ def get_plan(
configuration_path=configuration_path,
profile=profile,
ssh_identity_file=configurator_args.ssh_identity_file,
max_offers=configurator_args.max_offers,
full_offers=configurator_args.full_offers,
)
return run_plan, repo

Expand Down Expand Up @@ -387,6 +389,11 @@ def register_args(cls, parser: argparse.ArgumentParser):
type=int,
default=3,
)
configuration_group.add_argument(
"--full-offers",
action="store_true",
help="Show full offers not adjusted by requirements",
)
cls.register_env_args(configuration_group)
register_resources_args(configuration_group)
register_profile_args(parser)
Expand Down
4 changes: 3 additions & 1 deletion src/dstack/_internal/core/backends/aws/compute.py
Original file line number Diff line number Diff line change
Expand Up @@ -181,7 +181,9 @@ def get_all_offers_with_availability(self) -> List[InstanceOfferWithAvailability
)
return availability_offers

def get_offers_modifiers(self, requirements: Requirements) -> Iterable[OfferModifier]:
def get_offers_modifiers(
self, requirements: Requirements, full_offers: bool
) -> Iterable[OfferModifier]:
return [get_offers_disk_modifier(CONFIGURABLE_DISK_SIZE, requirements)]

def _get_offers_cached_key(self, requirements: Requirements) -> int:
Expand Down
4 changes: 3 additions & 1 deletion src/dstack/_internal/core/backends/azure/compute.py
Original file line number Diff line number Diff line change
Expand Up @@ -115,7 +115,9 @@ def get_all_offers_with_availability(self) -> List[InstanceOfferWithAvailability
)
return offers_with_availability

def get_offers_modifiers(self, requirements: Requirements) -> Iterable[OfferModifier]:
def get_offers_modifiers(
self, requirements: Requirements, full_offers: bool
) -> Iterable[OfferModifier]:
return [get_offers_disk_modifier(CONFIGURABLE_DISK_SIZE, requirements)]

def create_instance(
Expand Down
60 changes: 48 additions & 12 deletions src/dstack/_internal/core/backends/base/compute.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@
from enum import Enum
from functools import lru_cache
from pathlib import Path
from typing import Callable, Dict, List, Optional
from typing import Callable, ClassVar, Dict, List, Optional

import git
import requests
Expand Down Expand Up @@ -109,12 +109,23 @@ class Compute(ABC):
"""

@abstractmethod
def get_offers(self, requirements: Requirements) -> Iterator[InstanceOfferWithAvailability]:
def get_offers(
self, requirements: Requirements, full_offers: bool
) -> Iterator[InstanceOfferWithAvailability]:
"""
Returns offers with availability matching `requirements`.
If the provider is added to gpuhunt, typically gets offers using
`base.offers.get_catalog_offers()` and extends them with availability info.
It is called from async code in executor. It can block on call but not between yields.

if `full_offers` set to `True`, the method should not adjust offer's resources according to
`requirements`. For most backends, this flag has no meaning, as they work with predefined
provider offers (even configurable disk size reflects the actual disk created once the
instance is provisioned), but some backends such as Kubernetes and Slurm allocates flexible
slices of instances (nodes) according to the requested resources; such Computes usually
generate synthetic offers from discovered nodes on the fly; these synthetic offers should
reflect either resources that would be allocated based on `requirements`
(`full_offers=False`) or full allocatable node resources (`full_offers=True`).
"""
pass

Expand Down Expand Up @@ -190,11 +201,15 @@ def get_all_offers_with_availability(self) -> List[InstanceOfferWithAvailability
"""
pass

def get_offers_modifiers(self, requirements: Requirements) -> Iterable[OfferModifier]:
def get_offers_modifiers(
self, requirements: Requirements, full_offers: bool
) -> Iterable[OfferModifier]:
"""
Returns functions that modify offers before they are filtered by requirements.
A modifier function can return `None` to exclude the offer.
E.g. can be used to set appropriate disk size based on requirements.

See `Compute.get_offers()` for the `full_offers` argument description.
"""
return []

Expand All @@ -207,12 +222,16 @@ def get_offers_post_filter(
"""
return None

def get_offers(self, requirements: Requirements) -> Iterator[InstanceOfferWithAvailability]:
def get_offers(
self, requirements: Requirements, full_offers: bool
) -> Iterator[InstanceOfferWithAvailability]:
with self._offers_cache_execution_lock:
# Cache lock does not prevent concurrent execution.
# We use a separate lock to avoid requesting offers in parallel, re-doing the work and hitting rate limits.
cached_offers = self._get_all_offers_with_availability_cached()
offers = self.__apply_modifiers(cached_offers, self.get_offers_modifiers(requirements))
offers = self.__apply_modifiers(
cached_offers, self.get_offers_modifiers(requirements, full_offers)
)
offers = filter_offers_by_requirements(offers, requirements)
post_filter = self.get_offers_post_filter(requirements)
if post_filter is not None:
Expand Down Expand Up @@ -246,36 +265,53 @@ class ComputeWithFilteredOffersCached(ABC):
It caches offers using requirements as key.
"""

full_offers_argument_has_effect: ClassVar[bool] = False
"""
Set to `True` if `get_offers_by_requirements()` produces different results based on
the `full_offers` value. Doubles the amount of cached data.
"""

def __init__(self) -> None:
super().__init__()
self._offers_cache_lock = threading.Lock()
self._offers_cache = TTLCache(maxsize=10, ttl=180)

@abstractmethod
def get_offers_by_requirements(
self, requirements: Requirements
self,
requirements: Requirements,
full_offers: bool,
) -> List[InstanceOfferWithAvailability]:
"""
Returns backend offers with availability matching requirements.

See `Compute.get_offers()` for the `full_offers` argument description.
Set the class variable `full_offers_argument_has_effect` to `True` if the `full_offers`
value has an effect on the offers produced by this method.
"""
pass

def get_offers(self, requirements: Requirements) -> Iterator[InstanceOfferWithAvailability]:
return iter(self._get_offers_cached(requirements))
def get_offers(
self, requirements: Requirements, full_offers: bool
) -> Iterator[InstanceOfferWithAvailability]:
return iter(self._get_offers_cached(requirements, full_offers))

def _get_offers_cached_key(self, requirements: Requirements) -> int:
def _get_offers_cached_key(self, requirements: Requirements, full_offers: bool) -> int:
# Requirements is not hashable, so we use a hack to get arguments hash
return hash(requirements.json())
hashable_requirements = requirements.json()
if self.full_offers_argument_has_effect:
return hash((hashable_requirements, full_offers))
return hash(hashable_requirements)

@cachedmethod(
cache=lambda self: self._offers_cache,
key=_get_offers_cached_key,
lock=lambda self: self._offers_cache_lock,
)
def _get_offers_cached(
self, requirements: Requirements
self, requirements: Requirements, full_offers: bool
) -> List[InstanceOfferWithAvailability]:
return self.get_offers_by_requirements(requirements)
return self.get_offers_by_requirements(requirements, full_offers)


class ComputeWithCreateInstanceSupport(ABC):
Expand Down
4 changes: 3 additions & 1 deletion src/dstack/_internal/core/backends/crusoe/compute.py
Original file line number Diff line number Diff line change
Expand Up @@ -180,7 +180,9 @@ def _get_quota_map(self) -> dict[str, int]:
result[prog_name] = available
return result

def get_offers_modifiers(self, requirements: Requirements) -> Iterable[OfferModifier]:
def get_offers_modifiers(
self, requirements: Requirements, full_offers: bool
) -> Iterable[OfferModifier]:
# Only adjust disk size for types without ephemeral NVMe (disk_gb == 0).
# Types with ephemeral NVMe already have their disk_size set by gpuhunt.
base_modifier = get_offers_disk_modifier(CONFIGURABLE_DISK_SIZE, requirements)
Expand Down
4 changes: 3 additions & 1 deletion src/dstack/_internal/core/backends/gcp/compute.py
Original file line number Diff line number Diff line change
Expand Up @@ -168,7 +168,9 @@ def get_all_offers_with_availability(self) -> List[InstanceOfferWithAvailability
offer_with_availability.region = region
return offers_with_availability

def get_offers_modifiers(self, requirements: Requirements) -> Iterable[OfferModifier]:
def get_offers_modifiers(
self, requirements: Requirements, full_offers: bool
) -> Iterable[OfferModifier]:
modifiers = []

if requirements.reservation:
Expand Down
4 changes: 3 additions & 1 deletion src/dstack/_internal/core/backends/jarvislabs/compute.py
Original file line number Diff line number Diff line change
Expand Up @@ -91,7 +91,9 @@ def get_all_offers_with_availability(self) -> List[InstanceOfferWithAvailability
for offer in offers
]

def get_offers_modifiers(self, requirements: Requirements) -> Iterable[OfferModifier]:
def get_offers_modifiers(
self, requirements: Requirements, full_offers: bool
) -> Iterable[OfferModifier]:
return [get_offers_disk_modifier(CONFIGURABLE_DISK_SIZE, requirements)]

def create_instance(
Expand Down
6 changes: 5 additions & 1 deletion src/dstack/_internal/core/backends/kubernetes/compute.py
Original file line number Diff line number Diff line change
Expand Up @@ -169,7 +169,11 @@ def get_all_offers_with_availability(self) -> list[InstanceOfferWithAvailability
offers.extend(cluster_offers)
return offers

def get_offers_modifiers(self, requirements: Requirements) -> list[OfferModifier]:
def get_offers_modifiers(
self, requirements: Requirements, full_offers: bool
) -> list[OfferModifier]:
if full_offers:
return []
resource_requests = ResourceRequests.from_resources_spec(requirements.resources)
return [partial(_offer_modifier, resource_requests)]

Expand Down
4 changes: 3 additions & 1 deletion src/dstack/_internal/core/backends/nebius/compute.py
Original file line number Diff line number Diff line change
Expand Up @@ -133,7 +133,9 @@ def get_all_offers_with_availability(self) -> List[InstanceOfferWithAvailability
offer.with_availability(availability=InstanceAvailability.UNKNOWN) for offer in offers
]

def get_offers_modifiers(self, requirements: Requirements) -> Iterable[OfferModifier]:
def get_offers_modifiers(
self, requirements: Requirements, full_offers: bool
) -> Iterable[OfferModifier]:
return [get_offers_disk_modifier(CONFIGURABLE_DISK_SIZE, requirements)]

def create_instance(
Expand Down
4 changes: 3 additions & 1 deletion src/dstack/_internal/core/backends/oci/compute.py
Original file line number Diff line number Diff line change
Expand Up @@ -100,7 +100,9 @@ def get_all_offers_with_availability(self) -> List[InstanceOfferWithAvailability

return offers_with_availability

def get_offers_modifiers(self, requirements: Requirements) -> Iterable[OfferModifier]:
def get_offers_modifiers(
self, requirements: Requirements, full_offers: bool
) -> Iterable[OfferModifier]:
return [get_offers_disk_modifier(CONFIGURABLE_DISK_SIZE, requirements)]

def terminate_instance(
Expand Down
4 changes: 3 additions & 1 deletion src/dstack/_internal/core/backends/runpod/compute.py
Original file line number Diff line number Diff line change
Expand Up @@ -89,7 +89,9 @@ def get_all_offers_with_availability(self) -> List[InstanceOfferWithAvailability
]
return offers

def get_offers_modifiers(self, requirements: Requirements) -> Iterable[OfferModifier]:
def get_offers_modifiers(
self, requirements: Requirements, full_offers: bool
) -> Iterable[OfferModifier]:
gpu_disk_modifier = get_offers_disk_modifier(CONFIGURABLE_DISK_SIZE, requirements)

def disk_modifier(
Expand Down
6 changes: 5 additions & 1 deletion src/dstack/_internal/core/backends/slurm/compute.py
Original file line number Diff line number Diff line change
Expand Up @@ -114,7 +114,11 @@ def get_all_offers_with_availability(self) -> list[InstanceOfferWithAvailability
offers.extend(cluster_offers)
return offers

def get_offers_modifiers(self, requirements: Requirements) -> list[OfferModifier]:
def get_offers_modifiers(
self, requirements: Requirements, full_offers: bool
) -> list[OfferModifier]:
if full_offers:
return []
requested_resources = get_requested_resources_from_resources_spec(requirements.resources)
return [partial(self._offer_modifier, requested_resources)]

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -49,7 +49,7 @@ class {{ backend_name }}Compute(
self.config = config

def get_offers(
self, requirements: Requirements
self, requirements: Requirements, full_offers: bool
) -> Iterator[InstanceOfferWithAvailability]:
# If the provider is added to gpuhunt, you'd typically get offers
# using `get_catalog_offers()` and extend them with availability info.
Expand Down
2 changes: 1 addition & 1 deletion src/dstack/_internal/core/backends/vastai/compute.py
Original file line number Diff line number Diff line change
Expand Up @@ -90,7 +90,7 @@ def _make_catalog(self, options: VastAIProfileOptions) -> gpuhunt.Catalog:
return catalog

def get_offers_by_requirements(
self, requirements: Requirements
self, requirements: Requirements, full_offers: bool
) -> List[InstanceOfferWithAvailability]:
vastai_options = (
get_backend_profile_options(requirements.backend_options, VastAIProfileOptions)
Expand Down
4 changes: 3 additions & 1 deletion src/dstack/_internal/core/backends/verda/compute.py
Original file line number Diff line number Diff line change
Expand Up @@ -72,7 +72,9 @@ def get_all_offers_with_availability(self) -> List[InstanceOfferWithAvailability
offers_with_availability = self._get_offers_with_availability(offers)
return offers_with_availability

def get_offers_modifiers(self, requirements: Requirements) -> Iterable[OfferModifier]:
def get_offers_modifiers(
self, requirements: Requirements, full_offers: bool
) -> Iterable[OfferModifier]:
return [get_offers_disk_modifier(CONFIGURABLE_DISK_SIZE, requirements)]

def _get_offers_with_availability(
Expand Down
2 changes: 2 additions & 0 deletions src/dstack/_internal/core/compatibility/gpus.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,8 @@

def get_list_gpus_excludes(request: ListGpusRequest) -> Optional[IncludeExcludeDictType]:
list_gpus_excludes: IncludeExcludeDictType = {}
if not request.full_offers:
list_gpus_excludes["full_offers"] = True
run_spec_excludes = get_run_spec_excludes(request.run_spec)
if run_spec_excludes is not None:
list_gpus_excludes["run_spec"] = run_spec_excludes
Expand Down
2 changes: 2 additions & 0 deletions src/dstack/_internal/core/compatibility/runs.py
Original file line number Diff line number Diff line change
Expand Up @@ -70,6 +70,8 @@ def get_get_plan_excludes(request: GetRunPlanRequest) -> Optional[IncludeExclude
clients backward-compatibility with older servers.
"""
get_plan_excludes: IncludeExcludeDictType = {}
if not request.full_offers:
get_plan_excludes["full_offers"] = True
run_spec_excludes = get_run_spec_excludes(request.run_spec)
if run_spec_excludes is not None:
get_plan_excludes["run_spec"] = run_spec_excludes
Expand Down
1 change: 1 addition & 0 deletions src/dstack/_internal/server/routers/gpus.py
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,7 @@ async def list_gpus(
project=project,
run_spec=body.run_spec,
group_by=body.group_by,
full_offers=body.full_offers,
)
patch_list_gpus_response(resp, client_version)
return resp
1 change: 1 addition & 0 deletions src/dstack/_internal/server/routers/runs.py
Original file line number Diff line number Diff line change
Expand Up @@ -139,6 +139,7 @@ async def get_plan(
user=user,
run_spec=body.run_spec,
max_offers=body.max_offers,
full_offers=body.full_offers,
legacy_repo_dir=legacy_repo_dir,
)
patch_run_plan(run_plan, client_version)
Expand Down
5 changes: 4 additions & 1 deletion src/dstack/_internal/server/schemas/gpus.py
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
from typing import List, Literal, Optional
from typing import Annotated, List, Literal, Optional

from pydantic import Field

Expand All @@ -16,6 +16,9 @@ class ListGpusRequest(CoreModel):
description="List of fields to group by. Valid values: 'backend', 'region', 'count'. "
"Note: 'region' can only be used together with 'backend'.",
)
full_offers: Annotated[
bool, Field(description="Don't adjust backend offers by requirements")
] = False


class ListGpusResponse(CoreModel):
Expand Down
3 changes: 3 additions & 0 deletions src/dstack/_internal/server/schemas/runs.py
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,9 @@ class GetRunPlanRequest(CoreModel):
max_offers: Optional[int] = Field(
description="The maximum number of offers to return", ge=1, le=10000
)
full_offers: Annotated[
bool, Field(description="Return full offers not adjusted by requirements")
] = False


class SubmitRunRequest(CoreModel):
Expand Down
Loading
Loading