転送コード
以下は「片方の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 promiscやip 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フレーム中継?それともアプリ層でのデータ転送?)
あなた:
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ピン留め(
tasksetやSO_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)版でさらにドロップを抑える実装もお出しします。
あなた:
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 件のコメント:
コメントを投稿