Commit 2b1637ab authored by Lluis Gifre Renom's avatar Lluis Gifre Renom
Browse files

test: add live optical service and spectrum validation probe

parent ac973d4d
Loading
Loading
Loading
Loading
+29 −0
Changes for src/tests/ecoc26-hyr-agentic-optical/README.md: 29 added lines, 0 removed lines.
Original line number Diff line number Diff line
@@ -131,3 +131,32 @@ after teardown, following the checks in `src/opticalcontroller/README.md`.
Successful subsequent provisioning or an empty service inventory alone does
not prove spectrum synchronization. Preserve allocations belonging to other
services, and reconcile stale occupancy from earlier deployments separately.

### Automated MCP And Spectrum Checks

`check_spectrum.py` validates the controller without using an LLM. It requires
the Python `mcp` package, the supplied topology, and an Optical Controller URL
reachable from the test host. Service inventory must be empty and no active
reservations may exist; released audit records are allowed. The test
initializes Optical Controller's topology only in that state.
It preserves the existing spectrum baseline rather than resetting it.

```bash
python src/tests/ecoc26-hyr-agentic-optical/check_spectrum.py \
  --optical-url http://<optical-controller-host>:10060/OpticalTFS \
  --capacity 800 --output /tmp/ecoc-spectrum.json
```

The probe creates and deletes one unidirectional service, verifies slot reuse,
then creates two forward services sharing a ROADM link and a reverse service.
It removes them individually, checking retained ACTIVE services and disjoint
allocations after every step. Context and Optical Controller spectrum maps
must agree, and final occupancy must match the baseline. Cleanup targets only
the probe's uniquely named services. The JSON output stores per-step spectrum
snapshots and allocation details; keep live outputs outside the source tree.

Reservation creation, overlap/occupied-slot rejection, and release can be
checked with `src/tests/spectrum_negotiation/live_reservation_validation.py`.
That separate test retains a RELEASED reservation record. The reservation
DELETE endpoint also performs release, rather than deleting the audit record.
Check reservation status, not just list length, when verifying cleanup.
+234 −0
Changes for src/tests/ecoc26-hyr-agentic-optical/check_spectrum.py: 234 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.

"""Live ECOC service/spectrum validation through MCP; no LLM is needed."""

import argparse
import asyncio
import copy
import json
import logging
import urllib.request
import uuid
from pathlib import Path

from mcp import ClientSession
from mcp.client.sse import sse_client

LOGGER = logging.getLogger(__name__)
BANDS = ("c_slots", "l_slots", "s_slots")


def get_json(url):
    with urllib.request.urlopen(url, timeout=30) as reply:
        return json.load(reply)


class Probe:
    def __init__(self, args, session):
        self.args = args
        self.session = session
        self.owned = set()
        self.active = {}
        self.records = []
        self.baseline = None
        self.prefix = "spectrum-check-" + uuid.uuid4().hex[:8]

    async def call(self, name, **arguments):
        reply = await self.session.call_tool(name, arguments)
        if reply.isError:
            raise RuntimeError(f"{name}: {reply.content}")
        return json.loads(reply.content[0].text)

    async def services(self):
        reply = await self.call("tfs_list_services", context_uuid="admin")
        return {s["service_id"]["service_uuid"]["uuid"]: s
                for s in reply["services"]}

    async def snapshot(self, label):
        ctx = await self.call("tfs_list_optical_links")
        optical = get_json(self.args.optical_url + "/GetLinks")
        optical = {l["name"]: l["optical_details"]
                   for l in optical["optical_links"]}
        current = {}
        for link in ctx["optical_links"]:
            name = link["name"]
            current[name] = {}
            for band in BANDS:
                actual = optical[name].get(band, {})
                persisted = link["optical_details"].get(band, {})
                # Match the controller's sparse-map convention.
                normalized = {k: persisted.get(k, 1) for k in actual}
                assert normalized == actual, (label, name, band, "sync")
                current[name][band] = normalized
        if self.baseline is not None:
            expected = copy.deepcopy(self.baseline)
            for allocation in self.active.values():
                for link in allocation["links"]:
                    for slot in allocation["slots"]:
                        expected[link][allocation["band"]][str(slot)] = 0
            assert current == expected, (label, "unexpected spectrum change")
        services = await self.services()
        for sid in self.active:
            assert services[sid]["service_status"]["service_status"] == (
                "SERVICESTATUS_ACTIVE"
            ), (label, sid)
        self.records.append({"step": label, "active": copy.deepcopy(self.active),
                             "spectrum": current})
        self.args.output.write_text(json.dumps(self.records, indent=2))
        LOGGER.info("%s: %d ACTIVE services; spectrum synchronized",
                    label, len(self.active))
        return current

    async def create(self, suffix, source, destination):
        name = self.prefix + "-" + suffix
        sid = {"context_id": {"context_uuid": {"uuid": "admin"}},
               "service_uuid": {"uuid": name}}
        shell = {"service_id": sid, "name": name,
                 "service_type": "SERVICETYPE_OPTICAL_CONNECTIVITY"}
        # Track the name before submission so even a partial create is cleaned.
        self.owned.add(name)
        await self.call("tfs_create_services", context_uuid="admin",
                        services=[shell])
        endpoints = []
        for device in (source, destination):
            data = self.devices[device]
            assert len(data["device_endpoints"]) == 1, device
            endpoints.append({"device_id": data["device_id"],
                              "endpoint_uuid": data["device_endpoints"][0][
                                  "endpoint_id"]["endpoint_uuid"]})
        payload = dict(shell, service_endpoint_ids=endpoints,
                       service_config={"config_rules": []},
                       service_constraints=[
                           {"action": "CONSTRAINTACTION_SET", "custom": {
                               "constraint_type": key,
                               "constraint_value": value}}
                           for key, value in (
                               ("type", "flexi_grid"),
                               ("bidirectionality", "0"),
                               ("bandwidth[gbps]", str(self.args.capacity)),
                           )])
        await self.call("tfs_update_service", context_uuid="admin",
                        service_uuid=name, service=payload)
        service = await self.call("tfs_get_service", context_uuid="admin",
                                  service_uuid=name)
        canonical = service["service_id"]["service_uuid"]["uuid"]
        self.owned.discard(name)
        self.owned.add(canonical)
        assert service["service_status"]["service_status"] == (
            "SERVICESTATUS_ACTIVE"
        ), service
        settings = next(json.loads(r["custom"]["resource_value"])
                        for r in service["service_config"]["config_rules"]
                        if r.get("custom", {}).get("resource_key") == "/settings")
        band = {"C_BAND": "c_slots", "L_BAND": "l_slots",
                "S_BAND": "s_slots"}[settings["band_type"]]
        for link in settings["links"]:
            for slot in settings["slots"]:
                assert self.baseline[link][band][str(slot)] == 1, (
                    "Allocated a baseline-occupied slot", link, band, slot
                )
        for other in self.active.values():
            if other["band"] == band and (
                set(other["links"]) & set(settings["links"])
            ):
                assert not set(other["slots"]) & set(settings["slots"]), (
                    "Overlapping active allocations", name
                )
        self.active[canonical] = {"links": settings["links"],
                                  "slots": settings["slots"], "band": band}
        await self.snapshot("create-" + suffix)
        return canonical

    async def delete(self, sid):
        for attempt in range(3):
            services = await self.services()
            target = next((key for key, s in services.items()
                           if key == sid or s.get("name") == sid), None)
            if target is None:
                self.owned.discard(sid)
                self.active.pop(sid, None)
                return
            await self.call("tfs_delete_service", context_uuid="admin",
                            service_uuid=target)
            await asyncio.sleep(1)
        remaining = await self.services()
        assert not any(key == sid or s.get("name") == sid
                       for key, s in remaining.items()), ("delete failed", sid)
        self.owned.discard(sid)
        self.active.pop(sid, None)

    async def run(self):
        assert not await self.services(), "Requires an empty service inventory"
        reservations = await self.call("tfs_list_optical_spectrum_reservations",
                                       context_uuid="admin")
        active_reservations = [r for r in reservations["reservations"]
                               if r.get("status", "").split("_")[-1]
                               not in {"RELEASED", "EXPIRED"}]
        assert not active_reservations, "Active reservations must be absent"
        devices = await self.call("tfs_list_devices")
        self.devices = {d["name"]: d for d in devices["devices"]}
        # Initialize only while no services exist; never reload active flows.
        get_json(self.args.optical_url + "/GetTopology/admin/admin")
        self.baseline = await self.snapshot("baseline")
        first = await self.create("uni", "T1.1", "T2.1")
        first_allocation = copy.deepcopy(self.active[first])
        await self.delete(first)
        await self.snapshot("delete-uni")
        a = await self.create("forward-a", "T1.1", "T2.1")
        assert self.active[a] == first_allocation, "Released slots not reused"
        b = await self.create("forward-b", "T1.2", "T2.2")
        c = await self.create("reverse-b", "T2.2", "T1.2")
        await self.delete(a)
        await self.snapshot("delete-a-preserve-both-b-directions")
        await self.delete(b)
        await self.snapshot("delete-forward-b-preserve-reverse-b")
        await self.delete(c)
        await self.snapshot("all-deleted")
        assert not await self.services()
        connections = await self.call("tfs_list_all_connections",
                                      context_uuid="admin")
        assert not connections["connections"]


async def main():
    parser = argparse.ArgumentParser(description=__doc__)
    parser.add_argument("--mcp-url", default="http://127.0.0.1/mcp/sse")
    parser.add_argument("--optical-url", required=True,
                        help="Reachable Optical Controller /OpticalTFS URL")
    parser.add_argument("--capacity", type=int, default=800)
    parser.add_argument("--output", type=Path, required=True)
    args = parser.parse_args()
    args.output.parent.mkdir(parents=True, exist_ok=True)
    logging.basicConfig(level=logging.INFO)
    logging.getLogger("httpx").setLevel(logging.WARNING)
    async with sse_client(args.mcp_url, sse_read_timeout=300) as streams:
        async with ClientSession(*streams) as session:
            await session.initialize()
            probe = Probe(args, session)
            try:
                await probe.run()
            finally:
                failures = []
                for sid in list(probe.owned):
                    try:
                        await probe.delete(sid)
                    except Exception as exc:
                        failures.append((sid, str(exc)))
                if failures:
                    raise RuntimeError(f"Test cleanup failed: {failures}")


if __name__ == "__main__":
    asyncio.run(main())