Commit 6cffadd9 authored by Lluis Gifre Renom's avatar Lluis Gifre Renom
Browse files

Optical Controller component:

- Refactor RSA class methods for better slot management and error handling.
- Implement spectrum persistence logic in tools for optical link updates.
- Add comprehensive regression tests for spectrum allocation and release scenarios.
- Document spectrum persistence behavior in README.
parent cb4cfa52
Loading
Loading
Loading
Loading
+4 −4
Changes for src/opticalcontroller/OpticalController.py: 4 added lines, 4 removed lines.
Original line number Diff line number Diff line
@@ -519,10 +519,10 @@ class GetTopology(Resource):
                    
                links_dict["optical_links"].append(link_dict_type)
    
            rsa = RSA(node_dict, links_dict)
            if debug:
                print(f'rsa.init_link_slots2() {rsa}')
                print(rsa.init_link_slots2())
            new_rsa = RSA(node_dict, links_dict)
            slot_counts = new_rsa.init_link_slots2()
            rsa = new_rsa
            LOGGER.debug('Initialized optical spectrum: %s', slot_counts)
            return 'OK', 200
        except Exception:
            LOGGER.exception('Error while processing topology data')
+42 −0
Changes for src/opticalcontroller/README.md: 42 added lines, 0 removed lines.
Original line number Diff line number Diff line
@@ -18,3 +18,45 @@ pip install -r requirements_optical_ctrl_test.txt

python OpticalController.py
```

## Spectrum Persistence

Spectrum maps use `1` for available slots and `0` for occupied slots.
Allocation and release persist the affected optical link through Context before
publishing its new in-memory bitmap. Link identifiers, endpoints, ports, and
other metadata are retained. Releasing a range does not clear other occupied
ranges or bands; a present `used` flag remains set while occupancy remains.

Context write failures are logged and propagated. The optical deletion REST
request fails rather than reporting successful release, and Service propagates
Optical Controller server errors instead of discarding the service records.
Updates are per link, not an atomic transaction over the whole path. A failure
after earlier links were updated may therefore require a retry.

Topology loading preserves existing slot values and uses the existing sparse
map convention to fill missing slots as available. Empty bands remain empty.
Initialization is independent of DEBUG logging. This does not reconstruct
in-memory flow records after a controller restart or repair stale occupancy
left by earlier versions. Do not clear such slots automatically: first verify
service, connection, and reservation ownership.

## Regression Tests

From the repository root with component dependencies installed:

```bash
PYTHONPATH=src python -m pytest src/opticalcontroller/tests -q
PYTHONPATH=src python -m pytest \
  src/service/tests/test_optical_spectrum_reservation.py -q
```

`test_spectrum_sync.py` uses a Context test double to check allocation/release,
forward and reverse links, partial release, retained metadata, topology reload,
and write failure/retry. Its REST-boundary test requires the Optical Controller
Flask dependencies; run in the component image if unavailable locally.

After deploying a reviewed fix, validate one unidirectional service and a pair
of opposite-direction services against a known baseline. Compare allocated
slots in `/OpticalTFS/GetLinks` with `/tfs-api/optical_links` while ACTIVE and
after teardown. Keep another allocation active to verify that release affects
only the removed service. Existing stale data must be reconciled separately.
+30 −21
Changes for src/opticalcontroller/RSA.py: 30 added lines, 21 removed lines.
Original line number Diff line number Diff line
@@ -13,6 +13,7 @@
# limitations under the License.

import logging
from copy import deepcopy
from typing import Dict, Tuple
from opticalcontroller.dijkstra import *
from opticalcontroller.tools import *
@@ -75,8 +76,8 @@ class RSA():
            if not band_slots:
                return 0

            fib[band_name] = {str(slot_index): 1 for slot_index in range(width)}
            return width
            fib[band_name] = correct_slot(band_slots, width=width)
            return len(fib[band_name])

        if full_links:
            print("2026 initialize full spectrum")
@@ -327,23 +328,35 @@ class RSA():
        return c_sts, l_sts, s_sts

    def update_link(self, fib, slots, band):
        #print(fib)
        for i in slots:
            fib[band][str(i)] = 0
        if 'used' in fib:
            fib['used'] = True
        print(f"fib updated {fib}")    
        #print(fib)
        self._set_link_slots(fib, slots, band, available=False)
        
    def update_link_2(self, fib, slots, band, link):
        #print(fib)
        for i in slots:
            fib[band][str(i)] = 0
        if 'used' in fib:
            fib['used'] = True
        self._set_link_slots(fib, slots, band, available=False, link=link)

        set_link_update(fib,link)  
        #print(fib)    
    def _set_link_slots(self, fib, slots, band, available, link=None):
        if link is None:
            link = next(
                (item for item in self.links_dict["optical_links"]
                 if item["optical_details"] is fib), None
            )
        if link is None:
            raise ValueError("Cannot persist spectrum for an unknown link")
        updated = deepcopy(fib)
        for slot in slots:
            key = str(slot)
            if key not in updated.get(band, {}):
                raise ValueError(f"Unknown slot {slot} in band {band}")
            updated[band][key] = int(available)
        if "used" in updated:
            updated["used"] = any(
                value == 0
                for name in ("c_slots", "l_slots", "s_slots")
                for value in updated.get(name, {}).values()
            )
        # Publish locally only after Context accepts the new bitmap.
        set_link_update(updated, link)
        fib.clear()
        fib.update(updated)

    def update_optical_band(self, optical_band_id, slots, band):
        for i in slots:
@@ -354,11 +367,7 @@ class RSA():
            self.optical_bands[optical_band_id][band][str(i)] = 1

    def restore_link(self, fib, slots, band):
        for i in slots:
            fib[band][str(i)] = 1
        if 'used' in fib:
            fib['used'] = False
        #fib[band].sort()
        self._set_link_slots(fib, slots, band, available=True)
        
    # def restore_link_2(self, fib, slots, band, link):
    #     print("start restoring link")
+203 −0
Changes for src/opticalcontroller/tests/test_spectrum_sync.py: 203 added lines, 0 removed lines.
Original line number Diff line number Diff line
# Copyright 2022-2026 ETSI SDG TeraFlowSDN (TFS) (https://tfs.etsi.org/)
#
# Licensed under the Apache License, Version 2.0 (the "License");
# you may not use this file except in compliance with the License.
# You may obtain a copy of the License at
#
#      http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.

"""Controller spectrum publication with an in-memory Context test double."""

from copy import deepcopy
from importlib import import_module

import pytest
from google.protobuf.json_format import MessageToDict

from opticalcontroller.RSA import RSA
from opticalcontroller import tools


def make_link(name):
    return {
        "name": name,
        "link_id": {"link_uuid": {"uuid": name}},
        "link_endpoint_ids": [{
            "device_id": {"device_uuid": {"uuid": "device-a"}},
            "endpoint_uuid": {"uuid": "port-a"},
        }],
        "optical_details": {
            "length": 42.0, "src_port": "1", "dst_port": "6",
            "used": False,
            "c_slots": {str(i): 1 for i in range(8)},
            "l_slots": {"0": 1},
        },
    }


@pytest.fixture
def spectrum(monkeypatch):
    stored = {}
    clients = []

    class Context:
        fail = False

        def __init__(self):
            self.closed = False
            clients.append(self)

        def SetOpticalLink(self, request):
            if self.fail:
                raise RuntimeError("Context unavailable")
            stored[request.name] = MessageToDict(
                request, preserving_proto_field_name=True
            )

        def close(self):
            self.closed = True

    monkeypatch.setattr(tools, "ContextClient", Context)
    rsa = RSA.__new__(RSA)
    rsa.links_dict = {"optical_links": [make_link("A-B"), make_link("B-A")]}
    rsa.link_key__to__src_dst_node_names = {
        "A-B": ("A", "B"), "B-A": ("B", "A"),
    }
    rsa.optical_bands = {}
    rsa.c_slot_number = rsa.l_slot_number = rsa.s_slot_number = 0
    return rsa, stored, clients, Context


def make_flow(bidir=False):
    return {
        "flow_id": 1, "flows": [], "band_type": "c_slots",
        "slots": [1, 2], "fiber_forward": {}, "fiber_backward": {},
        "op-mode": 8, "n_slots": 2, "path": ["A", "B"],
        "links": ["A-B"], "bidir": bidir,
    }


@pytest.mark.parametrize("bidir", [False, True])
def test_flow_deletion_publishes_both_directions(spectrum, bidir):
    rsa, stored, clients, _ = spectrum
    names = ["A-B", "B-A"] if bidir else ["A-B"]
    for name in names:
        link = rsa.get_link_by_name(name)
        rsa.update_link_2(link["optical_details"], [1, 2], "c_slots", link)
        assert stored[name]["optical_details"]["c_slots"]["1"] == 0
    rsa.del_flow(make_flow(bidir), 1)
    assert set(stored) == set(names)
    for name in names:
        details = stored[name]["optical_details"]
        assert details["c_slots"]["1"] == details["c_slots"]["2"] == 1
        assert details["length"] == 42.0
        assert details["src_port"] == "1"
        assert stored[name]["link_endpoint_ids"]
        assert not details.get("used", False)
    assert all(client.closed for client in clients)


def test_partial_release_preserves_other_allocations_and_bands(spectrum):
    rsa, stored, _, _ = spectrum
    fib = rsa.get_link_by_name("A-B")["optical_details"]
    rsa.update_link(fib, [1, 2, 5], "c_slots")
    rsa.update_link(fib, [0], "l_slots")
    rsa.restore_link(fib, [1, 2], "c_slots")
    details = stored["A-B"]["optical_details"]
    assert details["c_slots"]["5"] == 0
    assert details["l_slots"]["0"] == 0
    assert details["used"]
    rsa.restore_link(fib, [5], "c_slots")
    assert stored["A-B"]["optical_details"]["used"]
    rsa.restore_link(fib, [0], "l_slots")
    assert not stored["A-B"]["optical_details"].get("used", False)


@pytest.mark.parametrize("operation", ["allocate", "release"])
def test_context_failure_does_not_publish_local_state(spectrum, operation):
    rsa, stored, clients, context = spectrum
    fib = rsa.get_link_by_name("A-B")["optical_details"]
    rsa.update_link(fib, [1, 2], "c_slots")
    before = deepcopy(fib)
    persisted = deepcopy(stored)
    context.fail = True
    call = rsa.restore_link if operation == "release" else rsa.update_link
    with pytest.raises(RuntimeError, match="Context unavailable"):
        call(fib, [1, 2], "c_slots")
    assert fib == before
    assert stored == persisted
    assert clients[-1].closed
    context.fail = False
    call(fib, [1, 2], "c_slots")
    assert stored["A-B"]["optical_details"]["c_slots"]["1"] == (
        1 if operation == "release" else 0
    )


def test_band_deletion_persists_released_slots(spectrum):
    rsa, stored, _, _ = spectrum
    link = rsa.get_link_by_name("A-B")
    rsa.update_link(link["optical_details"], [1, 2], "c_slots")
    rsa.optical_bands[7] = {
        "links": ["A-B"], "band_type": "c_slots", "n_slots": 2,
        "c_slots": {"1": 0, "2": 0}, "is_active": True,
    }
    rsa.del_band(make_flow(), o_b_id=7)
    assert stored["A-B"]["optical_details"]["c_slots"]["1"] == 1
    assert not rsa.optical_bands[7]["is_active"]


def test_reload_preserves_occupancy_and_unsupported_bands(spectrum, monkeypatch):
    rsa, _, _, _ = spectrum
    module = import_module("opticalcontroller.RSA")
    monkeypatch.setattr(module, "full_links", 1)
    fib = rsa.get_link_by_name("A-B")["optical_details"]
    fib["c_slots"] = {"1": 0, "4": 1}
    fib["s_slots"] = {}
    rsa.init_link_slots2()
    assert fib["c_slots"]["1"] == 0
    assert fib["c_slots"]["4"] == 1
    assert fib["c_slots"]["2"] == 1
    assert fib["s_slots"] == {}
    assert len(fib["c_slots"]) == module.Nc


def test_invalid_slot_does_not_change_state(spectrum):
    rsa, stored, _, _ = spectrum
    fib = rsa.get_link_by_name("A-B")["optical_details"]
    before = deepcopy(fib)
    with pytest.raises(ValueError, match="Unknown slot"):
        rsa.restore_link(fib, [9000], "c_slots")
    assert fib == before
    assert not stored


def test_rest_teardown_reports_persistence_failure_and_allows_retry(
    spectrum, monkeypatch
):
    pytest.importorskip("flask_restplus")
    api = import_module("opticalcontroller.OpticalController")
    rsa, stored, _, context = spectrum
    flow = dict(make_flow(), src="A", dst="B", bitrate=800, is_active=True)
    rsa.db_flows = {1: flow}
    fib = rsa.get_link_by_name("A-B")["optical_details"]
    rsa.update_link(fib, [1, 2], "c_slots")
    monkeypatch.setattr(api, "rsa", rsa)
    monkeypatch.setitem(api.app.config, "PROPAGATE_EXCEPTIONS", False)
    client = api.app.test_client()
    context.fail = True
    response = client.delete("/OpticalTFS/DelLightpath/A/B/800/1")
    assert response.status_code == 500
    assert flow["is_active"]
    assert stored["A-B"]["optical_details"]["c_slots"]["1"] == 0
    context.fail = False
    response = client.delete("/OpticalTFS/DelLightpath/A/B/800/1")
    assert response.status_code == 200
    assert not flow["is_active"]
    assert stored["A-B"]["optical_details"]["c_slots"]["1"] == 1
+20 −38
Changes for src/opticalcontroller/tools.py: 20 added lines, 38 removed lines.
Original line number Diff line number Diff line
@@ -13,12 +13,15 @@
# limitations under the License.

import json, logging, numpy as np
from google.protobuf.json_format import ParseDict
from typing import Dict, List, Optional, Tuple
from opticalcontroller.variables import  *
from context.client.ContextClient import ContextClient
from common.proto.context_pb2 import TopologyId , LinkId , OpticalLink , OpticalLinkDetails
from common.tools.object_factory.OpticalLink import correct_slot

LOGGER = logging.getLogger(__name__)


def common_slots(a, b):
    return list(np.intersect1d(a, b))
@@ -402,44 +405,23 @@ def update_optical_band (optical_bands,optical_band_id,band,link):
    set_link_update(fib,link,test=f"restoring_optical_band {link['link_id']}")
    
def set_link_update (fib:dict,link:dict,test="updating"):
    #print(link)
    print(f"invoked from {test}")
    print(f"fib updated {fib}")  
    optical_link = OpticalLink()
    linkId = LinkId()
    #linkId.link_uuid.uuid=link["link_id"]["link_uuid"]["uuid"]
    linkId.link_uuid.uuid=link["link_id"]["link_uuid"]["uuid"]
    optical_details = OpticalLinkDetails()
    optical_link.optical_details.length=0
    if "src_port" in fib :
       optical_link.optical_details.src_port=fib["src_port"]
    if "dst_port" in fib :   
       optical_link.optical_details.dst_port=fib["dst_port"]
    if "local_peer_port" in fib :   
       optical_link.optical_details.local_peer_port=fib['local_peer_port']
    if "remote_peer_port" in fib:   
       optical_link.optical_details.remote_peer_port=fib['remote_peer_port']
       
    optical_link.optical_details.used=fib['used'] if 'used' in fib else False
    if "c_slots" in fib :
        
        handle_slot( optical_link.optical_details.c_slots,fib["c_slots"])
    if "s_slots" in fib :
             
       handle_slot( optical_link.optical_details.s_slots,fib["s_slots"])
    if "l_slots" in fib :
           
       handle_slot( optical_link.optical_details.l_slots,fib["l_slots"])

    
    optical_link.name=link['name']
  
    optical_link.link_id.CopyFrom(linkId)
    
    """Persist spectrum without dropping link metadata or hiding failures."""
    optical_link = ParseDict(link, OpticalLink())
    details = optical_link.optical_details
    for field in ("length", "src_port", "dst_port", "local_peer_port",
                  "remote_peer_port", "used"):
        if field in fib:
            setattr(details, field, fib[field])
    for band in ("c_slots", "s_slots", "l_slots"):
        if band in fib:
            details.ClearField(band)
            getattr(details, band).update(fib[band])
    ctx_client = ContextClient()
    ctx_client.connect()
    try:
        ctx_client.SetOpticalLink(optical_link)
    except Exception as err:
        print (f"setOpticalLink {err}")    
    except Exception:
        LOGGER.exception("Cannot persist optical link %s (%s)",
                         optical_link.name, test)
        raise
    finally:
        ctx_client.close()