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
Original file line number Diff line number Diff line change
Expand Up @@ -514,6 +514,7 @@ async def _select_assignment(
preconditions: _ProcessedPreconditions,
candidate_fleet_models: list[FleetModel],
) -> _AssignmentResult:
creation_policy = context.run.run_spec.merged_profile.creation_policy
# Getting backend offers can be slow, so fleet selection must happen outside the DB transaction.
fleet_model, fleet_instances_with_offers, _ = await find_optimal_fleet_with_offers(
project=context.project,
Expand All @@ -524,6 +525,7 @@ async def _select_assignment(
master_job_provisioning_data=preconditions.master_job_provisioning_data,
volumes=preconditions.prepared_job_volumes.volumes,
exclude_not_available=True,
skip_backend_offers=creation_policy == CreationPolicy.REUSE,
skip_backend_offers_on_pool_capacity=True,
)

Expand All @@ -537,6 +539,12 @@ async def _select_assignment(
volumes=preconditions.prepared_job_volumes.volumes,
)

if creation_policy == CreationPolicy.REUSE:
return _TerminateSubmittedJobResult(
reason=JobTerminationReason.FAILED_TO_START_DUE_TO_NO_CAPACITY,
message="Could not reuse any instance for this job",
)

return _NewCapacityAssignment(fleet_id=fleet_model.id)


Expand Down Expand Up @@ -1235,13 +1243,6 @@ async def _process_provisioning(
prepared_job_volumes=preconditions.prepared_job_volumes,
)

if context.run.run_spec.merged_profile.creation_policy == CreationPolicy.REUSE:
logger.debug("%s: reuse instance failed", fmt(context.job_model))
return _TerminateSubmittedJobResult(
reason=JobTerminationReason.FAILED_TO_START_DUE_TO_NO_CAPACITY,
message="Could not reuse any instances for this job",
)

return await _process_new_capacity_provisioning(
item=item,
context=context,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
from dstack._internal.core.models.instances import InstanceStatus
from dstack._internal.core.models.placement import PlacementGroup
from dstack._internal.core.models.profiles import (
CreationPolicy,
FleetInstanceSelector,
InstanceHostnameSelector,
InstanceNameSelector,
Expand Down Expand Up @@ -1793,6 +1794,77 @@ async def test_assignment_creates_placeholder_instance_for_new_capacity(
assert placeholder.offer is None
assert placeholder.instance_num == 0

async def test_assigns_job_to_instance_with_reuse_creation_policy(
self, test_db, session: AsyncSession, worker: JobSubmittedWorker
):
project = await create_project(session=session)
user = await create_user(session=session)
repo = await create_repo(session=session, project_id=project.id)
fleet = await create_fleet(session=session, project=project)
instance = await create_instance(
session=session,
project=project,
fleet=fleet,
status=InstanceStatus.IDLE,
)
run_spec = get_run_spec(
repo_id=repo.name,
profile=Profile(creation_policy=CreationPolicy.REUSE),
)
run = await create_run(
session=session, project=project, repo=repo, user=user, run_spec=run_spec
)
job = await create_job(session=session, run=run)

await _process_job(session=session, worker=worker, job_model=job)

job = await _get_job(session, job.id)
assert job.status == JobStatus.SUBMITTED
assert job.instance_assigned
assert job.instance is not None and job.instance.id == instance.id
assert job.fleet_id == fleet.id

async def test_terminates_job_when_no_reusable_instances_with_reuse_creation_policy(
self, test_db, session: AsyncSession, worker: JobSubmittedWorker
):
project = await create_project(session=session)
user = await create_user(session=session)
repo = await create_repo(session=session, project_id=project.id)
fleet = await create_fleet(session=session, project=project)
await create_instance(
session=session,
project=project,
fleet=fleet,
status=InstanceStatus.BUSY,
)
run_spec = get_run_spec(
repo_id=repo.name,
profile=Profile(creation_policy=CreationPolicy.REUSE),
)
run = await create_run(
session=session, project=project, repo=repo, user=user, run_spec=run_spec
)
job = await create_job(session=session, run=run)

with patch("dstack._internal.server.services.backends.get_project_backends") as m:
await _process_job(session=session, worker=worker, job_model=job)

# Backend offers must not be requested with the reuse policy.
m.assert_not_called()
job = await _get_job(session, job.id)
assert job.status == JobStatus.TERMINATING
assert job.termination_reason == JobTerminationReason.FAILED_TO_START_DUE_TO_NO_CAPACITY
assert job.termination_reason_message == "Could not reuse any instance for this job"
assert not job.instance_assigned
# No placeholder must be created when reuse fails.
res = await session.execute(
select(InstanceModel).where(
InstanceModel.fleet_id == fleet.id,
InstanceModel.deleted == False,
)
)
assert len(res.scalars().all()) == 1

@pytest.mark.parametrize("fleet_type", ["cloud", "ssh"])
async def test_job_fails_when_fleet_is_full(
self, test_db, session: AsyncSession, worker: JobSubmittedWorker, fleet_type: str
Expand Down
Loading