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]

Reply via email to