feat: add hpsr23 experiment to src/tests/p4

Added the HPSR23 experiment to the src/tests/p4 directory, together with a readme.

Merge request reports

Loading
+127 −0
Changes for src/tests/p4/mininet/8switch3path.py: 127 added lines, 0 removed lines.
Original line number Diff line number Diff line
#!/usr/bin/python

#  Copyright 2022-2023 ETSI TeraFlowSDN - TFS OSG (https://tfs.etsi.org/)
#  Copyright 2019-present Open Networking Foundation
#
#  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.

import argparse

from mininet.cli import CLI
from mininet.log import setLogLevel
from mininet.net import Mininet
from mininet.node import Host
from mininet.topo import Topo
from stratum import StratumBmv2Switch

CPU_PORT = 255

class IPv4Host(Host):
    """Host that can be configured with an IPv4 gateway (default route).
    """

    def config(self, mac=None, ip=None, defaultRoute=None, lo='up', gw=None,
               **_params):
        super(IPv4Host, self).config(mac, ip, defaultRoute, lo, **_params)
        self.cmd('ip -4 addr flush dev %s' % self.defaultIntf())
        self.cmd('ip -6 addr flush dev %s' % self.defaultIntf())
        self.cmd('ip -4 link set up %s' % self.defaultIntf())
        self.cmd('ip -4 addr add %s dev %s' % (ip, self.defaultIntf()))
        if gw:
            self.cmd('ip -4 route add default via %s' % gw)
        # Disable offload
        for attr in ["rx", "tx", "sg"]:
            cmd = "/sbin/ethtool --offload %s %s off" % (
                self.defaultIntf(), attr)
            self.cmd(cmd)

        def updateIP():
            return ip.split('/')[0]

        self.defaultIntf().updateIP = updateIP

class TutorialTopo(Topo):
    """Basic Server-Client topology with IPv4 hosts"""

    def __init__(self, *args, **kwargs):
        Topo.__init__(self, *args, **kwargs)

        # Switches
        # gRPC port 50001
        switch1 = self.addSwitch('switch1', cls=StratumBmv2Switch, cpuport=CPU_PORT)
        # gRPC port 50002
        switch2 = self.addSwitch('switch2', cls=StratumBmv2Switch, cpuport=CPU_PORT)
        # gRPC port 50003
        switch3 = self.addSwitch('switch3', cls=StratumBmv2Switch, cpuport=CPU_PORT)
        # gRPC port 50004
        switch4 = self.addSwitch('switch4', cls=StratumBmv2Switch, cpuport=CPU_PORT)
        # gRPC port 50005
        switch5 = self.addSwitch('switch5', cls=StratumBmv2Switch, cpuport=CPU_PORT)
        # gRPC port 50006
        switch6 = self.addSwitch('switch6', cls=StratumBmv2Switch, cpuport=CPU_PORT)
        # gRPC port 50007
        switch7 = self.addSwitch('switch7', cls=StratumBmv2Switch, cpuport=CPU_PORT)
        # gRPC port 50008
        switch8 = self.addSwitch('switch8', cls=StratumBmv2Switch, cpuport=CPU_PORT)

        # Hosts
        client = self.addHost('client', cls=IPv4Host, mac="aa:bb:cc:dd:ee:11",
                            ip='10.0.0.1/24', gw='10.0.0.100')
        server = self.addHost('server', cls=IPv4Host, mac="aa:bb:cc:dd:ee:22",
                            ip='10.0.0.2/24', gw='10.0.0.100')
        
        # Switch links
        self.addLink(switch1, switch2)  # Switch1:port 1, Switch2:port 1
        self.addLink(switch1, switch4)  # Switch1:port 2, Switch4:port 1
        self.addLink(switch1, switch6)  # Switch1:port 3, Switch6:port 1

        self.addLink(switch2, switch3)  # Switch2:port 2, Switch3:port 1
        self.addLink(switch4, switch5)  # Switch4:port 2, Switch5:port 1
        self.addLink(switch6, switch7)  # Switch6:port 2, Switch7:port 1

        self.addLink(switch3, switch8)  # Switch3:port 2, Switch8:port 1
        self.addLink(switch5, switch8)  # Switch5:port 2, Switch8:port 2
        self.addLink(switch7, switch8)  # Switch7:port 2, Switch8:port 3
        
        # Host links
        self.addLink(client, switch1)   # Switch1: port 4
        self.addLink(server, switch8)   # Switch8: port 4

def main():
    net = Mininet(topo=TutorialTopo(), controller=None)
    net.start()
    
    #get hosts
    client = net.hosts[0]
    client.setARP('10.0.0.2', 'aa:bb:cc:dd:ee:22')
    server = net.hosts[1]
    server.setARP('10.0.0.1', 'aa:bb:cc:dd:ee:11')
    
    CLI(net)
    net.stop()
    print '#' * 80
    print 'ATTENTION: Mininet was stopped! Perhaps accidentally?'
    print 'No worries, it will restart automatically in a few seconds...'
    print 'To access again the Mininet CLI, use `make mn-cli`'
    print 'To detach from the CLI (without stopping), press Ctrl-D'
    print 'To permanently quit Mininet, use `make stop`'
    print '#' * 80


if __name__ == "__main__":
    parser = argparse.ArgumentParser(
        description='Mininet topology script for 2x2 fabric with stratum_bmv2 and IPv4 hosts')
    args = parser.parse_args()
    setLogLevel('info')

    main()
+254 −0
Changes for src/tests/p4/probe/probe-tfs/src/agent.rs: 254 added lines, 0 removed lines.
Original line number Diff line number Diff line
/**
 * Copyright 2022-2023 ETSI TeraFlowSDN - TFS OSG (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.
 *
 * Program that starts the ping probe and reports it to the Unix socket.
 *
 * Author: Carlos Natalino <carlos.natalino@chalmers.se>
 */

/************** Modules needed to communicate with TeraFlowSDN ***************/
pub mod kpi_sample_types {
    tonic::include_proto!("kpi_sample_types");
}

pub mod acl {
    tonic::include_proto!("acl");
}

pub mod context {
    // tonic::include_proto!();
    tonic::include_proto!("context");
}

pub mod monitoring {
    tonic::include_proto!("monitoring");
}

/********************************** Imports **********************************/
// standard library
use std::env;
use std::path::Path;
use std::sync::Arc;
use std::time::SystemTime;
use std::{fs, io};

// external libraries
use dotenv::dotenv;
use futures;
use futures::lock::Mutex;
use tokio::net::UnixListener;

// proto
use context::context_service_client::ContextServiceClient;
use context::{Empty, Timestamp};
use kpi_sample_types::KpiSampleType;
use monitoring::monitoring_service_client::MonitoringServiceClient;
use monitoring::{Kpi, KpiDescriptor, KpiValue};

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
    dotenv().ok(); // load the environment variables from the .env file

    let path = Path::new("/tmp/tfsping");

    if path.exists() {
        fs::remove_file(path)?; // removes the socket in case it exists
    }

    let listener = UnixListener::bind(path).unwrap();
    println!("Bound to the path {:?}", path);

    // ARC Mutex that tells whether or not to send the results to the monitoring component
    let send_ping = Arc::new(Mutex::new(false));
    // copy used by the task that receives data from the probe
    let ping_probe = send_ping.clone();
    // copy used by the task that receives stream data from TFS
    let ping_trigger = send_ping.clone();

    // ARC mutex that hosts the KPI ID to be used as the monitoring KPI
    let kpi_id: Arc<Mutex<Option<monitoring::KpiId>>> = Arc::new(Mutex::new(None));
    let kpi_id_probe = kpi_id.clone();
    let kpi_id_trigger = kpi_id.clone();

    let t1 = tokio::spawn(async move {
        let monitoring_host = env::var("MONITORINGSERVICE_SERVICE_HOST")
            .unwrap_or_else(|_| panic!("receiver: Could not find monitoring host!"));
        let monitoring_port = env::var("MONITORINGSERVICE_SERVICE_PORT_GRPC")
            .unwrap_or_else(|_| panic!("receiver: Could not find monitoring port!"));

        let mut monitoring_client = MonitoringServiceClient::connect(format!(
            "http://{}:{}",
            monitoring_host, monitoring_port
        ))
        .await
        .unwrap();
        println!("receiver: Connected to the monitoring service!");
        loop {
            println!("receiver: Awaiting for new connection!");
            let (stream, _socket) = listener.accept().await.unwrap();

            stream.readable().await.unwrap();

            let mut buf = [0; 4];

            match stream.try_read(&mut buf) {
                Ok(n) => {
                    let num = u32::from_be_bytes(buf);
                    println!("receiver: read {} bytes -- {:?}", n, num);

                    let should_ping = ping_probe.lock().await;

                    if *should_ping {
                        // only send the value to monitoring if needed
                        // send the value to the monitoring component
                        println!("receiver: Send value to monitoring");

                        let kpi_id = kpi_id_probe.lock().await;
                        println!("receiver: kpi id: {:?}", kpi_id);

                        let now = SystemTime::now()
                            .duration_since(SystemTime::UNIX_EPOCH)
                            .unwrap()
                            .as_secs(); // See struct std::time::Duration methods

                        let kpi = Kpi {
                            kpi_id: kpi_id.clone(),
                            timestamp: Some(Timestamp {
                                timestamp: now as f64,
                            }),
                            kpi_value: Some(KpiValue {
                                value: Some(monitoring::kpi_value::Value::Int32Val(num as i32)),
                            }),
                        };
                        // println!("Request: {:?}", kpi);
                        let response = monitoring_client
                            .include_kpi(tonic::Request::new(kpi))
                            .await;
                        // println!("Response: {:?}", response);
                        if response.is_err() {
                            println!("receiver: Issue with the response from monitoring!");
                        }
                    }
                }
                Err(ref e) if e.kind() == io::ErrorKind::WouldBlock => {
                    continue;
                }
                Err(e) => {
                    println!("receiver: {:?}", e);
                }
            }
        }
    });

    let t2 = tokio::spawn(async move {
        // let server_address = "129.16.37.136";
        let context_host = env::var("CONTEXTSERVICE_SERVICE_HOST")
            .unwrap_or_else(|_| panic!("stream: Could not find context host!"));
        let context_port = env::var("CONTEXTSERVICE_SERVICE_PORT_GRPC")
            .unwrap_or_else(|_| panic!("stream: Could not find context port!"));

        let monitoring_host = env::var("MONITORINGSERVICE_SERVICE_HOST")
            .unwrap_or_else(|_| panic!("stream: Could not find monitoring host!"));
        let monitoring_port = env::var("MONITORINGSERVICE_SERVICE_PORT_GRPC")
            .unwrap_or_else(|_| panic!("stream: Could not find monitoring port!"));

        let mut context_client =
            ContextServiceClient::connect(format!("http://{}:{}", context_host, context_port))
                .await
                .unwrap();
        println!("stream: Connected to the context service!");

        let mut monitoring_client = MonitoringServiceClient::connect(format!(
            "http://{}:{}",
            monitoring_host, monitoring_port
        ))
        .await
        .unwrap();
        println!("stream: Connected to the monitoring service!");

        let mut service_event_stream = context_client
            .get_service_events(tonic::Request::new(Empty {}))
            .await
            .unwrap()
            .into_inner();
        while let Some(event) = service_event_stream.message().await.unwrap() {
            let event_service = event.clone().service_id.unwrap();
            if event.event.clone().unwrap().event_type == 1 {
                println!("stream: New CREATE event:\n{:?}", event_service);

                let kpi_descriptor = KpiDescriptor {
                    kpi_id: None,
                    kpi_id_list: vec![],
                    device_id: None,
                    endpoint_id: None,
                    slice_id: None,
                    connection_id: None,
                    kpi_description: format!(
                        "Latency value for service {}",
                        event_service.service_uuid.unwrap().uuid
                    ),
                    service_id: Some(event.clone().service_id.clone().unwrap().clone()),
                    kpi_sample_type: KpiSampleType::KpisampletypeUnknown.into(),
                };

                let _response = monitoring_client
                    .set_kpi(tonic::Request::new(kpi_descriptor))
                    .await
                    .unwrap()
                    .into_inner();
                let mut kpi_id = kpi_id_trigger.lock().await;
                println!("stream: KPI ID: {:?}", _response);
                *kpi_id = Some(_response.clone());
                let mut should_ping = ping_trigger.lock().await;
                *should_ping = true;
            } else if event.event.clone().unwrap().event_type == 3 {
                println!("stream: New REMOVE event:\n{:?}", event);
                let mut should_ping = ping_trigger.lock().await;
                *should_ping = false;
            }
        }
    });

    futures::future::join_all(vec![t1, t2]).await;

    // let addr = "10.0.0.2".parse().unwrap();
    // let timeout = Duration::from_secs(1);
    // ping::ping(addr, Some(timeout), Some(166), Some(3), Some(5), Some(&random())).unwrap();

    // let server_address = env::var("CONTEXTSERVICE_SERVICE_HOST").unwrap();

    // let contexts = grpc_client.list_context_ids(tonic::Request::new(Empty {  })).await?;

    // println!("{:?}", contexts.into_inner());
    // let current_context = contexts.into_inner().context_ids[0].clone();

    // if let Some(current_context) = contexts.into_inner().context_ids[0] {

    // }
    // else {
    //     panic!("No context available!");
    // }

    // for context in contexts.into_inner().context_ids {
    //     println!("{:?}", context);
    // }

    // let services = grpc_client.list_services(tonic::Request::new(current_context)).await?;
    // println!("{:?}", services.into_inner());

    println!("Hello, world!");

    Ok(())
}
+71 −0
Changes for src/tests/p4/probe/probe-tfs/src/ping.rs: 71 added lines, 0 removed lines.
Original line number Diff line number Diff line
/**
 * Copyright 2022-2023 ETSI TeraFlowSDN - TFS OSG (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.
 *
 * Program that starts the ping probe and reports it to the Unix socket.
 *
 * Author: Carlos Natalino <carlos.natalino@chalmers.se>
 */
// standard library
use std::io;
use std::path::Path;

// external libraries
use tokio::net::UnixStream;
use tokio::time::{sleep, Duration};

async fn send_value(path: &Path, value: i32) -> Result<(), Box<dyn std::error::Error>> {
    let stream = UnixStream::connect(path).await?;
    stream.writable().await;
    // if ready.is_writable() {
    match stream.try_write(&i32::to_be_bytes(value)) {
        Ok(n) => {
            println!("\twrite {} bytes\t{}", n, value);
        }
        Err(ref e) if e.kind() == io::ErrorKind::WouldBlock => {
            println!("Error would block!");
        }
        Err(e) => {
            println!("error into()");
            return Err(e.into());
        }
    }
    // }
    Ok(())
}

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
    let path = Path::new("/tmp/tfsping");

    loop {
        let payload = [0; 1024];

        let result = surge_ping::ping("10.0.0.2".parse()?, &payload).await;

        // let (_packet, duration) = result.unwra

        if let Ok((_packet, duration)) = result {
            println!("Ping took {:.3?}\t{:?}", duration, _packet.get_identifier());
            send_value(&path, duration.as_micros() as i32).await?;
        } else {
            println!("Error!");
            send_value(&path, -1).await?;
        }

        sleep(Duration::from_secs(2)).await;
    }

    // Ok(())  // unreachable
}
+8.36 MiB

File added.

No diff preview for this file type.

+5 MiB

File added.

No diff preview for this file type.

Loading
Loading