Skip to content
Closed
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
Original file line number Diff line number Diff line change
Expand Up @@ -116,8 +116,8 @@
from dstack._internal.server.services.locking import get_locker
from dstack._internal.server.services.logging import fmt
from dstack._internal.server.services.offers import (
generate_shared_offer,
get_instance_offer_with_restricted_az,
get_offers_by_requirements,
)
from dstack._internal.server.services.pipelines import PipelineHinterProtocol
from dstack._internal.server.services.placement import (
Expand All @@ -128,16 +128,14 @@
)
from dstack._internal.server.services.runs import run_model_to_run
from dstack._internal.server.services.runs.plan import (
_get_backend_offers_in_fleet,
find_optimal_fleet_with_offers,
get_instance_offers_in_fleet,
get_run_candidate_fleet_models_filters,
get_run_profile_and_requirements_in_fleet,
get_targeted_instance_offers,
select_run_candidate_fleet_models_with_filters,
)
from dstack._internal.server.services.runs.spec import (
check_run_spec_requires_instance_mounts,
)
from dstack._internal.server.services.secrets import get_project_secrets_mapping
from dstack._internal.server.services.volumes import volume_model_to_volume
from dstack._internal.server.utils import tracing
Expand Down Expand Up @@ -1738,8 +1736,11 @@ async def _promote_or_create_instance_models_for_provisioned_jobs(
provisioned_job_model.used_instance_id = instance_model.id

instance_models.append(instance_model)
runtime_offer = offer
if offer.blocks < offer.total_blocks:
runtime_offer = generate_shared_offer(offer, offer.blocks, offer.total_blocks)
provisioned_job_model.job_runtime_data = _prepare_job_runtime_data(
offer, context.multinode
runtime_offer, context.multinode
).model_dump_json()
events.emit(
session,
Expand Down Expand Up @@ -1817,8 +1818,8 @@ def _create_instance_model_for_job(
price=offer.price,
region=offer.region,
volume_attachments=[],
total_blocks=1,
busy_blocks=1,
total_blocks=offer.total_blocks,
busy_blocks=offer.blocks,
)


Expand Down Expand Up @@ -1848,8 +1849,8 @@ def _promote_placeholder_instance(
instance_model.region = offer.region
instance_model.termination_policy = termination_policy
instance_model.termination_idle_time = termination_idle_time
instance_model.total_blocks = 1
instance_model.busy_blocks = 1
instance_model.total_blocks = offer.total_blocks
instance_model.busy_blocks = offer.blocks


async def _process_volume_attachments(
Expand Down Expand Up @@ -2405,18 +2406,15 @@ async def _provision_new_capacity(
placement_group_models=placement_group_models,
fleet_model=fleet_model,
)
multinode = requirements.multinode or is_multinode_job(job)
offers = await get_offers_by_requirements(
offers = await _get_backend_offers_in_fleet(
project=project,
profile=profile,
requirements=requirements,
exclude_not_available=True,
multinode=multinode,
master_job_provisioning_data=master_job_provisioning_data,
fleet_model=fleet_model,
run_spec=run.run_spec,
job=job,
volumes=volumes,
privileged=job.job_spec.privileged,
instance_mounts=check_run_spec_requires_instance_mounts(run.run_spec),
placement_group=placement_group_model_to_placement_group_optional(placement_group_model),
exclude_not_available=True,
master_job_provisioning_data=master_job_provisioning_data,
)
offers_iter = iter(offers)
offers_tried = 0
Expand Down
29 changes: 15 additions & 14 deletions src/dstack/_internal/server/services/instances.py
Original file line number Diff line number Diff line change
Expand Up @@ -59,7 +59,7 @@
from dstack._internal.server.schemas.runner import InstanceHealthResponse, TaskStatus
from dstack._internal.server.services import events
from dstack._internal.server.services.logging import fmt
from dstack._internal.server.services.offers import generate_shared_offer
from dstack._internal.server.services.offers import get_matching_shared_offer
from dstack._internal.server.services.projects import list_user_project_models
from dstack._internal.server.services.runner.client import ShimClient
from dstack._internal.utils import common as common_utils
Expand Down Expand Up @@ -693,7 +693,6 @@ def get_shared_instances_with_offers(
volumes: Optional[List[List[Volume]]] = None,
) -> list[tuple[InstanceModel, InstanceOfferWithAvailability]]:
instances_with_offers: list[tuple[InstanceModel, InstanceOfferWithAvailability]] = []
query_filter = requirements_to_query_filter(requirements)
filtered_instances = filter_instances(
instances=instances,
profile=profile,
Expand All @@ -711,18 +710,20 @@ def get_shared_instances_with_offers(
continue
total_blocks = common_utils.get_or_error(instance.total_blocks)
idle_blocks = total_blocks - instance.busy_blocks
min_blocks = total_blocks if multinode else 1
for blocks in range(min_blocks, total_blocks + 1):
shared_offer = generate_shared_offer(offer, blocks, total_blocks)
catalog_item = offer_to_catalog_item(shared_offer)
if gpuhunt.matches(catalog_item, query_filter):
if blocks <= idle_blocks:
shared_offer.availability = InstanceAvailability.IDLE
else:
shared_offer.availability = InstanceAvailability.BUSY
if shared_offer.availability == InstanceAvailability.IDLE or not idle_only:
instances_with_offers.append((instance, shared_offer))
break
shared_offer = get_matching_shared_offer(
offer,
requirements=requirements,
blocks=total_blocks,
multinode=multinode,
)
if shared_offer is None:
continue
if shared_offer.blocks <= idle_blocks:
shared_offer.availability = InstanceAvailability.IDLE
else:
shared_offer.availability = InstanceAvailability.BUSY
if shared_offer.availability == InstanceAvailability.IDLE or not idle_only:
instances_with_offers.append((instance, shared_offer))
return instances_with_offers


Expand Down
100 changes: 87 additions & 13 deletions src/dstack/_internal/server/services/offers.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,8 +6,13 @@
import gpuhunt

from dstack._internal.core.backends.base.backend import Backend
from dstack._internal.core.backends.base.compute import ComputeWithPlacementGroupSupport
from dstack._internal.core.backends.base.compute import (
ComputeWithCreateInstanceSupport,
ComputeWithPlacementGroupSupport,
)
from dstack._internal.core.backends.base.offers import filter_offers_by_requirements
from dstack._internal.core.backends.features import (
BACKENDS_WITH_CREATE_INSTANCE_SUPPORT,
BACKENDS_WITH_INSTANCE_VOLUMES_SUPPORT,
BACKENDS_WITH_MULTINODE_SUPPORT,
BACKENDS_WITH_PRIVILEGED_SUPPORT,
Expand All @@ -31,6 +36,7 @@ async def get_offers_by_requirements(
project: ProjectModel,
profile: Profile,
requirements: Requirements,
shared_offer_requirements: Optional[Requirements] = None,
exclude_not_available=False,
multinode: bool = False,
master_job_provisioning_data: Optional[JobProvisioningData] = None,
Expand Down Expand Up @@ -90,6 +96,11 @@ async def get_offers_by_requirements(
if backend_types is not None:
backends = [b for b in backends if b.TYPE in backend_types or b.TYPE == BackendType.DSTACK]

if blocks != 1 and shared_offer_requirements is not None:
backends = [
b for b in backends if isinstance(b.compute(), ComputeWithCreateInstanceSupport)
]

offers = await backends_services.get_backend_offers(
backends=backends,
requirements=requirements,
Expand All @@ -109,8 +120,16 @@ async def get_offers_by_requirements(
volumes_locations=volumes_locations,
)

if blocks != 1 and shared_offer_requirements is not None:
offers = _filter_shared_block_supported_offers(offers)

if blocks != 1:
offers = _get_shareable_offers(offers, blocks)
offers = _get_shareable_offers(
offers,
blocks,
requirements=shared_offer_requirements,
multinode=multinode,
)

if max_offers is not None:
offers = itertools.islice(offers, max_offers)
Expand Down Expand Up @@ -159,6 +178,7 @@ def generate_shared_offer(
) -> InstanceOfferWithAvailability:
full_resources = offer.instance.resources
resources = Resources(
cpu_arch=full_resources.cpu_arch,
cpus=full_resources.cpus // total_blocks * blocks,
memory_mib=full_resources.memory_mib // total_blocks * blocks,
gpus=full_resources.gpus[: len(full_resources.gpus) // total_blocks * blocks],
Expand All @@ -180,6 +200,30 @@ def generate_shared_offer(
)


def get_matching_shared_offer(
offer: InstanceOfferWithAvailability,
requirements: Requirements,
blocks: Union[int, Literal["auto"]],
*,
multinode: bool = False,
) -> Optional[InstanceOfferWithAvailability]:
resources = offer.instance.resources
gpu_count = len(resources.gpus)
if gpu_count > 0 and resources.gpus[0].vendor == gpuhunt.AcceleratorVendor.GOOGLE:
# TPUs cannot be shared.
gpu_count = 1
divisible, total_blocks = is_divisible_into_blocks(resources.cpus, gpu_count, blocks)
if not divisible:
return None

min_blocks = total_blocks if multinode else 1
shared_offers = (
generate_shared_offer(offer, selected_blocks, total_blocks)
for selected_blocks in range(min_blocks, total_blocks + 1)
)
return next(filter_offers_by_requirements(shared_offers, requirements), None)


def get_instance_offer_with_restricted_az(
instance_offer: InstanceOfferWithAvailability,
master_job_provisioning_data: Optional[JobProvisioningData],
Expand Down Expand Up @@ -254,20 +298,50 @@ def _filter_offers(
def _get_shareable_offers(
offers: Iterable[Tuple[Backend, InstanceOfferWithAvailability]],
blocks: Union[int, Literal["auto"]],
*,
requirements: Optional[Requirements] = None,
multinode: bool = False,
) -> Iterator[Tuple[Backend, InstanceOfferWithAvailability]]:
"""
Yields offers that can be shared with `total_blocks` set.
"""
for backend, offer in offers:
resources = offer.instance.resources
cpu_count = resources.cpus
gpu_count = len(resources.gpus)
if gpu_count > 0 and resources.gpus[0].vendor == gpuhunt.AcceleratorVendor.GOOGLE:
# TPUs cannot be shared
gpu_count = 1
divisible, total_blocks = is_divisible_into_blocks(cpu_count, gpu_count, blocks)
if not divisible:
if requirements is None:
resources = offer.instance.resources
cpu_count = resources.cpus
gpu_count = len(resources.gpus)
if gpu_count > 0 and resources.gpus[0].vendor == gpuhunt.AcceleratorVendor.GOOGLE:
# TPUs cannot be shared
gpu_count = 1
divisible, total_blocks = is_divisible_into_blocks(cpu_count, gpu_count, blocks)
if not divisible:
continue
new_offer = offer.model_copy()
new_offer.total_blocks = total_blocks
yield (backend, new_offer)
continue

shared_offer = get_matching_shared_offer(
offer,
requirements=requirements,
blocks=blocks,
multinode=multinode,
)
if shared_offer is None:
continue
annotated_offer = offer.model_copy()
annotated_offer.blocks = shared_offer.blocks
annotated_offer.total_blocks = shared_offer.total_blocks
yield backend, annotated_offer


def _filter_shared_block_supported_offers(
offers: Iterable[Tuple[Backend, InstanceOfferWithAvailability]],
) -> Iterator[Tuple[Backend, InstanceOfferWithAvailability]]:
for backend, offer in offers:
if backend.TYPE == offer.backend:
yield backend, offer
continue
if offer.backend not in BACKENDS_WITH_CREATE_INSTANCE_SUPPORT:
continue
new_offer = offer.model_copy()
new_offer.total_blocks = total_blocks
yield (backend, new_offer)
yield backend, offer
18 changes: 17 additions & 1 deletion src/dstack/_internal/server/services/requirements/combine.py
Original file line number Diff line number Diff line change
Expand Up @@ -62,10 +62,26 @@ def combine_fleet_and_run_profiles(

def combine_fleet_and_run_requirements(
fleet_requirements: Requirements, run_requirements: Requirements
) -> Optional[Requirements]:
try:
resources = _combine_resources(fleet_requirements.resources, run_requirements.resources)
except CombineError:
return None
return combine_fleet_and_run_requirements_with_resources(
fleet_requirements=fleet_requirements,
run_requirements=run_requirements,
resources=resources,
)


def combine_fleet_and_run_requirements_with_resources(
fleet_requirements: Requirements,
run_requirements: Requirements,
resources: ResourcesSpec,
) -> Optional[Requirements]:
try:
return Requirements(
resources=_combine_resources(fleet_requirements.resources, run_requirements.resources),
resources=resources.model_copy(deep=True),
max_price=_get_min_optional(fleet_requirements.max_price, run_requirements.max_price),
spot=_combine_spot_optional(fleet_requirements.spot, run_requirements.spot),
reservation=get_single_value_optional(
Expand Down
Loading