Skip to content

Commit 11dbf7d

Browse files
committed
Adding udp_listen
1 parent ffbbaef commit 11dbf7d

2 files changed

Lines changed: 80 additions & 1 deletion

File tree

README.md

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,8 @@ The following "special" binaries are included in the toolbox:
2828
Prepend a RFC3339 timestamp of nanosecond resolution to each line on stdin
2929
* **to_csv**
3030
Processes a crowsnest log file into a set of "topic-specific" csv files
31+
* **udp_listen**
32+
Eavesdrop onto UDP traffic and output to stdout, currently only supports multicast.
3133

3234

3335
## Recipes
@@ -36,7 +38,7 @@ The following are "recipes" for "run" commands that can be used with this image.
3638

3739
* Injecting data from "any" source into a mqtt broker using the standard brefv format (examplified by a multicast stream). Every UDP packet gets base64-encoded and packaged into a brefv envelope and then published to the broker:
3840
```
39-
socat -u UDP4-RECVFROM:60002,reuseaddr,ip-add-membership=239.192.0.2:enp2s0,fork SYSTEM:echo $$(base64 --wrap=0) \
41+
udp_listen --encode multicast 239.192.0.2 60002 --interface=172.16.6.1 \
4042
| raw_to_brefv \
4143
| mosquitto_pub -l -t '<topic>'
4244
```

bin/udp_listen

Lines changed: 77 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,77 @@
1+
#!/usr/bin/env python3.9
2+
3+
"""
4+
Command line utility tool for listening to udp traffic and
5+
outputting to stdout.
6+
"""
7+
8+
import sys
9+
import socket
10+
import struct
11+
import argparse
12+
from base64 import b64encode
13+
14+
15+
def process_multicast(args: argparse.Namespace):
16+
"""Listen in on multicast traffic and output to stdout
17+
18+
Args:
19+
args (argparse.Namespace): Command-line arguments
20+
"""
21+
22+
sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM, socket.IPPROTO_UDP)
23+
sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
24+
sock.setsockopt(socket.IPPROTO_IP, socket.IP_MULTICAST_TTL, args.TTL)
25+
sock.setsockopt(socket.IPPROTO_IP, socket.IP_MULTICAST_LOOP, int(args.loopback))
26+
sock.setsockopt(
27+
socket.IPPROTO_IP, socket.IP_MULTICAST_IF, socket.inet_aton(args.interface)
28+
)
29+
sock.setsockopt(
30+
socket.IPPROTO_IP,
31+
socket.IP_ADD_MEMBERSHIP,
32+
socket.inet_aton(args.group) + socket.inet_aton(args.interface),
33+
)
34+
sock.bind((args.group, args.port))
35+
sock.setblocking(True)
36+
37+
try:
38+
while True:
39+
data = sock.recvmsg(65535)[0]
40+
sys.stdout.write(
41+
(b64encode(data).decode() + "\n") if args.encode else data.decode()
42+
)
43+
sys.stdout.flush()
44+
finally:
45+
sock.setsockopt(
46+
socket.IPPROTO_IP,
47+
socket.IP_DROP_MEMBERSHIP,
48+
socket.inet_aton(args.group) + socket.inet_aton(args.interface),
49+
)
50+
sock.close()
51+
52+
53+
if __name__ == "__main__":
54+
parser = argparse.ArgumentParser(description="UDP listener")
55+
parser.add_argument(
56+
"--encode",
57+
action="store_true",
58+
default=False,
59+
help="base64 encode each packet content",
60+
)
61+
62+
sub_commands = parser.add_subparsers()
63+
64+
multicast_parser = sub_commands.add_parser("multicast")
65+
multicast_parser.add_argument("group", type=str)
66+
multicast_parser.add_argument("port", type=int)
67+
multicast_parser.add_argument("--interface", type=str, default="0.0.0.0")
68+
multicast_parser.add_argument("--loopback", type=bool, default=False)
69+
multicast_parser.add_argument(
70+
"--TTL",
71+
type=int,
72+
default=1,
73+
)
74+
multicast_parser.set_defaults(func=process_multicast)
75+
76+
args = parser.parse_args()
77+
args.func(args)

0 commit comments

Comments
 (0)