From 644565e132157359ce01ad5f83a72f457d9b80b5 Mon Sep 17 00:00:00 2001 From: justinpakzad <114518232+justinpakzad@users.noreply.github.com> Date: Fri, 24 Jul 2026 17:52:19 -0400 Subject: [PATCH] Fix wrong API paths for queued asset events in airflowctl --- airflow-ctl/src/airflowctl/api/operations.py | 4 ++-- .../tests/airflow_ctl/api/test_operations.py | 18 ++++++++++++++++++ 2 files changed, 20 insertions(+), 2 deletions(-) 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