mirror of https://lore.kernel.org/lkml/
 help / color / mirror / Atom feed
* [PATCH net v7 0/4] net: hsr: fix super-packet forwarding and ordering
@ 2026-10-09 20:13 Xin Xie
  2026-10-09 20:13 ` [PATCH net v7 1/4] net: hsr: keep GRO disabled on HSR/PRP ports Xin Xie
                   ` (3 more replies)
  0 siblings, 4 replies; 5+ messages in thread
From: Xin Xie @ 2026-10-09 20:13 UTC (permalink / raw)
  To: netdev, linux-kselftest, linux-kernel
  Cc: davem, edumazet, kuba, pabeni, horms, andrew+netdev, shuah, kees,
	petr.wozniak, qingfang.deng, fmaurer, luka.gejak, bigeasy,
	xiaoliang.yang_1, skhawaja, liuhangbin, stable, sdf.kernel,
	xiexinet

HSR/PRP requires per-wire-frame tags/RCTs and sequence numbers, and
duplicate discard is per frame. RX GRO and TX GSO can present multiple
frames as one skb and violate that assumption: a super-skb is either
rejected by a constrained lower device, or forwarded without valid
per-frame trailers and sequence numbers. An oversized PRP aggregate
can also truncate the RCT's 12-bit LSDU size field.

Patch 1 keeps software GRO off while devices are direct HSR/PRP
members, without rewriting wanted features. Fixed-on or driver-required
GRO_HW may remain enabled.

Patch 2 replaces the forwarding lock with one consumer to avoid
stacked-HSR deadlocks. It preserves per-lower submission order for
locally numbered frames, so old HSR peers that drop late sequence
numbers keep working. Only sequence allocation holds the short counter
lock; forwarding runs after it is released.

Patch 3 segments valid GSO before per-frame processing, so each wire
frame gets its own tag/RCT and sequence number. Patch 4 tests member
GRO policy, GSO forwarding and local submission order.

A typical GSO source is a container connected to a PRP RedBox
interlink through veth. An existing TCP/UDP GSO skb can reach the
interlink with GRO disabled:

  container/netns                  host PRP RedBox
  prp0-peer ---- veth ---- prp0-int (interlink RX)
                                  |
                           split GSO into frames
                                  |
                           sequence number + RCT
                              /          \
                           LAN A        LAN B

Compatibility with older PRP senders:

Valid PRP frames from older senders remain supported. Patch 2
preserves submission order for locally numbered frames without
changing the PRP frame format or requiring peers to upgrade.

Some old Linux PRP senders incorrectly append one RCT to an entire
GSO skb. This violates the per-frame PRP requirement. Patch 3 drops
such intact aggregates when trailing bytes extend past a known,
nonzero IP length; zero lengths and GSO_PARTIAL are not covered by
that check.

The mixed-version test uses two PRP VMs; the host runs no PRP.
On virtio-net/TAP paths, disabling GRO_HW in patch 1 can make the
host segment an old sender's malformed aggregate before it reaches
the receiving VM. The RCT then becomes transport payload, which the
receiver cannot reliably identify or remove. In our QEMU 8.2.2 UDP
test this produced an extra six-byte datagram; the base receiver
delivered the original data correctly. Upgrade the sender or disable
GSO/TSO on its HSR/PRP master; the latter avoided the failure in the
test.

Validation:

On the net test kernel (7.3.0-rc5-gb0fe53dd6370), the installed
ordered selftest passed 5/5 cases, including supervision, concurrent
producers and sequence wrap, with complete, loss-free captures.

GRO passed 11/11 with the preceding helper. The later cleanup-only
fix does not affect GRO cases. Failure-injection checks confirmed
error reporting and cleanup.

Earlier queue-stress and PREEMPT_RT results are retained. Paired W=1
builds had no additional diagnostics. These checks were not repeated
for the helper update.

Based on net 6dc989ea46b9 (2026-10-05). Integration preserves the
upstream interlink promiscuous-mode exception and setup-failure
cleanup. The series applies directly to this base.

Patch 3 depends on patch 2; no stable backport is requested here.
The remaining HSR dev->stats races are left to a separate series,
as Paolo suggested [1]. Pre-existing shared-skb mutations are also
handled separately.

[1] https://lore.kernel.org/netdev/4fc3b9f1-4bef-4b34-ae7a-e89037cce829@redhat.com/

Changes since v6, including the NIPA and Gemini review responses:

- Replace recursive, one-time GRO disabling with persistent member
  policy and restoration of the latest wanted state.
- Preserve submission order with one consumer; remove the bitmap
  prerequisite and use core per-CPU counters for new RX/TX drops.
- Handle bounded 802.1Q/802.1AD VLAN stacks and reject trailing data
  beyond a known nonzero IP length before segmentation.
- Replace child-process signalling and fixed readiness sleeps with
  namespace-bound sockets registered before traffic. Bound cleanup
  under signals and preserve failures.
- Check packet contents, retransmissions and capture completeness,
  rather than throughput or average sizes. Use a standalone Python
  helper, with no embedded Python in the shell entries.

Veth tests do not validate hardware GRO_HW. The generic GSO bit need
not be off, so tests now check GSO_MASK member types. The claim that
the old test always failed was not supported by its actual feature
state and recorded passing runs. Shared-skb mutation is a separate
existing issue; both reports about new drop accounting are addressed
by the per-CPU counters.

Previous postings (earlier design, newest first):
v6: https://lore.kernel.org/netdev/20260809121455.1745-1-xiexinet@gmail.com/
v5: https://lore.kernel.org/netdev/20260807140751.1351-1-xiexinet@gmail.com/
v4: https://lore.kernel.org/netdev/20260803222211.877-1-xiexinet@gmail.com/
v3: https://lore.kernel.org/netdev/20260731090224.18-1-xiexinet@gmail.com/
v2: https://lore.kernel.org/netdev/20260724161253.79-1-xiexinet@gmail.com/
v1: https://lore.kernel.org/netdev/20260722171836.196-1-xiexinet@gmail.com/

Xin Xie (4):
  net: hsr: keep GRO disabled on HSR/PRP ports
  net: hsr: preserve submission order without a forwarding lock
  net: hsr: segment GSO before per-frame forwarding
  selftests: net: hsr: verify GRO policy and ordered forwarding

 .../networking/net_cachelines/net_device.rst      |    1 +
 include/linux/netdevice.h                         |    6 +
 net/core/dev.c                                    |   10 +
 net/hsr/Makefile                                  |    3 +-
 net/hsr/hsr_device.c                              |   66 +-
 net/hsr/hsr_forward.c                             |  233 +++-
 net/hsr/hsr_forward.h                             |   37 +
 net/hsr/hsr_forward_queue.c                       |  294 +++++
 net/hsr/hsr_main.h                                |   18 +
 net/hsr/hsr_netlink.c                             |    4 +-
 net/hsr/hsr_slave.c                               |   60 +-
 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 +
 .../selftests/net/hsr/hsr_ordered_forwarding.sh   |    9 +
 16 files changed, 1784 insertions(+), 58 deletions(-)
 create mode 100644 net/hsr/hsr_forward_queue.c
 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


base-commit: 6dc989ea46b96ce170840174b4a38c4a387fb005
-- 
2.43.0

^ permalink raw reply	[flat|nested] 5+ messages in thread

* [PATCH net v7 1/4] net: hsr: keep GRO disabled on HSR/PRP ports
  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 ` Xin Xie
  2026-10-09 20:13 ` [PATCH net v7 2/4] net: hsr: preserve submission order without a forwarding lock Xin Xie
                   ` (2 subsequent siblings)
  3 siblings, 0 replies; 5+ messages in thread
From: Xin Xie @ 2026-10-09 20:13 UTC (permalink / raw)
  To: netdev, linux-kselftest, linux-kernel
  Cc: davem, edumazet, kuba, pabeni, horms, andrew+netdev, shuah, kees,
	petr.wozniak, qingfang.deng, fmaurer, luka.gejak, bigeasy,
	xiaoliang.yang_1, skhawaja, liuhangbin, stable, sdf.kernel,
	xiexinet

HSR/PRP add tags and sequence numbers per wire frame. GRO can merge
plain interlink traffic before those fields are added, and userspace
can re-enable GRO after the port is attached.

Mark slave A/B and interlink devices with a kernel role bit. Filter
software GRO and configurable GRO_HW requests before ndo_fix_features()
without changing wanted_features. Filtering before the driver keeps
feature dependencies intact; fixed-on or driver-required GRO_HW is
left enabled.

Apply the policy before registering the RX handler and reject attach
if software GRO remains enabled. On failure or detach, unlink the
upper before clearing the role and recomputing features, restoring
the user's last request.

GSO skbs may still arrive. A later patch segments valid GSO before
HSR/PRP processing.

Fixes: 5055cccfc2d1 ("net: hsr: Provide RedBox support (HSR-SAN)")
Signed-off-by: Xin Xie <xiexinet@gmail.com>
---
 .../networking/net_cachelines/net_device.rst  |  1 +
 include/linux/netdevice.h                     |  6 +++
 net/core/dev.c                                | 10 ++++
 net/hsr/hsr_slave.c                           | 47 +++++++++++++++++++
 4 files changed, 64 insertions(+)

diff --git a/Documentation/networking/net_cachelines/net_device.rst b/Documentation/networking/net_cachelines/net_device.rst
index 512f6d6fa3d8..114a29b2ff4c 100644
--- a/Documentation/networking/net_cachelines/net_device.rst
+++ b/Documentation/networking/net_cachelines/net_device.rst
@@ -168,6 +168,7 @@ unsigned_long:1                     see_all_hwtstamp_requests
 unsigned_long:1                     change_proto_down
 unsigned_long:1                     netns_immutable
 unsigned_long:1                     fcoe_mtu
+unsigned_long:1                     hsr_port
 struct list_head                    net_notifier_list
 struct macsec_ops*                  macsec_ops
 struct udp_tunnel_nic_info*         udp_tunnel_nic_info
diff --git a/include/linux/netdevice.h b/include/linux/netdevice.h
index 3cff2174dc03..0a917906fe2a 100644
--- a/include/linux/netdevice.h
+++ b/include/linux/netdevice.h
@@ -2102,6 +2102,7 @@ enum netdev_reg_state {
  *	@change_proto_down: device supports setting carrier via IFLA_PROTO_DOWN
  *	@netns_immutable: interface can't change network namespaces
  *	@fcoe_mtu:	device supports maximum FCoE MTU, 2158 bytes
+ *	@hsr_port:	direct HSR/PRP member (slave A/B or interlink)
  *
  *	@net_notifier_list:	List of per-net netdev notifier block
  *				that follow this device when it is moved
@@ -2522,6 +2523,11 @@ struct net_device {
 	unsigned long		change_proto_down:1;
 	unsigned long		netns_immutable:1;
 	unsigned long		fcoe_mtu:1;
+	/* Direct HSR/PRP member: slave A/B or interlink.
+	 * Software GRO is filtered off while set. Kernel role
+	 * state only; no user ABI or offload capability.
+	 */
+	unsigned long		hsr_port:1;
 
 	struct list_head	net_notifier_list;
 
diff --git a/net/core/dev.c b/net/core/dev.c
index 18dc88990510..0c53b59af439 100644
--- a/net/core/dev.c
+++ b/net/core/dev.c
@@ -11111,6 +11111,16 @@ int __netdev_update_features(struct net_device *dev)
 
 	features = netdev_get_wanted_features(dev);
 
+	/* A direct HSR/PRP port needs per-frame metadata: filter
+	 * software GRO and configurable GRO_HW requests before the
+	 * driver fix runs, keeping the existing driver and core
+	 * feature coupling.
+	 */
+	if (dev->hsr_port) {
+		features &= ~NETIF_F_GRO;
+		features &= ~(dev->hw_features & NETIF_F_GRO_HW);
+	}
+
 	if (dev->netdev_ops->ndo_fix_features)
 		features = dev->netdev_ops->ndo_fix_features(dev, features);
 
diff --git a/net/hsr/hsr_slave.c b/net/hsr/hsr_slave.c
index a546f70f9cc8..1afcacac6b3c 100644
--- a/net/hsr/hsr_slave.c
+++ b/net/hsr/hsr_slave.c
@@ -11,6 +11,7 @@
 #include <linux/etherdevice.h>
 #include <linux/if_arp.h>
 #include <linux/if_vlan.h>
+#include <net/netdev_lock.h>
 #include "hsr_main.h"
 #include "hsr_device.h"
 #include "hsr_forward.h"
@@ -137,6 +138,33 @@ static int hsr_check_dev_ok(struct net_device *dev,
 	return 0;
 }
 
+/* Member role and feature policy.
+ * Setting the role makes __netdev_update_features() filter GRO
+ * and GRO_HW for this device.  Only the role and feature update
+ * run under the device ops lock; unlink, promiscuous updates and
+ * the master recompute stay outside.
+ */
+static void hsr_portdev_role_set(struct net_device *dev)
+{
+	netdev_lock_ops(dev);
+	dev->hsr_port = true;
+	netdev_change_features(dev);
+	netdev_unlock_ops(dev);
+}
+
+static void hsr_portdev_role_clear(struct net_device *dev)
+{
+	netdev_lock_ops(dev);
+	dev->hsr_port = false;
+	/* Restore from the current wanted state only on a still-
+	 * registered device; a netns move stays registered and must
+	 * restore, a dying device just drops the role.
+	 */
+	if (dev->reg_state == NETREG_REGISTERED)
+		netdev_change_features(dev);
+	netdev_unlock_ops(dev);
+}
+
 /* Setup device to be added to the HSR bridge. */
 static int hsr_portdev_setup(struct hsr_priv *hsr, struct net_device *dev,
 			     struct hsr_port *port,
@@ -169,6 +197,19 @@ static int hsr_portdev_setup(struct hsr_priv *hsr, struct net_device *dev,
 	if (res)
 		goto fail_upper_dev_link;
 
+	/* Enable the member role and recompute before the RX handler
+	 * is published: software GRO must be off before frames can
+	 * arrive.  A still-active software GRO rejects the port, a
+	 * remaining GRO_HW alone does not.
+	 */
+	hsr_portdev_role_set(dev);
+	if (dev->features & NETIF_F_GRO) {
+		NL_SET_ERR_MSG_MOD(extack,
+				   "software GRO still on after feature update, cannot join");
+		res = -EBUSY;
+		goto fail_role_policy;
+	}
+
 	res = netdev_rx_handler_register(dev, hsr_handle_frame, port);
 	if (res)
 		goto fail_rx_handler;
@@ -177,7 +218,12 @@ static int hsr_portdev_setup(struct hsr_priv *hsr, struct net_device *dev,
 	return 0;
 
 fail_rx_handler:
+fail_role_policy:
+	/* Rollback order: unlink the upper, clear the role and
+	 * recompute, then undo this setup's promiscuous increment.
+	 */
 	netdev_upper_dev_unlink(dev, hsr_dev);
+	hsr_portdev_role_clear(dev);
 fail_upper_dev_link:
 	if (!port->hsr->fwd_offloaded || port->type == HSR_PT_INTERLINK)
 		dev_set_promiscuity(dev, -1);
@@ -248,6 +294,7 @@ void hsr_del_port(struct hsr_port *port)
 		if (port->type == HSR_PT_SLAVE_A || port->type == HSR_PT_SLAVE_B)
 			vlan_vids_del_by_dev(port->dev, master->dev);
 		netdev_upper_dev_unlink(port->dev, master->dev);
+		hsr_portdev_role_clear(port->dev);
 		if (hsr->prot_version == PRP_V1 &&
 		    port->type == HSR_PT_SLAVE_B) {
 			eth_hw_addr_set(port->dev, port->original_macaddress);
-- 
2.43.0


^ permalink raw reply	[flat|nested] 5+ messages in thread

* [PATCH net v7 2/4] net: hsr: preserve submission order without a forwarding lock
  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 ` 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 ` [PATCH net v7 4/4] selftests: net: hsr: verify GRO policy and ordered forwarding Xin Xie
  3 siblings, 0 replies; 5+ messages in thread
From: Xin Xie @ 2026-10-09 20:13 UTC (permalink / raw)
  To: netdev, linux-kselftest, linux-kernel
  Cc: davem, edumazet, kuba, pabeni, horms, andrew+netdev, shuah, kees,
	petr.wozniak, qingfang.deng, fmaurer, luka.gejak, bigeasy,
	xiaoliang.yang_1, skhawaja, liuhangbin, stable, sdf.kernel,
	xiexinet, syzbot+fbf74291c3b7e753b481

Holding seqnr_lock across dev_queue_xmit() serializes numbering and
submission, but creates transmit-lock dependencies that can deadlock
stacked HSR devices. Merely shrinking the lock to the counter lets
one CPU send frame N after another has sent N+1; old HSR peers discard
the late N as stale.

Use one consumer for master TX, interlink RX and internally generated
supervision frames. An idle caller processes its own input inline;
contended inputs wait in a FIFO drained by one BH work item. Allocate
sequence numbers only when the consumer executes a frame, and submit
it to every lower before numbering the next frame. The short counter
lock is released before forwarding.

Limit queued inputs to 1024 jobs and 8 MiB, reserving capacity for
internal supervision. Carry transmit recursion depth across deferred
execution, cancel the worker and purge the queue on teardown, and count
queue drops once through core per-CPU statistics. LAN A/B reception
stays synchronous and tagged frames retain their received numbers.
The ordering guarantee is per-lower submission, not physical wire
order across multiple TX queues.

On PREEMPT_RT, high-priority load preempted both the inline consumer
and BH worker while another CPU continued to enqueue inputs. All
2118 inputs completed without drops; the finite data workload
completed 4.07 seconds after the load ended.

Reported-by: syzbot+fbf74291c3b7e753b481@syzkaller.appspotmail.com
Link: https://syzkaller.appspot.com/bug?extid=fbf74291c3b7e753b481
Fixes: 06afd2c31d33 ("hsr: Synchronize sending frames to have always incremented outgoing seq nr.")
Fixes: 430d67bdcb04 ("net: hsr: Use the seqnr lock for frames received via interlink port.")
Signed-off-by: Xin Xie <xiexinet@gmail.com>
---
 net/hsr/Makefile            |   3 +-
 net/hsr/hsr_device.c        |  64 +++++---
 net/hsr/hsr_forward.c       |  87 +++++++---
 net/hsr/hsr_forward.h       |  32 ++++
 net/hsr/hsr_forward_queue.c | 306 ++++++++++++++++++++++++++++++++++++
 net/hsr/hsr_main.h          |  18 +++
 net/hsr/hsr_netlink.c       |   4 +-
 net/hsr/hsr_slave.c         |  13 +-
 8 files changed, 470 insertions(+), 57 deletions(-)
 create mode 100644 net/hsr/hsr_forward_queue.c

diff --git a/net/hsr/Makefile b/net/hsr/Makefile
index 34e581db5c41..223b66fe8577 100644
--- a/net/hsr/Makefile
+++ b/net/hsr/Makefile
@@ -6,7 +6,8 @@
 obj-$(CONFIG_HSR)	+= hsr.o
 
 hsr-y			:= hsr_main.o hsr_framereg.o hsr_device.o \
-			   hsr_netlink.o hsr_slave.o hsr_forward.o
+			   hsr_netlink.o hsr_slave.o hsr_forward.o \
+			   hsr_forward_queue.o
 hsr-$(CONFIG_DEBUG_FS) += hsr_debugfs.o
 
 obj-$(CONFIG_PRP_DUP_DISCARD_KUNIT_TEST) += prp_dup_discard_test.o
diff --git a/net/hsr/hsr_device.c b/net/hsr/hsr_device.c
index d14de44e14b7..68cd64a865fd 100644
--- a/net/hsr/hsr_device.c
+++ b/net/hsr/hsr_device.c
@@ -232,9 +232,7 @@ static netdev_tx_t hsr_dev_xmit(struct sk_buff *skb, struct net_device *dev)
 		skb->dev = master->dev;
 		skb_reset_mac_header(skb);
 		skb_reset_mac_len(skb);
-		spin_lock_bh(&hsr->seqnr_lock);
 		hsr_forward_skb(skb, master);
-		spin_unlock_bh(&hsr->seqnr_lock);
 	} else {
 		dev_core_stats_tx_dropped_inc(dev);
 		dev_kfree_skb_any(skb);
@@ -290,6 +288,31 @@ static struct sk_buff *hsr_init_skb(struct hsr_port *master, int extra)
 	return NULL;
 }
 
+/* Assign the supervision sequence number at execution time in the
+ * single consumer.  Internally built supervision frames carry
+ * ETH_P_PRP with the supervision tag right after the Ethernet
+ * header.  HSRv0 shares the data counter, later versions use the
+ * dedicated supervision counter.
+ */
+void hsr_assign_sup_seq(struct sk_buff *skb, struct hsr_priv *hsr,
+			enum hsr_exec_source source)
+{
+	struct hsr_sup_tag *hsr_stag;
+
+	WARN_ON_ONCE(source == HSR_EXEC_DIRECT_LAN);
+
+	hsr_stag = (struct hsr_sup_tag *)(skb_mac_header(skb) + ETH_HLEN);
+	spin_lock_bh(&hsr->seqnr_lock);
+	if (hsr->prot_version > 0) {
+		hsr_stag->sequence_nr = htons(hsr->sup_sequence_nr);
+		WRITE_ONCE(hsr->sup_sequence_nr, hsr->sup_sequence_nr + 1);
+	} else {
+		hsr_stag->sequence_nr = htons(hsr->sequence_nr);
+		WRITE_ONCE(hsr->sequence_nr, hsr->sequence_nr + 1);
+	}
+	spin_unlock_bh(&hsr->seqnr_lock);
+}
+
 static void send_hsr_supervision_frame(struct hsr_port *port,
 				       unsigned long *interval,
 				       const unsigned char *addr)
@@ -326,15 +349,10 @@ static void send_hsr_supervision_frame(struct hsr_port *port,
 	set_hsr_stag_path(hsr_stag, (hsr->prot_version ? 0x0 : 0xf));
 	set_hsr_stag_HSR_ver(hsr_stag, hsr->prot_version);
 
-	/* From HSRv1 on we have separate supervision sequence numbers. */
-	spin_lock_bh(&hsr->seqnr_lock);
-	if (hsr->prot_version > 0) {
-		hsr_stag->sequence_nr = htons(hsr->sup_sequence_nr);
-		hsr->sup_sequence_nr++;
-	} else {
-		hsr_stag->sequence_nr = htons(hsr->sequence_nr);
-		hsr->sequence_nr++;
-	}
+	/* The sequence number is assigned by the consumer at execution
+	 * time; zero placeholder until then.
+	 */
+	hsr_stag->sequence_nr = 0;
 
 	hsr_stag->tlv.HSR_TLV_type = type;
 	/* HSRv0 has 6 unused bytes after the MAC */
@@ -356,13 +374,10 @@ static void send_hsr_supervision_frame(struct hsr_port *port,
 		ether_addr_copy(hsr_sp->macaddress_A, hsr->macaddress_redbox);
 	}
 
-	if (skb_put_padto(skb, ETH_ZLEN)) {
-		spin_unlock_bh(&hsr->seqnr_lock);
+	if (skb_put_padto(skb, ETH_ZLEN))
 		return;
-	}
 
-	hsr_forward_skb(skb, port);
-	spin_unlock_bh(&hsr->seqnr_lock);
+	hsr_forward_sup_skb(skb, port);
 	return;
 }
 
@@ -397,10 +412,10 @@ static void send_prp_supervision_frame(struct hsr_port *master,
 	set_hsr_stag_path(hsr_stag, (hsr->prot_version ? 0x0 : 0xf));
 	set_hsr_stag_HSR_ver(hsr_stag, (hsr->prot_version ? 1 : 0));
 
-	/* From HSRv1 on we have separate supervision sequence numbers. */
-	spin_lock_bh(&hsr->seqnr_lock);
-	hsr_stag->sequence_nr = htons(hsr->sup_sequence_nr);
-	hsr->sup_sequence_nr++;
+	/* The sequence number is assigned by the consumer at execution
+	 * time; zero placeholder until then.
+	 */
+	hsr_stag->sequence_nr = 0;
 	hsr_stag->tlv.HSR_TLV_type = PRP_TLV_LIFE_CHECK_DD;
 	hsr_stag->tlv.HSR_TLV_length = sizeof(struct hsr_sup_payload);
 
@@ -424,13 +439,10 @@ static void send_prp_supervision_frame(struct hsr_port *master,
 		hsr_stlv->HSR_TLV_length = 0;
 	}
 
-	if (skb_put_padto(skb, ETH_ZLEN)) {
-		spin_unlock_bh(&hsr->seqnr_lock);
+	if (skb_put_padto(skb, ETH_ZLEN))
 		return;
-	}
 
-	hsr_forward_skb(skb, master);
-	spin_unlock_bh(&hsr->seqnr_lock);
+	hsr_forward_sup_skb(skb, master);
 }
 
 /* Announce (supervision frame) timer function
@@ -778,6 +790,7 @@ int hsr_dev_finalize(struct net_device *hsr_dev, struct net_device *slave[2],
 		return res;
 
 	spin_lock_init(&hsr->seqnr_lock);
+	hsr_forward_init(hsr);
 	/* Overflow soon to find bugs easier: */
 	hsr->sequence_nr = HSR_SEQNR_START;
 	hsr->sup_sequence_nr = HSR_SUP_SEQNR_START;
@@ -851,6 +864,7 @@ int hsr_dev_finalize(struct net_device *hsr_dev, struct net_device *slave[2],
 	return 0;
 
 err_unregister:
+	hsr_forward_stop(hsr);
 	hsr_del_ports(hsr);
 err_add_master:
 	hsr_del_self_node(hsr);
diff --git a/net/hsr/hsr_forward.c b/net/hsr/hsr_forward.c
index 7734a521a96c..e8b53367bf3e 100644
--- a/net/hsr/hsr_forward.c
+++ b/net/hsr/hsr_forward.c
@@ -627,11 +627,23 @@ static void check_local_dest(struct hsr_priv *hsr, struct sk_buff *skb,
 	}
 }
 
+/* Local numbering condition, shared with the consumer routing
+ * decision: exactly the master and interlink entries carry locally
+ * generated frames.
+ */
+static bool hsr_needs_local_numbering(const struct hsr_port *port)
+{
+	return port->type == HSR_PT_MASTER ||
+	       port->type == HSR_PT_INTERLINK;
+}
+
 static void handle_std_frame(struct sk_buff *skb,
-			     struct hsr_frame_info *frame)
+			     struct hsr_frame_info *frame,
+			     enum hsr_exec_source source)
 {
 	struct hsr_port *port = frame->port_rcv;
 	struct hsr_priv *hsr = port->hsr;
+	u16 seq;
 
 	frame->skb_hsr = NULL;
 	frame->skb_prp = NULL;
@@ -640,13 +652,20 @@ static void handle_std_frame(struct sk_buff *skb,
 	if (port->type != HSR_PT_MASTER)
 		frame->is_from_san = true;
 
-	if (port->type == HSR_PT_MASTER ||
-	    port->type == HSR_PT_INTERLINK) {
-		/* Sequence nr for the master/interlink node */
-		lockdep_assert_held(&hsr->seqnr_lock);
-		frame->sequence_nr = hsr->sequence_nr;
-		hsr->sequence_nr++;
-	}
+	if (!hsr_needs_local_numbering(port))
+		return;
+
+	/* Sequence nr for the master/interlink node.  Local sequence
+	 * numbers are assigned only by the single consumer; the
+	 * explicit call-chain source proves it, shared ownership state
+	 * does not.
+	 */
+	WARN_ON_ONCE(source == HSR_EXEC_DIRECT_LAN);
+	spin_lock_bh(&hsr->seqnr_lock);
+	seq = hsr->sequence_nr;
+	WRITE_ONCE(hsr->sequence_nr, seq + 1);
+	spin_unlock_bh(&hsr->seqnr_lock);
+	frame->sequence_nr = seq;
 }
 
 int hsr_fill_frame_info(__be16 proto, struct sk_buff *skb,
@@ -670,10 +689,10 @@ int hsr_fill_frame_info(__be16 proto, struct sk_buff *skb,
 		return 0;
 	}
 
-	/* Standard frame or PRP from master port */
-	handle_std_frame(skb, frame);
-
-	return 0;
+	/* Standard frame or PRP from master port: the caller completes
+	 * untagged input with its explicit execution source.
+	 */
+	return HSR_FRAME_PLAIN;
 }
 
 int prp_fill_frame_info(__be16 proto, struct sk_buff *skb,
@@ -690,13 +709,12 @@ int prp_fill_frame_info(__be16 proto, struct sk_buff *skb,
 		frame->sequence_nr = prp_get_skb_sequence_nr(rct);
 		return 0;
 	}
-	handle_std_frame(skb, frame);
-
-	return 0;
+	return HSR_FRAME_PLAIN;
 }
 
 static int fill_frame_info(struct hsr_frame_info *frame,
-			   struct sk_buff *skb, struct hsr_port *port)
+			   struct sk_buff *skb, struct hsr_port *port,
+			   enum hsr_exec_source source)
 {
 	struct hsr_priv *hsr = port->hsr;
 	struct hsr_vlan_ethhdr *vlan_hdr;
@@ -756,7 +774,9 @@ static int fill_frame_info(struct hsr_frame_info *frame,
 	frame->is_from_san = false;
 	frame->port_rcv = port;
 	ret = hsr->proto_ops->fill_frame_info(proto, skb, frame);
-	if (ret)
+	if (ret == HSR_FRAME_PLAIN)
+		handle_std_frame(skb, frame, source);
+	else if (ret)
 		return ret;
 
 	check_local_dest(port->hsr, skb, frame);
@@ -764,13 +784,16 @@ static int fill_frame_info(struct hsr_frame_info *frame,
 	return 0;
 }
 
-/* Must be called holding rcu read lock (because of the port parameter) */
-void hsr_forward_skb(struct sk_buff *skb, struct hsr_port *port)
+/* Per-frame forwarding path.  Must be called holding rcu read lock
+ * (because of the port parameter).
+ */
+void hsr_forward_frame(struct sk_buff *skb, struct hsr_port *port,
+		       enum hsr_exec_source source)
 {
 	struct hsr_frame_info frame;
 
 	rcu_read_lock();
-	if (fill_frame_info(&frame, skb, port) < 0)
+	if (fill_frame_info(&frame, skb, port, source) < 0)
 		goto out_drop;
 
 	hsr_register_frame_in(frame.node_src, port, frame.sequence_nr);
@@ -779,7 +802,7 @@ void hsr_forward_skb(struct sk_buff *skb, struct hsr_port *port)
 	/* Gets called for ingress frames as well as egress from master port.
 	 * So check and increment stats for master port only here.
 	 */
-	if (port->type == HSR_PT_MASTER || port->type == HSR_PT_INTERLINK) {
+	if (hsr_needs_local_numbering(port)) {
 		port->dev->stats.tx_packets++;
 		port->dev->stats.tx_bytes += skb->len;
 	}
@@ -794,3 +817,25 @@ void hsr_forward_skb(struct sk_buff *skb, struct hsr_port *port)
 	port->dev->stats.tx_dropped++;
 	kfree_skb(skb);
 }
+
+/* Submission entry.  Inputs that need a local sequence number go to
+ * the single consumer before any numbering happens; the synchronous
+ * LAN A/B receive path keeps its original behavior.
+ */
+void hsr_forward_skb(struct sk_buff *skb, struct hsr_port *port)
+{
+	if (hsr_needs_local_numbering(port)) {
+		hsr_queue_submit(skb, port, HSR_JOB_NORMAL);
+		return;
+	}
+	hsr_forward_frame(skb, port, HSR_EXEC_DIRECT_LAN);
+}
+
+/* Internally generated supervision frames always take the common
+ * submit point; their sequence numbers are assigned by the consumer
+ * at execution time.
+ */
+void hsr_forward_sup_skb(struct sk_buff *skb, struct hsr_port *port)
+{
+	hsr_queue_submit(skb, port, HSR_JOB_INTERNAL_SUP);
+}
diff --git a/net/hsr/hsr_forward.h b/net/hsr/hsr_forward.h
index 206636750b30..0fde1972c0a5 100644
--- a/net/hsr/hsr_forward.h
+++ b/net/hsr/hsr_forward.h
@@ -13,7 +13,39 @@
 #include <linux/netdevice.h>
 #include "hsr_main.h"
 
+/* Per-frame execution source, carried explicitly down the call chain. */
+enum hsr_exec_source {
+	HSR_EXEC_DIRECT_LAN,	/* synchronous LAN A/B receive path */
+	HSR_EXEC_INLINE,	/* inline consumer activation */
+	HSR_EXEC_WORKER,	/* BH worker consumer */
+};
+
+/* Input class, stored with the immutable queue charge. */
+enum hsr_job_class {
+	HSR_JOB_NORMAL,
+	HSR_JOB_INTERNAL_SUP,
+};
+
+/* Return codes of proto_ops->fill_frame_info(): frame completed from a
+ * wire tag/RCT (0), untagged input left for the caller to complete with
+ * its execution source (HSR_FRAME_PLAIN), or error (< 0).
+ */
+#define HSR_FRAME_PLAIN	1
+
 void hsr_forward_skb(struct sk_buff *skb, struct hsr_port *port);
+void hsr_forward_sup_skb(struct sk_buff *skb, struct hsr_port *port);
+void hsr_forward_frame(struct sk_buff *skb, struct hsr_port *port,
+		       enum hsr_exec_source source);
+void hsr_assign_sup_seq(struct sk_buff *skb, struct hsr_priv *hsr,
+			enum hsr_exec_source source);
+
+/* Ordered forwarding queue (hsr_forward_queue.c) */
+void hsr_queue_submit(struct sk_buff *skb, struct hsr_port *port,
+		      enum hsr_job_class class);
+void hsr_forward_init(struct hsr_priv *hsr);
+void hsr_forward_stop(struct hsr_priv *hsr);
+void hsr_forward_forget_port(struct hsr_priv *hsr, struct net_device *dev);
+
 struct sk_buff *prp_create_tagged_frame(struct hsr_frame_info *frame,
 					struct hsr_port *port);
 struct sk_buff *hsr_create_tagged_frame(struct hsr_frame_info *frame,
diff --git a/net/hsr/hsr_forward_queue.c b/net/hsr/hsr_forward_queue.c
new file mode 100644
index 000000000000..5e96c24dec00
--- /dev/null
+++ b/net/hsr/hsr_forward_queue.c
@@ -0,0 +1,306 @@
+// SPDX-License-Identifier: GPL-2.0
+/* Ordered forwarding queue for HSR and PRP.
+ *
+ * One bounded FIFO and one consumer per instance preserve the submission
+ * order of locally numbered frames to every lower device without a
+ * forwarding lock.  Master TX, interlink RX and internally generated
+ * supervision frames are submitted before any sequence number is
+ * assigned; an idle caller executes its own input inline, contended
+ * inputs are drained in submission order by one BH work item.
+ */
+
+#include <linux/netdevice.h>
+#include <linux/rcupdate.h>
+#include <net/dst.h>
+#include <linux/slab.h>
+#include <linux/workqueue.h>
+
+#include "hsr_main.h"
+#include "hsr_forward.h"
+
+struct hsr_job {
+	struct list_head	list;
+	struct sk_buff		*skb;
+	struct net_device	*dev;	/* held entry device */
+	enum hsr_port_type	type;	/* entry role at submit */
+	enum hsr_job_class	class;	/* ordinary / internal supervision */
+	int			depth;	/* dev_recursion_level() at submit */
+	unsigned int		charge;	/* immutable skb->truesize */
+};
+
+/* Queue-stage drop: one count per dropped original input on the held
+ * entry device.  Ordinary interlink RX is an RX drop; master TX and
+ * internally generated supervision frames are TX drops.
+ */
+static void hsr_fwd_drop_stat(struct net_device *dev, enum hsr_port_type type,
+			      enum hsr_job_class class)
+{
+	if (class == HSR_JOB_NORMAL && type == HSR_PT_INTERLINK)
+		dev_core_stats_rx_dropped_inc(dev);
+	else
+		dev_core_stats_tx_dropped_inc(dev);
+}
+
+static void hsr_job_drop(struct hsr_job *job)
+{
+	/* Drop statistics complete before dev_put(). */
+	hsr_fwd_drop_stat(job->dev, job->type, job->class);
+	kfree_skb(job->skb);
+	dev_put(job->dev);
+	kfree(job);
+}
+
+/* Admission under fwd_lock.  Ordinary inputs must fit both the total and
+ * the ordinary-subset bounds; internal supervision may use the reserve
+ * within the total.  Difference comparisons avoid overflow.
+ */
+static bool hsr_queue_fits(struct hsr_priv *hsr, struct hsr_job *job)
+{
+	if (hsr->fwd_jobs >= HSR_FWD_JOBS_MAX)
+		return false;
+	if (job->charge > HSR_FWD_BYTES_MAX - hsr->fwd_bytes)
+		return false;
+	if (job->class == HSR_JOB_NORMAL) {
+		if (hsr->fwd_ord_jobs >= HSR_FWD_ORD_JOBS_MAX)
+			return false;
+		if (job->charge > HSR_FWD_ORD_BYTES_MAX - hsr->fwd_ord_bytes)
+			return false;
+	}
+	return true;
+}
+
+static void hsr_queue_charge_add(struct hsr_priv *hsr, struct hsr_job *job)
+{
+	hsr->fwd_jobs++;
+	hsr->fwd_bytes += job->charge;
+	if (job->class == HSR_JOB_NORMAL) {
+		hsr->fwd_ord_jobs++;
+		hsr->fwd_ord_bytes += job->charge;
+	}
+}
+
+/* Exactly once, on dequeue to active and on purge. */
+static void hsr_queue_charge_del(struct hsr_priv *hsr, struct hsr_job *job)
+{
+	hsr->fwd_jobs--;
+	hsr->fwd_bytes -= job->charge;
+	if (job->class == HSR_JOB_NORMAL) {
+		hsr->fwd_ord_jobs--;
+		hsr->fwd_ord_bytes -= job->charge;
+	}
+}
+
+/* Re-validate role and device under RCU; no bare port/node is kept
+ * across queueing.
+ */
+static struct hsr_port *hsr_job_port(struct hsr_priv *hsr, struct hsr_job *job)
+{
+	struct hsr_port *port;
+
+	hsr_for_each_port(hsr, port)
+		if (port->type == job->type && port->dev == job->dev)
+			return port;
+	return NULL;
+}
+
+/* Complete one input: re-validate the entry under RCU, raise the xmit
+ * recursion depth to the level saved at submit (never lower it), assign
+ * the supervision sequence number if needed, then run the complete
+ * per-frame path.  Returns the budget units this input consumed.
+ */
+static unsigned int hsr_job_process(struct hsr_priv *hsr, struct hsr_job *job,
+				    enum hsr_exec_source source)
+{
+	struct hsr_port *port;
+	unsigned int raised = 0;
+
+	rcu_read_lock();
+	port = hsr_job_port(hsr, job);
+	if (!port) {
+		rcu_read_unlock();
+		hsr_job_drop(job);
+		return 1;
+	}
+
+	while (dev_recursion_level() < job->depth) {
+		dev_xmit_recursion_inc();
+		raised++;
+	}
+
+	if (job->class == HSR_JOB_INTERNAL_SUP)
+		hsr_assign_sup_seq(job->skb, hsr, source);
+	hsr_forward_frame(job->skb, port, source);
+
+	while (raised) {
+		dev_xmit_recursion_dec();
+		raised--;
+	}
+	rcu_read_unlock();
+
+	dev_put(job->dev);
+	kfree(job);
+
+	/* One input is one budget unit; GSO per-frame accounting is
+	 * added by the segmentation change.
+	 */
+	return 1;
+}
+
+static void hsr_queue_owned_release(struct hsr_priv *hsr)
+{
+	if (!hsr->fwd_stopped && !list_empty(&hsr->fwd_queue))
+		queue_work(system_bh_wq, &hsr->fwd_work);
+	else
+		hsr->fwd_owned = false;
+}
+
+static void hsr_queue_work(struct work_struct *work)
+{
+	struct hsr_priv *hsr = container_of(work, struct hsr_priv, fwd_work);
+	unsigned int used = 0;
+
+	while (used < HSR_FWD_BUDGET_MAX) {
+		struct hsr_job *job;
+		unsigned int cost;
+
+		spin_lock_bh(&hsr->fwd_lock);
+		if (hsr->fwd_stopped || list_empty(&hsr->fwd_queue)) {
+			hsr->fwd_owned = false;
+			spin_unlock_bh(&hsr->fwd_lock);
+			return;
+		}
+		job = list_first_entry(&hsr->fwd_queue, struct hsr_job, list);
+		list_del(&job->list);
+		hsr_queue_charge_del(hsr, job);
+		spin_unlock_bh(&hsr->fwd_lock);
+
+		cost = hsr_job_process(hsr, job, HSR_EXEC_WORKER);
+		if (cost >= HSR_FWD_BUDGET_MAX - used)
+			used = HSR_FWD_BUDGET_MAX;
+		else
+			used += cost;
+	}
+
+	/* Budget exhausted with a backlog: keep the ownership and hand
+	 * the queue to the same worker object again.
+	 */
+	spin_lock_bh(&hsr->fwd_lock);
+	hsr_queue_owned_release(hsr);
+	spin_unlock_bh(&hsr->fwd_lock);
+}
+
+static void hsr_queue_inline_one(struct hsr_priv *hsr, struct hsr_job *job)
+{
+	local_bh_disable();
+	hsr_job_process(hsr, job, HSR_EXEC_INLINE);
+	spin_lock_bh(&hsr->fwd_lock);
+	hsr_queue_owned_release(hsr);
+	spin_unlock_bh(&hsr->fwd_lock);
+	local_bh_enable();
+}
+
+void hsr_queue_submit(struct sk_buff *skb, struct hsr_port *port,
+		      enum hsr_job_class class)
+{
+	struct hsr_priv *hsr = port->hsr;
+	struct hsr_job *job;
+	bool inline_owner = false;
+
+	RCU_LOCKDEP_WARN(!rcu_read_lock_held(),
+			 "HSR queue submit outside RCU read-side");
+
+	job = kzalloc_obj(*job, GFP_ATOMIC);
+	if (!job) {
+		hsr_fwd_drop_stat(port->dev, port->type, class);
+		kfree_skb(skb);
+		return;
+	}
+	job->skb = skb;
+	job->dev = port->dev;
+	job->type = port->type;
+	job->class = class;
+	job->depth = dev_recursion_level();
+	job->charge = skb->truesize;
+	skb_dst_force(skb);
+	dev_hold(job->dev);
+
+	spin_lock_bh(&hsr->fwd_lock);
+	if (hsr->fwd_stopped) {
+		spin_unlock_bh(&hsr->fwd_lock);
+		hsr_job_drop(job);
+		return;
+	}
+	if (!hsr->fwd_owned) {
+		/* Idle: take the execution ownership and process this
+		 * input directly.  The job never enters the public queue
+		 * and does not consume queued charge.
+		 */
+		hsr->fwd_owned = true;
+		inline_owner = true;
+	} else if (hsr_queue_fits(hsr, job)) {
+		list_add_tail(&job->list, &hsr->fwd_queue);
+		hsr_queue_charge_add(hsr, job);
+	} else {
+		spin_unlock_bh(&hsr->fwd_lock);
+		hsr_job_drop(job);
+		return;
+	}
+	spin_unlock_bh(&hsr->fwd_lock);
+
+	if (inline_owner)
+		hsr_queue_inline_one(hsr, job);
+}
+
+void hsr_forward_forget_port(struct hsr_priv *hsr, struct net_device *dev)
+{
+	struct hsr_job *job, *next;
+	LIST_HEAD(purge);
+
+	spin_lock_bh(&hsr->fwd_lock);
+	list_for_each_entry_safe(job, next, &hsr->fwd_queue, list) {
+		if (job->dev != dev)
+			continue;
+		list_move_tail(&job->list, &purge);
+		hsr_queue_charge_del(hsr, job);
+	}
+	spin_unlock_bh(&hsr->fwd_lock);
+
+	list_for_each_entry_safe(job, next, &purge, list) {
+		list_del(&job->list);
+		hsr_job_drop(job);
+	}
+}
+
+void hsr_forward_stop(struct hsr_priv *hsr)
+{
+	struct hsr_job *job, *next;
+	LIST_HEAD(purge);
+
+	spin_lock_bh(&hsr->fwd_lock);
+	hsr->fwd_stopped = true;
+	spin_unlock_bh(&hsr->fwd_lock);
+
+	synchronize_net();
+	cancel_work_sync(&hsr->fwd_work);
+
+	spin_lock_bh(&hsr->fwd_lock);
+	list_splice_init(&hsr->fwd_queue, &purge);
+	hsr->fwd_jobs = 0;
+	hsr->fwd_ord_jobs = 0;
+	hsr->fwd_bytes = 0;
+	hsr->fwd_ord_bytes = 0;
+	hsr->fwd_owned = false;
+	spin_unlock_bh(&hsr->fwd_lock);
+
+	list_for_each_entry_safe(job, next, &purge, list) {
+		list_del(&job->list);
+		hsr_job_drop(job);
+	}
+}
+
+void hsr_forward_init(struct hsr_priv *hsr)
+{
+	spin_lock_init(&hsr->fwd_lock);
+	INIT_LIST_HEAD(&hsr->fwd_queue);
+	INIT_WORK(&hsr->fwd_work, hsr_queue_work);
+}
diff --git a/net/hsr/hsr_main.h b/net/hsr/hsr_main.h
index 53e95bae0ee2..294c61f5068d 100644
--- a/net/hsr/hsr_main.h
+++ b/net/hsr/hsr_main.h
@@ -14,6 +14,7 @@
 #include <linux/list.h>
 #include <linux/if_vlan.h>
 #include <linux/if_hsr.h>
+#include <linux/workqueue.h>
 
 /* Time constants as specified in the HSR specification (IEC-62439-3 2010)
  * Table 8.
@@ -25,6 +26,13 @@
 #define HSR_ANNOUNCE_INTERVAL		  100 /* ms */
 #define HSR_ENTRY_FORGET_TIME		  400 /* ms */
 
+/* Ordered-forwarding queue and work bounds. */
+#define HSR_FWD_JOBS_MAX		1024
+#define HSR_FWD_ORD_JOBS_MAX		960
+#define HSR_FWD_BYTES_MAX		(8 * 1024 * 1024)
+#define HSR_FWD_ORD_BYTES_MAX		(HSR_FWD_BYTES_MAX - 64 * 1024)
+#define HSR_FWD_BUDGET_MAX		64
+
 /* By how much may slave1 and slave2 timestamps of latest received frame from
  * each node differ before we notify of communication problem?
  */
@@ -202,6 +210,16 @@ struct hsr_priv {
 	enum hsr_version prot_version;	/* Indicate if HSRv0, HSRv1 or PRPv1 */
 	spinlock_t seqnr_lock;	/* locking for sequence_nr */
 	spinlock_t list_lock;	/* locking for node list */
+	/* Ordered-forwarding queue (hsr_forward_queue.c) */
+	spinlock_t fwd_lock;
+	struct list_head fwd_queue;
+	struct work_struct fwd_work;
+	unsigned int fwd_jobs;
+	unsigned int fwd_ord_jobs;
+	unsigned int fwd_bytes;
+	unsigned int fwd_ord_bytes;
+	bool fwd_owned;		/* consumer execution ownership */
+	bool fwd_stopped;
 	const struct hsr_proto_ops	*proto_ops;
 #define PRP_LAN_ID	0x5     /* 0x1010 for A and 0x1011 for B. Bit 0 is set
 				 * based on SLAVE_A or SLAVE_B
diff --git a/net/hsr/hsr_netlink.c b/net/hsr/hsr_netlink.c
index 88940e8014b2..971751bf38a1 100644
--- a/net/hsr/hsr_netlink.c
+++ b/net/hsr/hsr_netlink.c
@@ -12,6 +12,7 @@
 #include <net/rtnetlink.h>
 #include <net/genetlink.h>
 #include "hsr_main.h"
+#include "hsr_forward.h"
 #include "hsr_device.h"
 #include "hsr_framereg.h"
 
@@ -137,6 +138,7 @@ static void hsr_dellink(struct net_device *dev, struct list_head *head)
 	timer_delete_sync(&hsr->announce_timer);
 	timer_delete_sync(&hsr->announce_proxy_timer);
 
+	hsr_forward_stop(hsr);
 	hsr_debugfs_term(hsr);
 	hsr_del_ports(hsr);
 
@@ -173,7 +175,7 @@ static int hsr_fill_info(struct sk_buff *skb, const struct net_device *dev)
 
 	if (nla_put(skb, IFLA_HSR_SUPERVISION_ADDR, ETH_ALEN,
 		    hsr->sup_multicast_addr) ||
-	    nla_put_u16(skb, IFLA_HSR_SEQ_NR, hsr->sequence_nr))
+	    nla_put_u16(skb, IFLA_HSR_SEQ_NR, READ_ONCE(hsr->sequence_nr)))
 		goto nla_put_failure;
 	if (hsr->prot_version == PRP_V1)
 		proto = HSR_PROTOCOL_PRP;
diff --git a/net/hsr/hsr_slave.c b/net/hsr/hsr_slave.c
index 1afcacac6b3c..c93bd12a6749 100644
--- a/net/hsr/hsr_slave.c
+++ b/net/hsr/hsr_slave.c
@@ -74,16 +74,10 @@ static rx_handler_result_t hsr_handle_frame(struct sk_buff **pskb)
 	}
 	skb_reset_mac_len(skb);
 
-	/* Only the frames received over the interlink port will assign a
-	 * sequence number and require synchronisation vs other sender.
+	/* Interlink RX is locally numbered and takes the ordered
+	 * consumer; the LAN A/B receive path stays synchronous.
 	 */
-	if (port->type == HSR_PT_INTERLINK) {
-		spin_lock_bh(&hsr->seqnr_lock);
-		hsr_forward_skb(skb, port);
-		spin_unlock_bh(&hsr->seqnr_lock);
-	} else {
-		hsr_forward_skb(skb, port);
-	}
+	hsr_forward_skb(skb, port);
 
 finish_consume:
 	return RX_HANDLER_CONSUMED;
@@ -289,6 +283,7 @@ void hsr_del_port(struct hsr_port *port)
 		netdev_update_features(master->dev);
 		dev_set_mtu(master->dev, hsr_get_max_mtu(hsr));
 		netdev_rx_handler_unregister(port->dev);
+		hsr_forward_forget_port(hsr, port->dev);
 		if (!port->hsr->fwd_offloaded || port->type == HSR_PT_INTERLINK)
 			dev_set_promiscuity(port->dev, -1);
 		if (port->type == HSR_PT_SLAVE_A || port->type == HSR_PT_SLAVE_B)
-- 
2.43.0


^ permalink raw reply	[flat|nested] 5+ messages in thread

* [PATCH net v7 3/4] net: hsr: segment GSO before per-frame forwarding
  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 ` Xin Xie
  2026-10-09 20:13 ` [PATCH net v7 4/4] selftests: net: hsr: verify GRO policy and ordered forwarding Xin Xie
  3 siblings, 0 replies; 5+ messages in thread
From: Xin Xie @ 2026-10-09 20:13 UTC (permalink / raw)
  To: netdev, linux-kselftest, linux-kernel
  Cc: davem, edumazet, kuba, pabeni, horms, andrew+netdev, shuah, kees,
	petr.wozniak, qingfang.deng, fmaurer, luka.gejak, bigeasy,
	xiaoliang.yang_1, skhawaja, liuhangbin, stable, sdf.kernel,
	xiexinet

HSR/PRP require a tag/RCT and sequence number for each wire frame.
Adding them to a GSO aggregate before segmentation cannot provide
valid per-frame metadata.

Use the core GSO segmenter before normal per-frame processing.
Each segment follows the usual local-delivery and forwarding paths.
Remove GSO_MASK types from the master's hw_features so core TX also
segments local traffic before ndo_start_xmit().

Classify by inner protocol, not ingress port. Use the core's bounded
VLAN decoder for 802.1Q, 802.1AD and nested or accelerated tags.
Reject inner HSR/PRP aggregates, unreadable fragments and undecodable
headers. Also reject trailing bytes beyond a known nonzero IP length:
otherwise an old PRP sender's erroneous whole-aggregate RCT becomes
transport payload. This check also rejects locally destined malformed
inputs that the base kernel accepted. Zero IP lengths and GSO_PARTIAL
are left to the segmenter because they provide no full IP extent for
this check.

The preceding ordered-consumer patch is required for per-segment
local numbering. Charge the worker budget by the actual segment count,
or one on failure, and count each rejected input once through core
per-CPU statistics on its entry device.

Fixes: f421436a591d ("net/hsr: Add support for the High-availability Seamless Redundancy protocol (HSRv0)")
Signed-off-by: Xin Xie <xiexinet@gmail.com>
---
 net/hsr/hsr_device.c        |   2 +-
 net/hsr/hsr_forward.c       | 148 +++++++++++++++++++++++++++++++++++-
 net/hsr/hsr_forward.h       |   5 ++
 net/hsr/hsr_forward_queue.c |  22 ++----
 4 files changed, 158 insertions(+), 19 deletions(-)

diff --git a/net/hsr/hsr_device.c b/net/hsr/hsr_device.c
index 68cd64a865fd..a67d85551abe 100644
--- a/net/hsr/hsr_device.c
+++ b/net/hsr/hsr_device.c
@@ -698,7 +698,7 @@ void hsr_dev_setup(struct net_device *dev)
 	dev->needs_free_netdev = true;
 
 	dev->hw_features = NETIF_F_SG | NETIF_F_FRAGLIST | NETIF_F_HIGHDMA |
-			   NETIF_F_GSO_MASK | NETIF_F_HW_CSUM |
+			   NETIF_F_HW_CSUM |
 			   NETIF_F_HW_VLAN_CTAG_TX |
 			   NETIF_F_HW_VLAN_CTAG_FILTER;
 
diff --git a/net/hsr/hsr_forward.c b/net/hsr/hsr_forward.c
index e8b53367bf3e..0c227a20b7d8 100644
--- a/net/hsr/hsr_forward.c
+++ b/net/hsr/hsr_forward.c
@@ -12,6 +12,9 @@
 #include <linux/skbuff.h>
 #include <linux/etherdevice.h>
 #include <linux/if_vlan.h>
+#include <linux/ip.h>
+#include <linux/ipv6.h>
+#include <net/gso.h>
 #include "hsr_main.h"
 #include "hsr_framereg.h"
 
@@ -818,6 +821,149 @@ void hsr_forward_frame(struct sk_buff *skb, struct hsr_port *port,
 	kfree_skb(skb);
 }
 
+/* Shared drop mapping for the new queue and GSO gate stages: one
+ * count per dropped original input.  Ordinary master TX and all
+ * internal supervision (local and proxy) are TX drops on the entry
+ * device; ordinary slave A/B and interlink RX are RX drops.
+ * Internal supervision keeps the TX mapping regardless of its
+ * entry role.
+ */
+void hsr_fwd_drop_stat(struct net_device *dev, enum hsr_port_type type,
+		       enum hsr_job_class class)
+{
+	if (class == HSR_JOB_NORMAL && type != HSR_PT_MASTER)
+		dev_core_stats_rx_dropped_inc(dev);
+	else
+		dev_core_stats_tx_dropped_inc(dev);
+}
+
+/* Return the inner EtherType using the core's bounded 802.1Q/802.1AD
+ * decoder. An accelerated outer tag is metadata, not frame bytes.
+ */
+static __be16 hsr_gso_effective_proto(const struct sk_buff *skb, int *l3off)
+{
+	struct ethhdr eh;
+	const struct ethhdr *eth;
+
+	if (skb_vlan_tag_present(skb) &&
+	    !eth_type_vlan(skb->vlan_proto))
+		return 0;
+
+	eth = skb_header_pointer(skb, 0, sizeof(eh), &eh);
+	if (!eth)
+		return 0;
+
+	return __vlan_get_protocol(skb, eth->h_proto, l3off);
+}
+
+/* Bytes beyond the IP datagram must not become segment payload. Older
+ * PRP senders can append an RCT to a whole GSO skb. Zero IP lengths and
+ * GSO_PARTIAL headers do not provide the full aggregate length here.
+ */
+static bool hsr_gso_trailing_data(const struct sk_buff *skb, __be16 proto,
+				  int l3off)
+{
+	unsigned int l3len;
+
+	if (skb_shinfo(skb)->gso_type & SKB_GSO_PARTIAL)
+		return false;
+
+	if (proto == htons(ETH_P_IP)) {
+		struct iphdr buffer;
+		const struct iphdr *iph;
+
+		iph = skb_header_pointer(skb, l3off, sizeof(buffer), &buffer);
+		if (!iph)
+			return true;
+		l3len = ntohs(iph->tot_len);
+	} else if (proto == htons(ETH_P_IPV6)) {
+		struct ipv6hdr buffer;
+		const struct ipv6hdr *ip6h;
+
+		ip6h = skb_header_pointer(skb, l3off, sizeof(buffer), &buffer);
+		if (!ip6h)
+			return true;
+		l3len = ntohs(ip6h->payload_len);
+		if (l3len)
+			l3len += sizeof(*ip6h);
+	} else {
+		return false;
+	}
+
+	return l3len && l3off + l3len < skb->len;
+}
+
+/* GSO fan-out funnel: unfold super-packets before per-frame
+ * processing so each wire frame gets its own HSR/PRP tag and
+ * sequence number.  Returns the number of per-frame units actually
+ * processed (one for a single-frame input or a failed aggregate).
+ */
+unsigned int hsr_forward_input(struct sk_buff *skb, struct hsr_port *port,
+			       enum hsr_job_class class,
+			       enum hsr_exec_source source)
+{
+	struct sk_buff *segs, *next;
+	unsigned int count = 0;
+	__be16 proto;
+	int l3off;
+
+	if (likely(!skb_is_gso(skb))) {
+		hsr_forward_frame(skb, port, source);
+		return 1;
+	}
+
+	/* Conforming plain-protocol GSO super-packets carry
+	 * trailer-free sender payload and are segmented here: each
+	 * segment is delivered or forwarded as its own wire frame, on
+	 * any ingress role.
+	 *
+	 * The gate is content-based, not port-based.  An aggregate
+	 * whose effective protocol is ETH_P_HSR or ETH_P_PRP cannot be
+	 * safely segmented and is dropped, as is any skb whose headers
+	 * cannot be read or whose fragments are unreadable net_iov
+	 * (device-memory) pages: software segmentation cannot read
+	 * their payload, so each segment would inherit the unreadable
+	 * flag and carry uninitialized data.  Ordinary VLAN GSO
+	 * (802.1Q/802.1AD at any reachable tag depth, accelerated or
+	 * in-band) is classified to its inner protocol through the core
+	 * VLAN decoder and segmented like plain input.  With
+	 * NETIF_F_HW_HSR_TAG_RM the lower has already stripped the tag,
+	 * so such aggregates arrive plain and are segmented.
+	 */
+	if (!skb_frags_readable(skb))
+		goto drop_gso; /* net_iov frags are not host-readable */
+	proto = hsr_gso_effective_proto(skb, &l3off);
+	if (!proto)
+		goto drop_gso; /* classification failure, fail-safe */
+	if (proto == htons(ETH_P_HSR) || proto == htons(ETH_P_PRP))
+		goto drop_gso;
+	if (hsr_gso_trailing_data(skb, proto, l3off))
+		goto drop_gso;
+
+	/* features = 0: request full software segmentation.  tx_path is
+	 * true only for locally generated traffic on the master; ingress
+	 * follows RX checksum semantics.
+	 */
+	segs = __skb_gso_segment(skb, 0, port->type == HSR_PT_MASTER);
+	if (IS_ERR(segs) || unlikely(!segs))
+		goto drop_gso;
+
+	consume_skb(skb);
+	while (segs) {
+		next = segs->next;
+		segs->next = NULL;
+		hsr_forward_frame(segs, port, source);
+		count++;
+		segs = next;
+	}
+	return count;
+
+drop_gso:
+	hsr_fwd_drop_stat(port->dev, port->type, class);
+	kfree_skb(skb);
+	return 1;
+}
+
 /* Submission entry.  Inputs that need a local sequence number go to
  * the single consumer before any numbering happens; the synchronous
  * LAN A/B receive path keeps its original behavior.
@@ -828,7 +974,7 @@ void hsr_forward_skb(struct sk_buff *skb, struct hsr_port *port)
 		hsr_queue_submit(skb, port, HSR_JOB_NORMAL);
 		return;
 	}
-	hsr_forward_frame(skb, port, HSR_EXEC_DIRECT_LAN);
+	hsr_forward_input(skb, port, HSR_JOB_NORMAL, HSR_EXEC_DIRECT_LAN);
 }
 
 /* Internally generated supervision frames always take the common
diff --git a/net/hsr/hsr_forward.h b/net/hsr/hsr_forward.h
index 0fde1972c0a5..12fdf1bff4a3 100644
--- a/net/hsr/hsr_forward.h
+++ b/net/hsr/hsr_forward.h
@@ -36,6 +36,11 @@ void hsr_forward_skb(struct sk_buff *skb, struct hsr_port *port);
 void hsr_forward_sup_skb(struct sk_buff *skb, struct hsr_port *port);
 void hsr_forward_frame(struct sk_buff *skb, struct hsr_port *port,
 		       enum hsr_exec_source source);
+unsigned int hsr_forward_input(struct sk_buff *skb, struct hsr_port *port,
+			       enum hsr_job_class class,
+			       enum hsr_exec_source source);
+void hsr_fwd_drop_stat(struct net_device *dev, enum hsr_port_type type,
+		       enum hsr_job_class class);
 void hsr_assign_sup_seq(struct sk_buff *skb, struct hsr_priv *hsr,
 			enum hsr_exec_source source);
 
diff --git a/net/hsr/hsr_forward_queue.c b/net/hsr/hsr_forward_queue.c
index 5e96c24dec00..c97cc722eaca 100644
--- a/net/hsr/hsr_forward_queue.c
+++ b/net/hsr/hsr_forward_queue.c
@@ -28,18 +28,7 @@ struct hsr_job {
 	unsigned int		charge;	/* immutable skb->truesize */
 };
 
-/* Queue-stage drop: one count per dropped original input on the held
- * entry device.  Ordinary interlink RX is an RX drop; master TX and
- * internally generated supervision frames are TX drops.
- */
-static void hsr_fwd_drop_stat(struct net_device *dev, enum hsr_port_type type,
-			      enum hsr_job_class class)
-{
-	if (class == HSR_JOB_NORMAL && type == HSR_PT_INTERLINK)
-		dev_core_stats_rx_dropped_inc(dev);
-	else
-		dev_core_stats_tx_dropped_inc(dev);
-}
+/* Queue-stage drops share the mapping in hsr_forward.c. */
 
 static void hsr_job_drop(struct hsr_job *job)
 {
@@ -113,6 +102,7 @@ static unsigned int hsr_job_process(struct hsr_priv *hsr, struct hsr_job *job,
 {
 	struct hsr_port *port;
 	unsigned int raised = 0;
+	unsigned int cost;
 
 	rcu_read_lock();
 	port = hsr_job_port(hsr, job);
@@ -129,7 +119,7 @@ static unsigned int hsr_job_process(struct hsr_priv *hsr, struct hsr_job *job,
 
 	if (job->class == HSR_JOB_INTERNAL_SUP)
 		hsr_assign_sup_seq(job->skb, hsr, source);
-	hsr_forward_frame(job->skb, port, source);
+	cost = hsr_forward_input(job->skb, port, job->class, source);
 
 	while (raised) {
 		dev_xmit_recursion_dec();
@@ -140,10 +130,8 @@ static unsigned int hsr_job_process(struct hsr_priv *hsr, struct hsr_job *job,
 	dev_put(job->dev);
 	kfree(job);
 
-	/* One input is one budget unit; GSO per-frame accounting is
-	 * added by the segmentation change.
-	 */
-	return 1;
+	/* Actual per-frame units consumed; a failed aggregate costs 1. */
+	return cost;
 }
 
 static void hsr_queue_owned_release(struct hsr_priv *hsr)
-- 
2.43.0


^ permalink raw reply	[flat|nested] 5+ messages in thread

* [PATCH net v7 4/4] selftests: net: hsr: verify GRO policy and ordered forwarding
  2026-10-09 20:13 [PATCH net v7 0/4] net: hsr: fix super-packet forwarding and ordering Xin Xie
                   ` (2 preceding siblings ...)
  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
  3 siblings, 0 replies; 5+ messages in thread
From: Xin Xie @ 2026-10-09 20:13 UTC (permalink / raw)
  To: netdev, linux-kselftest, linux-kernel
  Cc: davem, edumazet, kuba, pabeni, horms, andrew+netdev, shuah, kees,
	petr.wozniak, qingfang.deng, fmaurer, luka.gejak, bigeasy,
	xiaoliang.yang_1, skhawaja, liuhangbin, stable, sdf.kernel,
	xiexinet

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


^ permalink raw reply	[flat|nested] 5+ messages in thread

end of thread, other threads:[~2026-10-09 20:13 UTC | newest]

Thread overview: 5+ messages (download: mbox.gz / follow: Atom feed)
-- links below jump to the message on this page --
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 ` [PATCH net v7 4/4] selftests: net: hsr: verify GRO policy and ordered forwarding Xin Xie

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®