Skip to content
Merged
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
71 changes: 39 additions & 32 deletions src/dstack/_internal/server/services/gpus.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,15 +7,16 @@
from dstack._internal.core.models.backends.base import BackendType
from dstack._internal.core.models.gpus import BackendGpu, BackendGpus, GpuGroup
from dstack._internal.core.models.instances import InstanceOfferWithAvailability
from dstack._internal.core.models.profiles import SpotPolicy
from dstack._internal.core.models.profiles import CreationPolicy
from dstack._internal.core.models.resources import Range
from dstack._internal.core.models.runs import Requirements, RunSpec, get_policy_map
from dstack._internal.server.models import ProjectModel
from dstack._internal.core.models.runs import RunSpec
from dstack._internal.server.models import InstanceModel, ProjectModel
from dstack._internal.server.schemas.gpus import ListGpusResponse
from dstack._internal.server.services.jobs import get_jobs_from_run_spec
from dstack._internal.server.services.offers import get_offers_by_requirements
from dstack._internal.server.services.runs.plan import (
get_backend_offers_in_run_candidate_fleets,
get_non_fleet_offers,
get_offers_in_run_candidate_fleets,
get_targeted_instance_offers,
)
from dstack._internal.utils.common import get_or_error

Expand Down Expand Up @@ -57,46 +58,52 @@ async def _get_gpu_offers(
session: AsyncSession,
project: ProjectModel,
run_spec: RunSpec,
) -> List[Tuple[Backend, InstanceOfferWithAvailability]]:
) -> list[InstanceOfferWithAvailability]:
"""Fetches all available instance offers that match the run spec's GPU requirements."""
# NOTE: Basically, this is a simplified version of get_job_plans(); keep them in sync
jobs = await get_jobs_from_run_spec(run_spec=run_spec, secrets={}, replica_num=0)
if len(jobs) == 0:
return []
job = jobs[0]
profile = run_spec.merged_profile
if profile.fleets is not None:
jobs = await get_jobs_from_run_spec(run_spec=run_spec, secrets={}, replica_num=0)
if len(jobs) == 0:
return []
return await get_backend_offers_in_run_candidate_fleets(
skip_backend_offers = profile.creation_policy == CreationPolicy.REUSE

instance_offers: list[tuple[InstanceModel, InstanceOfferWithAvailability]]
backend_offers: list[tuple[Backend, InstanceOfferWithAvailability]]
if profile.instances is not None:
instance_offers = await get_targeted_instance_offers(
session=session,
project=project,
run_spec=run_spec,
job=jobs[0],
volumes=None,
max_offers_per_fleet=None,
job=job,
)
requirements = Requirements(
resources=run_spec.configuration.resources,
max_price=profile.max_price,
spot=get_policy_map(profile.spot_policy, default=SpotPolicy.AUTO),
reservation=profile.reservation,
)
return await get_offers_by_requirements(
project=project,
profile=profile,
requirements=requirements,
exclude_not_available=False,
multinode=False,
volumes=None,
privileged=False,
instance_mounts=False,
)
backend_offers = []
elif profile.fleets is not None:
instance_offers, backend_offers = await get_offers_in_run_candidate_fleets(
session=session,
project=project,
run_spec=run_spec,
job=job,
skip_backend_offers=skip_backend_offers,
)
else:
instance_offers, backend_offers = await get_non_fleet_offers(
session=session,
project=project,
run_spec=run_spec,
job=job,
skip_backend_offers=skip_backend_offers,
)
return [offer for _, offer in instance_offers] + [offer for _, offer in backend_offers]


def _process_offers_into_backend_gpus(
offers: List[Tuple[Backend, InstanceOfferWithAvailability]],
offers: list[InstanceOfferWithAvailability],
) -> List[BackendGpus]:
"""Transforms raw offers into a structured list of BackendGpus, aggregating GPU info."""
backend_data: Dict[BackendType, Dict] = {}

for _, offer in offers:
for offer in offers:
backend_type = offer.backend
if backend_type not in backend_data:
backend_data[backend_type] = {"gpus": {}, "regions": set()}
Expand Down
24 changes: 12 additions & 12 deletions src/dstack/_internal/server/services/runs/plan.py
Original file line number Diff line number Diff line change
Expand Up @@ -166,7 +166,7 @@ async def get_job_plans(
)
backend_offers = []
elif run_spec.merged_profile.fleets is not None:
instance_offers, backend_offers = await _get_offers_in_run_candidate_fleets(
instance_offers, backend_offers = await get_offers_in_run_candidate_fleets(
session=session,
project=project,
run_spec=run_spec,
Expand All @@ -175,7 +175,7 @@ async def get_job_plans(
skip_backend_offers=skip_backend_offers,
)
else:
instance_offers, backend_offers = await _get_non_fleet_offers(
instance_offers, backend_offers = await get_non_fleet_offers(
session=session,
project=project,
run_spec=run_spec,
Expand Down Expand Up @@ -755,7 +755,7 @@ async def _get_pool_offers(
project: ProjectModel,
run_spec: RunSpec,
job: Job,
volumes: list[list[Volume]],
volumes: Optional[list[list[Volume]]],
) -> list[tuple[InstanceModel, InstanceOfferWithAvailability]]:
pool_offers: list[tuple[InstanceModel, InstanceOfferWithAvailability]] = []
detaching_instances_ids = await get_instances_ids_with_detaching_volumes(session)
Expand Down Expand Up @@ -786,12 +786,12 @@ async def _get_pool_offers(
return pool_offers


async def _get_non_fleet_offers(
async def get_non_fleet_offers(
session: AsyncSession,
project: ProjectModel,
run_spec: RunSpec,
job: Job,
volumes: list[list[Volume]],
volumes: Optional[list[list[Volume]]] = None,
skip_backend_offers: bool = False,
) -> tuple[
list[tuple[InstanceModel, InstanceOfferWithAvailability]],
Expand Down Expand Up @@ -836,7 +836,7 @@ async def get_backend_offers_in_run_candidate_fleets(
"""
Returns backend offers across the run's selected candidate fleets.

Used by `dstack offer --fleet ...` and `dstack offer --group-by ... --fleet ...`.
Helper of `get_offers_in_run_candidate_fleets()` that collects the backend part of its offers.
It resolves the selected fleets from `run_spec`, requests backend offers in each fleet,
merges them, and deduplicates identical backend offers across fleets.
"""
Expand Down Expand Up @@ -869,12 +869,12 @@ async def get_backend_offers_in_run_candidate_fleets(
return offers


async def _get_offers_in_run_candidate_fleets(
async def get_offers_in_run_candidate_fleets(
session: AsyncSession,
project: ProjectModel,
run_spec: RunSpec,
job: Job,
volumes: list[list[Volume]],
volumes: Optional[list[list[Volume]]] = None,
skip_backend_offers: bool = False,
) -> tuple[
list[tuple[InstanceModel, InstanceOfferWithAvailability]],
Expand All @@ -883,10 +883,10 @@ async def _get_offers_in_run_candidate_fleets(
"""
Returns existing-instance and backend offers across the run's candidate fleets.

Used by `dstack offer --fleet ...` without `--group-by`. Unlike normal `dstack apply`, it
does not choose a single best fleet. Instead, it gathers existing-instance and backend
offers from each selected fleet, keeps existing instances as separate reusable options, and
deduplicates identical backend offers across fleets.
Used by `dstack offer --fleet ...` (with or without `--group-by`). Unlike normal
`dstack apply`, it does not choose a single best fleet. Instead, it gathers existing-instance
and backend offers from each selected fleet, keeps existing instances as separate reusable
options, and deduplicates identical backend offers across fleets.
"""
candidate_fleet_models = await _select_candidate_fleet_models(
session=session,
Expand Down
118 changes: 117 additions & 1 deletion src/tests/_internal/server/routers/test_gpus.py
Original file line number Diff line number Diff line change
Expand Up @@ -15,17 +15,20 @@
InstanceType,
Resources,
)
from dstack._internal.core.models.profiles import Profile
from dstack._internal.core.models.profiles import CreationPolicy, Profile
from dstack._internal.core.models.runs import RunSpec
from dstack._internal.core.models.users import GlobalRole, ProjectRole
from dstack._internal.server.models import FleetModel, ProjectModel
from dstack._internal.server.services.projects import add_project_member
from dstack._internal.server.testing.common import (
create_fleet,
create_instance,
create_project,
create_repo,
create_user,
get_auth_headers,
get_fleet_spec,
get_instance_offer_with_availability,
get_run_spec,
)

Expand Down Expand Up @@ -79,6 +82,40 @@ def create_gpu_offer(
)


async def create_gpu_pool_instance(
session: AsyncSession,
project: ProjectModel,
fleet: FleetModel,
name: str,
gpu_name: str = "A100",
gpu_memory_gib: float = 80,
price: float = 5.0,
backend: BackendType = BackendType.AWS,
region: str = "us-west-2",
):
"""Helper to create an idle pool instance backed by a GPU offer."""
offer = get_instance_offer_with_availability(
backend=backend,
region=region,
gpu_count=1,
gpu_name=gpu_name,
gpu_memory_gib=gpu_memory_gib,
cpu_count=8,
memory_gib=64,
price=price,
)
return await create_instance(
session=session,
project=project,
fleet=fleet,
name=name,
backend=backend,
region=region,
offer=offer,
price=price,
)


def create_mock_backends_with_offers(
offers_by_backend: Dict[BackendType, List[InstanceOfferWithAvailability]],
) -> List[Mock]:
Expand Down Expand Up @@ -217,6 +254,85 @@ async def test_filters_gpus_by_multiple_specified_fleets(
response_data = response.json()
assert {gpu["backend"] for gpu in response_data["gpus"]} == {"aws", "runpod"}

@pytest.mark.asyncio
@pytest.mark.parametrize("test_db", ["sqlite", "postgres"], indirect=True)
async def test_includes_backend_offers_when_creation_policy_reuse_or_create(
self, test_db, session: AsyncSession, client: AsyncClient
):
user, project, repo, _ = await gpu_test_setup(session)
fleet = await create_fleet(session=session, project=project, name="pool-fleet")
await create_gpu_pool_instance(session, project, fleet, name="pool-instance")

offers_by_backend = {
BackendType.AWS: [create_gpu_offer(BackendType.AWS, "L4", 24576, 1.0)]
}
mocked_backends = create_mock_backends_with_offers(offers_by_backend)

with patch("dstack._internal.server.services.backends.get_project_backends") as m:
m.return_value = mocked_backends
run_spec = get_run_spec(run_name="test-run", repo_id=repo.name)
response = await call_gpus_api(client, project.name, user.token, run_spec)

assert response.status_code == 200
assert {gpu["name"] for gpu in response.json()["gpus"]} == {"A100", "L4"}

@pytest.mark.asyncio
@pytest.mark.parametrize("test_db", ["sqlite", "postgres"], indirect=True)
async def test_skips_backend_offers_when_creation_policy_reuse(
self, test_db, session: AsyncSession, client: AsyncClient
):
user, project, repo, _ = await gpu_test_setup(session)
fleet = await create_fleet(session=session, project=project, name="pool-fleet")
await create_gpu_pool_instance(session, project, fleet, name="pool-instance")

offers_by_backend = {
BackendType.AWS: [create_gpu_offer(BackendType.AWS, "L4", 24576, 1.0)]
}
mocked_backends = create_mock_backends_with_offers(offers_by_backend)

with patch("dstack._internal.server.services.backends.get_project_backends") as m:
m.return_value = mocked_backends
# reuse skips backend offers, keeping only the existing instance.
run_spec = get_run_spec(
run_name="test-run",
repo_id=repo.name,
profile=Profile(name="default", creation_policy=CreationPolicy.REUSE),
)
response = await call_gpus_api(client, project.name, user.token, run_spec)

assert response.status_code == 200
assert {gpu["name"] for gpu in response.json()["gpus"]} == {"A100"}

@pytest.mark.asyncio
@pytest.mark.parametrize("test_db", ["sqlite", "postgres"], indirect=True)
async def test_uses_targeted_instance_offers_when_instances_specified(
self, test_db, session: AsyncSession, client: AsyncClient
):
user, project, repo, _ = await gpu_test_setup(session)
fleet = await create_fleet(session=session, project=project, name="pool-fleet")
await create_gpu_pool_instance(session, project, fleet, name="targeted-instance")
await create_gpu_pool_instance(
session, project, fleet, name="other-instance", gpu_name="H100"
)
run_spec = get_run_spec(
run_name="test-run",
repo_id=repo.name,
profile=Profile(name="default", instances=["targeted-instance"]),
)

offers_by_backend = {
BackendType.AWS: [create_gpu_offer(BackendType.AWS, "L4", 24576, 1.0)]
}
mocked_backends = create_mock_backends_with_offers(offers_by_backend)

with patch("dstack._internal.server.services.backends.get_project_backends") as m:
m.return_value = mocked_backends
response = await call_gpus_api(client, project.name, user.token, run_spec)

assert response.status_code == 200
# Only the selected instance's GPU is listed: not the other instance, not backend offers.
assert {gpu["name"] for gpu in response.json()["gpus"]} == {"A100"}

@pytest.mark.asyncio
@pytest.mark.parametrize("test_db", ["sqlite", "postgres"], indirect=True)
async def test_returns_empty_gpus_when_no_offers(
Expand Down
20 changes: 10 additions & 10 deletions src/tests/_internal/server/routers/test_runs.py
Original file line number Diff line number Diff line change
Expand Up @@ -2628,14 +2628,14 @@ async def test_offer_without_fleets_uses_global_offer_collection(
global_offer = get_instance_offer_with_availability(price=1.0)
with (
patch(
"dstack._internal.server.services.runs.plan._get_non_fleet_offers",
"dstack._internal.server.services.runs.plan.get_non_fleet_offers",
new=AsyncMock(return_value=([(Mock(), global_offer)], [])),
) as get_non_fleet_offers_mock,
patch(
"dstack._internal.server.services.runs.plan._get_offers_in_run_candidate_fleets",
"dstack._internal.server.services.runs.plan.get_offers_in_run_candidate_fleets",
new=AsyncMock(
side_effect=AssertionError(
"_get_offers_in_run_candidate_fleets should not be called"
"get_offers_in_run_candidate_fleets should not be called"
)
),
) as get_offers_in_run_candidate_fleets_mock,
Expand Down Expand Up @@ -2693,13 +2693,13 @@ async def test_offer_with_fleets_uses_selected_fleet_offer_collection(
fleet_offer = get_instance_offer_with_availability(price=2.0)
with (
patch(
"dstack._internal.server.services.runs.plan._get_non_fleet_offers",
"dstack._internal.server.services.runs.plan.get_non_fleet_offers",
new=AsyncMock(
side_effect=AssertionError("_get_non_fleet_offers should not be called")
side_effect=AssertionError("get_non_fleet_offers should not be called")
),
) as get_non_fleet_offers_mock,
patch(
"dstack._internal.server.services.runs.plan._get_offers_in_run_candidate_fleets",
"dstack._internal.server.services.runs.plan.get_offers_in_run_candidate_fleets",
new=AsyncMock(return_value=([(Mock(), fleet_offer)], [])),
) as get_offers_in_run_candidate_fleets_mock,
patch(
Expand Down Expand Up @@ -2761,16 +2761,16 @@ async def test_regular_run_plan_uses_best_fleet_candidate_selection(
new=AsyncMock(return_value=(Mock(), [(Mock(), chosen_fleet_offer)], [])),
) as find_optimal_fleet_with_offers_mock,
patch(
"dstack._internal.server.services.runs.plan._get_non_fleet_offers",
"dstack._internal.server.services.runs.plan.get_non_fleet_offers",
new=AsyncMock(
side_effect=AssertionError("_get_non_fleet_offers should not be called")
side_effect=AssertionError("get_non_fleet_offers should not be called")
),
) as get_non_fleet_offers_mock,
patch(
"dstack._internal.server.services.runs.plan._get_offers_in_run_candidate_fleets",
"dstack._internal.server.services.runs.plan.get_offers_in_run_candidate_fleets",
new=AsyncMock(
side_effect=AssertionError(
"_get_offers_in_run_candidate_fleets should not be called"
"get_offers_in_run_candidate_fleets should not be called"
)
),
) as get_offers_in_run_candidate_fleets_mock,
Expand Down
Loading