2025年10月9日木曜日

Eth0->eth1転送

 


転送コード


以下は「片方のEthernetインターフェースで受け取ったフレームを、もう一方へそのまま転送する」最小のPython例です。Linux向け、root権限が必要です。手軽さ重視で Scapy 版と、標準ライブラリだけで動く raw socket(AF_PACKET) 版の2つを載せます。


1) Scapy版(簡単・分かりやすい)

#!/usr/bin/env python3 # bridge_scapy.py # 要: pip install scapy かつ root権限 from scapy.all import sniff, sendp import argparse import threading import signal import sys def forward_loop(in_if, out_if, bpf): # 受信 → そのまま送信 sniff( iface=in_if, prn=lambda pkt: sendp(pkt, iface=out_if, verbose=False), store=False, filter=bpf if bpf else None, promisc=True, ) def main(): ap = argparse.ArgumentParser(description="L2 frame forwarder (simple bridge) using Scapy") ap.add_argument("--if1", required=True, help="インターフェース1 (例: eth0)") ap.add_argument("--if2", required=True, help="インターフェース2 (例: eth1)") ap.add_argument("--bpf", default="", help="BPFフィルタ (例: 'arp or tcp')") args = ap.parse_args() # Ctrl+Cで両方止める def handle_sigint(signum, frame): sys.exit(0) signal.signal(signal.SIGINT, handle_sigint) t1 = threading.Thread(target=forward_loop, args=(args.if1, args.if2, args.bpf), daemon=True) t2 = threading.Thread(target=forward_loop, args=(args.if2, args.if1, args.bpf), daemon=True) t1.start(); t2.start() t1.join(); t2.join() if __name__ == "__main__": main()

使い方例:

sudo python3 bridge_scapy.py --if1 eth0 --if2 eth1 # 必要ならBPFで絞り込み sudo python3 bridge_scapy.py --if1 eth0 --if2 eth1 --bpf "arp or tcp"

ポイント

  • L2(イーサネット)フレームをそのまま転送します(ブリッジ動作)。

  • 速度より手軽さ重視。高スループットや低遅延が必要ならOSのブリッジ機能(ip link add br0 type bridge 等)やDPDK/ebpf等を検討してください。


2) 標準ライブラリだけ(AF_PACKET raw socket; Linux限定)

#!/usr/bin/env python3 # bridge_raw.py # 要: root権限 / Linux import socket import fcntl import struct import os import select import argparse import signal import sys SIOCGIFINDEX = 0x8933 ETH_P_ALL = 0x0003 # htons不要: socket作成で自動 def ifindex(sock, ifname: str) -> int: ifreq = struct.pack('16sH14s', ifname.encode('utf-8'), socket.AF_PACKET, b'\x00'*14) res = fcntl.ioctl(sock.fileno(), SIOCGIFINDEX, ifreq) return struct.unpack('16sH14s', res)[1] def bind_raw(ifname: str) -> socket.socket: s = socket.socket(socket.AF_PACKET, socket.SOCK_RAW, socket.htons(ETH_P_ALL)) s.setsockopt(socket.SOL_SOCKET, socket.SO_RCVBUF, 4 * 1024 * 1024) idx = ifindex(s, ifname) s.bind((ifname, 0)) # プロミスキャス有効化(best-effort) PACKET_ADD_MEMBERSHIP = 1 PACKET_MR_PROMISC = 1 mreq = struct.pack("IHH8s", idx, PACKET_MR_PROMISC, 0, b"\x00" * 8) try: s.setsockopt(socket.SOL_PACKET, PACKET_ADD_MEMBERSHIP, mreq) except OSError: pass return s def main(): ap = argparse.ArgumentParser(description="L2 frame forwarder (raw AF_PACKET)") ap.add_argument("--if1", required=True) ap.add_argument("--if2", required=True) ap.add_argument("--mtu", type=int, default=65535, help="受信バッファ(1フレーム)サイズ") args = ap.parse_args() def handle_sigint(signum, frame): sys.exit(0) signal.signal(signal.SIGINT, handle_sigint) s1 = bind_raw(args.if1) s2 = bind_raw(args.if2) # selectで双方向中継 while True: rlist, _, _ = select.select([s1, s2], [], []) for rs in rlist: try: data = rs.recv(args.mtu) if not data: continue # 受けた側と反対へ送る (s2 if rs is s1 else s1).send(data) except (BlockingIOError, InterruptedError): continue except OSError as e: # 一時的なエラーは無視 if e.errno in (11, 35): # EAGAIN/EWOULDBLOCK continue raise if __name__ == "__main__": main()

使い方:

sudo python3 bridge_raw.py --if1 eth0 --if2 eth1

運用メモ

  • どちらの方法も root権限 が必要です(生パケットアクセスのため)。

  • 片方/両方のIFでIPスタックによる処理を避けたい場合、ifconfig ethX promiscip link set ethX promisc on を確認。

  • 既にOS側でブリッジ(br0)を組んでいる場合は、本スクリプトと競合し得ます。片方にしてください。

  • フィルタしたい場合はScapy版の --bpf が手軽(例: --bpf "vlan and tcp")。

  • 高スループットが必要なら:OSブリッジ + ethtool でオフロード設定、あるいはeBPF/DPDK/専用スイッチを推奨。

必要なら、VLANタグの保持/書き換え、特定MACのみ通す、簡易L2学習(MACテーブル)付きに拡張したサンプルも用意できます。

あなた:
長いメッセージを送れるようバッファーを大きくできるコード
ChatGPT:

「長いメッセージを送れるようにバッファを大きくしたい」というのは、Ethernet転送コードで 一度に受け取れるパケットのサイズや送信バッファ容量を増やしたい という意味ですね。
以下のように、socket.setsockopt() を使って受信・送信バッファサイズを大きく設定できます。


🔧 バッファを拡張した例(AF_PACKET版)

#!/usr/bin/env python3 # bridge_buffered.py import socket, fcntl, struct, os, select, argparse, signal, sys SIOCGIFINDEX = 0x8933 ETH_P_ALL = 0x0003 # all protocols def ifindex(sock, ifname): ifreq = struct.pack('16sH14s', ifname.encode(), socket.AF_PACKET, b'\x00'*14) res = fcntl.ioctl(sock.fileno(), SIOCGIFINDEX, ifreq) return struct.unpack('16sH14s', res)[1] def bind_raw(ifname: str, bufsize=8*1024*1024): """raw socketを指定インターフェースにバインドし、バッファ拡張""" s = socket.socket(socket.AF_PACKET, socket.SOCK_RAW, socket.htons(ETH_P_ALL)) s.setsockopt(socket.SOL_SOCKET, socket.SO_RCVBUF, bufsize) # 受信バッファ s.setsockopt(socket.SOL_SOCKET, socket.SO_SNDBUF, bufsize) # 送信バッファ idx = ifindex(s, ifname) s.bind((ifname, 0)) return s def main(): ap = argparse.ArgumentParser(description="Ethernet bridge with larger buffers") ap.add_argument("--if1", required=True) ap.add_argument("--if2", required=True) ap.add_argument("--buf", type=int, default=8*1024*1024, help="受信・送信バッファサイズ (bytes)") ap.add_argument("--mtu", type=int, default=65535) args = ap.parse_args() s1 = bind_raw(args.if1, args.buf) s2 = bind_raw(args.if2, args.buf) print(f"Bridge running between {args.if1}{args.if2}") print(f"Buffer size: {args.buf:,} bytes") def handle_sigint(signum, frame): sys.exit(0) signal.signal(signal.SIGINT, handle_sigint) while True: r, _, _ = select.select([s1, s2], [], []) for rs in r: try: data = rs.recv(args.mtu) if not data: continue (s2 if rs is s1 else s1).send(data) except OSError as e: continue if __name__ == "__main__": main()

✅ ポイント

  • SO_RCVBUF / SO_SNDBUF により OS 側のソケットバッファを拡張。

  • --buf オプションで任意サイズを指定可能(例:--buf 16777216 → 16MB)。

  • Ethernetフレーム自体は最大 1518 バイト程度ですが、スクリプトの受信・転送キューを深くしてロスを防げます。

  • 実際のMTU(最大転送単位)を超えるサイズの「長いメッセージ」を送るには、アプリケーション層で分割するか、Jumbo Frame をNIC設定で有効にする必要があります(例:ip link set eth0 mtu 9000)。


💡 Jumbo Frameを併用する場合

sudo ip link set eth0 mtu 9000 sudo ip link set eth1 mtu 9000 sudo python3 bridge_buffered.py --if1 eth0 --if2 eth1 --buf 16777216 --mtu 9000

これで、大きなEthernetフレームを通しつつ、受信/送信バッファも拡大してパケットドロップを減らすことができます。


ご希望に応じて、

  • TCP/UDPレベルでの長メッセージ送受信(ソケット通信)版

  • 複数スレッド/非同期転送でのバッファリング強化版

なども示せますが、どちらの方向を想定されていますか?
(Ethernetフレーム中継?それともアプリ層でのデータ転送?)

あなた:
#!/usr/bin/env python3 # l2bridge.py import argparse import os import select import signal import socket import struct import sys import time ETH_P_ALL = 0x0003 # すべてのEtherType SOL_PACKET = 263 # from linux/sol_packet.h PACKET_ADD_MEMBERSHIP = 1 PACKET_DROP_MEMBERSHIP = 2 PACKET_MR_PROMISC = 1 PACKET_OUTGOING = 4 RUNNING = True def set_promisc(sock, ifindex: int, enable: bool = True): mreq = struct.pack("IHH8s", ifindex, PACKET_MR_PROMISC, 0, b"\x00"*8) opt = PACKET_ADD_MEMBERSHIP if enable else PACKET_DROP_MEMBERSHIP sock.setsockopt(SOL_PACKET, opt, mreq) def if_nametoindex(ifname: str) -> int: return socket.if_nametoindex(ifname) def open_raw_socket(ifname: str) -> socket.socket: s = socket.socket(socket.AF_PACKET, socket.SOCK_RAW, socket.htons(ETH_P_ALL)) # インターフェイスにバインド s.bind((ifname, 0)) # 受信バッファ拡大(任意) try: s.setsockopt(socket.SOL_SOCKET, socket.SO_RCVBUF, 4 * 1024 * 1024) except OSError: pass return s def forward_loop(if1: str, if2: str, print_stats: bool = True): s1 = open_raw_socket(if1) s2 = open_raw_socket(if2) # プロミスキャス有効化(これをしないと一部フレームを拾えない環境がある) set_promisc(s1, if_nametoindex(if1), True) set_promisc(s2, if_nametoindex(if2), True) cnt12 = cnt21 = 0 last_report = time.time() def _cleanup(*_): global RUNNING RUNNING = False signal.signal(signal.SIGINT, _cleanup) signal.signal(signal.SIGTERM, _cleanup) print(f"[+] Bridging {if1} <-> {if2} (Ctrl-C to stop)") # selectで両方監視 while RUNNING: r, _, _ = select.select([s1, s2], [], [], 1.0) for rs in r: try: # data: bytes, addr: (proto, ifindex, pkttype, halen, addrbytes) data, addr = rs.recvfrom(65535) pkttype = addr[2] if len(addr) >= 3 else None # 自分が送出した(または同一ホスト発)フレームは送り返さない if pkttype == PACKET_OUTGOING: continue if rs is s1: s2.send(data) cnt12 += 1 else: s1.send(data) cnt21 += 1 except BlockingIOError: continue except OSError as e: # インターフェイスがダウンした等 print(f"[!] socket error: {e}", file=sys.stderr) time.sleep(0.2) # 簡易統計 if print_stats and time.time() - last_report >= 5: print(f"[stats] {if1}→{if2}: {cnt12} frames, {if2}→{if1}: {cnt21} frames") last_report = time.time() # 後片付け try: set_promisc(s1, if_nametoindex(if1), False) set_promisc(s2, if_nametoindex(if2), False) except OSError: pass s1.close() s2.close() print("[+] Stopped.") def main(): if os.geteuid() != 0: print("This script must be run as root (RAW sockets).", file=sys.stderr) sys.exit(1) p = argparse.ArgumentParser(description="Simple L2 bridge between two Ethernet interfaces.") p.add_argument("-i1", "--iface1", required=True, help="First interface name (e.g., eth0)") p.add_argument("-i2", "--iface2", required=True, help="Second interface name (e.g., eth1)") p.add_argument("--no-stats", action="store_true", help="Disable periodic stats print") args = p.parse_args() try: forward_loop(args.iface1, args.iface2, print_stats=not args.no_stats) except KeyboardInterrupt: pass if __name__ == "__main__": main()   
ChatGPT:

いいコードです!そのままでも動きますが、「長いメッセージを落とさず捌く」目的なら ソケットバッファ拡張(任意) Jumbo MTU、それから 送信経路の遅延削減 を入れると安定します。
下の改訂版では次を追加しています:

  • --rcvbuf / --sndbuf で OS ソケットバッファを拡張(既定: 16MB)

  • 実際に設定された有効値を getsockopt で取得して表示(Linuxは内部的に2倍などの調整が入るため)

  • --mtu で1回の recvfrom() の読み出し上限を変更(Jumbo Frameを使うなら9000など)

  • 送信側で QdiscバイパスPACKET_QDISC_BYPASS)を有効化してカーネルキューの遅延を減少

  • 小さな後方互換の強化(例外処理など)

注: Ethernetの1フレーム自体はMTUを超えて巨大にはできません。長いアプリデータは上位プロトコル(TCP等)が分割します。ここでは「落としにくくするための受送信キュー拡張」と「Jumbo対応」を行っています。


改訂版 l2bridge.py

#!/usr/bin/env python3 # l2bridge.py (buffer-tuned) import argparse import os import select import signal import socket import struct import sys import time ETH_P_ALL = 0x0003 # すべてのEtherType SOL_PACKET = 263 # from linux/sol_packet.h PACKET_ADD_MEMBERSHIP = 1 PACKET_DROP_MEMBERSHIP = 2 PACKET_MR_PROMISC = 1 PACKET_OUTGOING = 4 PACKET_QDISC_BYPASS = 20 # bypass qdisc for lower TX latency (Linux >=3.14) RUNNING = True def set_promisc(sock, ifindex: int, enable: bool = True): mreq = struct.pack("IHH8s", ifindex, PACKET_MR_PROMISC, 0, b"\x00"*8) opt = PACKET_ADD_MEMBERSHIP if enable else PACKET_DROP_MEMBERSHIP sock.setsockopt(SOL_PACKET, opt, mreq) def if_nametoindex(ifname: str) -> int: return socket.if_nametoindex(ifname) def open_raw_socket(ifname: str, rcvbuf: int, sndbuf: int, qdisc_bypass: bool) -> socket.socket: s = socket.socket(socket.AF_PACKET, socket.SOCK_RAW, socket.htons(ETH_P_ALL)) # 受信/送信バッファ拡大(OSが内部で丸める/倍加することがある) if rcvbuf: try: s.setsockopt(socket.SOL_SOCKET, socket.SO_RCVBUF, rcvbuf) except OSError: pass if sndbuf: try: s.setsockopt(socket.SOL_SOCKET, socket.SO_SNDBUF, sndbuf) except OSError: pass # 送信経路の遅延を抑える(best-effort) if qdisc_bypass: try: s.setsockopt(SOL_PACKET, PACKET_QDISC_BYPASS, 1) except OSError: pass # インターフェイスにバインド s.bind((ifname, 0)) return s def human_bytes(n: int) -> str: for unit in ["B","KB","MB","GB"]: if n < 1024 or unit == "GB": return f"{n:.0f} {unit}" n /= 1024.0 def forward_loop(if1: str, if2: str, mtu: int, rcvbuf: int, sndbuf: int, qdisc_bypass: bool, print_stats: bool = True): s1 = open_raw_socket(if1, rcvbuf, sndbuf, qdisc_bypass) s2 = open_raw_socket(if2, rcvbuf, sndbuf, qdisc_bypass) # プロミスキャス有効化 set_promisc(s1, if_nametoindex(if1), True) set_promisc(s2, if_nametoindex(if2), True) # カーネルが実際に適用したサイズを表示(デバッグ/検証用) eff_rcv1 = s1.getsockopt(socket.SOL_SOCKET, socket.SO_RCVBUF) eff_snd1 = s1.getsockopt(socket.SOL_SOCKET, socket.SO_SNDBUF) eff_rcv2 = s2.getsockopt(socket.SOL_SOCKET, socket.SO_RCVBUF) eff_snd2 = s2.getsockopt(socket.SOL_SOCKET, socket.SO_SNDBUF) cnt12 = cnt21 = 0 last_report = time.time() def _cleanup(*_): global RUNNING RUNNING = False signal.signal(signal.SIGINT, _cleanup) signal.signal(signal.SIGTERM, _cleanup) print(f"[+] Bridging {if1} <-> {if2} (Ctrl-C to stop)") print(f" {if1}: RCV={human_bytes(eff_rcv1)}, SND={human_bytes(eff_snd1)} | " f"{if2}: RCV={human_bytes(eff_rcv2)}, SND={human_bytes(eff_snd2)}") print(f" recv buffer per call (mtu): {mtu} bytes" + (" | qdisc_bypass=on" if qdisc_bypass else "")) # selectで両方監視 while RUNNING: r, _, _ = select.select([s1, s2], [], [], 1.0) for rs in r: try: # data: bytes, addr: (proto, ifindex, pkttype, halen, addrbytes) data, addr = rs.recvfrom(mtu) pkttype = addr[2] if len(addr) >= 3 else None # 自分が送出した(または同一ホスト発)フレームは送り返さない if pkttype == PACKET_OUTGOING: continue if rs is s1: s2.send(data) cnt12 += 1 else: s1.send(data) cnt21 += 1 except (BlockingIOError, InterruptedError): continue except OSError as e: # インターフェイスがダウンした等 print(f"[!] socket error: {e}", file=sys.stderr) time.sleep(0.2) # 簡易統計 if print_stats and time.time() - last_report >= 5: print(f"[stats] {if1}{if2}: {cnt12} frames, {if2}{if1}: {cnt21} frames") last_report = time.time() # 後片付け try: set_promisc(s1, if_nametoindex(if1), False) set_promisc(s2, if_nametoindex(if2), False) except OSError: pass s1.close() s2.close() print("[+] Stopped.") def main(): if os.geteuid() != 0: print("This script must be run as root (RAW sockets).", file=sys.stderr) sys.exit(1) p = argparse.ArgumentParser(description="Simple L2 bridge between two Ethernet interfaces (buffer-tuned).") p.add_argument("-i1", "--iface1", required=True, help="First interface name (e.g., eth0)") p.add_argument("-i2", "--iface2", required=True, help="Second interface name (e.g., eth1)") p.add_argument("--mtu", type=int, default=65535, help="recv() buffer size per frame (Jumboなら9000など)") p.add_argument("--rcvbuf", type=int, default=16*1024*1024, help="SO_RCVBUF bytes") p.add_argument("--sndbuf", type=int, default=16*1024*1024, help="SO_SNDBUF bytes") p.add_argument("--qdisc-bypass", action="store_true", help="Enable PACKET_QDISC_BYPASS for low-latency TX") p.add_argument("--no-stats", action="store_true", help="Disable periodic stats print") args = p.parse_args() try: forward_loop( args.iface1, args.iface2, mtu=args.mtu, rcvbuf=args.rcvbuf, sndbuf=args.sndbuf, qdisc_bypass=args.qdisc_bypass, print_stats=not args.no_stats ) except KeyboardInterrupt: pass if __name__ == "__main__": main()

使い方例

# 標準: 受信/送信バッファを16MB、MTUは65535バイト読み(上限) sudo ./l2bridge.py -i1 eth0 -i2 eth1 # バッファを32MBに、Jumboフレーム(例: 9000B想定)に合わせて読み出しサイズも変更 sudo ip link set eth0 mtu 9000 sudo ip link set eth1 mtu 9000 sudo ./l2bridge.py -i1 eth0 -i2 eth1 --rcvbuf 33554432 --sndbuf 33554432 --mtu 9000 --qdisc-bypass

補足(性能をさらに攻めるなら)

  • TPACKET_V3(PACKET_RX_RING) を使ったリングバッファI/Oにするとドロップ耐性が大幅に上がります(実装は少し大きくなる)。

  • マルチキュー / CPUピン留めtasksetSO_BUSY_POLL)で遅延を詰められます。

  • 実転送が最優先なら OSのブリッジ(ip link add br0 type bridge が最も高速・堅牢です。

この改訂版で、バッファ不足による取りこぼしはかなり減るはず。必要なら TPACKET_V3 版のテンプレも出します。

あなた:
このコードをジャンボフレームに対応するように変更してください
ChatGPT:

了解です。頂いた l2bridge.pyジャンボフレーム(例: MTU=9000)対応に強化した版を示します。
ポイントは次のとおりです。

  • 受信バッファ長(recvfrom() の第1引数)を MTUに合わせて可変

  • SO_RCVBUF / SO_SNDBUF を十分に拡張(既定 16MB)

  • 任意で インターフェースMTUをioctlで設定--set-mtu

  • (任意)qdiscバイパスで TX 遅延を削減

重要: 実際にジャンボフレームを通すには、両方のNIC・スイッチ・対向機器すべてで MTU を揃える必要があります。
本スクリプトの --set-mtu は NIC 側の MTU 設定(root 必須)を行います。スイッチ側は別途設定してください。


変更済みコード(そのまま置き換え可)

#!/usr/bin/env python3 # l2bridge.py (jumbo-ready) import argparse import os import select import signal import socket import struct import sys import time ETH_P_ALL = 0x0003 # all EtherTypes SOL_PACKET = 263 PACKET_ADD_MEMBERSHIP = 1 PACKET_DROP_MEMBERSHIP = 2 PACKET_MR_PROMISC = 1 PACKET_OUTGOING = 4 PACKET_QDISC_BYPASS = 20 # optional, Linux >= 3.14 SIOCGIFMTU = 0x8921 SIOCSIFMTU = 0x8922 RUNNING = True def set_promisc(sock, ifindex: int, enable: bool = True): mreq = struct.pack("IHH8s", ifindex, PACKET_MR_PROMISC, 0, b"\x00"*8) opt = PACKET_ADD_MEMBERSHIP if enable else PACKET_DROP_MEMBERSHIP sock.setsockopt(SOL_PACKET, opt, mreq) def if_nametoindex(ifname: str) -> int: return socket.if_nametoindex(ifname) def get_mtu(ifname: str) -> int: # struct ifreq { char ifr_name[16]; union { int ifru_mtu; ... } ifr_ifru; } s = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) ifreq = struct.pack("16sH14s", ifname.encode(), 0, b"\x00"*14) try: res = fcntl_ioctl(s.fileno(), SIOCGIFMTU, ifreq) _, mtu, _ = struct.unpack("16sH14s", res) return int(mtu) except Exception: return -1 finally: s.close() def set_mtu_ioctl(ifname: str, mtu: int): s = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) # struct ifreq with ifr_mtu (int) packed into the last 4 bytes ifreq = struct.pack("16sI", ifname.encode(), mtu) # pad to sizeof(struct ifreq) ~ 40 bytes on 64-bit ifreq = ifreq.ljust(40, b"\x00") try: fcntl_ioctl(s.fileno(), SIOCSIFMTU, ifreq) finally: s.close() def fcntl_ioctl(fd, op, data): # 小さなヘルパ(import fcntl を避けるため) import fcntl as _fcntl return _fcntl.ioctl(fd, op, data) def open_raw_socket(ifname: str, rcvbuf: int, sndbuf: int, qdisc_bypass: bool) -> socket.socket: s = socket.socket(socket.AF_PACKET, socket.SOCK_RAW, socket.htons(ETH_P_ALL)) # OSソケットバッファを十分に拡張 if rcvbuf: try: s.setsockopt(socket.SOL_SOCKET, socket.SO_RCVBUF, rcvbuf) except OSError: pass if sndbuf: try: s.setsockopt(socket.SOL_SOCKET, socket.SO_SNDBUF, sndbuf) except OSError: pass # qdiscバイパス(任意) if qdisc_bypass: try: s.setsockopt(SOL_PACKET, PACKET_QDISC_BYPASS, 1) except OSError: pass s.bind((ifname, 0)) return s def human_bytes(n: int) -> str: for unit in ["B","KB","MB","GB"]: if n < 1024 or unit == "GB": return f"{n:.0f} {unit}" n /= 1024.0 def forward_loop(if1: str, if2: str, recv_len: int, rcvbuf: int, sndbuf: int, qdisc_bypass: bool, print_stats: bool = True): s1 = open_raw_socket(if1, rcvbuf, sndbuf, qdisc_bypass) s2 = open_raw_socket(if2, rcvbuf, sndbuf, qdisc_bypass) # プロミスキャス set_promisc(s1, if_nametoindex(if1), True) set_promisc(s2, if_nametoindex(if2), True) eff_rcv1 = s1.getsockopt(socket.SOL_SOCKET, socket.SO_RCVBUF) eff_snd1 = s1.getsockopt(socket.SOL_SOCKET, socket.SO_SNDBUF) eff_rcv2 = s2.getsockopt(socket.SOL_SOCKET, socket.SO_RCVBUF) eff_snd2 = s2.getsockopt(socket.SOL_SOCKET, socket.SO_SNDBUF) cnt12 = cnt21 = 0 last_report = time.time() def _cleanup(*_): global RUNNING RUNNING = False signal.signal(signal.SIGINT, _cleanup) signal.signal(signal.SIGTERM, _cleanup) print(f"[+] Bridging {if1} <-> {if2} (Ctrl-C to stop)") print(f" {if1}: RCV={human_bytes(eff_rcv1)}, SND={human_bytes(eff_snd1)} | " f"{if2}: RCV={human_bytes(eff_rcv2)}, SND={human_bytes(eff_snd2)}") print(f" recvfrom() size: {recv_len} bytes" + (" | qdisc_bypass=on" if qdisc_bypass else "")) while RUNNING: r, _, _ = select.select([s1, s2], [], [], 1.0) for rs in r: try: data, addr = rs.recvfrom(recv_len) pkttype = addr[2] if len(addr) >= 3 else None if pkttype == PACKET_OUTGOING: continue if rs is s1: s2.send(data); cnt12 += 1 else: s1.send(data); cnt21 += 1 except (BlockingIOError, InterruptedError): continue except OSError as e: print(f"[!] socket error: {e}", file=sys.stderr) time.sleep(0.2) if print_stats and time.time() - last_report >= 5: print(f"[stats] {if1}{if2}: {cnt12} frames, {if2}{if1}: {cnt21} frames") last_report = time.time() try: set_promisc(s1, if_nametoindex(if1), False) set_promisc(s2, if_nametoindex(if2), False) except OSError: pass s1.close(); s2.close() print("[+] Stopped.") def main(): if os.geteuid() != 0: print("This script must be run as root (RAW sockets).", file=sys.stderr) sys.exit(1) p = argparse.ArgumentParser(description="Simple L2 bridge between two Ethernet interfaces (jumbo-ready).") p.add_argument("-i1", "--iface1", required=True, help="First interface (e.g., eth0)") p.add_argument("-i2", "--iface2", required=True, help="Second interface (e.g., eth1)") # 受信バッファ長(1回のrecvfromで読み込む最大サイズ)。Jumboなら9000以上を推奨。 p.add_argument("--recv-len", type=int, default=16384, help="recvfrom() size in bytes (>= 実際のフレーム長; default 16384)") p.add_argument("--rcvbuf", type=int, default=16*1024*1024, help="SO_RCVBUF bytes (default 16MB)") p.add_argument("--sndbuf", type=int, default=16*1024*1024, help="SO_SNDBUF bytes (default 16MB)") p.add_argument("--qdisc-bypass", action="store_true", help="Enable PACKET_QDISC_BYPASS") # 任意: NICのMTUをioctlで設定 p.add_argument("--set-mtu", type=int, help="Set MTU on both interfaces via ioctl (root required)") p.add_argument("--no-stats", action="store_true", help="Disable periodic stats print") args = p.parse_args() # 任意: MTU設定(NIC側) if args.set_mtu: for ifn in (args.iface1, args.iface2): try: set_mtu_ioctl(ifn, args.set_mtu) print(f"[+] Set MTU {args.set_mtu} on {ifn}") except OSError as e: print(f"[!] Failed to set MTU on {ifn}: {e}", file=sys.stderr) try: forward_loop( args.iface1, args.iface2, recv_len=args.recv_len, rcvbuf=args.rcvbuf, sndbuf=args.sndbuf, qdisc_bypass=args.qdisc_bypass, print_stats=not args.no_stats ) except KeyboardInterrupt: pass if __name__ == "__main__": main()

使い方

# 1) スイッチ/対向もジャンボ対応済みの場合 sudo ./l2bridge.py -i1 eth0 -i2 eth1 --recv-len 16384 --qdisc-bypass # 2) NICのMTUをこのツールから設定してブリッジ(例: 9000) sudo ./l2bridge.py -i1 eth0 -i2 eth1 --set-mtu 9000 --recv-len 16384 --qdisc-bypass # 3) さらにバッファ増強(ドロップ耐性向上) sudo ./l2bridge.py -i1 eth0 -i2 eth1 --set-mtu 9000 --recv-len 16384 \ --rcvbuf 33554432 --sndbuf 33554432 --qdisc-bypass

実運用のコツ

  • --recv-len実際の最大フレーム長以上にしてください(例: 9000Bなら 16384 にして余裕を持たせる)。

  • 両IFとも 同じMTU に揃えること。経路途中のスイッチ/ルータもジャンボ許可が必要です。

  • 取りこぼしがある場合は --rcvbuf/--sndbuf を増やし、CPU使用率が高いなら taskset でCPUピン留めや ethtool -K でオフロード調整を検討してください。

必要なら、TPACKET_V3(PACKET_RX_RING)版でさらにドロップを抑える実装もお出しします。

あなた:
TPACKET_V3(PACKET_RX_RING)版をお願いします
ChatGPT:

やります!Linuxの TPACKET_V3(PACKET_RX_RING) を使って、受信をカーネル共有リング経由で行う L2 ブリッジの Python 実装を示します。
ジャンボフレーム(MTU 9000 など)にも対応しやすいようにパラメータ化しています。送信は通常の send() を使い、受信だけを V3 リングで高速化しています(TX_RING まで入れるとコードがかなり大きくなるため)。

必要: Linux, root, 64bit 環境想定(構造体サイズは x86_64 の UAPI に合わせています)


l2bridge_tpacketv3.py

#!/usr/bin/env python3 # -*- coding: utf-8 -*- # l2bridge_tpacketv3.py # # L2 bridge using AF_PACKET TPACKET_V3 (PACKET_RX_RING) for RX side. # - Linux / root required # - Jumbo-friendly (tunable block/frame sizes) # import argparse import ctypes as C import mmap import os import select import signal import socket import struct import sys import time # ---------- Linux constants ---------- ETH_P_ALL = 0x0003 SOL_PACKET = 263 PACKET_ADD_MEMBERSHIP = 1 PACKET_DROP_MEMBERSHIP = 2 PACKET_MR_PROMISC = 1 PACKET_OUTGOING = 4 PACKET_VERSION = 10 TPACKET_V3 = 3 PACKET_RX_RING = 5 PACKET_QDISC_BYPASS = 20 TP_STATUS_KERNEL = 0 TP_STATUS_USER = 1 # poll/select POLLIN = 0x0001 # ---------- ctypes structures from linux/uapi ---------- # struct tpacket_req3 class tpacket_req3(C.Structure): _fields_ = [ ("tp_block_size", C.c_uint), ("tp_block_nr", C.c_uint), ("tp_frame_size", C.c_uint), ("tp_frame_nr", C.c_uint), ("tp_retire_blk_tov", C.c_uint), ("tp_sizeof_priv", C.c_uint), ("tp_feature_req_word", C.c_uint), ] # struct tpacket_bd_ts class tpacket_bd_ts(C.Structure): _fields_ = [("ts_sec", C.c_uint), ("ts_nsec", C.c_uint)] # struct tpacket_hdr_v1 (block header core) class tpacket_hdr_v1(C.Structure): _fields_ = [ ("block_status", C.c_uint), ("num_pkts", C.c_uint), ("offset_to_first_pkt", C.c_uint), ("blk_len", C.c_uint), ("seq_num", C.c_ulonglong), ("ts_first_pkt", tpacket_bd_ts), ("ts_last_pkt", tpacket_bd_ts), ] # struct tpacket_block_desc: version & offset + hdr_v1 class tpacket_block_desc(C.Structure): _fields_ = [ ("version", C.c_uint), ("offset_to_priv", C.c_uint), ("hdr", tpacket_hdr_v1), ] # struct tpacket3_hdr (per-packet header inside a block) class tpacket3_hdr(C.Structure): _fields_ = [ ("tp_next_offset", C.c_uint), ("tp_sec", C.c_uint), ("tp_nsec", C.c_uint), ("tp_snaplen", C.c_uint), ("tp_len", C.c_uint), ("tp_status", C.c_uint), ("tp_mac", C.c_ushort), ("tp_net", C.c_ushort), ("hv1", C.c_ushort), # placeholder (vlan etc.), not used here ("hv2", C.c_ushort), ] # ---------- helpers ---------- def set_promisc(sock: socket.socket, ifindex: int, enable: bool = True): mreq = struct.pack("IHH8s", ifindex, PACKET_MR_PROMISC, 0, b"\x00"*8) opt = PACKET_ADD_MEMBERSHIP if enable else PACKET_DROP_MEMBERSHIP sock.setsockopt(SOL_PACKET, opt, mreq) def open_rx_ring_sock(ifname: str, block_size: int, block_nr: int, frame_size: int, retire_tov_ms: int, rcvbuf: int, qdisc_bypass: bool): s = socket.socket(socket.AF_PACKET, socket.SOCK_RAW, socket.htons(ETH_P_ALL)) # ring 前に version を V3 へ s.setsockopt(SOL_PACKET, PACKET_VERSION, struct.pack("I", TPACKET_V3)) if rcvbuf: try: s.setsockopt(socket.SOL_SOCKET, socket.SO_RCVBUF, int(rcvbuf)) except OSError: pass if qdisc_bypass: try: s.setsockopt(SOL_PACKET, PACKET_QDISC_BYPASS, 1) except OSError: pass # bind iface s.bind((ifname, 0)) # ring params # frame_nr は block_size/frame_size * block_nr frames_per_block = block_size // frame_size frame_nr = frames_per_block * block_nr req = tpacket_req3() req.tp_block_size = block_size req.tp_block_nr = block_nr req.tp_frame_size = frame_size req.tp_frame_nr = frame_nr req.tp_retire_blk_tov = retire_tov_ms req.tp_sizeof_priv = 0 req.tp_feature_req_word = 0 # apply RX_RING s.setsockopt(SOL_PACKET, PACKET_RX_RING, bytes(req)) # mmap ring ring_len = block_size * block_nr ring = mmap.mmap(s.fileno(), ring_len, flags=mmap.MAP_SHARED, prot=(mmap.PROT_READ | mmap.PROT_WRITE), offset=0) # promisc set_promisc(s, socket.if_nametoindex(ifname), True) return s, ring, frames_per_block, ring_len def close_rx_ring_sock(s: socket.socket, ring: mmap.mmap, ifname: str): try: set_promisc(s, socket.if_nametoindex(ifname), False) except OSError: pass try: ring.close() finally: s.close() def human(n: int): for u in ("B", "KB", "MB", "GB"): if n < 1024 or u == "GB": return f"{n:.0f}{u}" n /= 1024 # ---------- ring reader ---------- class RingReaderV3: def __init__(self, sock: socket.socket, ring: mmap.mmap, block_size: int, block_nr: int, frames_per_block: int): self.sock = sock self.ring = ring self.block_size = block_size self.block_nr = block_nr self.frames_per_block = frames_per_block self.cur_block = 0 self.block_desc_size = C.sizeof(tpacket_block_desc) self.pfd = select.poll() self.pfd.register(self.sock.fileno(), POLLIN) def _block_ptr(self, bidx: int): off = bidx * self.block_size return off def _as_block_desc(self, block_offset: int): # create ctypes view from mmap buffer buf = (C.c_char * self.block_desc_size).from_buffer(self.ring, block_offset) desc = tpacket_block_desc.from_buffer(buf) return desc def _packet_hdr_at(self, base_off: int, pkt_off: int): ptr = base_off + pkt_off hdr = tpacket3_hdr.from_buffer((C.c_char * C.sizeof(tpacket3_hdr)).from_buffer(self.ring, ptr)) return hdr, ptr def next_packets(self): # 1ブロックずつ処理。ユーザ側に渡す: [(data_mv1), (data_mv2), ...] while True: boff = self._block_ptr(self.cur_block) desc = self._as_block_desc(boff) if desc.hdr.block_status & TP_STATUS_USER == 0: # データ未到着。pollで待つ。 self.pfd.poll(1000) # 1s timeout # もう一度確認 if desc.hdr.block_status & TP_STATUS_USER == 0: continue # データあり pkts = [] pkt_off = desc.hdr.offset_to_first_pkt for _ in range(desc.hdr.num_pkts): ph, ph_off = self._packet_hdr_at(boff, pkt_off) data_off = ph_off + ph.tp_mac data_end = data_off + ph.tp_snaplen # memoryviewでゼロコピー参照 mv = memoryview(self.ring)[data_off:data_end] pkts.append(mv) if ph.tp_next_offset == 0: break pkt_off += ph.tp_next_offset # 処理済みとして KERNEL に返す desc.hdr.block_status = TP_STATUS_KERNEL # 次ブロックへ self.cur_block = (self.cur_block + 1) % self.block_nr return pkts # ---------- main bridge ---------- RUNNING = True def _cleanup(*_): global RUNNING RUNNING = False def main(): ap = argparse.ArgumentParser(description="L2 bridge using TPACKET_V3 RX ring (jumbo-friendly).") ap.add_argument("-i1", "--iface1", required=True) ap.add_argument("-i2", "--iface2", required=True) # --- ring params (tunable) --- ap.add_argument("--block-size", type=int, default=1<<20, help="block size bytes (default 1MB)") ap.add_argument("--block-nr", type=int, default=64, help="number of blocks (default 64)") ap.add_argument("--frame-size", type=int, default=2048, help="frame size bytes (default 2048; >= MTU + headroom)") ap.add_argument("--retire-tov", type=int, default=60, help="retire block timeout ms (default 60)") ap.add_argument("--rcvbuf", type=int, default=32*1024*1024, help="SO_RCVBUF") ap.add_argument("--sndbuf", type=int, default=32*1024*1024, help="SO_SNDBUF") ap.add_argument("--qdisc-bypass", action="store_true", help="Enable PACKET_QDISC_BYPASS") ap.add_argument("--no-stats", action="store_true") args = ap.parse_args() if os.geteuid() != 0: print("Run as root.", file=sys.stderr) sys.exit(1) # open TX sockets (普通のRAW) と RXリングソケットを分ける # RXリング: s1r/s2r, 送信用: s1tx/s2tx s1r, ring1, fpb1, rlen1 = open_rx_ring_sock( args.iface1, args.block_size, args.block_nr, args.frame_size, args.retire_tov, args.rcvbuf, args.qdisc_bypass) s2r, ring2, fpb2, rlen2 = open_rx_ring_sock( args.iface2, args.block_size, args.block_nr, args.frame_size, args.retire_tov, args.rcvbuf, args.qdisc_bypass) # TX ソケット(通常 RAW)。送信バッファ拡張。 s1tx = socket.socket(socket.AF_PACKET, socket.SOCK_RAW, socket.htons(ETH_P_ALL)) s1tx.bind((args.iface1, 0)) s2tx = socket.socket(socket.AF_PACKET, socket.SOCK_RAW, socket.htons(ETH_P_ALL)) s2tx.bind((args.iface2, 0)) try: s1tx.setsockopt(socket.SOL_SOCKET, socket.SO_SNDBUF, args.sndbuf) except OSError: pass try: s2tx.setsockopt(socket.SOL_SOCKET, socket.SO_SNDBUF, args.sndbuf) except OSError: pass if args.qdisc_bypass: for txs in (s1tx, s2tx): try: txs.setsockopt(SOL_PACKET, PACKET_QDISC_BYPASS, 1) except OSError: pass rr1 = RingReaderV3(s1r, ring1, args.block_size, args.block_nr, fpb1) rr2 = RingReaderV3(s2r, ring2, args.block_size, args.block_nr, fpb2) cnt12 = cnt21 = 0 last = time.time() signal.signal(signal.SIGINT, _cleanup) signal.signal(signal.SIGTERM, _cleanup) print(f"[+] TPACKET_V3 bridge {args.iface1} <-> {args.iface2}") print(f" block={human(args.block_size)} x {args.block_nr}, frame_size={args.frame_size}, rcvbuf={human(args.rcvbuf)} sndbuf={human(args.sndbuf)}") print(" HINT: For jumbo (MTU 9000), set frame_size >= 10240 to keep headroom.") while RUNNING: # それぞれのリングからまとめてパケットを受け取り、反対側へ送信 # poll は各 RingReader 内で実施(1ブロック単位) pkts1 = rr1.next_packets() for mv in pkts1: # 自ホスト発(PACKET_OUTGOING)の除外は RX_RING 側で既に行われる try: s2tx.send(mv) # memoryview をそのまま送る cnt12 += 1 except OSError: pass pkts2 = rr2.next_packets() for mv in pkts2: try: s1tx.send(mv) cnt21 += 1 except OSError: pass if not args.no_stats and time.time() - last >= 5: print(f"[stats] {args.iface1}{args.iface2}: {cnt12} | {args.iface2}{args.iface1}: {cnt21}") last = time.time() # cleanup close_rx_ring_sock(s1r, ring1, args.iface1) close_rx_ring_sock(s2r, ring2, args.iface2) s1tx.close(); s2tx.close() print("[+] Stopped.") if __name__ == "__main__": main()

使い方

# 標準設定で起動(1MBブロック x 64、フレーム 2048B) sudo ./l2bridge_tpacketv3.py -i1 eth0 -i2 eth1 # ジャンボ(MTU 9000)なら余裕を持って frame_size を大きめに sudo ./l2bridge_tpacketv3.py -i1 eth0 -i2 eth1 --frame-size 12288 --block-size $((1<<20)) --block-nr 64 # バッファ増強&低遅延(qdiscバイパス) sudo ./l2bridge_tpacketv3.py -i1 eth0 -i2 eth1 --frame-size 12288 --rcvbuf $((64*1024*1024)) \ --sndbuf $((64*1024*1024)) --qdisc-bypass

パラメータ設計の目安

  • frame_size: MTU + 512 程度に(アラインメントやメタデータ分の余裕)。MTU 9000 → 12288(12KB)推奨。

  • block_size: frame_size * N がページサイズの倍数になるように。1MB(1048576)などが扱いやすいです。

  • block_nr: トラフィックに合わせて。合計リングサイズ block_size * block_nr は数十〜数百MBまで増やすとドロップ耐性向上。

注意・補足

  • 経路上の全機器(NIC/スイッチ/対向)が同じ MTU で揃っていないとジャンボは通りません。

  • 本コードは RX側だけ V3 リング を使っています。TX もリングにするとさらにレイテンシを詰められますが、実装が大幅に複雑になります。

  • カーネルの UAPI に強く依存します。64bit Linux/x86_64 を想定した ctypes 定義です。環境差異がある場合は ctypes フィールドサイズを調整してください。

  • さらにスループットを上げたい場合:

    • CPUピン留め(taskset)、RPS/RFS チューニング、ethtool -G で NIC リング調整

    • マルチキュー(キューごとにソケットを分け、スレッドで並列受信)

    • VLAN/PPPoE 等のヘッダ有無に応じて frame_size を増やす

必要であれば TX_RING 併用版 や、マルチキュー対応(キューごとにスレッド)版 も作ります。気になる要件(目標スループット、遅延、CPU数)を教えてください。

0 件のコメント:

コメントを投稿