From bccda348ee4dd7ef653febb84c6bc3d417a466c5 Mon Sep 17 00:00:00 2001 From: Victor Skvortsov Date: Tue, 21 Jul 2026 14:09:02 +0500 Subject: [PATCH 1/2] Fix creation_policy ignored --- .../pipeline_tasks/jobs_submitted.py | 11 ++- .../pipeline_tasks/test_submitted_jobs.py | 72 +++++++++++++++++++ 2 files changed, 82 insertions(+), 1 deletion(-) diff --git a/src/dstack/_internal/server/background/pipeline_tasks/jobs_submitted.py b/src/dstack/_internal/server/background/pipeline_tasks/jobs_submitted.py index 1c782ef94..215635fd4 100644 --- a/src/dstack/_internal/server/background/pipeline_tasks/jobs_submitted.py +++ b/src/dstack/_internal/server/background/pipeline_tasks/jobs_submitted.py @@ -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, @@ -524,6 +525,8 @@ 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=context.run.run_spec.merged_profile.creation_policy + == CreationPolicy.REUSE, skip_backend_offers_on_pool_capacity=True, ) @@ -537,6 +540,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) @@ -1239,7 +1248,7 @@ async def _process_provisioning( 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", + message="Could not reuse any instance for this job", ) return await _process_new_capacity_provisioning( diff --git a/src/tests/_internal/server/background/pipeline_tasks/test_submitted_jobs.py b/src/tests/_internal/server/background/pipeline_tasks/test_submitted_jobs.py index b00ce5902..30c038945 100644 --- a/src/tests/_internal/server/background/pipeline_tasks/test_submitted_jobs.py +++ b/src/tests/_internal/server/background/pipeline_tasks/test_submitted_jobs.py @@ -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, @@ -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 From fc11c3b90cad5fbec25fe0d95d3f8b5426ab031f Mon Sep 17 00:00:00 2001 From: Victor Skvortsov Date: Tue, 21 Jul 2026 15:27:52 +0500 Subject: [PATCH 2/2] Drop legacy reuse check --- .../server/background/pipeline_tasks/jobs_submitted.py | 10 +--------- 1 file changed, 1 insertion(+), 9 deletions(-) diff --git a/src/dstack/_internal/server/background/pipeline_tasks/jobs_submitted.py b/src/dstack/_internal/server/background/pipeline_tasks/jobs_submitted.py index 215635fd4..1824e0070 100644 --- a/src/dstack/_internal/server/background/pipeline_tasks/jobs_submitted.py +++ b/src/dstack/_internal/server/background/pipeline_tasks/jobs_submitted.py @@ -525,8 +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=context.run.run_spec.merged_profile.creation_policy - == CreationPolicy.REUSE, + skip_backend_offers=creation_policy == CreationPolicy.REUSE, skip_backend_offers_on_pool_capacity=True, ) @@ -1244,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 instance for this job", - ) - return await _process_new_capacity_provisioning( item=item, context=context,