diff --git a/airflow-ctl/src/airflowctl/api/operations.py b/airflow-ctl/src/airflowctl/api/operations.py index 5684dfac8403f..a0c8a904e03ff 100644 --- a/airflow-ctl/src/airflowctl/api/operations.py +++ b/airflow-ctl/src/airflowctl/api/operations.py @@ -345,7 +345,7 @@ def delete_queued_events(self, asset_id: str) -> str | ServerResponseError: def delete_dag_queued_events(self, dag_id: str, before: str) -> str | ServerResponseError: """Delete a queued event for a dag.""" try: - self.client.delete(f"assets/dags/{dag_id}/queuedEvents", params={"before": before}) + self.client.delete(f"dags/{dag_id}/assets/queuedEvents", params={"before": before}) return dag_id except ServerResponseError as e: raise e @@ -353,7 +353,7 @@ def delete_dag_queued_events(self, dag_id: str, before: str) -> str | ServerResp def delete_queued_event(self, dag_id: str, asset_id: str) -> str | ServerResponseError: """Delete a queued event for a dag.""" try: - self.client.delete(f"assets/dags/{dag_id}/assets/{asset_id}/queuedEvents/") + self.client.delete(f"dags/{dag_id}/assets/{asset_id}/queuedEvents") return asset_id except ServerResponseError as e: raise e diff --git a/airflow-ctl/tests/airflow_ctl/api/test_operations.py b/airflow-ctl/tests/airflow_ctl/api/test_operations.py index ec6695077c8b2..fda8718545a91 100644 --- a/airflow-ctl/tests/airflow_ctl/api/test_operations.py +++ b/airflow-ctl/tests/airflow_ctl/api/test_operations.py @@ -449,6 +449,24 @@ def handle_request(request: httpx.Request) -> httpx.Response: response = client.assets.get_dag_queued_event(dag_id=self.dag_id, asset_id=self.asset_id) assert response == self.asset_queued_event_response + def test_delete_dag_queued_events(self): + def handle_request(request: httpx.Request) -> httpx.Response: + assert request.url.path == f"/api/v2/dags/{self.dag_id}/assets/queuedEvents" + return httpx.Response(204) + + client = make_api_client(transport=httpx.MockTransport(handle_request)) + response = client.assets.delete_dag_queued_events(dag_id=self.dag_id, before=self.before) + assert response == self.dag_id + + def test_delete_queued_event(self): + def handle_request(request: httpx.Request) -> httpx.Response: + assert request.url.path == f"/api/v2/dags/{self.dag_id}/assets/{self.asset_id}/queuedEvents" + return httpx.Response(204) + + client = make_api_client(transport=httpx.MockTransport(handle_request)) + response = client.assets.delete_queued_event(dag_id=self.dag_id, asset_id=self.asset_id) + assert response == self.asset_id + class TestBackfillOperations: backfill_id: NonNegativeInt = 1