I'm trying to use the multiprocessing module to implement a simple network traffic forwader.
My application listens on a port, and when it receives an inbound connection it makes an outgoing TCP connection to another server and then shuttles the data back and forth between the two connections.
I've been trying to test my code using the nc utility, but seems that my application goes inside the recv_bytes() call and blocks even though there is data there.
Here's a simplified version of my forwarding code (tested fully runnable):
from multiprocessing import Process
from multiprocessing.connection import Listener, Client, wait
def start_serving(listen_port, outbound_port):
with Listener(('', listen_port)) as server:
print(f"Waiting for connections on port {listen_port}")
with server.accept() as inbound_conn:
print(f"Connection accepted from {server.last_accepted}")
outbound_conn = Client(('localhost', outbound_port))
print(f"Connected to port {outbound_port}")
readers = [inbound_conn, outbound_conn]
print(f"inbound_reader = {inbound_conn}")
print(f"outbound_reader = {outbound_conn}")
while readers:
for r in wait(readers):
try:
print(f"Calling recv_bytes with reader {r}")
data = r.recv_bytes() # This blocks even when there's data
print(f"Out of recv_bytes with reader {r}")
except EOFError:
readers.remove(r)
else:
fwd_to_conn = None
if r is inbound_conn:
fwd_to_conn = outbound_conn
print("read from inbound connection")
elif r is outbound_conn:
fwd_to_conn = outbound_conn
print("read from outbound connection")
if fwd_to_conn is not None:
print(f"Forwarding {len(bytes)} bytes")
fwd_to_conn.send_bytes(data)
forwarder = Process(target=start_serving, daemon=True, args=(19001, 19002))
forwarder.start()
forwarder.join()
I run this script and in a separate consoles I run:
$ echo "Hi from outbound" | nc -l -p 19002
and
$ echo "Hi from inbound" | nc localhost 19001
This is the output I get:
Waiting for connections on port 19001
Connection accepted from ('127.0.0.1', 56874)
Connected to port 19002
inbound_reader = <multiprocessing.connection.Connection object at 0x7ff3730e2850>
outbound_reader = <multiprocessing.connection.Connection object at 0x7ff3730e2a60>
Calling recv_bytes with reader <multiprocessing.connection.Connection object at 0x7ff3730e2850>
As you can see the application is blocked inside the recv_bytes() call on the inbound connection even though there is data there.
Seems like a very simple application, so I'm hoping there's an obvious solution here.
EDIT
Digging into the multiprocessing.connection.Connection code a bit it seems that when using recv_bytes() the code expects the size of the message to be in the first 4 bytes. However, I don't want it to assume any kind of format and just read as much data as it can without trying to interpret it.
Thanks!