diff --git a/.github/workflows/tests.yaml b/.github/workflows/tests.yaml index 18e99e1..ccc6f55 100644 --- a/.github/workflows/tests.yaml +++ b/.github/workflows/tests.yaml @@ -24,9 +24,13 @@ jobs: tests: runs-on: ubuntu-24.04 strategy: - fail-fast: true + fail-fast: false matrix: - python-version: ["3.8"] + include: + - python-version: "3.8" + toxenv: py38 + - python-version: "3.14" + toxenv: py314 steps: - uses: actions/checkout@v3 - uses: actions/setup-python@v5 @@ -40,7 +44,5 @@ jobs: run: pip install tox - name: Unit tests run: | - tox -e ${{ matrix.python-version }} - # - name: Functional tests - # run: | - # tox -e ${{ matrix.python-version }}-functional + tox -e ${{ matrix.toxenv }} + diff --git a/evpn_connector/bgp/client.py b/evpn_connector/bgp/client.py index dc32ba1..84e9d4b 100644 --- a/evpn_connector/bgp/client.py +++ b/evpn_connector/bgp/client.py @@ -762,6 +762,44 @@ def list_peer(self, peer_address): resp = self._extract_response_from_grpc(list_peers_channel) return resp + @staticmethod + def _peer_converged(peer, settle_sec): + if peer.state.session_state != gobgp_pb2.PeerState.ESTABLISHED: + return False + if settle_sec: + uptime = peer.timers.state.uptime.seconds + if not uptime or time.time() - uptime < settle_sec: + return False + return all( + afi_safi.mp_graceful_restart.state.end_of_rib_received + for afi_safi in peer.afi_safis + if afi_safi.mp_graceful_restart.state.received + ) + + def count_peers(self, settle_sec=0): + """Return (total, converged) counts of configured BGP peers. + + Used by fail-static: a peer that is configured but has not + converged means the local RIB does not reflect the true remote + topology (e.g. the route reflector is gone), so flows derived + from it must not be trusted for deletion. + + Converged is stricter than ESTABLISHED, which only says the + session is up: the peer streams its table afterwards, and a RIB + read in that window is not the fabric yet. End-of-RIB (RFC 4724) + is the signal that it is over, reported per family but only + where the graceful-restart capability was negotiated. Without it + there is nothing to wait for but settle_sec, which also covers a + session that came up too recently to be trusted. + """ + peers = self.list_peer("") + if not peers: + return 0, 0 + converged = sum( + 1 for resp in peers if self._peer_converged(resp.peer, settle_sec) + ) + return len(peers), converged + def reset_peer(self, address, soft, direction): self.stub.ResetPeer( gobgp_pb2.ResetPeerRequest( diff --git a/evpn_connector/cmd/evpn.py b/evpn_connector/cmd/evpn.py index 68b0b0d..a082676 100644 --- a/evpn_connector/cmd/evpn.py +++ b/evpn_connector/cmd/evpn.py @@ -30,6 +30,7 @@ from evpn_connector.common import sentry from evpn_connector.ovs import client as ovs_client from evpn_connector.service import evpn +from evpn_connector.service import objects as evpnobj OBSENDER_APP_NAME = constants.GLOBAL_SERVICE_NAME @@ -70,6 +71,9 @@ def main(): router_mac_type5=CONF.gobgp.router_mac_type5, ) + # Read once, so the tunnel and the flows cannot disagree. + evpnobj.set_gbp(CONF.ovs.gbp) + # Init ovs client shell_ovs_client = ovs_client.OvSClient( sw_name=CONF.ovs.switch_name, @@ -77,6 +81,7 @@ def main(): enable_sudo=CONF.ovs.enable_sudo, ovsvsctl_bin=CONF.ovs.ovs_vsctl_bin_path, ovsofctl_bin=CONF.ovs.ovs_ofctl_bin_path, + gbp=CONF.ovs.gbp, ) # Start service @@ -87,6 +92,8 @@ def main(): vxlan_udp_port=CONF.ovs.vxlan_udp_port, as_number=CONF.gobgp.as_number, policy_enabled=CONF.gobgp.policy_enabled, + fail_static=CONF.gobgp.fail_static, + fail_static_min_peers=CONF.gobgp.fail_static_min_peers, configs_dir=CONF.daemon.configs_dir, router_mac_type5=CONF.gobgp.router_mac_type5, anycast_status_file=CONF.anycast.anycast_status_file, diff --git a/evpn_connector/common/conf_opts.py b/evpn_connector/common/conf_opts.py index 4c022d1..fdfc41e 100644 --- a/evpn_connector/common/conf_opts.py +++ b/evpn_connector/common/conf_opts.py @@ -87,12 +87,42 @@ default=constants.TYPE_5_DEFAULT_ROUTER_MAC_EXTENDED, help="Value of RouterMacExtended Ext Communities Attr for Type5", ), + cfg.BoolOpt( + name="fail_static", + default=False, + help=( + "Retain the last-known-good OVS flows while too few BGP " + "peers have converged (e.g. the route reflector is lost), " + "instead of deleting flows of routes gone from the RIB. " + "Additions and changes still land; deletions wait." + ), + ), + cfg.IntOpt( + name="fail_static_min_peers", + default=0, + min=0, + help=( + "Converged BGP peers (established, table received) needed " + "to trust the RIB for flow deletion. 0 or more than " + "configured means all; 1 suits redundant route reflectors." + ), + ), ] ovs_opts = [ cfg.StrOpt( name="switch_name", required=True, help="OpenvSwitch switch name" ), + cfg.BoolOpt( + name="gbp", + default=False, + help=( + "Carry the sender's group in VXLAN-GBP's Group Policy ID, " + "to and from skb mark bits 0..15 (reserved for it); the " + "underlay is trusted. Fabric-wide: OVS will not mix GBP and " + "non-GBP tunnels on one UDP port." + ), + ), cfg.StrOpt( name="tmp_flow_file_path", default="/tmp/evpn_tmp_flow_file", diff --git a/evpn_connector/common/constants.py b/evpn_connector/common/constants.py index 1f76a97..da75e11 100644 --- a/evpn_connector/common/constants.py +++ b/evpn_connector/common/constants.py @@ -51,6 +51,11 @@ # Traffic direction in REG1 REG_FROM_REMOTE = 0 REG_FROM_LOCAL = 1 + +# A sender's group: the GBP id on the wire, the low 16 bits of the skb mark on +# the host (a tunnel field does not survive a patch port). +GBP_TO_MARK = "move:NXM_NX_TUN_GBP_ID[]->NXM_NX_PKT_MARK[0..15]" +MARK_TO_GBP = "move:NXM_NX_PKT_MARK[0..15]->NXM_NX_TUN_GBP_ID[]" # ECMP Multipath hash algorithm (for details see man 7 ovs-actions: multipath) ECMP_HASH_ALGORITHM = "symmetric_l3l4+udp" diff --git a/evpn_connector/ovs/client.py b/evpn_connector/ovs/client.py index 105deea..63ba71c 100644 --- a/evpn_connector/ovs/client.py +++ b/evpn_connector/ovs/client.py @@ -40,8 +40,10 @@ def __init__( enable_sudo=False, ovsvsctl_bin=constants.OVSVSCTL_BIN, ovsofctl_bin=constants.OVSOFCTL_BIN, + gbp=False, ): self.vxlan_ofport = vxlan_ofport or constants.VXLAN_PORT_OFPORT + self.gbp = gbp self.sw_name = sw_name self.enable_sudo = enable_sudo self.tmp_flow_file_path = tmp_flow_file_path @@ -87,6 +89,12 @@ def create_tun_port( "options:dst_port={}".format(vxlan_udp_port), "ofport_request={}".format(self.vxlan_ofport), ] + if self.gbp: + cmd.append("options:exts=gbp") + else: + # Unmake an existing GBP tunnel; removing an absent key is + # a no-op. + cmd += ["--", "remove", "Interface", port_name, "options", "exts"] return shell.runsh(command=cmd, enable_sudo=self.enable_sudo) def sync_flows(self, flows): diff --git a/evpn_connector/service/evpn.py b/evpn_connector/service/evpn.py index 789e1e0..f47a172 100644 --- a/evpn_connector/service/evpn.py +++ b/evpn_connector/service/evpn.py @@ -31,6 +31,8 @@ LOG = logging.getLogger(__name__) COMMON_NAME = "local-vni-" +# Steps a session must outlive when nothing reports End-of-RIB. +PEER_SETTLE_STEPS = 2 def _rewrite_cfg_prefix(cfg, prefix): @@ -77,6 +79,8 @@ def __init__( event_type=None, error_event_type=None, policy_enabled=True, + fail_static=False, + fail_static_min_peers=0, ): super(EvpnConnectorService, self).__init__( step_period=step_period, @@ -100,9 +104,20 @@ def __init__( self.anycast_check_ofport = anycast_check_ofport self.anycast_check_mac = anycast_check_mac self.anycast_used = False + # Fail-static: the last flow set synced while enough peers had + # converged, retained if one is later lost. + self.fail_static = fail_static + self.fail_static_min_peers = fail_static_min_peers + self._peer_settle_sec = PEER_SETTLE_STEPS * step_period + self._last_good_flows = set() + self._peers_healthy = True + self._no_peers_since = None def _setup(self): super(EvpnConnectorService, self)._setup() + if self.fail_static: + # Otherwise a restart during an outage syncs the degraded RIB. + self._last_good_flows = self._read_synced_flows() LOG.info("Create if not exist ovs switch and vxlan port") self.ovs_client.create_bridge() self.ovs_client.create_tun_port( @@ -110,6 +125,19 @@ def _setup(self): ) self.update_peer(force=True) + def _read_synced_flows(self): + try: + with open(self.ovs_client.tmp_flow_file_path) as flow_file: + lines = flow_file.read().splitlines() + except OSError: + return set() + flows = set() + for line in lines: + match, _, action = line.partition(" ") + if match: + flows.add(evpnobj.OvsFlow(match, action)) + return flows + @staticmethod def _rt2lst(rt_list, local_asn=None): # Convert RouteTarget from str repr "ASN:LABEL" @@ -342,6 +370,86 @@ def update_peer( ) self.need_reset_peers = False + def _log_peer_health(self, healthy, reason): + # Checked every step, so only a change of state is worth a line. + if healthy == self._peers_healthy: + return + self._peers_healthy = healthy + if healthy: + LOG.info("BGP peers healthy again: %s", reason) + else: + LOG.warning( + "BGP peers degraded (%s); holding last-known flows " + "(fail-static)", + reason, + ) + + def _upstream_peers_healthy(self): + """Whether the RIB can be trusted for flow deletion. + + At least fail_static_min_peers (0: all) configured peers must + have converged, i.e. sent their table (see count_peers). No + peers is degraded for the settle period (a restarted gobgp), + then authoritative. An unreadable peer list is raised: the step + stops and the flows stay as they are. + """ + total, converged = self.gobgp_client.count_peers(self._peer_settle_sec) + if total == 0: + now = time.time() + if self._no_peers_since is None: + self._no_peers_since = now + if now - self._no_peers_since >= self._peer_settle_sec: + self._log_peer_health(True, "no peers are configured") + return True + self._log_peer_health(False, "no peers are configured yet") + return False + self._no_peers_since = None + # More required than configured means all of them. + required = min(self.fail_static_min_peers or total, total) + if converged < required: + self._log_peer_health( + False, + "only %d of %d peers converged, %d required" + % (converged, total, required), + ) + return False + self._log_peer_health( + True, "%d of %d peers converged" % (converged, total) + ) + return True + + def _apply_fail_static(self, target_flows, metrics, peers_were_healthy): + """Decide which flows to sync, retaining flows on peer loss. + + With enough peers converged the computed set is applied and + snapshotted; otherwise it is unioned with the snapshot, so + routes of a lost session are not deleted until it returns. + Health is sampled before the RIB read (peers_were_healthy) and + here, and both must hold: a session changing mid-step leaves a + RIB matching neither state. The union keeps the fresh side of + an equal match, so additions and changes land; deletions wait, + and a reused ofport meanwhile gets the stale flow. + """ + if not self.fail_static or ( + peers_were_healthy and self._upstream_peers_healthy() + ): + self._last_good_flows = target_flows + metrics["fail_static_active"] = 0 + metrics["fail_static_retained_cnt"] = 0 + return target_flows + + flows_to_sync = target_flows.union(self._last_good_flows) + metrics["fail_static_active"] = 1 + metrics["fail_static_retained_cnt"] = len(flows_to_sync) - len( + target_flows + ) + LOG.debug( + "Fail-static active: retaining %d flows on top of %d computed", + metrics["fail_static_retained_cnt"], + len(target_flows), + ) + return flows_to_sync + def get_anycast_status(self): if not os.path.isfile(self.anycast_status_file): LOG.warning( @@ -434,6 +542,10 @@ def _step(self): duration_metrics["update_policy_time"], ) + peers_were_healthy = ( + not self.fail_static or self._upstream_peers_healthy() + ) + LOG.debug("Start getting announces from gobgp") start_time = time.time() ( @@ -575,9 +687,15 @@ def _step(self): duration_metrics["prep_ovs_time"], ) + uni_flows = uni_flows.union(bum_flows) + + flows_to_sync = self._apply_fail_static( + uni_flows, metrics, peers_were_healthy + ) + start_time = time.time() LOG.debug("Sync flows in ovs") - self.ovs_client.sync_flows(uni_flows.union(bum_flows)) + self.ovs_client.sync_flows(flows_to_sync) duration_metrics["sync_ovs_time"] = time.time() - start_time LOG.info( "Sync ovs flows done for %0.4f sec", diff --git a/evpn_connector/service/objects.py b/evpn_connector/service/objects.py index b60c571..c985b2c 100644 --- a/evpn_connector/service/objects.py +++ b/evpn_connector/service/objects.py @@ -27,6 +27,19 @@ LOG = logging.getLogger(__name__) +# Whether the fabric carries a sender's group: switch-wide, set at start-up. +_GBP_ENABLED = False + + +def set_gbp(enabled): + global _GBP_ENABLED + _GBP_ENABLED = enabled + + +def _gbp(action): + """The GBP move with its separator, or nothing when gbp is off.""" + return "%s," % action if _GBP_ENABLED else "" + class BaseObj(object): def __init__(self): @@ -296,6 +309,9 @@ def _ovs_to_out_table_action(self, local=False): reg1_value = constants.REG_FROM_REMOTE if local: reg1_value = constants.REG_FROM_LOCAL + elif self.port_type == constants.EVPN_EDGE_TYPE_VXLAN: + # Off the wire on the switched path. + res += _gbp(constants.GBP_TO_MARK) res += "set_field:%d->reg0,set_field:%d->reg1,resubmit(,%d)" % ( self.vni, @@ -316,7 +332,9 @@ def ovs_output(self, for_group=False): self.ofport, ) elif self.port_type == constants.EVPN_EDGE_TYPE_VXLAN: - res += "set_field:%s->tun_id,set_field:%s->tun_dst,output:%d" % ( + # Onto the wire. + res += "%sset_field:%s->tun_id,set_field:%s->tun_dst,output:%d" % ( + _gbp(constants.MARK_TO_GBP), self.tun_id, self.next_hop, self.ofport, @@ -519,7 +537,9 @@ def ovs_output(self, local=False): self.ofport, ) elif self.port_type == constants.EVPN_EDGE_TYPE_VXLAN: - res += "set_field:%s->tun_id,set_field:%s->tun_dst,output:%d" % ( + # Onto the wire. + res += "%sset_field:%s->tun_id,set_field:%s->tun_dst,output:%d" % ( + _gbp(constants.MARK_TO_GBP), self.tun_id, self.next_hop, self.ofport, @@ -1058,7 +1078,9 @@ def import_announces(self, local_ce_prefixes, remote_ce_prefixes): @property def _ovs_to_out_table_action(self): - return "action=set_field:%d->reg2,resubmit(,%d)" % ( + # Off the wire on the routed path. + return "action=%sset_field:%d->reg2,resubmit(,%d)" % ( + _gbp(constants.GBP_TO_MARK), self.vrf_number, constants.OUTPUT_TABLE_NUM, ) @@ -1152,6 +1174,10 @@ def _ovs_to_out_table_action(self, local=False): reg1_value = constants.REG_FROM_REMOTE if local: reg1_value = constants.REG_FROM_LOCAL + else: + # Off the wire, only on tunnel ingress: local traffic has no + # header, and reading one would erase its sender's mark. + res += _gbp(constants.GBP_TO_MARK) res += "set_field:%d->reg0,set_field:%d->reg1,resubmit(,%d)" % ( self.vni, @@ -1164,7 +1190,9 @@ def _ovs_to_out_table_action(self, local=False): def ovs_output( self, for_group=False, tun_ofport=constants.VXLAN_PORT_OFPORT ): - return "set_field:%s->tun_id,set_field:%s->tun_dst,output:%d" % ( + # The flooded path carries it too. + return "%sset_field:%s->tun_id,set_field:%s->tun_dst,output:%d" % ( + _gbp(constants.MARK_TO_GBP), self.tun_id, self.next_hop, tun_ofport, diff --git a/evpn_connector/tests/unit/test_bgp_client.py b/evpn_connector/tests/unit/test_bgp_client.py index 3202784..69241bb 100644 --- a/evpn_connector/tests/unit/test_bgp_client.py +++ b/evpn_connector/tests/unit/test_bgp_client.py @@ -14,7 +14,10 @@ # License for the specific language governing permissions and limitations # under the License. +import mock + from evpn_connector.bgp import client +from evpn_connector.bgp.generated import gobgp_pb2 from evpn_connector.service import objects @@ -167,3 +170,64 @@ def test_filter_local_announces(self): assert not remote_edges assert remote_virt_nets == remote_vnet assert not remote_client_prefixes + + +class TestCountPeers(object): + """A session being up is not the table having arrived.""" + + @staticmethod + def _peer( + session_state=gobgp_pb2.PeerState.ESTABLISHED, + uptime=1000, + gr_received=True, + end_of_rib_received=True, + ): + peer = gobgp_pb2.Peer() + peer.state.session_state = session_state + peer.timers.state.uptime.seconds = uptime + afi_safi = peer.afi_safis.add() + afi_safi.mp_graceful_restart.state.received = gr_received + afi_safi.mp_graceful_restart.state.end_of_rib_received = ( + end_of_rib_received + ) + return peer + + def _count(self, peers, settle_sec=0): + responses = [mock.Mock(peer=peer) for peer in peers] + fake_self = mock.Mock( + list_peer=mock.Mock(return_value=responses), + _peer_converged=client.BGPClient._peer_converged, + ) + return client.BGPClient.count_peers(fake_self, settle_sec) + + def test_no_peers(self): + assert self._count([]) == (0, 0) + + def test_established_with_end_of_rib(self): + assert self._count([self._peer()]) == (1, 1) + + def test_established_while_the_table_is_still_coming(self): + """The window fail-static exists for: session up, RIB half-read.""" + assert self._count([self._peer(end_of_rib_received=False)]) == (1, 0) + + def test_not_established(self): + peer = self._peer(session_state=gobgp_pb2.PeerState.OPENSENT) + assert self._count([peer]) == (1, 0) + + def test_without_graceful_restart_nothing_reports_the_end(self): + """No End-of-RIB to wait for, so a fresh session is not trusted.""" + peer = self._peer(gr_received=False, end_of_rib_received=False) + with mock.patch("time.time", return_value=1005): + assert self._count([peer], settle_sec=10) == (1, 0) + peer.timers.state.uptime.seconds = 990 + assert self._count([peer], settle_sec=10) == (1, 1) + + def test_a_session_younger_than_the_settle_time(self): + peer = self._peer() + with mock.patch("time.time", return_value=1005): + assert self._count([peer], settle_sec=10) == (1, 0) + + def test_an_unset_uptime_is_not_trusted(self): + peer = self._peer(uptime=0) + with mock.patch("time.time", return_value=1005): + assert self._count([peer], settle_sec=10) == (1, 0) diff --git a/evpn_connector/tests/unit/test_evpn_objects.py b/evpn_connector/tests/unit/test_evpn_objects.py index 2a3d650..70cd387 100644 --- a/evpn_connector/tests/unit/test_evpn_objects.py +++ b/evpn_connector/tests/unit/test_evpn_objects.py @@ -16,11 +16,20 @@ import sys +import pytest + from evpn_connector.common import constants from evpn_connector.service import objects class TestEvpnConnectorObjects(object): + @pytest.fixture(autouse=True) + def gbp_off_afterwards(self): + # It is process-wide state; leaving it on would build every + # later flow in this run as a GBP one. + yield + objects.set_gbp(False) + def test_ovs_flow(self): match1 = "table=1,priority=100" match2 = "table=1,priority=200" @@ -333,3 +342,85 @@ def test_vnet_from_ce(self): assert len(vnet_vnis) == len(vnets) == 2 assert expected_vnets == vnets assert expected_vnis == vnet_vnis + + def _gbp_objects(self): + ce = objects.ClientEdge( + mac="d8:6b:5c:cd:97:ee", + vni=1, + ip="192.168.1.1", + rt=objects.RouteTarget(targets=[(65001, 100)]), + as_number=65001, + ofport=1, + port_type="vxlan", + tag=100, + next_hop="10.0.0.2", + ) + + prefix = objects.ClientEdgePrefix( + prefix="10.1.1.0", + prefix_len=24, + mac="d8:6b:5c:cd:97:ee", + router_mac="d8:6b:5c:cd:97:ef", + vni=1, + rt=objects.RouteTarget(targets=[(65001, 100)]), + as_number=65001, + ofport=1, + port_type="vxlan", + tag=100, + next_hop="10.0.0.2", + ) + + vnet = objects.VirtNet( + vni=1, + next_hop="10.0.0.2", + rt=objects.RouteTarget(targets=[(65001, 100)]), + as_number=65001, + ) + return ce, prefix, vnet + + def test_the_identity_is_carried_only_at_the_tunnel(self): + """The identity is put on the wire and taken off it, once each. + + Every path out has to carry it (a guest reached by a /32, or by + flooding, is otherwise a hole), and only tunnel ingress may take it + off: local traffic never had a header, and reading one would erase + its sender's mark. + """ + objects.set_gbp(True) + ce, prefix, vnet = self._gbp_objects() + + assert constants.MARK_TO_GBP in ce.ovs_output() + assert constants.MARK_TO_GBP in prefix.ovs_output() + assert constants.MARK_TO_GBP in vnet.ovs_output() + + # Both flows that a tunnel ingress can hit: the Type 2 one, and + # VirtNet's stand-in for a missing Type 2 announce. + for obj in (ce, vnet): + assert constants.GBP_TO_MARK in obj._ovs_to_out_table_action( + local=False + ) + assert constants.GBP_TO_MARK not in obj._ovs_to_out_table_action( + local=True + ) + + def test_no_identity_is_carried_when_gbp_is_off(self): + """Off by default means the flows are exactly the ones from before. + + OVS will not mix GBP and non-GBP tunnels on one UDP port, so a + fabric that has not asked for this must be left as it was. + """ + objects.set_gbp(False) + ce, prefix, vnet = self._gbp_objects() + + for flow in ( + ce.ovs_output(), + prefix.ovs_output(), + vnet.ovs_output(), + ce._ovs_to_out_table_action(local=False), + ce._ovs_to_out_table_action(local=True), + vnet._ovs_to_out_table_action(local=False), + vnet._ovs_to_out_table_action(local=True), + ): + assert constants.MARK_TO_GBP not in flow + assert constants.GBP_TO_MARK not in flow + assert ",," not in flow and not flow.endswith(",") diff --git a/evpn_connector/tests/unit/test_evpn_service.py b/evpn_connector/tests/unit/test_evpn_service.py index 7f41dd7..8df3c88 100644 --- a/evpn_connector/tests/unit/test_evpn_service.py +++ b/evpn_connector/tests/unit/test_evpn_service.py @@ -17,6 +17,7 @@ import json import mock import os +import pytest import shutil import tempfile @@ -42,10 +43,10 @@ def test_rt2lst_as_local_ovveride(self): == expected ) - def setup(self): + def setup_method(self, method): self.temp_dir = tempfile.mkdtemp() - def teardown(self): + def teardown_method(self, method): shutil.rmtree(self.temp_dir) def test_read_client_configs_with_no_folder(self): @@ -311,3 +312,294 @@ def test_read_l3_anycast_client_configs(self): assert expected_pr in res_pr assert expected_any_pr1 in res_pr assert expected_any_pr2 in res_pr + + +class TestFailStatic(object): + def _make_service( + self, + fail_static=True, + fail_static_min_peers=0, + step_period=5, + ): + return evpn.EvpnConnectorService( + source_ip="10.10.10.1", + as_number=1, + configs_dir="", + gobgp_client=mock.MagicMock(), + ovs_client=mock.MagicMock(), + sender=mock.MagicMock(), + vxlan_udp_port=4789, + router_mac_type5="11:22:33:44:55:66", + anycast_status_file="/tmp/anycast_status_file", + anycast_check_ofport=65277, + anycast_check_mac="12:34:56:78:90:aa", + fail_static=fail_static, + fail_static_min_peers=fail_static_min_peers, + step_period=step_period, + ) + + def test_snapshot_survives_a_restart(self, tmpdir): + """A restart while the RR is down keeps the synced flows.""" + flow_file = tmpdir.join("flows") + flow_file.write("match-a actions=NORMAL\nmatch-b actions=drop\n") + + service = self._make_service() + service.ovs_client.tmp_flow_file_path = str(flow_file) + service._setup() + + assert {f.to_string() for f in service._last_good_flows} == { + "match-a actions=NORMAL", + "match-b actions=drop", + } + + def test_no_flow_file_yet_is_an_empty_snapshot(self, tmpdir): + service = self._make_service() + service.ovs_client.tmp_flow_file_path = str(tmpdir.join("absent")) + service._setup() + + assert service._last_good_flows == set() + + def test_a_retained_flow_read_back_is_the_flow_it_was(self, tmpdir): + """A flow read back merges with the same flow from the RIB.""" + flow_file = tmpdir.join("flows") + flow_file.write("match-a actions=NORMAL\n") + + service = self._make_service() + service.ovs_client.tmp_flow_file_path = str(flow_file) + service._setup() + service.gobgp_client.count_peers.return_value = (2, 0) + + fresh = objects.OvsFlow("match-a", "actions=NORMAL") + result = service._apply_fail_static({fresh}, {}, True) + + assert result == {fresh} + + def test_peers_healthy_all_established(self): + service = self._make_service() + service.gobgp_client.count_peers.return_value = (2, 2) + + assert service._upstream_peers_healthy() is True + + def test_peers_degraded_no_peers(self): + """An empty peer list is a restarted gobgp, not a healthy node. + + It is indistinguishable from every session having been lost, and + it comes with an empty RIB, which is precisely the set of flows + that must not be trusted for deletion. + """ + service = self._make_service() + service.gobgp_client.count_peers.return_value = (0, 0) + + assert service._upstream_peers_healthy() is False + + def test_no_peers_past_settle_is_healthy(self): + service = self._make_service() + service.gobgp_client.count_peers.return_value = (0, 0) + + with mock.patch.object(evpn.time, "time", return_value=100): + assert service._upstream_peers_healthy() is False + with mock.patch.object(evpn.time, "time", return_value=110): + assert service._upstream_peers_healthy() is True + + def test_peers_appearing_restart_the_settle(self): + service = self._make_service() + service.gobgp_client.count_peers.return_value = (0, 0) + with mock.patch.object(evpn.time, "time", return_value=100): + service._upstream_peers_healthy() + service.gobgp_client.count_peers.return_value = (1, 1) + service._upstream_peers_healthy() + service.gobgp_client.count_peers.return_value = (0, 0) + + with mock.patch.object(evpn.time, "time", return_value=110): + assert service._upstream_peers_healthy() is False + + def test_restart_without_peers_drops_the_read_back_flows(self, tmpdir): + flow_file = tmpdir.join("flows") + flow_file.write("match-a actions=NORMAL\n") + service = self._make_service() + service.ovs_client.tmp_flow_file_path = str(flow_file) + service._setup() + service.gobgp_client.count_peers.return_value = (0, 0) + + with mock.patch.object(evpn.time, "time", return_value=100): + service._upstream_peers_healthy() + with mock.patch.object(evpn.time, "time", return_value=110): + result = service._apply_fail_static(set(), {}, True) + + assert result == set() + + def test_apply_without_a_snapshot_changes_nothing(self): + """A node that never had peers is not held back by fail-static.""" + service = self._make_service() + service.gobgp_client.count_peers.return_value = (0, 0) + metrics = {} + + result = service._apply_fail_static({"flow-a"}, metrics, True) + + assert result == {"flow-a"} + assert metrics["fail_static_retained_cnt"] == 0 + + def test_peers_degraded_partial(self): + service = self._make_service() + service.gobgp_client.count_peers.return_value = (2, 1) + + assert service._upstream_peers_healthy() is False + + def test_min_peers_allows_a_partial_fabric(self): + service = self._make_service(fail_static_min_peers=1) + service.gobgp_client.count_peers.return_value = (3, 1) + + assert service._upstream_peers_healthy() is True + + def test_min_peers_still_degrades_below_the_threshold(self): + service = self._make_service(fail_static_min_peers=2) + service.gobgp_client.count_peers.return_value = (3, 1) + + assert service._upstream_peers_healthy() is False + + def test_min_peers_above_configured_means_all_of_them(self): + service = self._make_service(fail_static_min_peers=5) + service.gobgp_client.count_peers.return_value = (2, 2) + + assert service._upstream_peers_healthy() is True + + def test_unreadable_peers_raise(self): + """A local gobgp fault is an error, not a held state. + + It is indistinguishable from a broken fabric, and the step that + raises never reaches sync_flows, so flows are retained anyway. + """ + service = self._make_service() + service.gobgp_client.count_peers.side_effect = RuntimeError("boom") + + with pytest.raises(RuntimeError): + service._upstream_peers_healthy() + + def test_apply_raises_and_keeps_the_snapshot_on_error(self): + service = self._make_service() + metrics = {} + service.gobgp_client.count_peers.return_value = (1, 1) + service._apply_fail_static({"flow-a"}, metrics, True) + service.gobgp_client.count_peers.side_effect = RuntimeError("boom") + + with pytest.raises(RuntimeError): + service._apply_fail_static({"flow-b"}, metrics, True) + + assert service._last_good_flows == {"flow-a"} + + def test_apply_healthy_snapshots_and_passes_through(self): + service = self._make_service() + service.gobgp_client.count_peers.return_value = (1, 1) + metrics = {} + target = {"flow-a", "flow-b"} + + result = service._apply_fail_static(target, metrics, True) + + assert result == target + assert service._last_good_flows == target + assert metrics["fail_static_active"] == 0 + assert metrics["fail_static_retained_cnt"] == 0 + + def test_apply_degraded_retains_last_good(self): + service = self._make_service() + metrics = {} + # Healthy step snapshots two flows + service.gobgp_client.count_peers.return_value = (1, 1) + service._apply_fail_static({"flow-a", "flow-remote"}, metrics, True) + # Peer lost: remote flow vanished from the computed set + service.gobgp_client.count_peers.return_value = (1, 0) + + result = service._apply_fail_static({"flow-a"}, metrics, True) + + assert result == {"flow-a", "flow-remote"} + # Snapshot must not be overwritten by the degraded set + assert service._last_good_flows == {"flow-a", "flow-remote"} + assert metrics["fail_static_active"] == 1 + assert metrics["fail_static_retained_cnt"] == 1 + + def test_apply_degraded_still_applies_local_changes(self): + service = self._make_service() + metrics = {} + service.gobgp_client.count_peers.return_value = (1, 1) + service._apply_fail_static({"flow-remote"}, metrics, True) + service.gobgp_client.count_peers.return_value = (1, 0) + + result = service._apply_fail_static({"flow-new-local"}, metrics, True) + + assert result == {"flow-remote", "flow-new-local"} + + def test_apply_recovery_drops_stale_flows(self): + service = self._make_service() + metrics = {} + service.gobgp_client.count_peers.return_value = (1, 1) + service._apply_fail_static({"flow-a", "flow-remote"}, metrics, True) + service.gobgp_client.count_peers.return_value = (1, 0) + service._apply_fail_static({"flow-a"}, metrics, True) + # Peer is back; RIB is authoritative again + service.gobgp_client.count_peers.return_value = (1, 1) + + result = service._apply_fail_static({"flow-a"}, metrics, True) + + assert result == {"flow-a"} + assert service._last_good_flows == {"flow-a"} + assert metrics["fail_static_active"] == 0 + + def test_apply_degraded_keeps_the_fresh_action_for_a_known_match(self): + """What makes retention safe: a stale flow never wins a match. + + Flows are equal by match alone, and a union keeps the side it + started from, so a match that is still computed is applied with + its current action and only matches that vanished are retained. + """ + service = self._make_service() + metrics = {} + stale = objects.OvsFlow( + match="table=0,priority=10,in_port=1", action="action=output:5" + ) + service.gobgp_client.count_peers.return_value = (1, 1) + service._apply_fail_static({stale}, metrics, True) + fresh = objects.OvsFlow( + match="table=0,priority=10,in_port=1", action="action=output:9" + ) + service.gobgp_client.count_peers.return_value = (1, 0) + + result = service._apply_fail_static({fresh}, metrics, True) + + assert [flow.to_string() for flow in result] == [fresh.to_string()] + assert metrics["fail_static_retained_cnt"] == 0 + + def test_health_before_the_rib_read_counts_too(self): + """A peer that converged mid-step leaves a half-read RIB. + + The second sample alone would call that set authoritative and + snapshot it, which is the very set fail-static exists to + distrust. + """ + service = self._make_service() + metrics = {} + service.gobgp_client.count_peers.return_value = (1, 1) + service._apply_fail_static({"flow-a", "flow-remote"}, metrics, True) + # Healthy now, but it was not when the RIB was read + result = service._apply_fail_static({"flow-a"}, metrics, False) + + assert result == {"flow-a", "flow-remote"} + assert service._last_good_flows == {"flow-a", "flow-remote"} + assert metrics["fail_static_active"] == 1 + + def test_settle_time_defaults_to_two_steps(self): + service = self._make_service(step_period=7) + service.gobgp_client.count_peers.return_value = (1, 1) + + service._upstream_peers_healthy() + + service.gobgp_client.count_peers.assert_called_once_with(14) + + def test_disabled_passes_through_when_degraded(self): + service = self._make_service(fail_static=False) + metrics = {} + service.gobgp_client.count_peers.return_value = (1, 0) + + result = service._apply_fail_static({"flow-a"}, metrics, True) + + assert result == {"flow-a"} + service.gobgp_client.count_peers.assert_not_called() diff --git a/evpn_connector/tests/unit/test_ovs_client.py b/evpn_connector/tests/unit/test_ovs_client.py index 1fc51a3..7d3fc22 100644 --- a/evpn_connector/tests/unit/test_ovs_client.py +++ b/evpn_connector/tests/unit/test_ovs_client.py @@ -88,10 +88,59 @@ def test_create_tun_port(self, mock_shell_run): 'options:local_ip="%s"' % local_ip, "options:dst_port=%d" % vxlan_udp_port, "ofport_request=%s" % vxlan_ofport, + "--", + "remove", + "Interface", + "vxlan_out", + "options", + "exts", ], enable_sudo=False, ) + def test_create_tun_port_with_gbp(self, mock_shell_run): + """With gbp on, the tunnel gains the extension and nothing else. + + OVS will not mix GBP and non-GBP tunnels on one UDP port, so this is + a property of a whole fabric, which is why it is opt-in. + """ + ovs_client = client.OvSClient( + "test_sw", "/tmp/flows.txt", vxlan_ofport=10, gbp=True + ) + + ovs_client.create_tun_port( + vxlan_source_ip="1.2.3.4", vxlan_udp_port=3423 + ) + + command = mock_shell_run.call_args[1]["command"] + assert "options:exts=gbp" in command + assert "remove" not in command + + def test_create_tun_port_without_gbp_unmakes_it(self, mock_shell_run): + """An existing tunnel is taken back, or the flag is one-way. + + The port outlives the daemon: --may-exist means turning gbp off + would otherwise leave a GBP tunnel behind forever. + """ + ovs_client = client.OvSClient( + "test_sw", "/tmp/flows.txt", vxlan_ofport=10, gbp=False + ) + + ovs_client.create_tun_port( + vxlan_source_ip="1.2.3.4", vxlan_udp_port=3423 + ) + + command = mock_shell_run.call_args[1]["command"] + assert "options:exts=gbp" not in command + assert command[-6:] == [ + "--", + "remove", + "Interface", + "vxlan_out", + "options", + "exts", + ] + def test_sync_flows(self, mock_shell_run): file_name = "/tmp/flows.txt" switch_name = "test_sw" diff --git a/requirements.txt b/requirements.txt index 7ad7ae3..0e972e8 100644 --- a/requirements.txt +++ b/requirements.txt @@ -1,16 +1,19 @@ -pbr>=1.10.0,<=5.8.1 # Apache-2.0 +pbr>=1.10.0,<=5.8.1 # Apache-2.0 (capped by loopster) +setuptools<81 # MIT (pbr 5.8.1 imports pkg_resources, gone in 81) six>=1.9.0,<=1.16.0 # MIT -oslo.config==3.22.0 # Apache-2.0 -#grpcio===1.41.1 # Apache-2.0 -grpcio===1.22.0; python_version<'3.8' # Apache-2.0 -grpcio===1.26.0; python_version>='3.8' # Apache-2.0 -protobuf===3.9.0; python_version<'3.8' # BSD -protobuf===3.14.0; python_version>='3.8' # BSD -netaddr>=0.7.18,<=0.8.0 # BSD +oslo.config==3.22.0; python_version<'3.10' # Apache-2.0 +oslo.config>=9.4.0,<11.0.0; python_version>='3.10' # Apache-2.0 +grpcio===1.39.0; python_version<'3.10' # Apache-2.0 +grpcio>=1.83.0,<2.0.0; python_version>='3.10' # Apache-2.0 +protobuf===3.17.3; python_version<'3.10' # BSD +protobuf>=7.35.1,<8.0.0; python_version>='3.10' # BSD +netaddr>=0.7.18,<=0.8.0; python_version<'3.10' # BSD +netaddr>=1.0.0,<2.0.0; python_version>='3.10' # BSD enum34===1.1.6; python_version<'3.8' subprocess32===3.2.6;python_version=='2.7' # PSF sentry-sdk==1.5.2;python_version<'3.8' # BSD -sentry-sdk==1.6.0;python_version>='3.8' # BSD +sentry-sdk==1.6.0;python_version>='3.8' and python_version<'3.10' # BSD +sentry-sdk>=2.0.0,<3.0.0; python_version>='3.10' # BSD loopster>=2.14.4,<3.0.0 # Apache-2.0 obsender>=5.0.0,<6.0.0 # Apache-2.0 pyyaml>=6.0 # MIT diff --git a/test-requirements.txt b/test-requirements.txt index 8f12f67..b1f8ae4 100644 --- a/test-requirements.txt +++ b/test-requirements.txt @@ -1,7 +1,9 @@ hacking<3 # Apache-2.0 typing-extensions<4.2.0;python_version=='3.6' # PSF-2.0 pytest-timer -pytest==4.6.1 +pytest==4.6.1; python_version<'3.10' +pytest>=8.0.0; python_version>='3.10' coverage>=4.0 # Apache-2.0 -mock==3.0.5 +mock==3.0.5; python_version<'3.10' +mock>=5.0.0; python_version>='3.10' flake8>=3.8.4 # Flake8 License (MIT) diff --git a/tox.ini b/tox.ini index 9dc635e..c4d2d4d 100644 --- a/tox.ini +++ b/tox.ini @@ -1,6 +1,6 @@ [tox] envlist = pep8 - begin,py27,py38,py313,end + begin,py27,py38,py313,py314,end skipsdist = True minversion = 2.0 @@ -14,7 +14,7 @@ deps = -r{toxinidir}/requirements.txt setenv = PYTHONDONTWRITEBYTECODE = 1 commands = - py27,py38: coverage run -m pytest {posargs} --timer-top-n=10 {[base]project_name}/tests/unit + py{27,38,39,310,311,312,313,314}: coverage run -m pytest {posargs} --timer-top-n=10 {[base]project_name}/tests/unit functional: pytest {posargs} --timer-top-n=10 {[base]project_name}/tests/functional [testenv:pep8]