Commit 89bbbea2 authored by Waleed Akbar's avatar Waleed Akbar
Browse files

Add telemetry subscription management scripts for MWC 2026 scenario

- Introduced run_telemetry-delete.sh to remove telemetry subscriptions.
- Enhanced telemetry-delete-slice1.py to manage subscription records.
- Created telemetry_subscriptions_store.py for on-disk subscription tracking.
- Updated telemetry-subscribe-slice1.py to record established subscriptions.
parent 59437ffe
Loading
Loading
Loading
Loading
+51 −0
Changes for src/tests/mwc26-f5ga/run_telemetry-delete.sh: 51 added lines, 0 removed lines.
Original line number Diff line number Diff line
#!/bin/bash
# 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.

# Counterpart of run_telemetry-subscribe.sh: deletes the telemetry subscriptions this
# controller established. Stopping telemetry-subscribe-slice1.py with Ctrl+C only closes
# the SSE stream; the subscription stays registered and the SIMAP connector keeps its
# collector/aggregator workers polling and publishing. Run this to actually remove them.
#
# Any extra arguments are passed through as explicit subscription ids, for subscriptions
# created before the local registry existed, e.g.:
#   ./run_telemetry-delete.sh 1210199060928987137 1210199508167884801

# Set working directory
cd "$(dirname "$0")" || exit 1

# Get the current hostname
HOSTNAME=$(hostname)
echo "Deleting telemetry subscriptions for ${HOSTNAME}..."

case "$HOSTNAME" in
    tfs-e2e-ctrl)
        echo "Unsubscribing from E2E Controller telemetry..."
        python3 telemetry-delete-slice1.py e2e E2E-L1 "$@"
        ;;
    tfs-agg-ctrl)
        echo "Unsubscribing from Aggregation Controller telemetry..."
        python3 telemetry-delete-slice1.py agg AggNet-L1 "$@"
        ;;
    tfs-ip-ctrl)
        echo "Unsubscribing from IP Controller telemetry..."
        python3 telemetry-delete-slice1.py trans-pkt Trans-L1 "$@"
        ;;
    *)
        echo "Unknown host: $HOSTNAME"
        echo "Usage: $0 [subscription_id ...]"
        echo "  This script must be run on tfs-e2e-ctrl, tfs-agg-ctrl, or tfs-ip-ctrl"
        exit 1
        ;;
esac
+114 −12
Changes for src/tests/mwc26-f5ga/telemetry-delete-slice1.py: 114 added lines, 12 removed lines.
Original line number Diff line number Diff line
@@ -13,34 +13,136 @@
# limitations under the License.


import sys
import requests
from requests.auth import HTTPBasicAuth
import telemetry_subscriptions_store as store


USAGE = '\n'.join([
    'Usage: {0:s} <network_name> <link_name> [--force] [subscription_id ...]',
    '',
    'Deletes the telemetry subscriptions established for <network_name>/<link_name>.',
    'Without explicit ids, every subscription recorded by telemetry-subscribe-slice1.py',
    'for that network/link is deleted. Pass ids explicitly to remove subscriptions that',
    'predate the local registry (for instance, ids read from the SIMAP connector logs).',
    '',
    '  --force  Stop tracking an id even if the controller rejected the deletion.',
    '           Useful because the NBI reports an unknown subscription as a generic',
    '           HTTP 500, so an already-deleted id cannot be told apart from a real',
    '           failure and would otherwise stay in the registry forever.',
    '',
    'Examples:',
    '  {0:s} trans-pkt Trans-L1',
    '  {0:s} trans-pkt Trans-L1 1210199060928987137 1210199508167884801',
])

ARGUMENTS = [argument for argument in sys.argv[1:] if argument != '--force']
FORCE     = '--force' in sys.argv[1:]

if len(ARGUMENTS) < 2:
    print(USAGE.format(sys.argv[0]))
    sys.exit(1)

RESTCONF_ADDRESS  = '127.0.0.1'
RESTCONF_PORT     = 80
TELEMETRY_ID = 1109405947767160833
TARGET_SIMAP_NAME = ARGUMENTS[0]
TARGET_LINK_NAME  = ARGUMENTS[1]

try:
    EXPLICIT_IDS = [int(argument) for argument in ARGUMENTS[2:]]
except ValueError as e:
    print('Subscription ids must be integers: {:s}'.format(str(e)))
    sys.exit(1)

UNSUBSCRIBE_URI = '/restconf/operations/subscriptions:delete-subscription'
UNSUBSCRIBE_URL = 'http://{:s}:{:d}{:s}'.format(RESTCONF_ADDRESS, RESTCONF_PORT, UNSUBSCRIBE_URI)
REQUEST = {
    'ietf-subscribed-notifications:input': {
        'id': TELEMETRY_ID,
    }
}


def main() -> None:
    print('[E2E] Delete Telemetry slice1...')
def delete_subscription(subscription_id : int) -> bool:
    """Delete one subscription. Returns True if it is gone (deleted or already absent)."""
    headers = {'accept': 'application/json'}
    auth    = HTTPBasicAuth('admin', 'admin')
    print(UNSUBSCRIBE_URL)
    print(REQUEST)
    request = {'ietf-subscribed-notifications:input': {'id': subscription_id}}

    print('  - Deleting subscription {:d}...'.format(subscription_id))
    try:
        reply = requests.post(
        UNSUBSCRIBE_URL, headers=headers, json=REQUEST, auth=auth,
            UNSUBSCRIBE_URL, headers=headers, json=request, auth=auth,
            verify=False, allow_redirects=True, timeout=30
        )
    reply.raise_for_status()
    except requests.RequestException as e:
        print('    FAILED: {:s}'.format(str(e)))
        return False

    if reply.ok:
        print('    OK')
        return True

    body = reply.content.decode('UTF-8', errors='replace')

    # The connector raises NotFoundException for an unknown subscription. That is the
    # desired end state, so treat it as success and stop tracking the id.
    if 'not found' in body.lower() or reply.status_code == 404:
        print('    Already gone (HTTP {:d})'.format(reply.status_code))
        return True

    print('    FAILED: HTTP {:d}: {:s}'.format(reply.status_code, body.strip()))
    if reply.status_code == 500:
        # The NBI does not propagate the gRPC NOT_FOUND detail, it answers a bare
        # "Internal Server Error", so an already-deleted subscription lands here too.
        # Check the connector log to tell the two cases apart:
        #   kubectl logs -n tfs deploy/nbiservice | grep -i subscription
        print('    HINT: HTTP 500 is also what an already-deleted subscription returns.')
        print('          Check the nbiservice log, or re-run with --force to stop tracking it.')
    return False


def main() -> None:
    print('[{:s}] Delete Telemetry subscriptions for link {:s}...'.format(
        TARGET_SIMAP_NAME.upper(), TARGET_LINK_NAME
    ))

    if len(EXPLICIT_IDS) > 0:
        subscription_ids = EXPLICIT_IDS
        print('Deleting {:d} subscription(s) given on the command line.'.format(
            len(subscription_ids)
        ))
    else:
        subscription_ids = store.list_ids(TARGET_SIMAP_NAME, TARGET_LINK_NAME)
        if len(subscription_ids) == 0:
            print('No subscriptions recorded for {:s}/{:s} in {:s}.'.format(
                TARGET_SIMAP_NAME, TARGET_LINK_NAME, store.get_store_path()
            ))
            print('Nothing to do. Pass the subscription ids explicitly if you know them.')
            return
        print('Found {:d} recorded subscription(s) in {:s}.'.format(
            len(subscription_ids), store.get_store_path()
        ))

    print(UNSUBSCRIBE_URL)

    failed = list()
    for subscription_id in subscription_ids:
        if delete_subscription(subscription_id):
            store.remove(TARGET_SIMAP_NAME, TARGET_LINK_NAME, subscription_id)
        else:
            failed.append(subscription_id)
            if FORCE:
                print('    Dropping {:d} from the registry anyway (--force).'.format(
                    subscription_id
                ))
                store.remove(TARGET_SIMAP_NAME, TARGET_LINK_NAME, subscription_id)

    deleted = len(subscription_ids) - len(failed)
    print('Deleted {:d}/{:d} subscription(s).'.format(deleted, len(subscription_ids)))

    if len(failed) > 0:
        print('Failed subscription ids: {:s}'.format(
            ', '.join(str(subscription_id) for subscription_id in failed)
        ))
        sys.exit(1)


if __name__ == '__main__':
    main()
+14 −0
Changes for src/tests/mwc26-f5ga/telemetry-subscribe-slice1.py: 14 added lines, 0 removed lines.
Original line number Diff line number Diff line
@@ -16,6 +16,7 @@
import sys
import requests
from requests.auth import HTTPBasicAuth
import telemetry_subscriptions_store as store


if len(sys.argv) < 3:
@@ -67,8 +68,21 @@ def main() -> None:
        raise Exception('Unexpected Reply: {:s}'.format(str(reply_data)))
    subscription_uri = reply_data['uri']

    if 'id' not in reply_data:
        raise Exception('Unexpected Reply: {:s}'.format(str(reply_data)))
    subscription_id = reply_data['id']

    # Record the subscription so run_telemetry-delete.sh can remove it later. Ctrl+C below
    # only closes the SSE stream; the subscription itself stays registered and the SIMAP
    # connector keeps polling and publishing until it is explicitly deleted.
    store.add(TARGET_SIMAP_NAME, TARGET_LINK_NAME, subscription_id, subscription_uri)
    print('Established subscription {:d} (recorded in {:s})'.format(
        subscription_id, store.get_store_path()
    ))

    stream_url = 'http://{:s}:{:d}{:s}'.format(RESTCONF_ADDRESS, RESTCONF_PORT, subscription_uri)
    print('Opening stream "{:s}" (press Ctrl+C to stop)...'.format(stream_url))
    print('NOTE: Ctrl+C does NOT delete the subscription; run ./run_telemetry-delete.sh')

    with requests.get(stream_url, stream=True, auth=auth) as resp:
        for i, line in enumerate(resp.iter_lines(decode_unicode=True), 1):
+105 −0
Changes for src/tests/mwc26-f5ga/telemetry_subscriptions_store.py: 105 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.

"""Tiny on-disk registry of the telemetry subscriptions established by this scenario.

The RESTCONF NBI exposes only establish-subscription and delete-subscription; there is
no endpoint to list the subscriptions that are currently active. Without a local record,
a subscription established by telemetry-subscribe-slice1.py can only be removed if its
id was copied by hand from the console, which is why stale subscriptions accumulate in
the SIMAP connector and keep its collector/aggregator workers producing telemetry long
after the TFS services are gone.

This module keeps that record in tmp/telemetry-subscriptions.json so that
telemetry-delete-slice1.py can clean up every subscription the scenario created.
"""

import datetime, json, os
from typing import Dict, List


SCRIPT_DIR = os.path.dirname(os.path.abspath(__file__))
STORE_PATH = os.path.join(SCRIPT_DIR, 'tmp', 'telemetry-subscriptions.json')


def get_store_path() -> str:
    return STORE_PATH


def make_key(network_name : str, link_name : str) -> str:
    return '{:s}:{:s}'.format(network_name, link_name)


def load() -> Dict[str, List[Dict]]:
    if not os.path.exists(STORE_PATH): return dict()
    try:
        with open(STORE_PATH, 'r', encoding='UTF-8') as file:
            data = json.load(file)
    except (OSError, ValueError) as e:
        print('WARNING: unable to read {:s} ({:s}), starting from an empty store'.format(
            STORE_PATH, str(e)
        ))
        return dict()
    if not isinstance(data, dict): return dict()
    return data


def save(data : Dict[str, List[Dict]]) -> None:
    os.makedirs(os.path.dirname(STORE_PATH), exist_ok=True)
    # Write through a temporary file so an interrupted run cannot truncate the store.
    tmp_path = STORE_PATH + '.tmp'
    with open(tmp_path, 'w', encoding='UTF-8') as file:
        json.dump(data, file, indent=4, sort_keys=True)
        file.write('\n')
    os.replace(tmp_path, STORE_PATH)


def add(network_name : str, link_name : str, subscription_id : int, subscription_uri : str) -> None:
    """Record a newly established subscription."""
    data = load()
    key = make_key(network_name, link_name)
    entries = data.setdefault(key, list())

    # Never record the same id twice; refresh the existing entry instead.
    entries = [entry for entry in entries if entry.get('id') != subscription_id]
    entries.append({
        'id'         : subscription_id,
        'uri'        : subscription_uri,
        'network'    : network_name,
        'link'       : link_name,
        'created_at' : datetime.datetime.now(datetime.timezone.utc).isoformat(),
    })

    data[key] = entries
    save(data)


def list_ids(network_name : str, link_name : str) -> List[int]:
    """Return the recorded subscription ids for a network/link, in the order recorded."""
    entries = load().get(make_key(network_name, link_name), list())
    return [entry['id'] for entry in entries if 'id' in entry]


def remove(network_name : str, link_name : str, subscription_id : int) -> None:
    """Forget a subscription that is no longer active."""
    data = load()
    key = make_key(network_name, link_name)
    if key not in data: return

    entries = [entry for entry in data[key] if entry.get('id') != subscription_id]
    if len(entries) == 0:
        data.pop(key, None)
    else:
        data[key] = entries
    save(data)