Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 8 additions & 6 deletions .github/workflows/tests.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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 }}

38 changes: 38 additions & 0 deletions evpn_connector/bgp/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down
7 changes: 7 additions & 0 deletions evpn_connector/cmd/evpn.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -70,13 +71,17 @@ 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,
tmp_flow_file_path=CONF.ovs.tmp_flow_file_path,
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
Expand All @@ -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,
Expand Down
30 changes: 30 additions & 0 deletions evpn_connector/common/conf_opts.py
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
5 changes: 5 additions & 0 deletions evpn_connector/common/constants.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"

Expand Down
8 changes: 8 additions & 0 deletions evpn_connector/ovs/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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):
Expand Down
120 changes: 119 additions & 1 deletion evpn_connector/service/evpn.py
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand Down Expand Up @@ -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,
Expand All @@ -100,16 +104,40 @@ 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()
Comment thread
gmelikov marked this conversation as resolved.
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(
vxlan_source_ip=self.source_ip, vxlan_udp_port=self.vxlan_udp_port
)
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"
Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -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()
(
Expand Down Expand Up @@ -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",
Expand Down
Loading
Loading