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