henry3260 commented on code in PR #70376:
URL: https://github.com/apache/airflow/pull/70376#discussion_r3645395959
##########
airflow-ctl/src/airflowctl/api/operations.py:
##########
@@ -289,106 +280,73 @@ def create_event(
self, asset_event_body: CreateAssetEventsBody
) -> AssetEventResponse | ServerResponseError:
"""Create an asset event."""
- try:
- # Ensure extra is initialised before sent to API
- if asset_event_body.extra is None:
- asset_event_body.extra = {}
- self.response = self.client.post(
- "assets/events", json=asset_event_body.model_dump(mode="json",
exclude_none=True)
- )
- return
AssetEventResponse.model_validate_json(self.response.content)
- except ServerResponseError as e:
- raise e
+ # Ensure extra is initialised before sent to API
+ if asset_event_body.extra is None:
+ asset_event_body.extra = {}
+ self.response = self.client.post(
+ "assets/events", json=asset_event_body.model_dump(mode="json",
exclude_none=True)
+ )
+ return AssetEventResponse.model_validate_json(self.response.content)
def materialize(self, asset_id: str) -> DAGRunResponse |
ServerResponseError:
"""Materialize an asset."""
- try:
- self.response = self.client.post(f"assets/{asset_id}/materialize")
- return DAGRunResponse.model_validate_json(self.response.content)
- except ServerResponseError as e:
- raise e
+ self.response = self.client.post(f"assets/{asset_id}/materialize")
+ return DAGRunResponse.model_validate_json(self.response.content)
def get_queued_events(self, asset_id: str) ->
QueuedEventCollectionResponse | ServerResponseError:
"""Get queued events for an asset."""
- try:
- self.response = self.client.get(f"assets/{asset_id}/queuedEvents")
- return
QueuedEventCollectionResponse.model_validate_json(self.response.content)
- except ServerResponseError as e:
- raise e
+ self.response = self.client.get(f"assets/{asset_id}/queuedEvents")
+ return
QueuedEventCollectionResponse.model_validate_json(self.response.content)
def get_dag_queued_events(
self, dag_id: str, before: str
) -> QueuedEventCollectionResponse | ServerResponseError:
"""Get queued events for a dag."""
- try:
- self.response =
self.client.get(f"dags/{dag_id}/assets/queuedEvents", params={"before": before})
- return
QueuedEventCollectionResponse.model_validate_json(self.response.content)
- except ServerResponseError as e:
- raise e
+ self.response = self.client.get(f"dags/{dag_id}/assets/queuedEvents",
params={"before": before})
+ return
QueuedEventCollectionResponse.model_validate_json(self.response.content)
def get_dag_queued_event(self, dag_id: str, asset_id: str) ->
QueuedEventResponse | ServerResponseError:
"""Get a queued event for a dag."""
- try:
- self.response =
self.client.get(f"dags/{dag_id}/assets/{asset_id}/queuedEvents")
- return
QueuedEventResponse.model_validate_json(self.response.content)
- except ServerResponseError as e:
- raise e
+ self.response =
self.client.get(f"dags/{dag_id}/assets/{asset_id}/queuedEvents")
+ return QueuedEventResponse.model_validate_json(self.response.content)
def delete_queued_events(self, asset_id: str) -> str | ServerResponseError:
"""Delete a queued event for an asset."""
- try:
- self.client.delete(f"assets/{asset_id}/queuedEvents/")
- return asset_id
- except ServerResponseError as e:
- raise e
+ 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."""
- try:
- self.client.delete(f"assets/dags/{dag_id}/queuedEvents",
params={"before": before})
- return dag_id
- except ServerResponseError as e:
- raise e
+ self.client.delete(f"assets/dags/{dag_id}/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."""
Review Comment:
```suggestion
"""Delete a queued event for a Dag."""
```
##########
airflow-ctl/src/airflowctl/api/operations.py:
##########
@@ -289,106 +280,73 @@ def create_event(
self, asset_event_body: CreateAssetEventsBody
) -> AssetEventResponse | ServerResponseError:
"""Create an asset event."""
- try:
- # Ensure extra is initialised before sent to API
- if asset_event_body.extra is None:
- asset_event_body.extra = {}
- self.response = self.client.post(
- "assets/events", json=asset_event_body.model_dump(mode="json",
exclude_none=True)
- )
- return
AssetEventResponse.model_validate_json(self.response.content)
- except ServerResponseError as e:
- raise e
+ # Ensure extra is initialised before sent to API
+ if asset_event_body.extra is None:
+ asset_event_body.extra = {}
+ self.response = self.client.post(
+ "assets/events", json=asset_event_body.model_dump(mode="json",
exclude_none=True)
+ )
+ return AssetEventResponse.model_validate_json(self.response.content)
def materialize(self, asset_id: str) -> DAGRunResponse |
ServerResponseError:
"""Materialize an asset."""
- try:
- self.response = self.client.post(f"assets/{asset_id}/materialize")
- return DAGRunResponse.model_validate_json(self.response.content)
- except ServerResponseError as e:
- raise e
+ self.response = self.client.post(f"assets/{asset_id}/materialize")
+ return DAGRunResponse.model_validate_json(self.response.content)
def get_queued_events(self, asset_id: str) ->
QueuedEventCollectionResponse | ServerResponseError:
"""Get queued events for an asset."""
- try:
- self.response = self.client.get(f"assets/{asset_id}/queuedEvents")
- return
QueuedEventCollectionResponse.model_validate_json(self.response.content)
- except ServerResponseError as e:
- raise e
+ self.response = self.client.get(f"assets/{asset_id}/queuedEvents")
+ return
QueuedEventCollectionResponse.model_validate_json(self.response.content)
def get_dag_queued_events(
self, dag_id: str, before: str
) -> QueuedEventCollectionResponse | ServerResponseError:
"""Get queued events for a dag."""
- try:
- self.response =
self.client.get(f"dags/{dag_id}/assets/queuedEvents", params={"before": before})
- return
QueuedEventCollectionResponse.model_validate_json(self.response.content)
- except ServerResponseError as e:
- raise e
+ self.response = self.client.get(f"dags/{dag_id}/assets/queuedEvents",
params={"before": before})
+ return
QueuedEventCollectionResponse.model_validate_json(self.response.content)
def get_dag_queued_event(self, dag_id: str, asset_id: str) ->
QueuedEventResponse | ServerResponseError:
"""Get a queued event for a dag."""
- try:
- self.response =
self.client.get(f"dags/{dag_id}/assets/{asset_id}/queuedEvents")
- return
QueuedEventResponse.model_validate_json(self.response.content)
- except ServerResponseError as e:
- raise e
+ self.response =
self.client.get(f"dags/{dag_id}/assets/{asset_id}/queuedEvents")
+ return QueuedEventResponse.model_validate_json(self.response.content)
def delete_queued_events(self, asset_id: str) -> str | ServerResponseError:
"""Delete a queued event for an asset."""
- try:
- self.client.delete(f"assets/{asset_id}/queuedEvents/")
- return asset_id
- except ServerResponseError as e:
- raise e
+ 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."""
Review Comment:
```suggestion
"""Delete a queued event for a Dag."""
```
##########
airflow-ctl/src/airflowctl/api/operations.py:
##########
@@ -460,86 +400,62 @@ def create(
connection: ConnectionBody,
) -> ConnectionResponse | ServerResponseError:
"""Create a connection."""
- try:
- self.response = self.client.post(
- "connections", json=connection.model_dump(mode="json",
by_alias=True, exclude_none=True)
- )
- return
ConnectionResponse.model_validate_json(self.response.content)
- except ServerResponseError as e:
- raise e
+ self.response = self.client.post(
+ "connections", json=connection.model_dump(mode="json",
by_alias=True, exclude_none=True)
+ )
+ return ConnectionResponse.model_validate_json(self.response.content)
def bulk(self, connections: BulkBodyConnectionBody) -> BulkResponse |
ServerResponseError:
"""CRUD multiple connections."""
- try:
- self.response = self.client.patch(
- "connections", json=connections.model_dump(mode="json",
by_alias=True)
- )
- return BulkResponse.model_validate_json(self.response.content)
- except ServerResponseError as e:
- raise e
+ self.response = self.client.patch(
+ "connections", json=connections.model_dump(mode="json",
by_alias=True)
+ )
+ return BulkResponse.model_validate_json(self.response.content)
def create_defaults(self) -> None | ServerResponseError:
"""Create default connections."""
- try:
- self.response = self.client.post("connections/defaults")
- return None
- except ServerResponseError as e:
- raise e
+ self.response = self.client.post("connections/defaults")
+ return None
def delete(self, conn_id: str) -> str | ServerResponseError:
"""Delete a connection."""
- try:
- self.client.delete(f"connections/{conn_id}")
- return conn_id
- except ServerResponseError as e:
- raise e
+ self.client.delete(f"connections/{conn_id}")
+ return conn_id
def update(
self,
connection: ConnectionBody,
) -> ConnectionResponse | ServerResponseError:
"""Update a connection."""
- try:
- self.response = self.client.patch(
- f"connections/{connection.connection_id}",
- json=connection.model_dump(mode="json", by_alias=True),
- )
- return
ConnectionResponse.model_validate_json(self.response.content)
- except ServerResponseError as e:
- raise e
+ self.response = self.client.patch(
+ f"connections/{connection.connection_id}",
+ json=connection.model_dump(mode="json", by_alias=True),
+ )
+ return ConnectionResponse.model_validate_json(self.response.content)
def test(
self,
connection: ConnectionBody,
) -> ConnectionTestResponse | ServerResponseError:
"""Test a connection."""
- try:
- self.response = self.client.post(
- "connections/test", json=connection.model_dump(mode="json",
by_alias=True)
- )
- return
ConnectionTestResponse.model_validate_json(self.response.content)
- except ServerResponseError as e:
- raise e
+ self.response = self.client.post(
+ "connections/test", json=connection.model_dump(mode="json",
by_alias=True)
+ )
+ return
ConnectionTestResponse.model_validate_json(self.response.content)
class DagsOperations(BaseOperations):
"""Dags operations."""
def get(self, dag_id: str) -> DAGResponse | ServerResponseError:
"""Get a Dag."""
- try:
- self.response = self.client.get(f"dags/{dag_id}")
- return DAGResponse.model_validate_json(self.response.content)
- except ServerResponseError as e:
- raise e
+ self.response = self.client.get(f"dags/{dag_id}")
+ return DAGResponse.model_validate_json(self.response.content)
def get_details(self, dag_id: str) -> DAGDetailsResponse |
ServerResponseError:
- """Get a DAG details."""
- try:
- self.response = self.client.get(f"dags/{dag_id}/details")
- return
DAGDetailsResponse.model_validate_json(self.response.content)
- except ServerResponseError as e:
- raise e
+ """Get a Dag details."""
+ self.response = self.client.get(f"dags/{dag_id}/details")
+ return DAGDetailsResponse.model_validate_json(self.response.content)
def get_tags(self) -> DAGTagCollectionResponse | ServerResponseError:
"""Get all DAG tags."""
Review Comment:
```suggestion
"""Get all Dag tags."""
```
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]