Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions airflow-ctl/src/airflowctl/api/operations.py
Original file line number Diff line number Diff line change
Expand Up @@ -345,15 +345,15 @@ 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

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
Expand Down
18 changes: 18 additions & 0 deletions airflow-ctl/tests/airflow_ctl/api/test_operations.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down