* [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 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