This is an automated email from the ASF dual-hosted git repository.

henry3260 pushed a commit to branch v3-3-test
in repository https://gitbox.apache.org/repos/asf/airflow.git


The following commit(s) were added to refs/heads/v3-3-test by this push:
     new ac61af57e56 Fix incorrect asset queued-event endpoint paths in 
airflowctl (#70461) (#70463)
ac61af57e56 is described below

commit ac61af57e56e4f2201092fd30ca1c734b94d53b4
Author: Henry Chen <[email protected]>
AuthorDate: Sun Jul 26 17:22:09 2026 +0800

    Fix incorrect asset queued-event endpoint paths in airflowctl (#70461) 
(#70463)
    
    (cherry picked from commit 15df3aaac07284aa4d4afecd242aa65cd9685761)
---
 airflow-ctl/src/airflowctl/api/operations.py       |  6 ++---
 .../tests/airflow_ctl/api/test_operations.py       | 31 ++++++++++++++++++++++
 2 files changed, 34 insertions(+), 3 deletions(-)

diff --git a/airflow-ctl/src/airflowctl/api/operations.py 
b/airflow-ctl/src/airflowctl/api/operations.py
index 3b840845f41..61366004676 100644
--- a/airflow-ctl/src/airflowctl/api/operations.py
+++ b/airflow-ctl/src/airflowctl/api/operations.py
@@ -300,17 +300,17 @@ class AssetsOperations(BaseOperations):
 
     def delete_queued_events(self, asset_id: str) -> str | ServerResponseError:
         """Delete a queued event for an asset."""
-        self.client.delete(f"assets/{asset_id}/queuedEvents/")
+        self.client.delete(f"assets/{asset_id}/queuedEvents")
         return asset_id
 
     def delete_dag_queued_events(self, dag_id: str, before: str) -> str | 
ServerResponseError:
         """Delete a queued event for a Dag."""
-        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
 
     def delete_queued_event(self, dag_id: str, asset_id: str) -> str | 
ServerResponseError:
         """Delete a queued event for a Dag."""
-        
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
 
 
diff --git a/airflow-ctl/tests/airflow_ctl/api/test_operations.py 
b/airflow-ctl/tests/airflow_ctl/api/test_operations.py
index 52faecee73e..cf61ed5dafa 100644
--- a/airflow-ctl/tests/airflow_ctl/api/test_operations.py
+++ b/airflow-ctl/tests/airflow_ctl/api/test_operations.py
@@ -431,6 +431,37 @@ class TestAssetsOperations:
         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_queued_events(self):
+        def handle_request(request: httpx.Request) -> httpx.Response:
+            assert request.method == "DELETE"
+            assert request.url.path == 
f"/api/v2/assets/{self.asset_id}/queuedEvents"
+            return httpx.Response(204)
+
+        client = make_api_client(transport=httpx.MockTransport(handle_request))
+        response = client.assets.delete_queued_events(asset_id=self.asset_id)
+        assert response == self.asset_id
+
+    def test_delete_dag_queued_events(self):
+        def handle_request(request: httpx.Request) -> httpx.Response:
+            assert request.method == "DELETE"
+            assert request.url.path == 
f"/api/v2/dags/{self.dag_id}/assets/queuedEvents"
+            assert request.url.params["before"] == self.before
+            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.method == "DELETE"
+            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

Reply via email to