Commit b533cf9301 for frr
commit b533cf93011d6c56f776c9f828f59c4a95831354
Author: Jafar Al-Gharaibeh <jafar@atcorp.com>
Date: Mon Oct 5 10:30:01 2026 -0500
tests: retry IGMP include reports if the group stalls
The fragmentation test sends its INCLUDE burst once. A burst that
arrives before the IGMP socket exists, or that is dropped while pimd
is not scheduled, never shows up, and polling show ip igmp groups
cannot recover it. Wait for the socket, and resend only when the
installed source count stops growing.
Signed-off-by: Jafar Al-Gharaibeh <jafar@atcorp.com>
diff --git a/tests/topotests/pim_igmp_specific_query_frag/igmp_send.py b/tests/topotests/pim_igmp_specific_query_frag/igmp_send.py
index 4ac8dad10a..10fcf6eb63 100644
--- a/tests/topotests/pim_igmp_specific_query_frag/igmp_send.py
+++ b/tests/topotests/pim_igmp_specific_query_frag/igmp_send.py
@@ -19,36 +19,46 @@ from scapy.layers.inet import IP, IPOption_Router_Alert
from scapy.sendrecv import sendp
# pylint: disable=C0413
-load_contrib('igmpv3')
+load_contrib("igmpv3")
from scapy.contrib.igmpv3 import IGMPv3, IGMPv3mr, IGMPv3gr
+
def get_packet(rtype, group, src_ip, sources):
"Builds a single Ethernet/IP/IGMPv3 Membership Report packet"
record = IGMPv3gr(rtype=rtype, maddr=group, srcaddrs=sources)
pkt = (
- Ether(dst="01:00:5e:00:00:16") /
- IP(src=src_ip, dst="224.0.0.22", options=[IPOption_Router_Alert()]) /
- IGMPv3() /
- IGMPv3mr(records=[record])
+ Ether(dst="01:00:5e:00:00:16")
+ / IP(src=src_ip, dst="224.0.0.22", options=[IPOption_Router_Alert()])
+ / IGMPv3()
+ / IGMPv3mr(records=[record])
)
return pkt
+
def main():
parser = argparse.ArgumentParser(description="IGMPv3 Host Script")
- parser.add_argument("--mode", choices=["include", "block"], required=True,
- help="Operation mode: 'include' (rtype 1) or 'block' (rtype 6)")
+ parser.add_argument(
+ "--mode",
+ choices=["include", "block"],
+ required=True,
+ help="Operation mode: 'include' (rtype 1) or 'block' (rtype 6)",
+ )
parser.add_argument("--group", required=True, help="Target Multicast Group IP")
parser.add_argument("--src", required=True, help="Source IP (Packet Sender)")
- parser.add_argument("--json", required=True, help="Path to JSON file containing source IP list")
+ parser.add_argument(
+ "--json", required=True, help="Path to JSON file containing source IP list"
+ )
parser.add_argument("--iface", default="h1-eth0", help="Network interface to use")
- parser.add_argument("--step", type=int, default=250, help="Number of source IPs per packet")
+ parser.add_argument(
+ "--step", type=int, default=250, help="Number of source IPs per packet"
+ )
args = parser.parse_args()
# Load source IPs from JSON file
try:
- with open(args.json, 'r') as f:
+ with open(args.json, "r") as f:
data = json.load(f)
sources = data if isinstance(data, list) else data.get("sources", [])
except Exception as e:
@@ -77,5 +87,6 @@ def main():
print(f"Success")
+
if __name__ == "__main__":
main()
diff --git a/tests/topotests/pim_igmp_specific_query_frag/test_pim_igmp_specific_query_frag.py b/tests/topotests/pim_igmp_specific_query_frag/test_pim_igmp_specific_query_frag.py
index 5662581b3e..c73d371243 100644
--- a/tests/topotests/pim_igmp_specific_query_frag/test_pim_igmp_specific_query_frag.py
+++ b/tests/topotests/pim_igmp_specific_query_frag/test_pim_igmp_specific_query_frag.py
@@ -40,22 +40,30 @@ from lib.topotest import json_cmp
from lib.topogen import Topogen, TopoRouter, get_topogen
from lib.topolog import logger
-HOST_IP="192.168.10.2"
-IGMP_GROUP="232.1.1.12"
+HOST_IP = "192.168.10.2"
+ROUTER_IP = "192.168.10.1"
+IGMP_GROUP = "232.1.1.12"
TOTAL_SOURCES = 5000
SOURCE_PREFIX = "10.0"
STEP = 250
+# INCLUDE reports are not retransmitted. A burst that arrives before the
+# IGMP socket exists, or that the kernel drops while pimd is not scheduled,
+# is gone. Resend a few times if the installed source count stops growing,
+# but only after it has stayed flat for several polls. Sending all retries
+# during one stall drops them together.
+MAX_INCLUDE_RESENDS = 3
+INCLUDE_RESEND_POLLS = 10
+
@pytest.fixture
def sources_db():
db = {
- f"{SOURCE_PREFIX}.{(i // 256) % 256}.{i % 256}": {
- "query_count": 0
- }
+ f"{SOURCE_PREFIX}.{(i // 256) % 256}.{i % 256}": {"query_count": 0}
for i in range(1, TOTAL_SOURCES + 1)
}
return db
+
def build_topo(tgen):
"Build function"
@@ -67,21 +75,24 @@ def build_topo(tgen):
s1.add_link(tgen.gears["r1"])
s1.add_link(tgen.gears["h1"])
+
def setup_module(mod):
"Sets up the pytest environment"
tgen = Topogen(build_topo, mod.__name__)
tgen.start_topology()
- r1 = tgen.gears['r1']
+ r1 = tgen.gears["r1"]
r1.load_frr_config()
tgen.start_router()
+
def teardown_module(mod):
"Teardown the pytest environment"
tgen = get_topogen()
tgen.stop_topology()
+
def igmpv3_sources_save_to_json(tgen, sources_db):
sources_list = list(sources_db.keys())
sources_file = os.path.join(tgen.logdir, "sources.json")
@@ -89,6 +100,7 @@ def igmpv3_sources_save_to_json(tgen, sources_db):
with open(sources_file, "w") as f:
json.dump(sources_list, f)
+
def igmpv3_send_membership_reports(tgen, mode):
sources_file_path = os.path.join(tgen.logdir, "sources.json")
@@ -102,41 +114,90 @@ def igmpv3_send_membership_reports(tgen, mode):
f"--step {STEP}"
)
- h1 = tgen.gears['h1']
+ h1 = tgen.gears["h1"]
output = h1.run(cmd)
assert "Success" in output, f"Script failed. Output: {output}"
+
def igmpv3_send_membership_reports_include(tgen):
logger.info("Send IGMPv3 Membership Reports (MODE_IS_INCLUDE)")
igmpv3_send_membership_reports(tgen, "include")
+
+def igmpv3_wait_for_socket(tgen):
+ "Wait until pimd has an IGMP socket on r1-eth0."
+ logger.info("Wait for IGMP socket on r1-eth0")
+
+ # show ip igmp interface only lists interfaces that already have a socket.
+ expected = {"r1-eth0": {"querierIp": ROUTER_IP}}
+
+ def verify_socket():
+ r1 = tgen.gears["r1"]
+ data = r1.vtysh_cmd("show ip igmp interface r1-eth0 json", isjson=True)
+ return json_cmp(data, expected)
+
+ success, result = topotest.run_and_expect(verify_socket, None, count=30, wait=1)
+ assert success, f"IGMP socket on r1-eth0 did not come up! Last output:\n{result}"
+
+
+def igmpv3_group_source_count(data):
+ if not isinstance(data, dict):
+ return 0
+ iface = data.get("r1-eth0") or {}
+ for grp in iface.get("groups") or []:
+ if grp.get("group") == IGMP_GROUP:
+ return grp.get("sourcesCount") or 0
+ return 0
+
+
def igmpv3_send_membership_reports_block(tgen):
logger.info("Send IGMPv3 Membership Reports (BLOCK_OLD_SOURCES)")
igmpv3_send_membership_reports(tgen, "block")
+
def igmpv3_verify_group(tgen):
expected = {
- "r1-eth0": {
- "groups": [
- {
- "group": IGMP_GROUP,
- "sourcesCount": TOTAL_SOURCES
- }
- ]
- }
+ "r1-eth0": {"groups": [{"group": IGMP_GROUP, "sourcesCount": TOTAL_SOURCES}]}
}
+ # count starts at -1 so the first poll does not resend the burst the
+ # test just transmitted. A resend waits until the count has been flat
+ # for INCLUDE_RESEND_POLLS polls, so the retries are not spent inside
+ # one stall.
+ progress = {"count": -1, "resends": 0, "stalls": 0}
def verify_group():
r1 = tgen.gears["r1"]
data = r1.vtysh_cmd("show ip igmp groups json", isjson=True)
- return json_cmp(data, expected)
+ result = json_cmp(data, expected)
+ count = igmpv3_group_source_count(data)
+ if result is not None and count <= progress["count"]:
+ progress["stalls"] += 1
+ if (
+ progress["resends"] < MAX_INCLUDE_RESENDS
+ and progress["stalls"] >= INCLUDE_RESEND_POLLS
+ ):
+ progress["resends"] += 1
+ progress["stalls"] = 0
+ logger.info(
+ "IGMP group stalled at %s/%s sources, resending INCLUDE reports (%s/%s)",
+ count,
+ TOTAL_SOURCES,
+ progress["resends"],
+ MAX_INCLUDE_RESENDS,
+ )
+ igmpv3_send_membership_reports_include(tgen)
+ else:
+ progress["stalls"] = 0
+ progress["count"] = count
+ return result
success, result = topotest.run_and_expect(verify_group, None, count=60, wait=1)
assert success, f"IGMP group did not converge! Last output:\n{result}"
+
def igmpv3_verify_group_sources(tgen, sources_db):
r1 = tgen.gears["r1"]
data = r1.vtysh_cmd("show ip igmp sources json", isjson=True)
@@ -156,12 +217,14 @@ def igmpv3_verify_group_sources(tgen, sources_db):
assert not error_msg, f"IGMP Sources mismatch."
+
def igmpv3_verify_group_and_sources(tgen, sources_db):
logger.info("Verify IGMPv3 Group and Source States")
igmpv3_verify_group(tgen)
igmpv3_verify_group_sources(tgen, sources_db)
+
def igmpv3_start_traffic_capture(tgen):
logger.info("Start IGMPv3 traffic capture")
@@ -175,10 +238,11 @@ def igmpv3_start_traffic_capture(tgen):
h1.run(cmd)
time.sleep(2)
+
def igmpv3_wait_for_expiration(tgen):
logger.info("Wait for expiration of IGMPv3 sources")
- expected = {'totalGroups': 0 }
+ expected = {"totalGroups": 0}
def verify_group():
r1 = tgen.gears["r1"]
@@ -188,6 +252,7 @@ def igmpv3_wait_for_expiration(tgen):
success, result = topotest.run_and_expect(verify_group, None, count=60, wait=1)
assert success, f"IGMP group did not converge! Last output:\n{result}"
+
def igmpv3_stop_traffic_capture(tgen):
logger.info("Stop IGMPv3 traffic capture")
@@ -204,6 +269,7 @@ def igmpv3_stop_traffic_capture(tgen):
h1.run(cmd)
time.sleep(2)
+
def igmpv3_verify_specific_query(tgen, sources_db):
logger.info("Verify processing of IGMPv3 Group and Source Specific Queries")
@@ -227,7 +293,7 @@ def igmpv3_verify_specific_query(tgen, sources_db):
lmqc = data["r1-eth0"]["lastMemberQueryCount"]
for src, stats in sources_db.items():
- count = stats['query_count']
+ count = stats["query_count"]
assert count == lmqc, (
f"Source {src} mismatch! "
f"Expected {lmqc} Specific Queries, "
@@ -239,6 +305,7 @@ def igmpv3_verify_specific_query(tgen, sources_db):
f"(matches LMQC in config)"
)
+
def test_igmpv3_specific_query_frag(tgen, sources_db):
"Tests IGMP specefic query fragmentation"
@@ -248,6 +315,7 @@ def test_igmpv3_specific_query_frag(tgen, sources_db):
logger.info(f"Starting IGMPv3 fragmentation test ({TOTAL_SOURCES} sources)")
igmpv3_sources_save_to_json(tgen, sources_db)
+ igmpv3_wait_for_socket(tgen)
igmpv3_send_membership_reports_include(tgen)
igmpv3_verify_group_and_sources(tgen, sources_db)
igmpv3_start_traffic_capture(tgen)
@@ -257,6 +325,7 @@ def test_igmpv3_specific_query_frag(tgen, sources_db):
igmpv3_verify_specific_query(tgen, sources_db)
logger.info(f"Success!")
+
def test_memory_leak():
"Run the memory leak test and report results."
tgen = get_topogen()
@@ -265,6 +334,7 @@ def test_memory_leak():
tgen.report_memory_leaks()
+
if __name__ == "__main__":
args = ["-s"] + sys.argv[1:]
sys.exit(pytest.main(args))