Loading src/automation/service/__main__.py +11 −5 Changes for src/automation/service/__main__.py: 11 added lines, 5 removed lines. Original line number Diff line number Diff line Loading @@ -85,11 +85,17 @@ def main(): grpc_service.stop() # event_engine.stop() # Drop the database try: Engine.drop_database(db_engine) except Exception as ex: LOGGER.exception(f"Failed to check/create the database: {str(db_engine.url)}") # Deliberately NOT dropping the database here. It used to be dropped on # every shutdown on the assumption that the next startup's # Engine.create_database() call would always recreate it first -- but in # a Kubernetes rolling restart, the old pod's shutdown and the new pod's # startup briefly overlap, and if this DROP ran after the new pod's # CREATE, the database was left missing with no automatic recovery: the # new pod had already passed its create-database step and never retries # it, so every DB operation afterward failed with "database ... does not # exist" (observed in practice). No other TFS component drops its own # database on shutdown; CREATE DATABASE IF NOT EXISTS on startup is # already the correct, safe, idempotent pattern on its own. LOGGER.info("Bye") return 0 Loading src/common/tools/database/GenericDatabase.py +1 −1 Changes for src/common/tools/database/GenericDatabase.py: 1 added line, 1 removed line. Original line number Diff line number Diff line Loading @@ -76,7 +76,7 @@ class Database: extra_details=["Unique key voilation: {:}".format(e)] ) else: LOGGER.error(f"Failed to insert new row into {row.__class__.__name__} table. {str(e)}") raise OperationFailedException ("Deletion by column id", extra_details=["unable to delete row {:}".format(e)]) raise OperationFailedException ("Insertion of row", extra_details=["unable to insert row {:}".format(e)]) finally: session.close() Loading src/opticalcontroller/Dockerfile +3 −0 Changes for src/opticalcontroller/Dockerfile: 3 added lines, 0 removed lines. Original line number Diff line number Diff line Loading @@ -69,6 +69,9 @@ COPY src/kpi_manager/client/. kpi_manager/client/ COPY src/telemetry/__init__.py telemetry/__init__.py COPY src/telemetry/frontend/__init__.py telemetry/frontend/__init__.py COPY src/telemetry/frontend/client/. telemetry/frontend/client/ COPY src/analytics/__init__.py analytics/__init__.py COPY src/analytics/frontend/__init__.py analytics/frontend/__init__.py COPY src/analytics/frontend/client/. analytics/frontend/client/ COPY src/automation/__init__.py automation/__init__.py COPY src/automation/client/. automation/client/ COPY src/opticalcontroller/. opticalcontroller/ Loading src/opticalcontroller/OpticalController.py +327 −102 File changed.Preview size limit exceeded, changes collapsed. Show changes src/opticalcontroller/RSA_P2MP.py +96 −0 Changes for src/opticalcontroller/RSA_P2MP.py: 96 added lines, 0 removed lines. Original line number Diff line number Diff line Loading @@ -163,6 +163,52 @@ class RSA_P2MP: self.roadm_subcarriers.clear() self.flow_id = 0 def merge_topology(self, new_nodes, new_links): """Incrementally merge newly-discovered devices/links into this already-running instance, WITHOUT touching any active-flow state (db_flows, global_slots_used, global_subcarriers_used, optical_bands, slots_used, subcarriers_used, etc.). Replacing rsa_pmp wholesale on every topology refresh (the previous approach) silently wiped every active P2MP flow's bookkeeping the moment topology changed, since a brand new instance starts with all of that state empty. This instead only adds genuinely new nodes (by their dict key, i.e. device UUID) and links (by "name", matching _initialize_graph's own join key) that aren't already present; anything already known -- including every active flow's slots, subcarriers, and optical bands -- is left completely untouched. Safe to call at any time, even with flows actively in progress. Returns (added_node_count, added_link_count). """ added_nodes = 0 for node_id, node_info in new_nodes.items(): if node_id not in self.nodes_dict: self.nodes_dict[node_id] = node_info self.graph.add_vertex(node_id) added_nodes += 1 existing_link_names = { link.get("name") for link in self.links_dict.get("optical_links", []) } added_links = 0 for link in new_links.get("optical_links", []): name = link.get("name") if not name or name in existing_link_names: continue self.links_dict.setdefault("optical_links", []).append(link) details = link.get("optical_details", {}) # Matches _initialize_graph's own parsing exactly (same # assumption: exactly one "-" separating src and dst). src, dst = name.split('-') self.graph.add_edge( src, dst, details.get("src_port"), details.get("dst_port"), 1 ) existing_link_names.add(name) added_links += 1 return added_nodes, added_links def _initialize_graph(self): # Create an empty graph graph = Graph() Loading Loading @@ -665,6 +711,56 @@ class RSA_P2MP: f"from flow {flow_id}" ) def del_subflow_by_dest(self, flow_id, dest): """Release just dest's sub-flow entries from flow_id, leaving every other destination's sub-flow (and its reserved spectrum) untouched. Mirrors del_subflow exactly, but scoped by destination instead of source: needed because a single P2MP flow_id can carry multiple destinations (see the "dst" field flows_per_link entries each carry -- see rsa_multi_computation), and deleting just one of them must not release spectrum that the flow's other, still-active destinations depend on. Without this, TapiDeleteP2MP's only option for a partial delete was releasing the flow's slots/subcarriers/ optical bands wholesale (or not at all), even though other destinations of that same flow_id were still marked active in _active_path_computations -- their spectrum could then be silently re-allocated to a different flow while supposedly still in service. Returns (True, message) on success, (False, message) if flow_id doesn't exist or dest isn't part of it. """ if flow_id not in self.db_flows: return False, f"Flow {flow_id} not found" flow = self.db_flows[flow_id] flow_dst = flow.get("dst", []) if isinstance(flow_dst, str): flow_dst = [flow_dst] if dest not in flow_dst: return False, f"Destination {dest} is not part of flow {flow_id}" kept, released = [], [] for entry in flow.get("flows_per_link", []): (released if entry.get("dst") == dest else kept).append(entry) for entry in released: for slot in entry.get("assigned_slots") or []: self.global_slots_used.discard(slot) for subcarrier in entry.get("subcarriers") or []: self.global_subcarriers_used.discard(subcarrier) sub_band_id = entry.get("sub_flow_id") if sub_band_id in self.optical_bands: del self.optical_bands[sub_band_id] flow["flows_per_link"] = kept flow["dst"] = [d for d in flow_dst if d != dest] return True, ( f"Released {len(released)} sub-flow entry/entries for destination " f"{dest} from flow {flow_id}" ) def get_fibers_forward(self, links, slots, band): Loading Loading
src/automation/service/__main__.py +11 −5 Changes for src/automation/service/__main__.py: 11 added lines, 5 removed lines. Original line number Diff line number Diff line Loading @@ -85,11 +85,17 @@ def main(): grpc_service.stop() # event_engine.stop() # Drop the database try: Engine.drop_database(db_engine) except Exception as ex: LOGGER.exception(f"Failed to check/create the database: {str(db_engine.url)}") # Deliberately NOT dropping the database here. It used to be dropped on # every shutdown on the assumption that the next startup's # Engine.create_database() call would always recreate it first -- but in # a Kubernetes rolling restart, the old pod's shutdown and the new pod's # startup briefly overlap, and if this DROP ran after the new pod's # CREATE, the database was left missing with no automatic recovery: the # new pod had already passed its create-database step and never retries # it, so every DB operation afterward failed with "database ... does not # exist" (observed in practice). No other TFS component drops its own # database on shutdown; CREATE DATABASE IF NOT EXISTS on startup is # already the correct, safe, idempotent pattern on its own. LOGGER.info("Bye") return 0 Loading
src/common/tools/database/GenericDatabase.py +1 −1 Changes for src/common/tools/database/GenericDatabase.py: 1 added line, 1 removed line. Original line number Diff line number Diff line Loading @@ -76,7 +76,7 @@ class Database: extra_details=["Unique key voilation: {:}".format(e)] ) else: LOGGER.error(f"Failed to insert new row into {row.__class__.__name__} table. {str(e)}") raise OperationFailedException ("Deletion by column id", extra_details=["unable to delete row {:}".format(e)]) raise OperationFailedException ("Insertion of row", extra_details=["unable to insert row {:}".format(e)]) finally: session.close() Loading
src/opticalcontroller/Dockerfile +3 −0 Changes for src/opticalcontroller/Dockerfile: 3 added lines, 0 removed lines. Original line number Diff line number Diff line Loading @@ -69,6 +69,9 @@ COPY src/kpi_manager/client/. kpi_manager/client/ COPY src/telemetry/__init__.py telemetry/__init__.py COPY src/telemetry/frontend/__init__.py telemetry/frontend/__init__.py COPY src/telemetry/frontend/client/. telemetry/frontend/client/ COPY src/analytics/__init__.py analytics/__init__.py COPY src/analytics/frontend/__init__.py analytics/frontend/__init__.py COPY src/analytics/frontend/client/. analytics/frontend/client/ COPY src/automation/__init__.py automation/__init__.py COPY src/automation/client/. automation/client/ COPY src/opticalcontroller/. opticalcontroller/ Loading
src/opticalcontroller/OpticalController.py +327 −102 File changed.Preview size limit exceeded, changes collapsed. Show changes
src/opticalcontroller/RSA_P2MP.py +96 −0 Changes for src/opticalcontroller/RSA_P2MP.py: 96 added lines, 0 removed lines. Original line number Diff line number Diff line Loading @@ -163,6 +163,52 @@ class RSA_P2MP: self.roadm_subcarriers.clear() self.flow_id = 0 def merge_topology(self, new_nodes, new_links): """Incrementally merge newly-discovered devices/links into this already-running instance, WITHOUT touching any active-flow state (db_flows, global_slots_used, global_subcarriers_used, optical_bands, slots_used, subcarriers_used, etc.). Replacing rsa_pmp wholesale on every topology refresh (the previous approach) silently wiped every active P2MP flow's bookkeeping the moment topology changed, since a brand new instance starts with all of that state empty. This instead only adds genuinely new nodes (by their dict key, i.e. device UUID) and links (by "name", matching _initialize_graph's own join key) that aren't already present; anything already known -- including every active flow's slots, subcarriers, and optical bands -- is left completely untouched. Safe to call at any time, even with flows actively in progress. Returns (added_node_count, added_link_count). """ added_nodes = 0 for node_id, node_info in new_nodes.items(): if node_id not in self.nodes_dict: self.nodes_dict[node_id] = node_info self.graph.add_vertex(node_id) added_nodes += 1 existing_link_names = { link.get("name") for link in self.links_dict.get("optical_links", []) } added_links = 0 for link in new_links.get("optical_links", []): name = link.get("name") if not name or name in existing_link_names: continue self.links_dict.setdefault("optical_links", []).append(link) details = link.get("optical_details", {}) # Matches _initialize_graph's own parsing exactly (same # assumption: exactly one "-" separating src and dst). src, dst = name.split('-') self.graph.add_edge( src, dst, details.get("src_port"), details.get("dst_port"), 1 ) existing_link_names.add(name) added_links += 1 return added_nodes, added_links def _initialize_graph(self): # Create an empty graph graph = Graph() Loading Loading @@ -665,6 +711,56 @@ class RSA_P2MP: f"from flow {flow_id}" ) def del_subflow_by_dest(self, flow_id, dest): """Release just dest's sub-flow entries from flow_id, leaving every other destination's sub-flow (and its reserved spectrum) untouched. Mirrors del_subflow exactly, but scoped by destination instead of source: needed because a single P2MP flow_id can carry multiple destinations (see the "dst" field flows_per_link entries each carry -- see rsa_multi_computation), and deleting just one of them must not release spectrum that the flow's other, still-active destinations depend on. Without this, TapiDeleteP2MP's only option for a partial delete was releasing the flow's slots/subcarriers/ optical bands wholesale (or not at all), even though other destinations of that same flow_id were still marked active in _active_path_computations -- their spectrum could then be silently re-allocated to a different flow while supposedly still in service. Returns (True, message) on success, (False, message) if flow_id doesn't exist or dest isn't part of it. """ if flow_id not in self.db_flows: return False, f"Flow {flow_id} not found" flow = self.db_flows[flow_id] flow_dst = flow.get("dst", []) if isinstance(flow_dst, str): flow_dst = [flow_dst] if dest not in flow_dst: return False, f"Destination {dest} is not part of flow {flow_id}" kept, released = [], [] for entry in flow.get("flows_per_link", []): (released if entry.get("dst") == dest else kept).append(entry) for entry in released: for slot in entry.get("assigned_slots") or []: self.global_slots_used.discard(slot) for subcarrier in entry.get("subcarriers") or []: self.global_subcarriers_used.discard(subcarrier) sub_band_id = entry.get("sub_flow_id") if sub_band_id in self.optical_bands: del self.optical_bands[sub_band_id] flow["flows_per_link"] = kept flow["dst"] = [d for d in flow_dst if d != dest] return True, ( f"Released {len(released)} sub-flow entry/entries for destination " f"{dest} from flow {flow_id}" ) def get_fibers_forward(self, links, slots, band): Loading