From: Xin Xie <xiexinet@gmail.com>
To: netdev@vger.kernel.org, linux-kselftest@vger.kernel.org,
linux-kernel@vger.kernel.org
Cc: davem@davemloft.net, edumazet@google.com, kuba@kernel.org,
pabeni@redhat.com, horms@kernel.org, andrew+netdev@lunn.ch,
shuah@kernel.org, kees@kernel.org, petr.wozniak@gmail.com,
qingfang.deng@linux.dev, fmaurer@redhat.com,
luka.gejak@linux.dev, bigeasy@linutronix.de,
xiaoliang.yang_1@nxp.com, skhawaja@google.com,
liuhangbin@gmail.com, stable@vger.kernel.org,
sdf.kernel@gmail.com, xiexinet@gmail.com
Subject: [PATCH net v7 4/4] selftests: net: hsr: verify GRO policy and ordered forwarding
Date: Fri, 9 Oct 2026 22:13:24 +0200 [thread overview]
Message-ID: <20261009201324.17-5-xiexinet@gmail.com> (raw)
In-Reply-To: <20261009201324.17-1-xiexinet@gmail.com>
Add member GRO policy and ordered-forwarding tests with two shell
entry points and one installed Python helper.
Exercise attach, feature requests, detach and failed attachment,
including a busy RX handler. Check wanted-state restoration, the
master's GSO type mask and continued delivery through the busy handler.
Capture actual TCP and UDP GSO traffic and compare every input with
both LAN outputs. Count retransmitted TCP bytes separately so a retry
cannot hide a dropped aggregate. Check payloads, tag/RCT, VLAN stack
and sequence numbers, requiring complete captures without socket loss
or truncation before any packet verdict.
Cover PRP LAN-ingress local delivery, plain and VLAN traffic, local
and proxy supervision, concurrent producers and sequence wrap.
Concurrency tests check per-LAN ordering and contents, without claiming
proof of a worker handoff.
One process owns the namespace-bound sockets and capture. Work and
cleanup are bounded; INT/TERM release owned resources and preserve
failures. Missing prerequisites SKIP the affected case; feature-query
errors and missing packet evidence FAIL.
Signed-off-by: Xin Xie <xiexinet@gmail.com>
---
tools/testing/selftests/net/hsr/Makefile | 3 +
tools/testing/selftests/net/hsr/config | 1 +
.../selftests/net/hsr/hsr_gro_superpacket.py | 1088 +++++++++++++++++
.../selftests/net/hsr/hsr_gro_superpacket.sh | 9 +
.../net/hsr/hsr_ordered_forwarding.sh | 9 +
5 files changed, 1110 insertions(+)
create mode 100755 tools/testing/selftests/net/hsr/hsr_gro_superpacket.py
create mode 100755 tools/testing/selftests/net/hsr/hsr_gro_superpacket.sh
create mode 100755 tools/testing/selftests/net/hsr/hsr_ordered_forwarding.sh
diff --git a/tools/testing/selftests/net/hsr/Makefile b/tools/testing/selftests/net/hsr/Makefile
index 2150e487ac7d..81e55309ed20 100644
--- a/tools/testing/selftests/net/hsr/Makefile
+++ b/tools/testing/selftests/net/hsr/Makefile
@@ -8,8 +8,11 @@ TEST_PROGS := \
hsr_redbox.sh \
link_faults.sh \
prp_ping.sh \
+ hsr_gro_superpacket.sh \
+ hsr_ordered_forwarding.sh \
# end of TEST_PROGS
TEST_FILES += hsr_common.sh
+TEST_FILES += hsr_gro_superpacket.py
include ../../lib.mk
diff --git a/tools/testing/selftests/net/hsr/config b/tools/testing/selftests/net/hsr/config
index 205cc4d3d64b..724fbe3a6eab 100644
--- a/tools/testing/selftests/net/hsr/config
+++ b/tools/testing/selftests/net/hsr/config
@@ -2,5 +2,6 @@ CONFIG_BRIDGE=y
CONFIG_HSR=y
CONFIG_IPV6=y
CONFIG_NET_SCH_NETEM=m
+CONFIG_PACKET=y
CONFIG_VETH=y
CONFIG_VLAN_8021Q=m
diff --git a/tools/testing/selftests/net/hsr/hsr_gro_superpacket.py b/tools/testing/selftests/net/hsr/hsr_gro_superpacket.py
new file mode 100755
index 000000000000..d39bf1af291f
--- /dev/null
+++ b/tools/testing/selftests/net/hsr/hsr_gro_superpacket.py
@@ -0,0 +1,1088 @@
+#!/usr/bin/env python3
+# SPDX-License-Identifier: GPL-2.0
+"""HSR/PRP GRO policy, segmentation and functional ordering regressions."""
+
+from array import array
+from collections import Counter
+from contextlib import contextmanager
+import ctypes
+import errno
+import json
+import os
+import re
+import selectors
+import shutil
+import signal
+import socket
+import struct
+import subprocess
+import sys
+import threading
+import time
+import uuid
+
+
+class Skip(Exception):
+ pass
+
+
+def require(condition, message):
+ if not condition:
+ raise ValueError(message)
+
+
+def features(text):
+ result = {}
+ for line in text.splitlines():
+ if ":" not in line or line.startswith("Features for"):
+ continue
+ name, state = line.strip().split(":", 1)
+ match = re.fullmatch(r"\s*(on|off)((?:\s+\[[^\]]*\])*)", state)
+ require(match is not None, "malformed feature line: " + line)
+ result[name] = (match[1], "[fixed]" in match[2])
+ require(result, "empty feature listing")
+ return result
+
+
+def decode(raw, ancillary=()):
+ require(len(raw) >= 14, "short Ethernet capture")
+ proto = struct.unpack_from("!H", raw, 12)[0]
+ off, stack, aux, checksum_partial = 14, [], (), False
+ for level, kind, data in ancillary:
+ if level == 263 and kind == 8:
+ require(len(data) >= 20, "short PACKET_AUXDATA")
+ status, length, snap, mac, net, tci, tpid = struct.unpack("IIIHHHH", data[:20])
+ require(length == snap, "truncated packet auxiliary length")
+ checksum_partial = bool(status & 8)
+ if status & 0x10:
+ aux = ((tpid if status & 0x40 else 0x8100, tci),)
+ while proto in (0x8100, 0x88a8):
+ require(off + 4 <= len(raw), "short VLAN header")
+ tci, inner = struct.unpack_from("!HH", raw, off)
+ stack.append((proto, tci))
+ proto, off = inner, off + 4
+ row = dict(raw=raw, size=len(raw) + 4 * len(aux), vlans=aux + tuple(stack),
+ aux=aux, inband=tuple(stack), src=raw[6:12], dst=raw[:6],
+ proto=proto, seq=None, path=None, payload=b"", ports=(), tcpseq=None,
+ supervision=None, supseq=None, tag=False, encoding="plain", checksum_partial=checksum_partial)
+ end = len(raw)
+ if proto == 0x892f:
+ require(off + 6 <= end, "short HSR tag")
+ field, seq, proto = struct.unpack_from("!HHH", raw, off)
+ row.update(seq=seq, path=field >> 12, tag=True, encoding="hsr")
+ expected = len(raw) - 14 - (4 if raw[12:14] == b"\x81\x00" else 0)
+ require(field & 0xfff == expected, "HSR LSDU size differs")
+ off += 6
+ elif end >= off + 6 and raw[-2:] == b"\x88\xfb":
+ seq, field, suffix = struct.unpack_from("!HHH", raw, end - 6)
+ expected = len(raw) - 14 - (4 if raw[12:14] == b"\x81\x00" else 0)
+ require(field & 0xfff == expected, "PRP LSDU size differs")
+ row.update(seq=seq, path=field >> 12, tag=True, encoding="prp")
+ end -= 6
+ row["proto"] = proto
+ row["network"] = raw[off:end]
+ if raw[:3] == b"\x01\x15\x4e" and proto == 0x88fb:
+ require(off + 12 <= end, "short supervision payload")
+ version, seq, kind, size = struct.unpack_from("!HHBB", raw, off)
+ require(kind in (20, 21, 22, 23) and size in (6, 12), "invalid supervision TLV")
+ row.update(supervision=raw[off + 6:off + 12], supseq=seq)
+ return row
+ if proto != 0x0800:
+ return row
+ require(off + 20 <= end and raw[off] >> 4 == 4, "invalid IPv4 header")
+ ihl, total = (raw[off] & 15) * 4, struct.unpack_from("!H", raw, off + 2)[0]
+ require(ihl >= 20 and total >= ihl and off + total <= end, "invalid IP extent")
+ protocol, l4 = raw[off + 9], off + ihl
+ ipend = off + total
+ row.update(ip_src=raw[off + 12:off + 16], ip_dst=raw[off + 16:off + 20])
+ if protocol == 6:
+ require(l4 + 20 <= ipend, "short TCP header")
+ size = (raw[l4 + 12] >> 4) * 4
+ require(size >= 20 and l4 + size <= ipend, "invalid TCP extent")
+ row.update(ports=struct.unpack_from("!HH", raw, l4),
+ tcpseq=struct.unpack_from("!I", raw, l4 + 4)[0],
+ tcpflags=raw[l4 + 13], payload=raw[l4 + size:ipend])
+ elif protocol == 17:
+ require(l4 + 8 <= ipend, "short UDP header")
+ size = struct.unpack_from("!H", raw, l4 + 4)[0]
+ require(size >= 8 and l4 + size == ipend, "invalid UDP extent")
+ row.update(ports=struct.unpack_from("!HH", raw, l4), payload=raw[l4 + 8:ipend])
+ return row
+
+
+class Capture:
+ def __init__(self, lab, ns, dev, pkttype=None):
+ self.lab, self.rows, self.seen = lab, [], 0
+ self.ns, self.dev, self.packets = ns, dev, None
+ self.pkttype, self.name = pkttype, ns + "/" + dev
+ self.complete, self.drops, self.truncated = False, None, 0
+ self.last_seen = time.monotonic()
+ self.socket = lab.sock(ns, socket.AF_PACKET, socket.SOCK_RAW, socket.htons(3))
+ self.socket.setsockopt(socket.SOL_SOCKET, socket.SO_RCVBUF, 8 * 1024 * 1024)
+ self.socket.setsockopt(263, 8, 1)
+ self.socket.bind((dev, 0))
+ self.socket.setblocking(False)
+ lab.selector.register(self.socket, selectors.EVENT_READ, self)
+
+ def drain(self):
+ count = 0
+ while count < 256:
+ self.lab.remaining()
+ try:
+ raw, anc, flags, addr = self.socket.recvmsg(131072, 256)
+ except BlockingIOError:
+ return count, True
+ self.seen += 1
+ self.last_seen = time.monotonic()
+ count += 1
+ if flags & (socket.MSG_TRUNC | socket.MSG_CTRUNC):
+ self.truncated += 1
+ raise ValueError(self.name + ": capture truncated")
+ if self.pkttype is None or addr[2] == self.pkttype:
+ self.rows.append(decode(raw, anc))
+ return count, False
+
+ def validate(self):
+ require(self.complete and self.drops == 0 and self.truncated == 0,
+ self.name + ": incomplete/dropped/truncated capture")
+
+
+class Lab:
+ def __init__(self, deadline):
+ self.deadline = deadline
+ self.names, self.sockets, self.captures, self.threads = [], [], [], []
+ self.selector = selectors.DefaultSelector()
+ self.original = os.open("/proc/self/ns/net", os.O_RDONLY)
+ self.token, self.restore_failed = uuid.uuid4().hex[:10], False
+ self.stop = threading.Event()
+ self.libc = ctypes.CDLL(None, use_errno=True)
+
+ def remaining(self):
+ value = self.deadline - time.monotonic()
+ require(value > 0, "120-second work deadline expired")
+ return value
+
+ def run(self, args, check=True):
+ require(not self.restore_failed, "network namespace restoration failed")
+ result = subprocess.run([str(arg) for arg in args], text=True, capture_output=True,
+ timeout=min(5, self.remaining()))
+ require(not check or result.returncode == 0,
+ "%s: rc%d %s" % (args, result.returncode, result.stderr.strip()))
+ return result
+
+ def ip(self, ns, *args, check=True):
+ return self.run(["ip", "-n", ns, *args], check)
+
+ def cmd(self, ns, *args, check=True):
+ return self.run(["ip", "netns", "exec", ns, *args], check)
+
+ def features(self, ns, dev):
+ return features(self.cmd(ns, "ethtool", "-k", dev).stdout)
+
+ def ns(self, label):
+ name = "hsrp4-%s-%d" % (self.token, len(self.names))
+ require(not os.path.lexists("/var/run/netns/" + name), "namespace collision")
+ self.names.append(name)
+ self.run(["ip", "netns", "add", name])
+ self.ip(name, "link", "set", "dev", "lo", "up")
+ return name
+
+ def veth(self, ns, dev, peer_ns, peer):
+ self.ip(ns, "link", "add", "name", dev, "type", "veth", "peer", "name", peer,
+ "netns", peer_ns)
+
+ def setns(self, descriptor):
+ if self.libc.setns(descriptor, 0x40000000):
+ error = ctypes.get_errno()
+ raise OSError(error, os.strerror(error))
+
+ @contextmanager
+ def enter(self, ns):
+ descriptor = os.open("/var/run/netns/" + ns, os.O_RDONLY)
+ try:
+ self.setns(descriptor)
+ yield
+ finally:
+ try:
+ self.setns(self.original)
+ except BaseException:
+ self.restore_failed = True
+ raise
+ finally:
+ os.close(descriptor)
+
+ def sock(self, ns, family, kind, proto=0):
+ with self.enter(ns):
+ result = socket.socket(family, kind, proto)
+ self.sockets.append(result)
+ result.settimeout(min(3, self.remaining()))
+ return result
+
+ def capture(self, ns, dev, pkttype=None):
+ result = Capture(self, ns, dev, pkttype)
+ self.captures.append(result)
+ return result
+
+ def observe(self, seconds, pump=None):
+ end = time.monotonic() + seconds
+ while True:
+ self.remaining()
+ if pump is not None:
+ pump()
+ for key, events in self.selector.select(min(0.01, max(0, end - time.monotonic()))):
+ key.data.drain()
+ if time.monotonic() >= end:
+ return
+
+ def finish(self):
+ active = [cap for cap in self.captures if not cap.complete]
+ end = min(self.deadline, time.monotonic() + 2)
+ quiet = time.monotonic()
+ while time.monotonic() < end:
+ self.observe(0.01)
+ empty = True
+ for cap in active:
+ count, drained = cap.drain()
+ if count or not drained:
+ empty = False
+ quiet = max([quiet] + [cap.last_seen for cap in active])
+ if empty and time.monotonic() - quiet >= 0.1:
+ break
+ else:
+ raise ValueError("capture failed to drain within fixed deadline")
+ for ns in sorted({cap.ns for cap in active}):
+ require(time.monotonic() < end, "capture cutoff/drain deadline expired")
+ require(ns in self.names, "capture namespace is not owned")
+ self.ip(ns, "link", "set", "dev", "lo", "down")
+ require(time.monotonic() < end, "capture cutoff/drain deadline expired")
+ state = link_state(self, ns, "lo")
+ require(time.monotonic() < end, "capture cutoff/drain deadline expired")
+ require("UP" not in state["flags"],
+ "capture cutoff loopback is still up")
+ # Rebinding to a different down device synchronizes the old receive hook.
+ # Its ENETDOWN socket error must be cleared before draining queued frames.
+ for cap in active:
+ require(time.monotonic() < end, "capture cutoff/drain deadline expired")
+ require(cap.dev != "lo", "capture already uses the cutoff device")
+ cap.socket.bind(("lo", 0))
+ require(cap.socket.getsockopt(socket.SOL_SOCKET, socket.SO_ERROR) == errno.ENETDOWN,
+ cap.name + ": capture cutoff did not report ENETDOWN")
+ for cap in active:
+ while True:
+ require(time.monotonic() < end, "capture cutoff/drain deadline expired")
+ if cap.drain()[1]:
+ break
+ require(time.monotonic() < end, "capture cutoff/drain deadline expired")
+ packets, cap.drops = struct.unpack("II", cap.socket.getsockopt(263, 6, 8))
+ cap.packets = packets
+ require(time.monotonic() < end, "capture cutoff/drain deadline expired")
+ require(packets == cap.seen + cap.drops,
+ "%s: queued capture remains packets=%d seen=%d drops=%d" %
+ (cap.name, packets, cap.seen, cap.drops))
+ cap.complete = True
+ cap.validate()
+ self.selector.unregister(cap.socket)
+ cap.socket.close()
+ print("# capture %s packets=%d drops=0 truncated=0 complete=1" %
+ (cap.name, cap.seen), flush=True)
+
+ def close(self):
+ self.stop.set()
+ end, errors = time.monotonic() + 20, []
+ for sock in self.sockets:
+ sock.close()
+ for thread in self.threads:
+ if thread.ident is not None:
+ thread.join(max(0, end - time.monotonic()))
+ if thread.is_alive():
+ errors.append("sender thread did not stop")
+ self.selector.close()
+ if self.restore_failed:
+ errors.append("namespace restoration failed; no namespace commands attempted")
+ else:
+ for name in reversed(self.names):
+ if time.monotonic() >= end:
+ errors.append("cleanup deadline expired with owned namespaces remaining")
+ break
+ try:
+ result = subprocess.run(["ip", "netns", "del", name],
+ capture_output=True, text=True,
+ timeout=min(3, end - time.monotonic()))
+ require(result.returncode == 0 and not os.path.lexists("/var/run/netns/" + name),
+ "namespace cleanup failed: " + name + " " + result.stderr)
+ except BaseException as error:
+ errors.append(str(error))
+ os.close(self.original)
+ directory = os.getenv("HSR_CAPTURE_EVIDENCE")
+ if directory:
+ try:
+ for number, cap in enumerate(self.captures):
+ require(time.monotonic() < end, "cleanup/evidence deadline expired")
+ path = os.path.join(directory, "%s-%d.json" % (self.token, number))
+ record = dict(name=cap.name, seen=cap.seen, packets=cap.packets, drops=cap.drops,
+ truncated=cap.truncated, complete=cap.complete,
+ pkttype=cap.pkttype, rows=cap.rows)
+ with open(path, "x") as stream:
+ json.dump(record, stream, default=lambda value: value.hex())
+ except BaseException as error:
+ errors.append("capture evidence: " + str(error))
+ require(not errors, "; ".join(errors))
+
+
+def coverage(frames, expected):
+ counts = array("I", [0]) * len(expected)
+ for offset, payload in frames:
+ require(payload and offset >= 0 and offset + len(payload) <= len(expected),
+ "invalid stream extent")
+ require(payload == expected[offset:offset + len(payload)], "stream payload differs")
+ for position in range(offset, offset + len(payload)):
+ counts[position] += 1
+ return counts
+
+
+def protocol_order(rows, wrap=False, key="seq"):
+ numbers = [row[key] for row in rows]
+ require(numbers and all(value is not None for value in numbers), "missing protocol sequence")
+ require(all(0 < ((right - left) & 0xffff) < 0x8000
+ for left, right in zip(numbers, numbers[1:])), "protocol sequence out of order")
+ if wrap:
+ require(any(left > right for left, right in zip(numbers, numbers[1:])), "no actual sequence wrap")
+
+
+def wire_rows(rows, port, stack, tagged=True):
+ data = [row for row in rows if row["ports"] and row["ports"][1] == port and row["payload"]]
+ require(data, "no captured flow payload")
+ for row in data:
+ require(row["tag"] == tagged, "flow tag/RCT missing or unexpected")
+ require(row["encoding"] == ("hsr" if tagged else "plain"), "wrong protocol tag kind")
+ require(row["vlans"] == stack, "flow VLAN stack differs")
+ require(row["size"] <= 1520 + 4 * len(stack), "unsegmented output frame")
+ require(len({row["src"] for row in data}) == 1 and
+ len({row["dst"] for row in data}) == 1, "flow MAC addresses differ within lane")
+ if tagged:
+ protocol_order(data)
+ return data
+
+
+def copy_identity(row):
+ # Lane/path and source substitution differ legitimately between the two ports.
+ return (row["seq"], row["vlans"], row["dst"], row["proto"], row["network"])
+
+
+def verify_stream(captures, port, payload, stack=(), local=False, inband=False, ingress=None):
+ for cap in captures:
+ cap.validate()
+ source = [row for row in captures[0].rows if row["ports"] and
+ row["ports"][1] == port and row["payload"]]
+ require(source, "sender TCP capture empty")
+ syn = [row for row in captures[0].rows if row["ports"] and
+ row["ports"][1] == port and row.get("tcpflags", 0) & 2]
+ require(syn, "missing captured TCP initial sequence")
+ isn = (syn[0]["tcpseq"] + 1) & 0xffffffff
+ source_counts = coverage([((row["tcpseq"] - isn) & 0xffffffff, row["payload"])
+ for row in source], payload)
+ require(0 not in source_counts, "sender did not send complete stream")
+ aggregate = [row for row in source if len(row["payload"]) > 1500]
+ if ingress is not None:
+ ingress.validate()
+ inputs = [row for row in ingress.rows if row["ports"] and
+ row["ports"][1] == port and row["payload"]]
+ require(all(row["vlans"] == stack and row["encoding"] == "plain" for row in inputs),
+ "HSR/PRP ingress classification/VLAN differs")
+ require(coverage([((row["tcpseq"] - isn) & 0xffffffff, row["payload"])
+ for row in inputs], payload) == source_counts,
+ "sender/actual ingress stream occurrence differs")
+ require(not aggregate or any(len(row["payload"]) > 1500 for row in inputs),
+ "no actual ingress aggregate")
+ require(all(row["vlans"] == stack for row in source), "sender VLAN stack differs")
+ if stack:
+ require(all(row["aux"] == (() if inband else stack[:1]) and
+ row["inband"] == (stack if inband else stack[1:]) for row in aggregate),
+ "sender aggregate VLAN encoding differs")
+ if local:
+ outputs = [wire_rows(captures[1].rows, port, stack, tagged=False)]
+ require(all(row["ports"][1] == port for row in outputs[0]), "unexpected local flow direction")
+ else:
+ outputs = [wire_rows(cap.rows, port, stack) for cap in captures[1:]]
+ for index, rows in enumerate(outputs):
+ require(all(row["path"] == index for row in rows), "wrong HSR lane/path")
+ require([copy_identity(row) for row in outputs[0]] ==
+ [copy_identity(row) for row in outputs[1]], "LAN copies differ")
+ for rows in outputs:
+ observed = coverage([((row["tcpseq"] - isn) & 0xffffffff, row["payload"])
+ for row in rows], payload)
+ require(observed == source_counts, "lost/duplicate stream occurrence (including retransmission)")
+ # Missing aggregation is a prerequisite gap only after transfer and all other checks.
+ if not aggregate:
+ raise Skip("no sender GSO aggregate observed")
+ print("# TCP bytes=%d sent-occurrences=%d aggregates=%d outputs=%s" %
+ (len(payload), sum(source_counts), len(aggregate), [len(rows) for rows in outputs]), flush=True)
+
+
+def verify_udp(case, captures, expected, ports, vlans=(), expect_wrap=False,
+ input_order=False, producer_order=False):
+ lanes = []
+ for index, cap in enumerate(captures):
+ cap.validate()
+ rows = [row for row in cap.rows if row["ports"] and row["ports"][1] in ports]
+ require(Counter(row["payload"] for row in rows) == Counter(expected),
+ case + ": UDP payload/ID set differs")
+ if input_order:
+ require([row["payload"] for row in rows] == expected, case + ": input order differs")
+ if producer_order:
+ for producer in (1, 2):
+ observed = [row["payload"] for row in rows if row["payload"][6] == producer]
+ require(observed == [body for body in expected if body[6] == producer],
+ case + ": producer order differs")
+ for row in rows:
+ require(row["tag"] and row["size"] <= 1520 + 4 * len(vlans),
+ case + ": missing tag or unsegmented frame")
+ require(row["vlans"] == vlans, case + ": VLAN stack differs")
+ require(row["path"] == index, case + ": wrong HSR lane/path")
+ require(row["encoding"] == "hsr", case + ": wrong protocol tag kind")
+ protocol_order(rows, expect_wrap)
+ lanes.append(rows)
+ require([copy_identity(row) for row in lanes[0]] == [copy_identity(row) for row in lanes[1]],
+ case + ": unequal LAN copies")
+ print("# %s expected=%d lanes=%s" % (case, len(expected), [len(rows) for rows in lanes]), flush=True)
+ return lanes
+
+
+def local_pair(lab, ns, dev):
+ lab.veth(ns, dev, ns, dev + "p")
+ for name in (dev, dev + "p"):
+ lab.ip(ns, "link", "set", "dev", name, "up")
+
+
+def link_state(lab, ns, dev):
+ result = lab.ip(ns, "-j", "-d", "link", "show", "dev", dev)
+ rows = json.loads(result.stdout)
+ require(len(rows) == 1, "%s: missing link description" % dev)
+ state = rows[0]
+ require(isinstance(state.get("ifindex"), int), "missing ifindex")
+ require(isinstance(state.get("promiscuity"), int),
+ "%s: missing promiscuity count" % dev)
+ return state
+
+
+def set_gro(lab, ns, dev, wanted, active=None):
+ result = lab.cmd(ns, "ethtool", "-K", dev, "gro", wanted, check=False)
+ forced = (active == "off" and wanted == "on" and result.returncode == 1
+ and result.stdout == "Actual changes:\nrx-gro: off [requested on]\n"
+ and result.stderr == "Could not change any device features\n")
+ require(result.returncode == 0 or forced,
+ "%s: GRO setter rc%d %s" % (dev, result.returncode, result.stderr.strip()))
+ features = lab.features(ns, dev)
+ require("generic-receive-offload" in features,
+ "%s: missing software GRO feature" % dev)
+ actual = features["generic-receive-offload"][0]
+ require(actual == (active if active is not None else wanted),
+ "%s: requested GRO %s, observed %s" % (dev, wanted, actual))
+
+
+def hsr_attach(lab, ns, master, first, second, interlink=None):
+ args = ["link", "add", "name", master, "type", "hsr",
+ "slave1", first, "slave2", second, "version", "1"]
+ if interlink is not None:
+ args.extend(("interlink", interlink))
+ result = lab.ip(ns, *args, check=False)
+ if result.returncode and interlink is not None:
+ # Tool-option absence needs a successful ordinary-HSR control.
+ control = lab.ip(ns, "link", "add", "ctl", "type", "hsr",
+ "slave1", first, "slave2", second,
+ "version", "1", check=False)
+ if control.returncode == 0:
+ lab.ip(ns, "link", "del", "ctl")
+ error = result.stderr.lower()
+ absent = ("interlink" in error and any(word in error for word in
+ ("unknown command", "unknown option", "unknown argument", 'what is "interlink"')))
+ require(control.returncode == 0,
+ "interlink failed and ordinary HSR control failed: %s / %s"
+ % (result.stderr.strip(), control.stderr.strip()))
+ if absent and "rtnetlink" not in error:
+ raise Skip("iproute2 has no interlink option")
+ require(result.returncode == 0,
+ "HSR attach failed: " + result.stderr.strip())
+ lab.ip(ns, "link", "set", "dev", master, "up")
+
+
+def policy(lab):
+ """All roles, initial preferences and last requests, twice each."""
+ ns = lab.ns("policy")
+ for dev in ("m", "p", "q"):
+ local_pair(lab, ns, dev)
+ for role in ("A", "B", "interlink"):
+ for initial in ("on", "off"):
+ for last in ("on", "off"):
+ label = "%s initial=%s last=%s" % (role, initial, last)
+ before = link_state(lab, ns, "m")
+ for cycle in range(2):
+ set_gro(lab, ns, "m", initial)
+ first = "m" if role == "A" else "p"
+ second = "m" if role == "B" else "q"
+ hsr_attach(lab, ns, "hsr0", first, second,
+ "m" if role == "interlink" else None)
+ active = lab.features(ns, "m")
+ require(active["generic-receive-offload"][0] == "off",
+ label + ": GRO remained enabled after attach")
+ for wanted in ("on", "off", last):
+ set_gro(lab, ns, "m", wanted, active="off")
+ lab.ip(ns, "link", "del", "hsr0")
+ after = link_state(lab, ns, "m")
+ require("master" not in after,
+ label + ": upper remains after detach")
+ require(after["promiscuity"] == before["promiscuity"],
+ label + ": promiscuity changed after detach")
+ require(lab.features(ns, "m")[
+ "generic-receive-offload"][0] == last,
+ label + ": last GRO preference was not restored")
+ print("# policy %s cycles=2" % label, flush=True)
+
+
+def master_gso(lab):
+ ns = lab.ns("master-mask")
+ for dev in ("a", "b"):
+ local_pair(lab, ns, dev)
+ hsr_attach(lab, ns, "hsr0", "a", "b")
+ features = lab.features(ns, "hsr0")
+ names = [name for name in features if
+ ("segmentation" in name or name in
+ ("tx-gso-robust", "tx-gso-partial", "tx-gso-list"))
+ and name != "generic-segmentation-offload"]
+ require(names, "master GSO type mask is absent from feature report")
+ enabled = [name for name in names if features[name][0] != "off"]
+ require(not enabled, "master retains GSO types: " + ", ".join(enabled))
+ print("# master GSO types off: " + ", ".join(sorted(names)), flush=True)
+
+
+def rollback_bridge(lab):
+ """B's existing bridge causes A's earlier attachment to unwind."""
+ ns = lab.ns("rollback-bridge")
+ for dev in ("a", "b", "c"):
+ local_pair(lab, ns, dev)
+ set_gro(lab, ns, "a", "on")
+ before = link_state(lab, ns, "a")
+ lab.ip(ns, "link", "add", "br0", "type", "bridge")
+ lab.ip(ns, "link", "set", "dev", "b", "master", "br0")
+ lab.ip(ns, "link", "set", "dev", "br0", "up")
+ result = lab.ip(ns, "link", "add", "hsr0", "type", "hsr",
+ "slave1", "a", "slave2", "b", "version", "1",
+ check=False)
+ require(result.returncode != 0, "second-port attachment unexpectedly passed")
+ require(link_state(lab, ns, "b").get("master") == "br0",
+ "failed second port lost its original bridge")
+ after = link_state(lab, ns, "a")
+ require("master" not in after, "rolled-back A retains an upper")
+ require(after["promiscuity"] == before["promiscuity"],
+ "rolled-back A retains a promiscuity reference")
+ require(lab.features(ns, "a")["generic-receive-offload"][0] == "on",
+ "rolled-back A did not restore GRO")
+ # This positive control detects a residual role or RX handler on A.
+ hsr_attach(lab, ns, "hsr1", "a", "c")
+ lab.ip(ns, "link", "del", "hsr1")
+ require(lab.features(ns, "a")["generic-receive-offload"][0] == "on",
+ "GRO not restored after the rollback positive control")
+ print("# bridge second-port rollback restored A; reattach passed", flush=True)
+
+
+def multicast_flow(lab, ns, phase):
+ """Exact multicast receipt on both lower and macvlan, with no workers."""
+ captures = [lab.capture(ns, dev) for dev in ("m", "mv")]
+ sender = lab.sock(ns, socket.AF_PACKET, socket.SOCK_RAW, 0)
+ sender.bind(("mp", 0))
+ sender.setblocking(False)
+ source = bytes.fromhex(link_state(lab, ns, "mp")["address"].replace(":", ""))
+ prefix = ("hsr-macvlan-%s:" % phase).encode()
+ frames = [bytes.fromhex("01005e000009") + source + b"\x88\xb5" +
+ prefix + struct.pack("!H", index) + b"m" * 96
+ for index in range(50)]
+ for frame in frames:
+ require(sender.send(frame) == len(frame), "short multicast send")
+ lab.observe(0.002)
+ # finish() establishes real empty queues and valid loss/truncation stats.
+ lab.finish()
+ expected = Counter(frames)
+ for dev, capture in zip(("m", "mv"), captures):
+ observed = Counter(row["raw"] for row in capture.rows
+ if row["raw"][12:14] == b"\x88\xb5" and
+ row["raw"][14:].startswith(prefix))
+ require(observed == expected,
+ "%s %s multicast copies differ: missing=%d excess=%d" %
+ (phase, dev, sum((expected - observed).values()),
+ sum((observed - expected).values())))
+ print("# macvlan %s multicast lower=50/50 upper=50/50" % phase, flush=True)
+
+
+def rollback_macvlan(lab):
+ """A non-master upper occupies the RX handler; delivery survives."""
+ ns = lab.ns("rollback-macvlan")
+ for dev in ("m", "p"):
+ local_pair(lab, ns, dev)
+ result = lab.ip(ns, "link", "add", "link", "m", "name", "mv",
+ "type", "macvlan", "mode", "bridge", check=False)
+ if result.returncode:
+ error = result.stderr.lower()
+ if "unknown device type" in error or "operation not supported" in error:
+ raise Skip("macvlan device kind unavailable: " + result.stderr.strip())
+ require(False, "macvlan creation failed: " + result.stderr.strip())
+ lab.ip(ns, "link", "set", "dev", "mv", "up")
+ lab.ip(ns, "maddress", "add", "01:00:5e:00:00:09", "dev", "mv")
+ set_gro(lab, ns, "m", "on")
+ before = link_state(lab, ns, "m")
+ require("master" not in before, "macvlan lower already has a master")
+
+ def association():
+ index = lab.cmd(ns, "cat", "/sys/class/net/mv/iflink").stdout.strip()
+ require(index == str(link_state(lab, ns, "m")["ifindex"]),
+ "macvlan no longer belongs to the original lower")
+ membership = lab.ip(ns, "maddress", "show", "dev", "mv").stdout.lower()
+ require("01:00:5e:00:00:09" in membership,
+ "macvlan multicast membership missing")
+
+ association()
+ multicast_flow(lab, ns, "before")
+ result = lab.ip(ns, "link", "add", "hsr0", "type", "hsr",
+ "slave1", "m", "slave2", "p", "version", "1",
+ check=False)
+ require(result.returncode != 0 and "busy" in result.stderr.lower(),
+ "expected busy RX-handler failure: " + result.stderr.strip())
+ after = link_state(lab, ns, "m")
+ require("master" not in after, "failed attachment left an HSR upper")
+ require(after["promiscuity"] == before["promiscuity"],
+ "failed attachment changed promiscuity")
+ require(lab.features(ns, "m")["generic-receive-offload"][0] == "on",
+ "failed RX-handler attachment did not restore GRO")
+ association()
+ multicast_flow(lab, ns, "after")
+
+
+GSO_ROWS = (
+ ("plain", (), False, "A"),
+ ("s-tag-inband", ((0x88a8, 100),), True, "A"),
+ ("c-tag-accelerated", ((0x8100, 50),), False, "A"),
+ ("s-tag-accelerated", ((0x88a8, 100),), False, "B"),
+ ("sc-nested-accelerated", ((0x88a8, 200), (0x8100, 30)), False, "B"),
+ ("cc-nested-inband", ((0x8100, 40), (0x8100, 20)), True, "A"),
+)
+
+
+def vlan_chain(lab, ns, parent, stack, prefix):
+ for depth, (protocol, identifier) in enumerate(stack):
+ dev = "%s%d" % (prefix, depth)
+ lab.ip(ns, "link", "add", "link", parent, "name", dev,
+ "type", "vlan", "id", str(identifier), "protocol",
+ "802.1AD" if protocol == 0x88a8 else "802.1Q")
+ lab.ip(ns, "link", "set", "dev", dev, "up")
+ parent = dev
+ return parent
+
+
+def vlan_tx_mode(lab, ns, lower, stack, inband):
+ if not stack:
+ return
+ aliases = (("tx-vlan-stag-hw-insert", "tx-vlan-stag-offload")
+ if stack[0][0] == 0x88a8 else
+ ("tx-vlan-offload", "tx-vlan-hw-insert"))
+ features = lab.features(ns, lower)
+ name = next((alias for alias in aliases if alias in features), None)
+ if name is None:
+ raise Skip("no VLAN insertion feature: " + ", ".join(aliases))
+ wanted = "off" if inband else "on"
+ if features[name][1] and features[name][0] != wanted:
+ raise Skip("fixed VLAN insertion feature cannot be " + wanted)
+ lab.cmd(ns, "ethtool", "-K", lower, name, wanted)
+ require(lab.features(ns, lower)[name][0] == wanted,
+ "VLAN insertion mode did not match its requested state")
+
+
+def ordered_payload(producer, index, size=64):
+ prefix = b"HSRORD" + bytes([producer]) + struct.pack("!I", index)
+ return prefix + bytes([producer ^ 0xA5]) * (size - len(prefix))
+
+
+def ordered_socket(lab, ns, address):
+ sock = lab.sock(ns, socket.AF_INET, socket.SOCK_DGRAM)
+ sock.bind((address, 0))
+ sock.setsockopt(socket.IPPROTO_IP, socket.IP_MULTICAST_TTL, 1)
+ sock.settimeout(1)
+ return sock
+
+
+def ordered_lanes(lab, topology):
+ return [lab.capture(topology["dut"], lane, pkttype=4)
+ for lane in ("da", "db")]
+
+
+def ordered_send(lab, sock, producer, first, count, port, segment=0):
+ expected = []
+ if segment:
+ sock.setsockopt(socket.IPPROTO_UDP, 103, segment)
+ blocks = 16 if segment else 1
+ size = segment or 64
+ for number in range(count):
+ payloads = [ordered_payload(producer, first + number * blocks + n,
+ size) for n in range(blocks)]
+ body = b"".join(payloads)
+ require(sock.sendto(body, ("239.0.0.9", port)) == len(body),
+ "short UDP send")
+ expected.extend(payloads)
+ # Keep capture receive queues bounded throughout long wrap warmups.
+ if number % 8 == 0:
+ lab.observe(0)
+ lab.observe(0.1)
+ return expected
+
+
+def ordered_data(lab, gso=False):
+ topology = topo(lab)
+ caps = ordered_lanes(lab, topology)
+ local = ordered_socket(lab, topology["dut"], "100.64.3.1")
+ san = ordered_socket(lab, topology["tx"], "100.64.2.1")
+ # Observe interlink aggregates before the two tagged LAN copies expand.
+ sent = lab.capture(topology["tx"], "t0", pkttype=4)
+ ingress = lab.capture(topology["dut"], "dil")
+ expected = ordered_send(lab, local, 1, 0, 32, 5100)
+ if gso:
+ try:
+ san.setsockopt(socket.IPPROTO_UDP, 103, 512)
+ except OSError as error:
+ if error.errno in (errno.ENOPROTOOPT, errno.EOPNOTSUPP):
+ raise Skip("UDP GSO unavailable") from error
+ raise
+ expected += ordered_send(lab, san, 2, 0, 16, 5101, segment=512 if gso else 0)
+ expected += ordered_send(lab, local, 1, 32, 32, 5100)
+ lab.finish()
+ verify_udp("local and aggregate boundaries", caps, expected,
+ ports={5100, 5101}, input_order=True)
+ inputs = [row for row in sent.rows if row["ports"] and row["ports"][1] == 5101]
+ expected_inputs = [b"".join(ordered_payload(2, n * 16 + part, 512)
+ for part in range(16)) for n in range(16)] if gso else expected[32:48]
+ aggregated = any(len(row["payload"]) > 1500 for row in inputs)
+ if gso and not aggregated:
+ expected_inputs = expected[32:-32]
+ require([row["payload"] for row in inputs] == expected_inputs,
+ "UDP sender input set/order differs")
+ incoming = [row for row in ingress.rows if row["ports"] and row["ports"][1] == 5101]
+ require([row["payload"] for row in incoming] == expected_inputs,
+ "actual UDP ingress input set/order differs")
+ require(all(row["encoding"] == "plain" and not row["vlans"] for row in inputs + incoming),
+ "UDP input classification/VLAN differs")
+ if gso and not aggregated:
+ raise Skip("interlink UDP GSO aggregate not observed")
+
+
+def ordered_supervision(lab):
+ topology = topo(lab)
+ caps = ordered_lanes(lab, topology)
+ san = ordered_socket(lab, topology["tx"], "100.64.2.1")
+ expected = ordered_send(lab, san, 3, 0, 8, 5110)
+ # Both timers have time for multiple life checks; no fake worker state.
+ lab.observe(10)
+ lab.finish()
+ verify_udp("proxy-learning inputs", caps, expected, ports={5110},
+ input_order=True)
+ verify_supervision(caps, topology["mac_master"], topology["mac_tx"])
+
+
+def ordered_concurrent(lab):
+ cpus = sorted(os.sched_getaffinity(0))[:2]
+ if len(cpus) != 2:
+ raise Skip("two functional producers need two allowed CPUs")
+ topology = topo(lab)
+ caps = ordered_lanes(lab, topology)
+ sockets = [ordered_socket(lab, topology["tx"], "100.64.2.1"),
+ ordered_socket(lab, topology["dut"], "100.64.3.1")]
+ barrier = threading.Barrier(3)
+ stop = threading.Event()
+ errors = []
+ windows = [None, None]
+ until = min(lab.deadline, time.monotonic() + 15)
+
+ def producer(index):
+ try:
+ # Socket namespace binding was completed before threads existed.
+ os.sched_setaffinity(0, {cpus[index]})
+ barrier.wait(timeout=2)
+ start = time.monotonic()
+ for number in range(128):
+ require(not stop.is_set() and not lab.stop.is_set() and time.monotonic() < until,
+ "producer work deadline")
+ body = ordered_payload(index + 1, number)
+ require(sockets[index].sendto(body,
+ ("239.0.0.9", 5120 + index)) == len(body),
+ "short producer send")
+ time.sleep(0.001)
+ windows[index] = (start, time.monotonic())
+ except Exception as error:
+ errors.append(error)
+ stop.set()
+ barrier.abort()
+
+ threads = [threading.Thread(target=producer, args=(index,))
+ for index in range(2)]
+ unfinished = []
+ try:
+ for thread in threads:
+ lab.threads.append(thread)
+ thread.start()
+ barrier.wait(timeout=2)
+ while any(thread.is_alive() for thread in threads):
+ require(time.monotonic() < until, "producer join deadline")
+ lab.observe(0.01)
+ finally:
+ stop.set()
+ barrier.abort()
+ join_end = min(lab.deadline, time.monotonic() + 2)
+ for thread in threads:
+ if thread.ident is not None:
+ thread.join(max(0, join_end - time.monotonic()))
+ unfinished = [thread for thread in threads if thread.is_alive()]
+ if unfinished:
+ print("# cleanup failure: sender thread remained alive")
+ if errors:
+ raise errors[0]
+ require(not unfinished, "sender thread remained alive")
+ require(all(windows) and windows[0][0] < windows[1][1] and
+ windows[1][0] < windows[0][1], "producer send windows do not overlap")
+ print("# functional producers CPUs=%s windows=%s" % (cpus, windows))
+ lab.finish()
+ expected = [ordered_payload(producer, number)
+ for producer in (1, 2) for number in range(128)]
+ verify_udp("concurrent master/interlink inputs", caps, expected,
+ ports={5120, 5121}, producer_order=True)
+
+
+def ordered_wrap(lab):
+ topology = topo(lab)
+ caps = ordered_lanes(lab, topology)
+ sender = ordered_socket(lab, topology["tx"], "100.64.2.1")
+ probe = ordered_send(lab, sender, 4, 0, 1, 5130)
+ # This only selects the warmup count; final loss-free captures certify it.
+ rows = [[row for row in cap.rows if row["ports"] and row["ports"][1] == 5130]
+ for cap in caps]
+ require(all(len(lane) == 1 for lane in rows) and
+ rows[0][0]["seq"] == rows[1][0]["seq"], "wrap probe missing or unequal")
+ start = rows[0][0]["seq"]
+ warm_count = (65400 - ((start + 1) & 0xFFFF)) & 0xFFFF
+ warm = ordered_send(lab, sender, 5, 0, warm_count, 5131)
+ # Separate IDs/ports prevent sequence-space reuse from hiding a loss.
+ expected = ordered_send(lab, sender, 6, 0, 256, 5132)
+ lab.finish()
+ verify_udp("wrap start", caps, probe, ports={5130}, input_order=True)
+ if warm:
+ verify_udp("wrap warmup", caps, warm, ports={5131}, input_order=True)
+ verify_udp("actual protocol wrap", caps, expected, ports={5132},
+ expect_wrap=True, input_order=True)
+
+
+ORDERED_CASES = (("local-data-order", ordered_data),
+ ("local-data-aggregate-boundaries", lambda lab: ordered_data(lab, True)),
+ ("local-proxy-supervision", ordered_supervision),
+ ("two-functional-producers", ordered_concurrent),
+ ("actual-sequence-wrap", ordered_wrap))
+
+
+def topo(lab):
+ tx, dut, rx = [lab.ns(label) for label in ("tx", "dut", "rx")]
+ for ns, dev, peer_ns, peer in ((tx, "t0", dut, "dil"),
+ (dut, "da", rx, "ra"), (dut, "db", rx, "rb")):
+ lab.veth(ns, dev, peer_ns, peer)
+ for owner, link in ((ns, dev), (peer_ns, peer)):
+ lab.ip(owner, "link", "set", "dev", link, "up")
+ hsr_attach(lab, dut, "hm", "da", "db", "dil")
+ attributes = link_state(lab, dut, "hm")["linkinfo"]["info_data"]
+ require(attributes.get("proto") == 0 and attributes.get("version") == 1 and
+ all(attributes.get(key) == value for key, value in
+ (("slave1", "da"), ("slave2", "db"), ("interlink", "dil"))),
+ "HSR topology has unexpected protocol/version")
+ lab.ip(tx, "addr", "add", "100.64.2.1/24", "dev", "t0")
+ lab.ip(dut, "addr", "add", "100.64.3.1/24", "dev", "hm")
+ for ns, dev in ((tx, "t0"), (dut, "hm")):
+ lab.ip(ns, "route", "add", "239.0.0.0/8", "dev", dev)
+ return dict(tx=tx, dut=dut, rx=rx, master="hm",
+ mac_master=bytes.fromhex(link_state(lab, dut, "hm")["address"].replace(":", "")),
+ mac_tx=bytes.fromhex(link_state(lab, tx, "t0")["address"].replace(":", "")))
+
+
+def verify_supervision(captures, master, proxy):
+ observed = []
+ all_frames = []
+ for lane, cap in enumerate(captures):
+ cap.validate()
+ tagged = [row for row in cap.rows if row["encoding"] == "hsr"]
+ protocol_order(tagged)
+ all_frames.append([copy_identity(row) for row in tagged])
+ rows = [row for row in cap.rows if row["supervision"] is not None]
+ require(all(row["encoding"] == "hsr" and row["path"] == lane
+ for row in rows), "wrong supervision tag kind/path")
+ require(all(row["supervision"] in (master, proxy) for row in rows), "unexpected supervised MAC")
+ for mac in (master, proxy):
+ group = [row for row in rows if row["supervision"] == mac]
+ require(len(group) >= 2, "missing repeated local/proxy supervision TLV")
+ protocol_order(group, key="supseq")
+ protocol_order(rows)
+ observed.append([copy_identity(row) for row in rows])
+ require(observed[0] == observed[1], "supervision LAN copies differ")
+ require(all_frames[0] == all_frames[1], "data/supervision LAN order differs")
+ print("# supervision exact local/proxy TLV sets lanes=%s" %
+ [len(rows) for rows in observed], flush=True)
+
+
+def tcp_transfer(lab, sender_ns, receiver_ns, dev, source, destination, port, captures):
+ payload = bytes(range(256)) * 512
+ reply = b"HSR-TCP-REPLY"
+ listener = lab.sock(receiver_ns, socket.AF_INET, socket.SOCK_STREAM)
+ listener.bind((destination, port))
+ listener.listen(1)
+ client = lab.sock(sender_ns, socket.AF_INET, socket.SOCK_STREAM)
+ client.setsockopt(socket.SOL_SOCKET, socket.SO_BINDTODEVICE, dev.encode() + b"\0")
+ client.bind((source, 0))
+ client.connect((destination, port))
+ server, address = listener.accept()
+ lab.sockets.append(server)
+ for sock in (client, server):
+ sock.setblocking(False)
+ sent, received, echoed = 0, bytearray(), bytearray()
+ state = dict(reply=False, write_closed=False)
+ end = min(lab.deadline, time.monotonic() + 12)
+
+ def pump():
+ nonlocal sent
+ if sent < len(payload):
+ try:
+ sent += client.send(payload[sent:sent + 65536])
+ except BlockingIOError:
+ pass
+ elif not state["write_closed"]:
+ client.shutdown(socket.SHUT_WR)
+ state["write_closed"] = True
+ try:
+ data = server.recv(65536)
+ received.extend(data)
+ except BlockingIOError:
+ pass
+ if len(received) == len(payload) and not state["reply"]:
+ require(server.send(reply) == len(reply), "short TCP reply")
+ state["reply"] = True
+ try:
+ echoed.extend(client.recv(65536))
+ except BlockingIOError:
+ pass
+
+ while len(echoed) < len(reply) and time.monotonic() < end:
+ lab.observe(0.01, pump)
+ require(bytes(received) == payload and bytes(echoed) == reply,
+ "TCP application did not receive complete expected bytes")
+ for sock in (server, client, listener):
+ sock.close()
+ lab.observe(0.2)
+ lab.finish()
+ return payload, reply
+
+
+def tcp_interlink(lab):
+ topology = topo(lab)
+ tx, dut, rx = (topology[key] for key in ("tx", "dut", "rx"))
+ lab.ip(rx, "link", "add", "hr", "type", "hsr",
+ "slave1", "ra", "slave2", "rb", "version", "1")
+ lab.ip(rx, "link", "set", "dev", "hr", "up")
+ lab.ip(tx, "addr", "add", "10.67.0.2/24", "dev", "t0")
+ lab.ip(rx, "addr", "add", "10.67.0.3/24", "dev", "hr")
+ lab.cmd(tx, "ethtool", "-K", "t0", "tso", "on", "gso", "on")
+ caps = [lab.capture(tx, "t0", pkttype=4)] + ordered_lanes(lab, topology)
+ ingress = lab.capture(dut, "dil")
+ payload, reply = tcp_transfer(lab, tx, rx, "t0", "10.67.0.2", "10.67.0.3", 5200, caps + [ingress])
+ verify_stream(caps, 5200, payload, ingress=ingress)
+
+
+def tcp_prp(lab, row):
+ name, stack, inband, side = row
+ tx, dut = lab.ns("prp-tx"), lab.ns("prp-dut")
+ for lane in ("a", "b"):
+ lab.veth(tx, "t" + lane, dut, "d" + lane)
+ lab.ip(tx, "link", "set", "dev", "t" + lane, "up")
+ lab.ip(dut, "link", "set", "dev", "d" + lane, "up")
+ lab.ip(dut, "link", "add", "pm", "type", "hsr", "slave1", "da", "slave2", "db", "proto", "1")
+ lab.ip(dut, "link", "set", "dev", "pm", "up")
+ attributes = link_state(lab, dut, "pm")["linkinfo"]["info_data"]
+ require(attributes.get("proto") == 1 and
+ attributes.get("slave1") == "da" and attributes.get("slave2") == "db",
+ "PRP topology has unexpected protocol/ports")
+ lower = "t" + side.lower()
+ sender = vlan_chain(lab, tx, lower, stack, "txv")
+ receiver = vlan_chain(lab, dut, "pm", stack, "rxv")
+ lab.ip(tx, "addr", "add", "10.68.0.2/24", "dev", sender)
+ lab.ip(dut, "addr", "add", "10.68.0.1/24", "dev", receiver)
+ vlan_tx_mode(lab, tx, lower, stack, inband)
+ lab.cmd(tx, "ethtool", "-K", lower, "tso", "on", "gso", "on")
+ caps = [lab.capture(tx, lower, pkttype=4), lab.capture(dut, "pm", pkttype=0)]
+ ingress = lab.capture(dut, "d" + side.lower())
+ lanes = [lab.capture(dut, dev, pkttype=4) for dev in ("da", "db")]
+ payload, reply = tcp_transfer(lab, tx, dut, sender, "10.68.0.2", "10.68.0.1", 5201, caps + lanes + [ingress])
+ verify_stream(caps, 5201, payload, stack, local=True, inband=inband, ingress=ingress)
+ copies = []
+ for index, cap in enumerate(lanes):
+ cap.validate()
+ rows = [frame for frame in cap.rows if frame["ports"] and
+ frame["ports"][0] == 5201 and frame["payload"]]
+ require(len(rows) == 1 and rows[0]["payload"] == reply, "PRP reply not exactly once per LAN")
+ require(rows[0]["tag"] and rows[0]["vlans"] == stack and
+ rows[0]["size"] <= 1520 + 4 * len(stack), "invalid PRP reply tag/size/VLAN")
+ require(rows[0]["path"] == 10 + index, "wrong PRP LAN identifier")
+ require(rows[0]["encoding"] == "prp", "wrong PRP reply tag kind")
+ copies.append(copy_identity(rows[0]))
+ require(copies[0] == copies[1], "PRP reply LAN copies differ")
+ print("# PRP LAN %s local data segmented; reply each LAN=1/1" % name, flush=True)
+
+
+def interrupt(signum, frame):
+ raise InterruptedError("signal %d" % signum)
+
+
+def main():
+ mode = sys.argv[1] if len(sys.argv) == 2 else ""
+ require(mode in ("gro", "ordered"), "usage: hsr_gro_superpacket.py gro|ordered")
+ cases = [("GRO policy", policy), ("master GSO mask", master_gso),
+ ("bridge rollback", rollback_bridge), ("busy-handler rollback", rollback_macvlan),
+ ("natural TCP interlink", tcp_interlink)]
+ cases += [("PRP LAN " + row[0], lambda lab, row=row: tcp_prp(lab, row)) for row in GSO_ROWS]
+ if mode == "ordered":
+ cases = ORDERED_CASES
+ print("TAP version 13\n1..%d" % len(cases), flush=True)
+ missing = [tool for tool in ("ip", "ethtool") if shutil.which(tool) is None]
+ if os.geteuid() or missing:
+ for number, (name, function) in enumerate(cases, 1):
+ print("ok %d - %s # SKIP need root/ip/ethtool %s" % (number, name, missing), flush=True)
+ return 4
+ for sig in (signal.SIGINT, signal.SIGTERM):
+ signal.signal(sig, interrupt)
+ lab = Lab(time.monotonic() + 120)
+ failure, passed = None, 0
+ try:
+ for number, (name, function) in enumerate(cases, 1):
+ try:
+ function(lab)
+ print("ok %d - %s" % (number, name), flush=True)
+ passed += 1
+ except Skip as error:
+ print("ok %d - %s # SKIP %s" % (number, name, error), flush=True)
+ except BaseException as error:
+ failure = error
+ print("not ok %d - %s: %s" % (number, name, error), flush=True)
+ for remaining, (later, unused) in enumerate(cases[number:], number + 1):
+ print("not ok %d - %s: stopped after prior failure" % (remaining, later), flush=True)
+ break
+ finally:
+ for sig in (signal.SIGINT, signal.SIGTERM):
+ signal.signal(sig, signal.SIG_IGN)
+ try:
+ lab.close()
+ except BaseException as error:
+ print("# cleanup failure: " + str(error), flush=True)
+ if failure is None:
+ failure = error
+ return 1 if failure is not None else (0 if passed else 4)
+
+
+if __name__ == "__main__":
+ raise SystemExit(main())
diff --git a/tools/testing/selftests/net/hsr/hsr_gro_superpacket.sh b/tools/testing/selftests/net/hsr/hsr_gro_superpacket.sh
new file mode 100755
index 000000000000..9da711ad2ab2
--- /dev/null
+++ b/tools/testing/selftests/net/hsr/hsr_gro_superpacket.sh
@@ -0,0 +1,9 @@
+#!/bin/bash
+# SPDX-License-Identifier: GPL-2.0
+
+if ! command -v python3 > /dev/null; then
+ echo "SKIP: missing python3"
+ exit 4
+fi
+
+exec python3 -B "$(dirname "$0")/hsr_gro_superpacket.py" gro "$@"
diff --git a/tools/testing/selftests/net/hsr/hsr_ordered_forwarding.sh b/tools/testing/selftests/net/hsr/hsr_ordered_forwarding.sh
new file mode 100755
index 000000000000..b43218709455
--- /dev/null
+++ b/tools/testing/selftests/net/hsr/hsr_ordered_forwarding.sh
@@ -0,0 +1,9 @@
+#!/bin/bash
+# SPDX-License-Identifier: GPL-2.0
+
+if ! command -v python3 > /dev/null; then
+ echo "SKIP: missing python3"
+ exit 4
+fi
+
+exec python3 -B "$(dirname "$0")/hsr_gro_superpacket.py" ordered "$@"
--
2.53.0
prev parent reply other threads:[~2026-10-09 20:13 UTC|newest]
Thread overview: 5+ messages / expand[flat|nested] mbox.gz Atom feed top
2026-10-09 20:13 [PATCH net v7 0/4] net: hsr: fix super-packet forwarding and ordering Xin Xie
2026-10-09 20:13 ` [PATCH net v7 1/4] net: hsr: keep GRO disabled on HSR/PRP ports Xin Xie
2026-10-09 20:13 ` [PATCH net v7 2/4] net: hsr: preserve submission order without a forwarding lock Xin Xie
2026-10-09 20:13 ` [PATCH net v7 3/4] net: hsr: segment GSO before per-frame forwarding Xin Xie
2026-10-09 20:13 ` Xin Xie [this message]
Reply instructions:
You may reply publicly to this message via plain-text email
using any one of the following methods:
* Save the following mbox file, import it into your mail client,
and reply-to-all from there: mbox
Avoid top-posting and favor interleaved quoting:
https://en.wikipedia.org/wiki/Posting_style#Interleaved_style
* Reply using the --to, --cc, and --in-reply-to
switches of git-send-email(1):
git send-email \
--in-reply-to=20261009201324.17-5-xiexinet@gmail.com \
--to=xiexinet@gmail.com \
--cc=andrew+netdev@lunn.ch \
--cc=bigeasy@linutronix.de \
--cc=davem@davemloft.net \
--cc=edumazet@google.com \
--cc=fmaurer@redhat.com \
--cc=horms@kernel.org \
--cc=kees@kernel.org \
--cc=kuba@kernel.org \
--cc=linux-kernel@vger.kernel.org \
--cc=linux-kselftest@vger.kernel.org \
--cc=liuhangbin@gmail.com \
--cc=luka.gejak@linux.dev \
--cc=netdev@vger.kernel.org \
--cc=pabeni@redhat.com \
--cc=petr.wozniak@gmail.com \
--cc=qingfang.deng@linux.dev \
--cc=sdf.kernel@gmail.com \
--cc=shuah@kernel.org \
--cc=skhawaja@google.com \
--cc=stable@vger.kernel.org \
--cc=xiaoliang.yang_1@nxp.com \
/path/to/YOUR_REPLY
https://kernel.org/pub/software/scm/git/docs/git-send-email.html
* If your mail client supports setting the In-Reply-To header
via mailto: links, try the mailto: link
Be sure your reply has a Subject: header at the top and a blank line
before the message body.
This is a public inbox, see mirroring instructions
for how to clone and mirror all data and code used for this inbox
all inboxes | Powered by JetHome®