Commit 9a06b597 authored by Andrea Sgambelluri's avatar Andrea Sgambelluri
Browse files

RSA P2MP and Optical contoller cleanup

parent f0666d46
Loading
Loading
Loading
Loading
+136 −84
Changes for src/opticalcontroller/OpticalController.py: 136 added lines, 84 removed lines.
Original line number Diff line number Diff line
@@ -55,6 +55,97 @@ from context.client.ContextClient import ContextClient



# ══════════════════════════════════════════════════════════════════════════
# OpticalController -- map of this file, for anyone debugging it later.
#
# This is a single Flask-RESTX app (Api `optical`) that does three
# unrelated-looking but related jobs:
#
#   1. LEGACY point-to-point RSA (Route & Spectrum Assignment):
#      AddLightpath / AddFlexLightpath / DelLightpath / DelFLightpath /
#      GetOpticalBand(s) and friends, backed by the global `rsa`
#      (opticalcontroller.RSA.RSA) instance. Single source, single
#      destination only. Mostly untouched legacy code.
#
#   2. P2MP (point-to-multipoint) RSA, backed by the global `rsa_pmp`
#      (opticalcontroller.RSA_P2MP.RSA_P2MP) instance -- see that module's
#      own header comment for its internal data structures. Key endpoints:
#      AddMultiLightpath(Comp) / ComputeP2MP (TAPI) / PathRecomp
#      (== /recompute-p2mp) / RecomputeSubflow / TapiDeleteP2MP /
#      GetActivePathComputations. `_active_path_computations` (below) is
#      this file's OWN bookkeeping of which sources/destinations are
#      currently "in service" per destination hub -- kept in sync with
#      rsa_pmp's db_flows by hand at every call site; the two are not the
#      same data structure and can drift if a new mutation site forgets to
#      update both (see _register_path_computation / the release logic in
#      TapiDeleteP2MP and PathRecomp for the pattern to follow).
#      GetTopologyMulti (re)reads topology from Context and either builds
#      `rsa_pmp` fresh (first call) or merges into the existing instance
#      (RSA_P2MP.merge_topology) so active flows survive a topology refresh.
#
#   3. The ECOC26 demo's alarm/failover/monitoring engine -- the part of
#      this file that has actually been under active development. Roughly,
#      in the order a real alarm sample flows through:
#        a. Policy (Java, CommonPolicyServiceImpl) forwards EVERY sample it
#           sees (breach or not) to PostAlarmNotification below, which:
#        b. resolves the raw KPI UUID to a device name + role
#           (_resolve_kpi_context),
#        c. computes low/high breach flags itself, directly from
#           measured_value against RECEIVED_POWER_RANGE_DBM /
#           PRE_FEC_BER_RANGE (it does NOT trust Policy's own breach flags),
#        d. runs the device through the "confirmed up" baseline-health gate
#           (_record_breach_state_and_check_confirmed_up /
#           _confirmed_up_devices) so alarms aren't trusted until the device
#           has shown one genuine healthy reading pair -- see the big
#           comment block above that function for the full rationale,
#        e. for the T1.2 failover device specifically, also runs the
#           separate T1.2 mute/confirm/revert state machine
#           (_update_t1_2_confirmation_state) -- see the comment block above
#           _T1_2_RESTING_POWER_RANGE_DBM,
#        f. dedupes against the last notification for this service+KPI
#           (_last_pushed_alarm) so Analytics re-publishing an unchanged
#           aggregated value doesn't spam duplicate notifications,
#        g. stores the notification (_alarm_notifications) and, if genuinely
#           new, calls _switch_monitored_device_and_recompute (the real
#           failover: recomputes the path via _recompute_subflow AND starts
#           real monitoring on the failover device via
#           _start_leaf_monitoring) UNCONDITIONALLY -- not gated on whether
#           any webhook push below succeeds,
#        h. pushes the notification to the static webhook
#           (_ALARM_WEBHOOK_URL) and/or any registered subscriber
#           (_alarm_subscriptions) on background threads.
#      Reverting back to the original device (T1.1) happens from three
#      independent triggers, all funneled through the same
#      _revert_to_original_leaf_device: T1.2 settling back to its own quiet
#      baseline (_update_t1_2_confirmation_state), an explicit path
#      delete that removed the failed-over path
#      (_restore_monitored_device_after_delete), or a new path computation
#      explicitly sourced from T1.1 again (_maybe_revert_on_source_recompute).
#      All failover/revert state (_LEAF_DEVICE_NAME, _leaf_failover_done,
#      _t1_2_confirmed_up) is guarded by the single _failover_state_lock
#      (an RLock -- see its own comment for why reentrant) since a single
#      real alarm typically arrives as two near-simultaneous samples
#      (RECEIVED_POWER and PRE_FEC_BER), each processed on its own request/
#      thread.
#
# Cross-cutting concerns worth knowing about before touching any of this:
#   - _RECOMPUTE_COOLDOWN_S / _reserve_recompute_slot / per-destination
#     _last_recompute_time_by_dest: rate-limits actually-PERFORMED path
#     recomputes (own failover, /RecomputeSubflow, /pathrecomp,
#     /recompute-p2mp) per destination hub, so a flood of external calls
#     reacting to the same fault can't re-churn the path repeatedly.
#   - _active_leaf_monitoring / _stop_leaf_monitoring: every call to
#     _start_leaf_monitoring first tears down whatever collector(s)/Analyzer
#     it previously started for that same device, so repeated failover/
#     revert cycles don't leave old monitoring pipelines running forever.
#   - Almost every one of these mechanisms exists because of a specific,
#     previously-observed failure mode during hardening for the ECOC26 demo
#     -- read the comment directly above each piece of state for the
#     concrete "here's what broke without this" rationale before changing
#     it; a change that looks like simplification may silently reopen one
#     of those.
# ══════════════════════════════════════════════════════════════════════════

logging.basicConfig(level=logging.INFO)
LOGGER = logging.getLogger(__name__)

@@ -113,6 +204,15 @@ def index():
    return render_template('index.html')


# ══════════════════════════════════════════════════════════════════════════
# LEGACY point-to-point RSA endpoints (job #1 in the file map above).
# Backed by the global `rsa` (opticalcontroller.RSA.RSA) instance -- a
# single-source/single-destination sibling of RSA_P2MP, initialized lazily
# by GetTopology/process_topology below. Mostly untouched, pre-dates the
# P2MP/alarm work; skip to "New addition: P2MP failover endpoints" further
# down for the actively-developed part of this file.
# ══════════════════════════════════════════════════════════════════════════

#@optical.route('/AddLightpath/<string:src>/<string:dst>/<int:bitrate>')
@optical.route('/AddLightpath/<string:src>/<string:dst>/<int:bitrate>/<int:bidir>')
@optical.response(200, 'Success')
@@ -897,29 +997,15 @@ class GetFlows(Resource):



"""@optical.route('/GetLinks')
@optical.response(200, 'Success')
@optical.response(404, 'Error, not found')
class GetMultiFlows(Resource):
    @staticmethod
    def get():
        try:
            if debug:
                print(rsa_pmp.db_flows)

            # Convert numpy int64 to native Python int
            def convert_int64(obj):
                if isinstance(obj, np.int64):
                    return int(obj)
                raise TypeError

            # Serialize the data with the custom converter
            json_data = json.dumps(rsa_pmp.db_flows, default=convert_int64)
            return json_data, 200
        except Exception as e:
            return f"Error: {str(e)}", 404"""


# NOTE: this GET /GetLinks resource used to be defined FOUR times in this
# file (two of them as inert triple-quoted docstrings, harmless; the other
# two as real, back-to-back, byte-for-byte identical classes). Flask-RESTX's
# @optical.route() keeps whichever registration for a given URL rule happens
# FIRST and silently ignores every later one for the same rule (verified
# directly against this deployment's flask_restplus: a second class
# registered on an already-used route never runs) -- so the second live
# definition, and both commented-out ones, were dead weight, not a second
# code path. Only the one real, actually-serving definition is kept below.
@optical.route('/GetLinks')
@optical.response(200, 'Success')
@optical.response(404, 'Error, not found')
@@ -958,66 +1044,6 @@ class GetMultiFlows(Resource):
            return f"Error: {str(e)}", 404


"""@optical.route('/GetLinks')
@optical.response(200, 'Success')
@optical.response(404, 'Error, not found')
class GetMultiFlows(Resource):
    @staticmethod
    def get():
        try:
            if debug:
                print(rsa_pmp.db_flows)

            # Convert numpy int64 to native Python int
            def convert_int64(obj):
                if isinstance(obj, np.int64):
                    return int(obj)
                raise TypeError

            # Serialize the data with the custom converter
            json_data = json.dumps(rsa_pmp.db_flows, default=convert_int64)
            return json_data, 200
        except Exception as e:
            return f"Error: {str(e)}", 404"""


@optical.route('/GetLinks')
@optical.response(200, 'Success')
@optical.response(404, 'Error, not found')
class GetMultiFlows(Resource):
    @staticmethod
    def get():
        try:
            # Initialize context and topology like in AddMultiLightpathComp
            context_client = ContextClient()
            context_id_x = json_context_id(DEFAULT_CONTEXT_NAME)
            topology_id_x = json_topology_id(DEFAULT_TOPOLOGY_NAME, context_id_x)
            topology_details = context_client.GetTopologyDetails(TopologyId(**topology_id_x))
            
            topo_id_str = topology_id_x["topology_uuid"]["uuid"]
            cxt_id_str = topology_id_x["context_id"]["context_uuid"]["uuid"]
            process_topology(topo_id_str, cxt_id_str)
            
            # Ensure rsa_pmp is initialized like in AddMultiLightpathComp
            global rsa_pmp
            if rsa_pmp is None:
                rsa_pmp = RSA_P2MP(node_dict, links_dict)
            
            if debug:
                print(rsa_pmp.links_dict)

            # Convert numpy int64 to native Python int
            def convert_int64(obj):
                if isinstance(obj, np.int64):
                    return int(obj)
                raise TypeError

            # Serialize the data with the custom converter
            json_data = json.dumps(rsa_pmp.links_dict)
            return json_data, 200
        except Exception as e:
            return f"Error: {str(e)}", 404

@optical.route('/GetOpticalBands')
@optical.response(200, 'Success')
@optical.response(404, 'Error, not found')
@@ -1701,7 +1727,18 @@ _alarm_subscriptions: dict = {}
# ──────────────────────────────────────────────────────────────────────────────

class TapiClient:
    """Manual TAPI REST client test"""
    """Builds the TAPI-formatted topology response for GET
    /restconf/operations/tapi-topology:get-topology (and is reused as a
    topology source inside ComputeP2MP._create_tapi_connectivity_response,
    see the `topology_raw` fallback there).

    NOTE: topology_extract()/read_DSC_only() below talk to a HARDCODED
    external HTTP address (http://10.30.7.65/tfs-api/... and
    http://deviceservice:10065/...) via plain `requests` calls, instead of
    going through ContextClient/DeviceClient like the rest of this file.
    That IP is specific to whichever deployment this was last pointed at --
    if TAPI topology retrieval breaks after moving to a different cluster/
    demo rig, this hardcoded address is almost certainly why."""

    def __init__(self):
        self.headers = {
@@ -2103,6 +2140,21 @@ class ComputeP2MP(Resource):
        return int(v)

    def _create_tapi_connectivity_response(self, dscm_flow, sources, destinations, bitrate, bidirectional):
        """Turn an already-computed RSA_P2MP flow (dscm_flow -- the
        DSCM-converted form of rsa_pmp.db_flows[flow_id], see
        flow_to_DSCM_message) into a TAPI tapi-connectivity:connectivity-service
        response: one end-point per destination (first) and per source
        digital-subcarrier-group (after), each carrying the
        service-interface-point UUID looked up from the live topology
        (tapi_client.topology_extract(), cached on self.topology_raw if a
        caller set it first) plus per-endpoint spectrum config, and one
        top-level optical-connection-attributes block
        (_create_optical_attributes) describing modulation/power/subcarrier
        layout. Device -> physical-port/channel/frequency mappings are
        governed by the fixed lookup tables above
        (_DEVICE_CHANNEL_MAP/_GROUP_PORT_MAP/_DEVICE_FREQUENCY_MAP) rather
        than derived purely from the topology, since this is a fixed demo
        rig rather than a general deployment."""

        # Main frequency (no unit conversion here if your dscm gives Hz already)
        freq_raw = dscm_flow.get('frequency') or dscm_flow.get('central-frequency') or dscm_flow.get('central_frequency')
+305 −17

File changed.

Preview size limit exceeded, changes collapsed.

+34 −24
Changes for src/opticalcontroller/heuristicMST.py: 34 added lines, 24 removed lines.
Original line number Diff line number Diff line
@@ -21,26 +21,36 @@ import heapq
from collections import defaultdict

class Graph:
    """Undirected adjacency-list graph used for P2MP/point-to-point path
    computation (see spf_find_path / heuristic_mst below).

    Storage: self.graph is a dict keyed by node name -> list of
    (neighbor, src_port, dst_port, weight) tuples. add_edge() appends the
    edge in BOTH directions (src->dst and dst->src, with ports swapped) so
    the graph can be walked either way -- there is no separate "reverse"
    structure.

    NOTE: this class used to define add_vertex/add_edge/printGraph twice.
    Because a class body executes top-to-bottom and each `def` simply
    rebinds the method name, the second definition of each silently
    replaced the first -- the first add_vertex/add_edge (which referenced
    self.vertices/self.edges, attributes that are never initialized
    anywhere in __init__) were dead code, unreachable from the moment the
    class was defined. They have been removed; only the versions that were
    actually ever called (the ones operating on self.graph) remain.
    """
    def __init__(self):
        self.graph = defaultdict(list)
        self.visited = {}  # Track visited nodes



    def add_vertex(self, node):
        self.vertices[node] = []
        if node not in self.graph:
            self.graph[node] = []

    def add_edge(self, src, dst, src_port, dst_port, weight):
        self.edges.append((src, dst, src_port, dst_port, weight))
        self.vertices[src].append((dst, src_port, dst_port, weight))

    def printGraph(self):
        print("Vertices:")
        for vertex, edges in self.vertices.items():
            print(f"{vertex}: {edges}")
        print("Edges:")
        for edge in self.edges:
            print(edge)
        # Store port information along with the edges
        self.graph[src].append((dst, src_port, dst_port, weight))
        self.graph[dst].append((src, dst_port, src_port, weight))  # Assuming undirected graph

    def printGraph(self):
        print("Vertices:")
@@ -51,15 +61,6 @@ class Graph:
            for edge in edges:
                print(f"{vertex} -> {edge[0]} (src_port: {edge[1]}, dst_port: {edge[2]}, weight: {edge[3]})")

    def add_vertex(self, node):
        if node not in self.graph:
            self.graph[node] = []

    def add_edge(self, src, dst, src_port, dst_port, weight):
        # Store port information along with the edges
        self.graph[src].append((dst, src_port, dst_port, weight))
        self.graph[dst].append((src, dst_port, src_port, weight))  # Assuming undirected graph

    def get_edge_weight(self, src, dst):
        # Return the weight of the edge between src and dst nodes
        for neighbor, _, _, weight in self.graph.get(src, []):
@@ -78,7 +79,12 @@ class Graph:
        return self.graph.get(node, [])

def spf_find_path(graph, start, end, visited_edges=None):
    # Find the shortest path between start and end node using Dijkstra's algorithm
    """Dijkstra shortest path from start to end over `graph` (a Graph
    instance). Returns a list of (from_node, to_node, src_port, dst_port)
    hops, or None if end is unreachable. visited_edges, when given, is a
    set of (from_node, to_node) pairs -- an edge NOT in that set is skipped
    (used to keep a path from immediately walking back over an edge
    reserved by a different direction of the same computation)."""
    priority_queue = [(0, start, [])]  # (cost, node, path)
    visited_nodes = set()
    
@@ -103,7 +109,11 @@ def spf_find_path(graph, start, end, visited_edges=None):


def heuristic_mst(graph, start):
    # A heuristic-based MST construction (e.g., using edge weights).
    """Build a minimum-spanning-tree-like Graph rooted at `start` via a
    Prim's-algorithm-style expansion (cheapest-edge-first from the visited
    set). Used by compute_path_pmp (RSA_P2MP.py) as a cheap way to get a
    single tree that spf_find_path can then extract a src->dst route from,
    rather than re-running Dijkstra from scratch for every destination."""
    mst = Graph()
    visited = set()
    min_heap = [(0, start, None, None)]  # (cost, node, parent, port)