diff --git a/README.md b/README.md index 0259151c..ea059c16 100644 --- a/README.md +++ b/README.md @@ -199,27 +199,30 @@ $ docker run --rm -d --name rabbitmq -p 5672:5672 -p 15672:15672 rabbitmq:latest $ docker run --rm -d --name mongo -p 27017:27017 -e MONGO_INITDB_ROOT_USERNAME=guest -e MONGO_INITDB_ROOT_PASSWORD=guest mongo:7.0.11 ``` -Some environment variables are expected to be set for the tests to -work as expected, so you may want to copy `env.template` to `.env` and -edit it according to your environment, and make sure the env vars are -present in your shell: +Load environment variables from `env.template` directly into your shell +(no need to copy it to `.env`), then override `MQ_HOST` and `MQ_PORT` +to point at the local RabbitMQ container started above: ```console -$ cp env.template .env -$ # and then edit .env to suit your environment -$ source .env +$ set -a; source env.template; set +a +$ export MQ_HOST=localhost MQ_PORT=5672 ``` -And now, activate a virtual environment, install the requirements, and -then run `pytest`: +Activate a virtual environment, install the requirements, and run `pytest`: -``` +```console $ python3 -m venv venv --upgrade-deps $ source ./venv/bin/activate $ pip3 install --editable .[test] $ pytest ``` +To run a specific test file, pass its path to `pytest`: + +```console +$ pytest sdx_controller/test/test_l2vpn_controller_patch.py +``` + diff --git a/pyproject.toml b/pyproject.toml index 8a3a4187..a3d81cc1 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "flit_core.buildapi" [project] name = "sdx-controller" -version = "3.2.0" +version = "3.2.2" description = "AtlanticWave-SDX project's main controller" authors = [ { name = "Yufeng Xin", email = "yxin@renci.org" }, @@ -29,7 +29,7 @@ dependencies = [ "pika >= 1.2.0", "dataset", "pymongo > 3.0", - "sdx-pce @ git+https://github.com/atlanticwave-sdx/pce@v3.2.1", + "sdx-pce @ git+https://github.com/atlanticwave-sdx/pce@v3.2.2", ] [project.optional-dependencies] diff --git a/sdx_controller/controllers/l2vpn_controller.py b/sdx_controller/controllers/l2vpn_controller.py index ee64ef9c..5f0a020c 100644 --- a/sdx_controller/controllers/l2vpn_controller.py +++ b/sdx_controller/controllers/l2vpn_controller.py @@ -1,5 +1,7 @@ +import copy import logging import os +import time import uuid import connexion @@ -90,7 +92,15 @@ def delete_connection(service_id): logger.info(f"Removing connection: {service_id} {connection.get('status')}") - connection_handler.remove_connection(current_app.te_manager, service_id, "API") + remove_reason, remove_code = connection_handler.remove_connection( + current_app.te_manager, service_id, "API" + ) + if remove_code // 100 != 2: + logger.info( + f"Delete failed (connection id: {service_id}): " + f"reason='{remove_reason}', code={remove_code}" + ) + # return remove_reason, remove_code db_instance.mark_deleted(MongoCollections.CONNECTIONS, f"{service_id}") db_instance.mark_deleted(MongoCollections.BREAKDOWNS, f"{service_id}") except Exception as e: @@ -188,14 +198,12 @@ def place_connection(body): body["id"] = service_id logger.info(f"Request has no ID. Generated ID: {service_id}") - body["status"] = str(ConnectionStateMachine.State.REQUESTED) + conn_status = ConnectionStateMachine.State.REQUESTED + body["status"] = str(conn_status) # used in lc_message_handler to count the oxp success response body["oxp_success_count"] = 0 - conn_status = ConnectionStateMachine.State.UNDER_PROVISIONING - body, _ = connection_state_machine(body, conn_status) - db_instance.add_key_value_pair_to_db(MongoCollections.CONNECTIONS, service_id, body) logger.info( @@ -203,8 +211,19 @@ def place_connection(body): ) reason, code = connection_handler.place_connection(current_app.te_manager, body) - if code // 100 != 2: + if code // 100 == 2: + # conn_status = ConnectionStateMachine.State.UNDER_PROVISIONING + # body, _ = connection_state_machine(body, conn_status) + # db_instance.update_field_in_json( + # MongoCollections.CONNECTIONS, + # service_id, + # "status", + # str(conn_status), + # ) + logger.info(f"place_connection succeeds: ID: {service_id} body='{body}'") + else: conn_status = ConnectionStateMachine.State.REJECTED + body, _ = connection_state_machine(body, conn_status) db_instance.update_field_in_json( MongoCollections.CONNECTIONS, service_id, @@ -215,9 +234,16 @@ def place_connection(body): f"place_connection result: ID: {service_id} reason='{reason}', code={code}" ) + current_conn = db_instance.get_value_from_db( + MongoCollections.CONNECTIONS, f"{service_id}" + ) response = { "service_id": service_id, - "status": parse_conn_status(str(conn_status)), + "status": parse_conn_status( + current_conn.get("status", str(conn_status)) + if current_conn + else str(conn_status) + ), "reason": reason, } @@ -258,26 +284,62 @@ def patch_connection(service_id, body=None): # noqa: E501 logger.info(f"Gathered connexion JSON: {new_body}") - body.update(new_body) - - body, _ = connection_state_machine(body, ConnectionStateMachine.State.MODIFYING) + if "id" not in new_body: + new_body["id"] = service_id - body["oxp_success_count"] = 0 + # Validate the new request body before making any change to the existing connection. + # This is to avoid the case where we have already removed the original connection but the new request body is invalid, which will cause the connection to be deleted but not re-created. + # We can reuse the same validation function used in place_connection since the request body for patch_connection has the same schema as place_connection. + # + te_manager = current_app.te_manager # Assuming te_manager is accessible like this + try: + # Validate the new request body + te_manager.generate_traffic_matrix(connection_request=new_body) + except Exception as request_err: + logger.error("ERROR: invalid patch request: " + str(request_err)) + error_code = getattr(request_err, "request_code", None) + if not isinstance(error_code, int): + # Backward-compatible fallback for exception strings like "... (Code: 400)". + error_code = 400 + err_text = str(request_err) + if "Code:" in err_text: + candidate = err_text.split("Code:")[-1].replace(")", "").strip() + try: + error_code = int(candidate) + except (TypeError, ValueError): + logger.warning( + f"Could not parse error code from patch validation error: {err_text}" + ) + return f"Error: patch request is not valid: {request_err}", error_code + + logger.info("Modifying connection") + # Get roll back connection before removing connection + rollback_conn_body = copy.deepcopy(body) + body.update(new_body) - db_instance.add_key_value_pair_to_db(MongoCollections.CONNECTIONS, service_id, body) + conn_status = ConnectionStateMachine.State.MODIFYING + body, _ = connection_state_machine(body, conn_status) + db_instance.update_field_in_json( + MongoCollections.CONNECTIONS, + service_id, + "status", + str(conn_status), + ) try: logger.info("Removing connection") - # Get roll back connection before removing connection - rollback_conn_body = body remove_conn_reason, remove_conn_code = connection_handler.remove_connection( current_app.te_manager, service_id, "API" ) if remove_conn_code // 100 != 2: - body, _ = connection_state_machine(body, ConnectionStateMachine.State.DOWN) - db_instance.add_key_value_pair_to_db( - MongoCollections.CONNECTIONS, service_id, body + conn_status = ConnectionStateMachine.State.DOWN + body, _ = connection_state_machine(body, conn_status) + db_instance.update_field_in_json( + MongoCollections.CONNECTIONS, + service_id, + "status", + str(conn_status), ) response = { "service_id": service_id, @@ -289,20 +351,35 @@ def patch_connection(service_id, body=None): # noqa: E501 logger.info(f"Removed connection: {service_id}") except Exception as e: logger.info(f"Delete failed (connection id: {service_id}): {e}") + conn_status = ConnectionStateMachine.State.DOWN + body, _ = connection_state_machine(body, conn_status) + db_instance.update_field_in_json( + MongoCollections.CONNECTIONS, + service_id, + "status", + str(conn_status), + ) return f"Failed, reason: {e}", 500 - + time.sleep(10) logger.info( - f"Placing new connection {service_id} with te_manager: {current_app.te_manager}" - ) - - body, _ = connection_state_machine( - body, ConnectionStateMachine.State.UNDER_PROVISIONING + f"Modifying: Placing new connection {service_id} with te_manager: {current_app.te_manager}" ) + # Reset: remove_connection archives/deletes the original entry, + # so persist the patched request before re-placement. + conn_status = ConnectionStateMachine.State.REQUESTED + body["status"] = str(conn_status) + body["oxp_success_count"] = 0 + body["oxp_response"] = {} db_instance.add_key_value_pair_to_db(MongoCollections.CONNECTIONS, service_id, body) reason, code = connection_handler.place_connection(current_app.te_manager, body) if code // 100 == 2: # Service created successfully + # conn_status = ConnectionStateMachine.State.UNDER_PROVISIONING + # body, _ = connection_state_machine(body, conn_status) + # db_instance.add_key_value_pair_to_db( + # MongoCollections.CONNECTIONS, service_id, body + # ) code = 201 logger.info(f"Placed: ID: {service_id} reason='{reason}', code={code}") response = { @@ -311,47 +388,79 @@ def patch_connection(service_id, body=None): # noqa: E501 "reason": reason, } return response, code - else: - body, _ = connection_state_machine(body, ConnectionStateMachine.State.DOWN) logger.info( - f"Failed to place new connection. ID: {service_id} reason='{reason}', code={code}" + f"Modifying: Failed to place new connection. ID: {service_id} reason='{reason}', code={code}" ) logger.info("Rolling back to old connection.") - if not rollback_conn_body: - response = { - "service_id": service_id, - "status": parse_conn_status(body["status"]), - "reason": f"Failure, unable to rollback to last successful L2VPN: {reason}", - } - return response, code - # because above placement failed, so re-place the original connection request. + + rollback_conn_body["status"] = str(ConnectionStateMachine.State.REQUESTED) + # used in lc_message_handler to count the oxp success response + rollback_conn_body["oxp_success_count"] = 0 + rollback_conn_body["oxp_response"] = {} + conn_request = rollback_conn_body conn_request["id"] = service_id + db_instance.add_key_value_pair_to_db( + MongoCollections.CONNECTIONS, service_id, conn_request + ) + rollback_conn_reason = "Rollback attempt did not complete" try: rollback_conn_reason, rollback_conn_code = connection_handler.place_connection( current_app.te_manager, conn_request ) if rollback_conn_code // 100 == 2: - db_instance.add_key_value_pair_to_db( - MongoCollections.CONNECTIONS, service_id, conn_request + # conn_status = ConnectionStateMachine.State.UNDER_PROVISIONING + # rollback_conn_body, _ = connection_state_machine( + # rollback_conn_body, conn_status + # ) + # db_instance.update_field_in_json( + # MongoCollections.CONNECTIONS, + # service_id, + # "status", + # str(conn_status), + # ) + # still return 400 to indicate the patch request is not successful, since we have already rolled back to original connection, which is under provisioning state, so the connection is not down and not failed. + rollback_conn_code = code + else: + conn_status = ConnectionStateMachine.State.REJECTED + body, _ = connection_state_machine(body, conn_status) + db_instance.update_field_in_json( + MongoCollections.CONNECTIONS, + service_id, + "status", + str(conn_status), ) + rollback_conn_code = 500 logger.info( f"Roll back connection result: ID: {service_id} reason='{rollback_conn_reason}', code={rollback_conn_code}" ) except Exception as e: + conn_status = ConnectionStateMachine.State.REJECTED + db_instance.update_field_in_json( + MongoCollections.CONNECTIONS, + service_id, + "status", + str(conn_status), + ) logger.info(f"Rollback failed (connection id: {service_id}): {e}") - return f"Rollback failed, reason: {e}", 500 + rollback_conn_reason = f"Rollback failed: {e}" + rollback_conn_code = 500 + current_conn = db_instance.get_value_from_db( + MongoCollections.CONNECTIONS, f"{service_id}" + ) response = { "service_id": service_id, "reason": f"Failure, rolled back to last successful L2VPN: {reason}", - "status": parse_conn_status(conn_request["status"]), + "status": parse_conn_status( + current_conn.get("status", "") if current_conn else "" + ), } - return response, code + return response, rollback_conn_code def get_archived_connections_by_id(service_id): diff --git a/sdx_controller/handlers/connection_handler.py b/sdx_controller/handlers/connection_handler.py index 6fc1e18e..7c33921b 100644 --- a/sdx_controller/handlers/connection_handler.py +++ b/sdx_controller/handlers/connection_handler.py @@ -200,21 +200,29 @@ def _send_breakdown_to_lc(self, breakdown, operation, connection_request): } if operation == "delete": - oxp_response = connection_request.get("oxp_response") - - # evc_id is the service_id in the OXP response, it differs from the service_id in the connection. - evc_id = ( - oxp_response.get(domain_name, [None, {}])[1].get("service_id") - if oxp_response - else None + logger.debug( + f"Handling delete operation for connection {connection_request}" ) - - if not oxp_response or not evc_id: - return ( - "Connection does not have OXP response, cannot remove connection", - 404, + oxp_response = None + evc_id = None + try: + oxp_response = connection_request.get("oxp_response") + # evc_id is the service_id in the OXP response, it differs from the service_id in the connection. + evc_id = ( + oxp_response.get(domain_name, [None, {}])[1].get("service_id") + if oxp_response + else None + ) + if not oxp_response or not evc_id: + return ( + "Connection does not have OXP response, cannot remove connection", + 404, + ) + mq_link["evc_id"] = evc_id + except Exception as e: + logger.error( + f"Error occurred while processing OXP response in delete: {e}" ) - mq_link["evc_id"] = evc_id producer = TopicQueueProducer( timeout=5, exchange_name=exchange_name, routing_key=domain_name @@ -286,6 +294,17 @@ def place_connection( MongoCollections.BREAKDOWNS, connection_request["id"], breakdown ) self._process_port(connection_request["id"], ctx.ingress_port, "post") + conn_status = ConnectionStateMachine.State.UNDER_PROVISIONING + connection_request, _ = connection_state_machine( + connection_request, conn_status + ) + self.db_instance.update_field_in_json( + MongoCollections.CONNECTIONS, + connection_request["id"], + "status", + str(conn_status), + ) + status, code = self._send_breakdown_to_lc( breakdown, "post", connection_request ) @@ -343,6 +362,16 @@ def place_connection( operation="post", connection_request=connection_request, ) + conn_status = ConnectionStateMachine.State.UNDER_PROVISIONING + connection_request, _ = connection_state_machine( + connection_request, conn_status + ) + self.db_instance.update_field_in_json( + MongoCollections.CONNECTIONS, + connection_request["id"], + "status", + str(conn_status), + ) status, code = self._send_breakdown_to_lc( breakdown, "post", connection_request ) @@ -426,6 +455,12 @@ def remove_connection( status, code = self._send_breakdown_to_lc( breakdown, "delete", connection_request ) + if code // 100 != 2: + logger.error( + f"Could not publish delete breakdown for {service_id}: " + f"reason='{status}', code={code}" + ) + return status, code self._process_path_to_db( te_manager, operation="delete", connection_request=connection_request ) @@ -519,27 +554,59 @@ def handle_link_failure(self, te_manager, failed_links): logger.debug("Removed connection:") logger.debug(connection) + # time.sleep(10) + connection, _ = connection_state_machine( connection, ConnectionStateMachine.State.RECOVERING ) connection["oxp_success_count"] = 0 + connection["oxp_response"] = {} self.db_instance.add_key_value_pair_to_db( MongoCollections.CONNECTIONS, service_id, connection ) _reason, code = self.place_connection(te_manager, connection) - if code // 100 != 2: + + if code // 100 == 2: + # Service created successfully + # conn_status = ConnectionStateMachine.State.UNDER_PROVISIONING + # connection, _ = connection_state_machine( + # connection, conn_status + # ) + # self.db_instance.update_field_in_json( + # MongoCollections.CONNECTIONS, + # service_id, + # "status", + # str(conn_status), + # ) + logger.info( + f"link failure rerouting: place_connection succeeds: ID: {service_id} connection='{connection}'" + ) + code = 201 + else: + conn_status = ConnectionStateMachine.State.ERROR connection, _ = connection_state_machine( - connection, ConnectionStateMachine.State.ERROR + connection, conn_status ) - self.db_instance.add_key_value_pair_to_db( + self.db_instance.update_field_in_json( MongoCollections.CONNECTIONS, service_id, - connection, + "status", + str(conn_status), ) + _reason = ( + "place_connection failed during link failure rerouting" + ) + code = 400 logger.info( f"place_connection result: ID: {service_id} reason='{_reason}', code={code}" ) + # response = { + # "service_id": service_id, + # "status": parse_conn_status(connection["status"]), + # "reason": _reason, + # } + # return response, code def handle_uni_ports_up_to_down(self, uni_ports_up_to_down): """ @@ -572,9 +639,13 @@ def handle_uni_ports_up_to_down(self, uni_ports_up_to_down): logger.debug(f"Cannot find connection {service_id} in DB.") continue logger.info(f"Updating connection {service_id} status to 'down'.") - connection["status"] = "DOWN" - self.db_instance.add_key_value_pair_to_db( - MongoCollections.CONNECTIONS, service_id, connection + conn_status = "DOWN" + connection["status"] = conn_status + self.db_instance.update_field_in_json( + MongoCollections.CONNECTIONS, + service_id, + "status", + str(conn_status), ) logger.debug(f"Connection status updated for {service_id}") else: @@ -613,9 +684,13 @@ def handle_uni_ports_down_to_up(self, uni_ports_down_to_up): continue logger.info(f"Updating connection {service_id} status to 'up'.") - connection["status"] = "UP" - self.db_instance.add_key_value_pair_to_db( - MongoCollections.CONNECTIONS, service_id, connection + conn_status = "UP" + connection["status"] = conn_status + self.db_instance.update_field_in_json( + MongoCollections.CONNECTIONS, + service_id, + "status", + str(conn_status), ) logger.debug(f"Connection status updated for {service_id}") diff --git a/sdx_controller/handlers/lc_message_handler.py b/sdx_controller/handlers/lc_message_handler.py index 76230361..54821d71 100644 --- a/sdx_controller/handlers/lc_message_handler.py +++ b/sdx_controller/handlers/lc_message_handler.py @@ -93,6 +93,7 @@ def process_lc_json_msg( logger.info(f"Could not find breakdown for {service_id}") return None + conn_status = connection.get("status") oxp_number = len(breakdown) oxp_success_count = connection.get("oxp_success_count", 0) lc_domain = msg_json.get("lc_domain") @@ -108,46 +109,46 @@ def process_lc_json_msg( if msg_json.get("operation") != "delete": oxp_success_count += 1 connection["oxp_success_count"] = oxp_success_count + logger.info( + f"Update oxp_success_count: {oxp_success_count}; oxp_number: {oxp_number}" + ) if oxp_success_count == oxp_number: - if connection.get("status") and ( - connection.get("status") - == str(ConnectionStateMachine.State.RECOVERING) - ): - connection, _ = connection_state_machine( - connection, - ConnectionStateMachine.State.UNDER_PROVISIONING, - ) + conn_status = ConnectionStateMachine.State.UP connection, _ = connection_state_machine( - connection, ConnectionStateMachine.State.UP + connection, conn_status ) else: if connection.get("status") and ( connection.get("status") - == str(ConnectionStateMachine.State.RECOVERING) - ): - connection, _ = connection_state_machine( - connection, ConnectionStateMachine.State.ERROR - ) - elif ( - connection.get("status") - and connection.get("status") - != str(ConnectionStateMachine.State.DOWN) - and connection.get("status") - != str(ConnectionStateMachine.State.ERROR) + == str(ConnectionStateMachine.State.MODIFYING) + or connection.get("status") + == str(ConnectionStateMachine.State.UNDER_PROVISIONING) ): - connection, _ = connection_state_machine( - connection, ConnectionStateMachine.State.DOWN - ) + conn_status = ConnectionStateMachine.State.DOWN + connection, _ = connection_state_machine(connection, conn_status) # ToDo: eg: if 3 oxps in the breakdowns: (1) all up: up (2) parital down: remove_connection() # release successful oxp circuits if some are down: remove_connection() (3) count the responses # to finalize the status of the connection. - self.db_instance.add_key_value_pair_to_db( + self.db_instance.update_field_in_json( MongoCollections.CONNECTIONS, service_id, - connection, + "status", + str(conn_status), + ) + self.db_instance.update_field_in_json( + MongoCollections.CONNECTIONS, + service_id, + "oxp_response", + oxp_response, + ) + self.db_instance.update_field_in_json( + MongoCollections.CONNECTIONS, + service_id, + "oxp_success_count", + oxp_success_count, ) - logger.info("Connection updated: " + service_id) + logger.info("Connection updated: " + str(connection)) return # topology message RPC from OXP: no exchange name is defined. diff --git a/sdx_controller/messaging/rpc_queue_consumer.py b/sdx_controller/messaging/rpc_queue_consumer.py index 299e58aa..a564d710 100644 --- a/sdx_controller/messaging/rpc_queue_consumer.py +++ b/sdx_controller/messaging/rpc_queue_consumer.py @@ -141,7 +141,10 @@ def __init__(self, thread_queue, exchange_name, te_manager): self.channel = self.connection.channel() self.exchange_name = exchange_name - self.channel.queue_declare(queue=SUB_QUEUE) + # RabbitMQ no longer permits transient non-exclusive queues by default. + # This shared controller queue should be durable so it remains compatible + # with newer broker defaults. + self.channel.queue_declare(queue=SUB_QUEUE, durable=True) self._thread_queue = thread_queue self.te_manager = te_manager