diff --git a/src/dstack/_internal/server/background/pipeline_tasks/jobs_terminating.py b/src/dstack/_internal/server/background/pipeline_tasks/jobs_terminating.py index aac0fd8d64..3f75070fe4 100644 --- a/src/dstack/_internal/server/background/pipeline_tasks/jobs_terminating.py +++ b/src/dstack/_internal/server/background/pipeline_tasks/jobs_terminating.py @@ -332,6 +332,7 @@ class _UnregisterReplicaResult: class _ProcessResult: job_update_map: _JobUpdateMap = field(default_factory=_JobUpdateMap) instance_update_map: Optional[_InstanceUpdateMap] = None + delete_instance: bool = False volume_update_rows: list[_VolumeUpdateRow] = field(default_factory=list) detached_volume_ids: set[uuid.UUID] = field(default_factory=set) unassign_event_message: Optional[str] = None @@ -476,7 +477,11 @@ async def _apply_process_result( ) -> None: set_processed_update_map_fields(result.job_update_map) set_unlock_update_map_fields(result.job_update_map) - if instance_model is not None and result.instance_update_map is None: + if ( + instance_model is not None + and result.instance_update_map is None + and not result.delete_instance + ): result.instance_update_map = _InstanceUpdateMap() if result.instance_update_map is not None: set_processed_update_map_fields(result.instance_update_map) @@ -514,7 +519,24 @@ async def _apply_process_result( ) return - if instance_model is not None and instance_update_map is not None: + if instance_model is not None and result.delete_instance: + res = await session.execute( + delete(InstanceModel) + .where( + InstanceModel.id == instance_model.id, + InstanceModel.lock_token == item.lock_token, + InstanceModel.lock_owner == related_instance_lock_owner, + ) + .returning(InstanceModel.id) + ) + deleted_ids = list(res.scalars().all()) + if len(deleted_ids) == 0: + logger.error( + "Failed to delete placeholder instance %s for terminating job %s.", + instance_model.id, + item.id, + ) + elif instance_model is not None and instance_update_map is not None: res = await session.execute( update(InstanceModel) .where( @@ -647,13 +669,11 @@ async def _process_terminating_job( return result if is_placeholder_instance(instance_model): - # Placeholder has no VM and no provisioning data. Skip graceful stop, - # container stop, and volume detach. - instance_update_map = get_or_error(result.instance_update_map) - if instance_model.status != InstanceStatus.TERMINATING: - instance_update_map["status"] = InstanceStatus.TERMINATING - instance_update_map["skip_min_processing_interval"] = True - instance_update_map["termination_reason"] = InstanceTerminationReason.JOB_FINISHED + # Placeholder has no VM and no provisioning data, so there is nothing to terminate. + # Skip graceful stop, container stop, and volume detach. Delete the instance right away + # to avoid placeholders accumulating for runs retrying on no capacity. + result.instance_update_map = None + result.delete_instance = True result.job_update_map["instance_id"] = None await _unregister_replica_and_update_result(result=result, job_model=job_model) result.job_update_map["status"] = _get_job_termination_status(job_model) diff --git a/src/dstack/_internal/server/background/pipeline_tasks/runs/__init__.py b/src/dstack/_internal/server/background/pipeline_tasks/runs/__init__.py index 40aae39bc9..727535f5bc 100644 --- a/src/dstack/_internal/server/background/pipeline_tasks/runs/__init__.py +++ b/src/dstack/_internal/server/background/pipeline_tasks/runs/__init__.py @@ -25,6 +25,9 @@ set_processed_update_map_fields, set_unlock_update_map_fields, ) +from dstack._internal.server.background.pipeline_tasks.runs.common import ( + delete_superseded_no_capacity_job_submissions, +) from dstack._internal.server.db import get_db, get_session_ctx from dstack._internal.server.models import InstanceModel, JobModel, ProjectModel, RunModel from dstack._internal.server.services import events @@ -415,14 +418,17 @@ async def _apply_pending_result( actor=events.SystemActor(), targets=[events.Target.from_model(job_model)], ) - + await delete_superseded_no_capacity_job_submissions( + session=session, + run_id=item.id, + new_job_models=result.new_job_models, + ) emit_run_status_change_event( session=session, run_model=context.run_model, old_status=context.run_model.status, new_status=result.run_update_map.get("status", context.run_model.status), ) - await _unlock_related_jobs( session=session, item=item, @@ -596,6 +602,11 @@ async def _apply_active_result( actor=events.SystemActor(), targets=[events.Target.from_model(job_model)], ) + await delete_superseded_no_capacity_job_submissions( + session=session, + run_id=item.id, + new_job_models=result.new_job_models, + ) old_status = run_model.status new_status = result.run_update_map.get("status", old_status) diff --git a/src/dstack/_internal/server/background/pipeline_tasks/runs/active.py b/src/dstack/_internal/server/background/pipeline_tasks/runs/active.py index 8448c52c32..aafe132eab 100644 --- a/src/dstack/_internal/server/background/pipeline_tasks/runs/active.py +++ b/src/dstack/_internal/server/background/pipeline_tasks/runs/active.py @@ -475,8 +475,8 @@ async def _build_retry_job_models( run_model=context.run_model, job=new_job, status=JobStatus.SUBMITTED, + submission_num=old_job_model.submission_num + 1, ) - job_model.submission_num = old_job_model.submission_num + 1 new_job_models.append(job_model) return new_job_models diff --git a/src/dstack/_internal/server/background/pipeline_tasks/runs/common.py b/src/dstack/_internal/server/background/pipeline_tasks/runs/common.py index 0c6c9730c4..3f27e3b2a8 100644 --- a/src/dstack/_internal/server/background/pipeline_tasks/runs/common.py +++ b/src/dstack/_internal/server/background/pipeline_tasks/runs/common.py @@ -1,12 +1,16 @@ import json +import uuid from datetime import datetime from typing import Optional +from sqlalchemy import delete +from sqlalchemy.ext.asyncio import AsyncSession + from dstack._internal.core.models.configurations import ( DEFAULT_REPLICA_GROUP_NAME, ServiceConfiguration, ) -from dstack._internal.core.models.runs import JobStatus, RunSpec +from dstack._internal.core.models.runs import JobStatus, JobTerminationReason, RunSpec from dstack._internal.proxy.gateway.schemas.stats import PerWindowStats from dstack._internal.server.models import JobModel, RunModel from dstack._internal.server.services.jobs import get_job_spec, get_jobs_from_run_spec @@ -86,8 +90,8 @@ async def build_scale_up_job_models( run_model=run_model, job=new_job, status=JobStatus.SUBMITTED, + submission_num=old_job_model.submission_num + 1, ) - job_model.submission_num = old_job_model.submission_num + 1 new_job_models.append(job_model) scheduled_replicas += 1 @@ -115,3 +119,32 @@ async def build_scale_up_job_models( new_job_models.append(job_model) return new_job_models + + +async def delete_superseded_no_capacity_job_submissions( + session: AsyncSession, + run_id: uuid.UUID, + new_job_models: list[JobModel], +) -> None: + """ + Delete previous job submissions that failed to start due to no capacity without + ever provisioning. Such submissions are created each time a run is retrying + on no capacity and carry no information beyond the events. + The direct predecessor of each new submission is kept so that the run's + `status_message` can still show `retrying` while the new submission is in flight. + """ + for job_model in new_job_models: + if job_model.submission_num < 2: + continue + await session.execute( + delete(JobModel).where( + JobModel.run_id == run_id, + JobModel.replica_num == job_model.replica_num, + JobModel.job_num == job_model.job_num, + JobModel.submission_num < job_model.submission_num - 1, + JobModel.status.in_(JobStatus.finished_statuses()), + JobModel.termination_reason + == JobTerminationReason.FAILED_TO_START_DUE_TO_NO_CAPACITY, + JobModel.job_provisioning_data.is_(None), + ) + ) diff --git a/src/dstack/_internal/server/migrations/versions/2026/07_22_1256_ad348ea93493_add_eventtargetmodel_entity_run_id.py b/src/dstack/_internal/server/migrations/versions/2026/07_22_1256_ad348ea93493_add_eventtargetmodel_entity_run_id.py new file mode 100644 index 0000000000..1b5e2481f9 --- /dev/null +++ b/src/dstack/_internal/server/migrations/versions/2026/07_22_1256_ad348ea93493_add_eventtargetmodel_entity_run_id.py @@ -0,0 +1,51 @@ +"""Add EventTargetModel.entity_run_id + +Revision ID: ad348ea93493 +Revises: e9c5e7e26c78 +Create Date: 2026-07-22 12:56:04.527041+00:00 + +""" + +import sqlalchemy as sa +import sqlalchemy_utils +from alembic import op + +# revision identifiers, used by Alembic. +revision = "ad348ea93493" +down_revision = "e9c5e7e26c78" +branch_labels = None +depends_on = None + + +def upgrade() -> None: + # ### commands auto generated by Alembic - please adjust! ### + with op.batch_alter_table("event_targets", schema=None) as batch_op: + batch_op.add_column( + sa.Column( + "entity_run_id", sqlalchemy_utils.types.uuid.UUIDType(binary=False), nullable=True + ) + ) + batch_op.create_index( + batch_op.f("ix_event_targets_entity_run_id"), ["entity_run_id"], unique=False + ) + batch_op.create_foreign_key( + batch_op.f("fk_event_targets_entity_run_id_runs"), + "runs", + ["entity_run_id"], + ["id"], + ondelete="CASCADE", + ) + + # ### end Alembic commands ### + + +def downgrade() -> None: + # ### commands auto generated by Alembic - please adjust! ### + with op.batch_alter_table("event_targets", schema=None) as batch_op: + batch_op.drop_constraint( + batch_op.f("fk_event_targets_entity_run_id_runs"), type_="foreignkey" + ) + batch_op.drop_index(batch_op.f("ix_event_targets_entity_run_id")) + batch_op.drop_column("entity_run_id") + + # ### end Alembic commands ### diff --git a/src/dstack/_internal/server/migrations/versions/2026/07_22_1331_87d4312605e5_backfill_eventtargetmodel_entity_run_id.py b/src/dstack/_internal/server/migrations/versions/2026/07_22_1331_87d4312605e5_backfill_eventtargetmodel_entity_run_id.py new file mode 100644 index 0000000000..20bc68bf28 --- /dev/null +++ b/src/dstack/_internal/server/migrations/versions/2026/07_22_1331_87d4312605e5_backfill_eventtargetmodel_entity_run_id.py @@ -0,0 +1,45 @@ +"""Backfill EventTargetModel.entity_run_id + +Revision ID: 87d4312605e5 +Revises: ad348ea93493 +Create Date: 2026-07-22 13:31:47.070886+00:00 + +""" + +from alembic import op + +# revision identifiers, used by Alembic. +revision = "87d4312605e5" +down_revision = "ad348ea93493" +branch_labels = None +depends_on = None + + +def upgrade() -> None: + # Events recorded before entity_run_id was introduced have it unset. + # The within_runs events filter relies on entity_run_id, so backfill it. + # Job targets are backfilled via the jobs table, so this migration must run + # before the job models the events reference are deleted + # (e.g. superseded no capacity submissions deleted on resubmission). + # Old replicas can still record events without entity_run_id while + # this migration is being deployed. Such events won't be backfilled and + # won't be returned by the within_runs events filter. + op.execute( + """ + UPDATE event_targets SET entity_run_id = entity_id + WHERE entity_type = 'run' AND entity_run_id IS NULL + """ + ) + op.execute( + """ + UPDATE event_targets SET entity_run_id = jobs.run_id + FROM jobs + WHERE jobs.id = event_targets.entity_id + AND event_targets.entity_type = 'job' + AND event_targets.entity_run_id IS NULL + """ + ) + + +def downgrade() -> None: + pass diff --git a/src/dstack/_internal/server/models.py b/src/dstack/_internal/server/models.py index fc73010263..d5452f0590 100644 --- a/src/dstack/_internal/server/models.py +++ b/src/dstack/_internal/server/models.py @@ -1183,6 +1183,11 @@ class EventTargetModel(BaseModel): ) entity_project: Mapped[Optional["ProjectModel"]] = relationship() + entity_run_id: Mapped[Optional[uuid.UUID]] = mapped_column( + ForeignKey("runs.id", ondelete="CASCADE"), nullable=True, index=True + ) + entity_run: Mapped[Optional["RunModel"]] = relationship() + entity_type: Mapped[EventTargetType] = mapped_column( EnumAsString(EventTargetType, 100), index=True ) diff --git a/src/dstack/_internal/server/services/events.py b/src/dstack/_internal/server/services/events.py index dd7b33dc7f..5eaf73b064 100644 --- a/src/dstack/_internal/server/services/events.py +++ b/src/dstack/_internal/server/services/events.py @@ -76,6 +76,7 @@ class Target: project_id: Optional[uuid.UUID] id: uuid.UUID name: str + run_id: Optional[uuid.UUID] = None def __post_init__(self): if self.type == EventTargetType.USER and self.project_id is not None: @@ -84,6 +85,10 @@ def __post_init__(self): raise ValueError(f"{self.type} target must have project_id") if self.type == EventTargetType.PROJECT and self.id != self.project_id: raise ValueError("Project target id must be equal to project_id") + if self.type in [EventTargetType.RUN, EventTargetType.JOB] and self.run_id is None: + raise ValueError(f"{self.type} target must have run_id") + if self.type == EventTargetType.RUN and self.id != self.run_id: + raise ValueError("Run target id must be equal to run_id") @staticmethod def from_model( @@ -126,6 +131,7 @@ def from_model( project_id=model.project_id or model.project.id, id=model.id, name=model.job_name, + run_id=model.run_id, ) if isinstance(model, ProjectModel): return Target( @@ -140,6 +146,7 @@ def from_model( project_id=model.project_id or model.project.id, id=model.id, name=model.run_name, + run_id=model.id, ) if isinstance(model, SecretModel): return Target( @@ -225,6 +232,7 @@ def emit(session: AsyncSession, message: str, actor: AnyActor, targets: list[Tar entity_project_id=target.project_id, entity_id=target.id, entity_name=target.name, + entity_run_id=target.run_id, ) ) session.add(event) @@ -354,23 +362,7 @@ async def list_events( ) ) if within_runs is not None: - query = select(JobModel.id).where(JobModel.run_id.in_(within_runs)) - res = await session.execute(query) - # In Postgres, fetching job IDs separately is orders of magnitude faster - # than using a subquery. - job_ids = list(res.unique().scalars().all()) - target_filters.append( - or_( - and_( - EventTargetModel.entity_type == EventTargetType.RUN, - EventTargetModel.entity_id.in_(within_runs), - ), - and_( - EventTargetModel.entity_type == EventTargetType.JOB, - EventTargetModel.entity_id.in_(job_ids), - ), - ) - ) + target_filters.append(EventTargetModel.entity_run_id.in_(within_runs)) if include_target_types is not None: target_filters.append(EventTargetModel.entity_type.in_(include_target_types)) diff --git a/src/dstack/_internal/server/services/runs/__init__.py b/src/dstack/_internal/server/services/runs/__init__.py index 39d4815d08..4d9d523845 100644 --- a/src/dstack/_internal/server/services/runs/__init__.py +++ b/src/dstack/_internal/server/services/runs/__init__.py @@ -833,6 +833,7 @@ def create_job_model_for_new_submission( run_model: RunModel, job: Job, status: JobStatus, + submission_num: int = 0, ) -> JobModel: """ Create a new job. @@ -849,7 +850,7 @@ def create_job_model_for_new_submission( job_name=f"{job.job_spec.job_name}", replica_num=job.job_spec.replica_num, deployment_num=run_model.deployment_num, - submission_num=len(job.job_submissions), + submission_num=submission_num, submitted_at=now, last_processed_at=now, status=status, diff --git a/src/tests/_internal/server/background/pipeline_tasks/test_runs/test_active.py b/src/tests/_internal/server/background/pipeline_tasks/test_runs/test_active.py index 551a957768..8c91fce0ba 100644 --- a/src/tests/_internal/server/background/pipeline_tasks/test_runs/test_active.py +++ b/src/tests/_internal/server/background/pipeline_tasks/test_runs/test_active.py @@ -296,6 +296,86 @@ async def test_retries_no_capacity_replica_and_keeps_service_running( assert retried_job.status == JobStatus.SUBMITTED assert len(jobs) == 3 + async def test_replica_retry_deletes_superseded_no_capacity_submissions( + self, test_db, session: AsyncSession, worker: RunWorker + ) -> None: + project = await create_project(session=session) + user = await create_user(session=session) + repo = await create_repo(session=session, project_id=project.id) + run_spec = get_run_spec( + repo_id=repo.name, + profile=Profile( + name="default", + retry=ProfileRetry(duration=3600, on_events=[RetryEvent.INTERRUPTION]), + ), + configuration=ServiceConfiguration( + port=8080, + commands=["echo Hi!"], + replicas=Range[int](min=2, max=2), + ), + ) + run = await create_run( + session=session, + project=project, + repo=repo, + user=user, + run_spec=run_spec, + status=RunStatus.RUNNING, + ) + await create_job( + session=session, + run=run, + status=JobStatus.RUNNING, + submitted_at=run.submitted_at, + last_processed_at=run.last_processed_at, + replica_num=1, + job_provisioning_data=get_job_provisioning_data(), + ) + superseded_job = await create_job( + session=session, + run=run, + status=JobStatus.FAILED, + termination_reason=JobTerminationReason.FAILED_TO_START_DUE_TO_NO_CAPACITY, + submitted_at=run.submitted_at, + last_processed_at=run.last_processed_at, + replica_num=0, + submission_num=0, + ) + interrupted_job = await create_job( + session=session, + run=run, + status=JobStatus.TERMINATING, + termination_reason=JobTerminationReason.INTERRUPTED_BY_NO_CAPACITY, + submitted_at=run.submitted_at, + last_processed_at=run.last_processed_at, + replica_num=0, + submission_num=1, + job_provisioning_data=get_job_provisioning_data(), + ) + lock_run(run) + await session.commit() + + now = run.submitted_at + timedelta(minutes=3) + with patch( + "dstack._internal.server.background.pipeline_tasks.runs.active.get_current_datetime", + return_value=now, + ): + await worker.process(run_to_pipeline_item(run)) + + jobs = list( + ( + await session.execute( + select(JobModel) + .where(JobModel.run_id == run.id, JobModel.replica_num == 0) + .order_by(JobModel.submission_num) + ) + ).scalars() + ) + assert [j.submission_num for j in jobs] == [1, 2] + assert superseded_job.id not in [j.id for j in jobs] + assert interrupted_job.id in [j.id for j in jobs] + assert jobs[-1].status == JobStatus.SUBMITTED + async def test_retries_scheduled_run_no_capacity_from_trigger_time( self, test_db, session: AsyncSession, worker: RunWorker ) -> None: diff --git a/src/tests/_internal/server/background/pipeline_tasks/test_runs/test_pending.py b/src/tests/_internal/server/background/pipeline_tasks/test_runs/test_pending.py index 6cfe5fe00e..ca772ca24e 100644 --- a/src/tests/_internal/server/background/pipeline_tasks/test_runs/test_pending.py +++ b/src/tests/_internal/server/background/pipeline_tasks/test_runs/test_pending.py @@ -11,6 +11,7 @@ from dstack._internal.core.models.resources import Range from dstack._internal.core.models.runs import ( JobStatus, + JobTerminationReason, RunStatus, ) from dstack._internal.server.background.pipeline_tasks.runs import RunWorker @@ -21,6 +22,7 @@ create_repo, create_run, create_user, + get_job_provisioning_data, get_run_spec, ) from dstack._internal.utils.common import get_current_datetime @@ -147,6 +149,77 @@ async def test_resubmits_retrying_run_after_delay( assert new_job.replica_num == old_job.replica_num assert new_job.submission_num == old_job.submission_num + 1 + async def test_resubmission_deletes_superseded_no_capacity_submissions( + self, test_db, session: AsyncSession, worker: RunWorker + ) -> None: + """ + On resubmission, older submissions that failed to start due to no capacity + without provisioning are deleted. The direct predecessor, provisioned + submissions, and submissions failed for other reasons are kept. + """ + project = await create_project(session=session) + user = await create_user(session=session) + repo = await create_repo(session=session, project_id=project.id) + run = await create_run( + session=session, + project=project, + repo=repo, + user=user, + status=RunStatus.PENDING, + resubmission_attempt=4, + ) + old_time = get_current_datetime() - timedelta(minutes=10) + superseded_job = await create_job( + session=session, + run=run, + status=JobStatus.FAILED, + last_processed_at=old_time, + submission_num=0, + termination_reason=JobTerminationReason.FAILED_TO_START_DUE_TO_NO_CAPACITY, + ) + provisioned_job = await create_job( + session=session, + run=run, + status=JobStatus.FAILED, + last_processed_at=old_time, + submission_num=1, + termination_reason=JobTerminationReason.FAILED_TO_START_DUE_TO_NO_CAPACITY, + job_provisioning_data=get_job_provisioning_data(), + ) + other_reason_job = await create_job( + session=session, + run=run, + status=JobStatus.FAILED, + last_processed_at=old_time, + submission_num=2, + termination_reason=JobTerminationReason.TERMINATED_BY_USER, + ) + predecessor_job = await create_job( + session=session, + run=run, + status=JobStatus.FAILED, + last_processed_at=old_time, + submission_num=3, + termination_reason=JobTerminationReason.FAILED_TO_START_DUE_TO_NO_CAPACITY, + ) + lock_run(run) + await session.commit() + + await worker.process(run_to_pipeline_item(run)) + + await session.refresh(run) + assert run.status == RunStatus.SUBMITTED + res = await session.execute( + select(JobModel).where(JobModel.run_id == run.id).order_by(JobModel.submission_num) + ) + jobs = list(res.scalars().all()) + assert [j.submission_num for j in jobs] == [1, 2, 3, 4] + assert superseded_job.id not in [j.id for j in jobs] + assert provisioned_job.id in [j.id for j in jobs] + assert other_reason_job.id in [j.id for j in jobs] + assert predecessor_job.id in [j.id for j in jobs] + assert jobs[-1].status == JobStatus.SUBMITTED + async def test_noops_when_run_lock_changes_after_processing( self, test_db, session: AsyncSession, worker: RunWorker ) -> None: diff --git a/src/tests/_internal/server/background/pipeline_tasks/test_terminating_jobs.py b/src/tests/_internal/server/background/pipeline_tasks/test_terminating_jobs.py index 0ecb64cd24..f6f93613dd 100644 --- a/src/tests/_internal/server/background/pipeline_tasks/test_terminating_jobs.py +++ b/src/tests/_internal/server/background/pipeline_tasks/test_terminating_jobs.py @@ -844,10 +844,14 @@ async def test_terminates_job_with_placeholder_instance( await worker.process(_job_to_pipeline_item(job)) await session.refresh(job) - await session.refresh(placeholder) assert job.status == JobStatus.TERMINATED assert job.instance_id is None - assert placeholder.status == InstanceStatus.TERMINATING + assert job.used_instance_id == placeholder.id + # The placeholder never provisioned anything, so it is deleted, not terminated + res = await session.execute( + select(InstanceModel.id).where(InstanceModel.id == placeholder.id) + ) + assert res.scalar_one_or_none() is None async def test_retries_detaching_when_used_instance_is_missing( self, test_db, session: AsyncSession, worker: JobTerminatingWorker diff --git a/src/tests/_internal/server/routers/test_events.py b/src/tests/_internal/server/routers/test_events.py index cb8e44b85a..7c47cc724b 100644 --- a/src/tests/_internal/server/routers/test_events.py +++ b/src/tests/_internal/server/routers/test_events.py @@ -6,9 +6,11 @@ import pytest_asyncio from freezegun import freeze_time from httpx import AsyncClient +from sqlalchemy import delete from sqlalchemy.ext.asyncio import AsyncSession from dstack._internal.core.models.users import GlobalRole, ProjectRole +from dstack._internal.server.models import JobModel from dstack._internal.server.services import events from dstack._internal.server.services.projects import add_project_member from dstack._internal.server.testing.common import ( @@ -929,6 +931,33 @@ async def test_within_runs(self, session: AsyncSession, client: AsyncClient) -> resp.raise_for_status() assert len(resp.json()) == 3 + async def test_within_runs_finds_events_of_deleted_jobs( + self, session: AsyncSession, client: AsyncClient + ) -> None: + user = await create_user(session=session) + project = await create_project(session=session, owner=user) + repo = await create_repo(session=session, project_id=project.id) + run = await create_run(session=session, project=project, repo=repo, user=user) + job = await create_job(session=session, run=run) + events.emit( + session, + "Job created on new submission", + actor=events.SystemActor(), + targets=[events.Target.from_model(job)], + ) + await session.commit() + # Superseded no-capacity submissions are deleted on resubmission + await session.execute(delete(JobModel).where(JobModel.id == job.id)) + await session.commit() + + resp = await client.post( + "/api/events/list", + headers=get_auth_headers(user.token), + json={"within_runs": [str(run.id)]}, + ) + resp.raise_for_status() + assert len(resp.json()) == 1 + async def test_include_target_types(self, session: AsyncSession, client: AsyncClient) -> None: user = await create_user(session=session) project = await create_project(session=session, owner=user)