Loading device 0 → 100644 +220 −0 File added.Preview size limit exceeded, changes collapsed. Show changes src/device/service/drivers/OpenFlow/OpenFlowDriver.py +44 −2 Original line number Diff line number Diff line Loading @@ -13,6 +13,7 @@ # limitations under the License. import json import logging, requests, threading import resource from requests.auth import HTTPBasicAuth from typing import Any, Iterator, List, Optional, Tuple, Union from common.method_wrappers.Decorator import MetricsPool, metered_subclass_method Loading Loading @@ -149,12 +150,53 @@ class OpenFlowDriver(_Driver): # @metered_subclass_method(METRICS_POOL) def SetConfig(self, resources: List[Tuple[str, Any]]) -> List[Union[bool, Exception]]: url = f"{self.__base_url}/stats/flowentry/add" results = [] LOGGER.info(f'SetConfig_resources:{resources}') LOGGER.info(f"SetConfig_resources: {resources}") if not resources: return results with self.__lock: for resource in resources: try: resource_key, resource_value = resource resource_value_dict = json.loads(resource_value) dpid = int(resource_value_dict["dpid"], 16) in_port = int(resource_value_dict["in-port"].split("-")[1][3:]) out_port = int(resource_value_dict["out-port"].split("-")[1][3:]) flow_entry = { "dpid": dpid, "cookie": 0, "priority": 32768, "match": { "in_port": in_port }, "actions": [ { "type": "OUTPUT", "port": out_port } ] } try: response = requests.post(url,json=flow_entry,timeout=self.__timeout,verify=False,auth=self.__auth) response.raise_for_status() results.append(True) LOGGER.info(f"Successfully posted flow entry: {flow_entry}") except requests.exceptions.Timeout: LOGGER.error(f"Timeout connecting to {url}") results.append(Exception(f"Timeout connecting to {url}")) except requests.exceptions.RequestException as e: LOGGER.error(f"Error posting to {url}: {e}") results.append(e) else: LOGGER.warning(f"Skipped invalid resource_key: {resource_key}") results.append(Exception(f"Invalid resource_key: {resource_key}")) except Exception as e: LOGGER.error(f"Error processing resource {resource_key}: {e}", exc_info=True) results.append(e) return results # # # Loading src/service/service/service_handlers/l3nm_ryu/L3NMryuServiceHandler.py +22 −14 Original line number Diff line number Diff line Loading @@ -88,6 +88,7 @@ class RYUServiceHandler(_ServiceHandler): src_device, src_endpoint, = self._get_endpoint_details(endpoints[0]) dst_device, dst_endpoint, = self._get_endpoint_details(endpoints[-1]) src_controller = self.__task_executor.get_device_controller(src_device) del src_controller.device_config.config_rules[:] for index in range(len(endpoints) - 1): current_device, current_endpoint = self._get_endpoint_details(endpoints[index]) Loading @@ -100,30 +101,37 @@ class RYUServiceHandler(_ServiceHandler): flow_split = service_name.split('-') flow_rule_forward = f"{flow_split[0]}-{flow_split[2]}" flow_rule_reverse = f"{flow_split[2]}-{flow_split[0]}" forward_resource_value = json.dumps({"dpid": current_device.name, "in-port": in_port_forward, "out-port": out_port_forward}) forward_rule = ConfigRule( custom=ConfigRule_Custom( forward_resource_value = ({"dpid": current_device.name, "in-port": in_port_forward, "out-port": out_port_forward}) forward_rule = json_config_rule_set ( resource_key=f"/device[{current_endpoint.name.split('-')[0]}]/flow[{flow_rule_forward}]", resource_value=forward_resource_value ) ) LOGGER.debug(f"Forward configuration rule: {forward_rule}") src_controller.device_config.config_rules.append(forward_rule) src_controller.device_config.config_rules.append(ConfigRule(**forward_rule)) in_port_reverse = next_endpoint.name out_port_reverse = current_endpoint.name reverse_resource_value = json.dumps({"dpid": current_device.name, "in-port": in_port_reverse, "out-port": out_port_reverse}) reverse_rule = ConfigRule( custom=ConfigRule_Custom( reverse_resource_value = ({"dpid": current_device.name, "in-port": in_port_reverse, "out-port": out_port_reverse}) reverse_rule = json_config_rule_set( resource_key=f"/device[{current_endpoint.name.split('-')[0]}]/flow[{flow_rule_reverse}]", resource_value=reverse_resource_value ) ) LOGGER.debug(f"Reverse configuration rule: {reverse_rule}") src_controller.device_config.config_rules.append(reverse_rule) src_controller.device_config.config_rules.append(ConfigRule(**reverse_rule)) self.__task_executor.configure_device(src_controller) results.append(True) def get_config_rules(controller): try: config_rules = controller.device_config.config_rules for rule in config_rules: if rule.HasField("custom"): resource_key = rule.custom.resource_key resource_value = rule.custom.resource_value LOGGER.debug(f"Resource key in config: {resource_key}, Resource value in config: {resource_value}") except Exception as e: print(f"Error accessing config rules: {e}") get_config_rules(src_controller) LOGGER.debug(f"Configuration rules: {src_controller.device_config.config_rules}") return results except Exception as e: Loading Loading
src/device/service/drivers/OpenFlow/OpenFlowDriver.py +44 −2 Original line number Diff line number Diff line Loading @@ -13,6 +13,7 @@ # limitations under the License. import json import logging, requests, threading import resource from requests.auth import HTTPBasicAuth from typing import Any, Iterator, List, Optional, Tuple, Union from common.method_wrappers.Decorator import MetricsPool, metered_subclass_method Loading Loading @@ -149,12 +150,53 @@ class OpenFlowDriver(_Driver): # @metered_subclass_method(METRICS_POOL) def SetConfig(self, resources: List[Tuple[str, Any]]) -> List[Union[bool, Exception]]: url = f"{self.__base_url}/stats/flowentry/add" results = [] LOGGER.info(f'SetConfig_resources:{resources}') LOGGER.info(f"SetConfig_resources: {resources}") if not resources: return results with self.__lock: for resource in resources: try: resource_key, resource_value = resource resource_value_dict = json.loads(resource_value) dpid = int(resource_value_dict["dpid"], 16) in_port = int(resource_value_dict["in-port"].split("-")[1][3:]) out_port = int(resource_value_dict["out-port"].split("-")[1][3:]) flow_entry = { "dpid": dpid, "cookie": 0, "priority": 32768, "match": { "in_port": in_port }, "actions": [ { "type": "OUTPUT", "port": out_port } ] } try: response = requests.post(url,json=flow_entry,timeout=self.__timeout,verify=False,auth=self.__auth) response.raise_for_status() results.append(True) LOGGER.info(f"Successfully posted flow entry: {flow_entry}") except requests.exceptions.Timeout: LOGGER.error(f"Timeout connecting to {url}") results.append(Exception(f"Timeout connecting to {url}")) except requests.exceptions.RequestException as e: LOGGER.error(f"Error posting to {url}: {e}") results.append(e) else: LOGGER.warning(f"Skipped invalid resource_key: {resource_key}") results.append(Exception(f"Invalid resource_key: {resource_key}")) except Exception as e: LOGGER.error(f"Error processing resource {resource_key}: {e}", exc_info=True) results.append(e) return results # # # Loading
src/service/service/service_handlers/l3nm_ryu/L3NMryuServiceHandler.py +22 −14 Original line number Diff line number Diff line Loading @@ -88,6 +88,7 @@ class RYUServiceHandler(_ServiceHandler): src_device, src_endpoint, = self._get_endpoint_details(endpoints[0]) dst_device, dst_endpoint, = self._get_endpoint_details(endpoints[-1]) src_controller = self.__task_executor.get_device_controller(src_device) del src_controller.device_config.config_rules[:] for index in range(len(endpoints) - 1): current_device, current_endpoint = self._get_endpoint_details(endpoints[index]) Loading @@ -100,30 +101,37 @@ class RYUServiceHandler(_ServiceHandler): flow_split = service_name.split('-') flow_rule_forward = f"{flow_split[0]}-{flow_split[2]}" flow_rule_reverse = f"{flow_split[2]}-{flow_split[0]}" forward_resource_value = json.dumps({"dpid": current_device.name, "in-port": in_port_forward, "out-port": out_port_forward}) forward_rule = ConfigRule( custom=ConfigRule_Custom( forward_resource_value = ({"dpid": current_device.name, "in-port": in_port_forward, "out-port": out_port_forward}) forward_rule = json_config_rule_set ( resource_key=f"/device[{current_endpoint.name.split('-')[0]}]/flow[{flow_rule_forward}]", resource_value=forward_resource_value ) ) LOGGER.debug(f"Forward configuration rule: {forward_rule}") src_controller.device_config.config_rules.append(forward_rule) src_controller.device_config.config_rules.append(ConfigRule(**forward_rule)) in_port_reverse = next_endpoint.name out_port_reverse = current_endpoint.name reverse_resource_value = json.dumps({"dpid": current_device.name, "in-port": in_port_reverse, "out-port": out_port_reverse}) reverse_rule = ConfigRule( custom=ConfigRule_Custom( reverse_resource_value = ({"dpid": current_device.name, "in-port": in_port_reverse, "out-port": out_port_reverse}) reverse_rule = json_config_rule_set( resource_key=f"/device[{current_endpoint.name.split('-')[0]}]/flow[{flow_rule_reverse}]", resource_value=reverse_resource_value ) ) LOGGER.debug(f"Reverse configuration rule: {reverse_rule}") src_controller.device_config.config_rules.append(reverse_rule) src_controller.device_config.config_rules.append(ConfigRule(**reverse_rule)) self.__task_executor.configure_device(src_controller) results.append(True) def get_config_rules(controller): try: config_rules = controller.device_config.config_rules for rule in config_rules: if rule.HasField("custom"): resource_key = rule.custom.resource_key resource_value = rule.custom.resource_value LOGGER.debug(f"Resource key in config: {resource_key}, Resource value in config: {resource_value}") except Exception as e: print(f"Error accessing config rules: {e}") get_config_rules(src_controller) LOGGER.debug(f"Configuration rules: {src_controller.device_config.config_rules}") return results except Exception as e: Loading