|
import asyncio |
|
import json |
|
import sys |
|
|
|
import stun # pip install pystun3 |
|
import miniupnpc # pip install miniupnpc |
|
|
|
RENDEZVOUS = ("YOUR.SERVER.IP", 55555) |
|
LOCAL_PORT = 50000 |
|
|
|
class P2PClient(asyncio.DatagramProtocol): |
|
def __init__(self, username): |
|
self.username = username |
|
self.peers = {} # peer_username -> (ip, port) |
|
self.transport = None |
|
|
|
async def start(self): |
|
# 1. STUN lookup |
|
nat, ext_ip, ext_port = stun.get_ip_info() |
|
print(f"[STUN] NAT type={nat}, public={ext_ip}:{ext_port}") |
|
|
|
# 2. UPnP port mapping (TCP) |
|
upnp = miniupnpc.UPnP() |
|
upnp.discoverdelay = 200 |
|
upnp.discover() |
|
upnp.selectigd() |
|
upnp.addportmapping(LOCAL_PORT, "TCP", upnp.lanaddr, LOCAL_PORT, |
|
f"p2p-{self.username}", "" ) |
|
print(f"[UPnP] mapped {LOCAL_PORT} → {upnp.externalipaddress()}:{LOCAL_PORT}") |
|
|
|
# 3. Open our UDP socket and register |
|
loop = asyncio.get_running_loop() |
|
await loop.create_datagram_endpoint( |
|
lambda: self, local_addr=("0.0.0.0", LOCAL_PORT) |
|
) |
|
self.send_udp({"action":"register", "username":self.username}) |
|
print("[UDP] registration sent") |
|
|
|
# 4. Start listening for incoming TCP |
|
server = await asyncio.start_server( |
|
self.on_tcp_connection, "0.0.0.0", LOCAL_PORT |
|
) |
|
asyncio.create_task(server.serve_forever()) |
|
print(f"[TCP] listening on port {LOCAL_PORT}") |
|
|
|
#–– UDP protocol callbacks ––# |
|
def connection_made(self, transport): |
|
self.transport = transport |
|
|
|
def datagram_received(self, data, addr): |
|
msg = json.loads(data.decode()) |
|
act = msg.get("action") |
|
|
|
if act == "registered": |
|
print("[UDP] registered OK") |
|
elif act == "connection_info": |
|
peer = msg["peer_username"] |
|
ip, port = msg["peer_ip"], msg["peer_port"] |
|
print(f"[UDP] got peer {peer} → {ip}:{port}") |
|
self.peers[peer] = (ip, port) |
|
# punch and simultaneous TCP open: |
|
asyncio.create_task(self.holepunch_and_connect(peer, ip, port)) |
|
elif act == "error": |
|
print("[UDP] error:", msg.get("message")) |
|
|
|
def send_udp(self, msg, addr=None): |
|
target = addr or RENDEZVOUS |
|
self.transport.sendto(json.dumps(msg).encode(), target) |
|
|
|
#–– simultaneous TCP open ––# |
|
async def holepunch_and_connect(self, peer, ip, port): |
|
await asyncio.sleep(0.1) # give TCP server a moment |
|
print(f"[PUNCH] simultaneous TCP open to {peer}") |
|
# 1) outbound connect attempt |
|
try: |
|
r, w = await asyncio.open_connection(ip, port) |
|
print(f"[TCP] outbound to {peer} succeeded") |
|
asyncio.create_task(self.handle_stream(r, w)) |
|
except Exception as e: |
|
print(f"[TCP] outbound to {peer} failed: {e}") |
|
|
|
#–– incoming TCP ––# |
|
async def on_tcp_connection(self, reader, writer): |
|
peer = writer.get_extra_info("peername") |
|
print(f"[TCP] inbound connection from {peer}") |
|
asyncio.create_task(self.handle_stream(reader, writer)) |
|
|
|
#–– unified handler ––# |
|
async def handle_stream(self, reader, writer): |
|
peer = writer.get_extra_info("peername") |
|
try: |
|
while True: |
|
data = await reader.read(1024) |
|
if not data: |
|
break |
|
print(f"[{peer}] {data.decode().strip()}") |
|
writer.write(f"ECHO> {data.decode()}".encode()) |
|
await writer.drain() |
|
except Exception as e: |
|
print("Stream error:", e) |
|
finally: |
|
writer.close() |
|
await writer.wait_closed() |
|
|
|
#–– user commands ––# |
|
def request_peer(self, peer_username): |
|
print(f"[CMD] requesting {peer_username}") |
|
self.send_udp({ |
|
"action": "request_connection", |
|
"username": self.username, |
|
"peer_username": peer_username |
|
}) |
|
|
|
async def send_message(self, peer_username, msg): |
|
if peer_username not in self.peers: |
|
print("No info for peer", peer_username) |
|
return |
|
ip, port = self.peers[peer_username] |
|
try: |
|
r, w = await asyncio.open_connection(ip, port) |
|
w.write(msg.encode()) |
|
await w.drain() |
|
resp = await r.read(1024) |
|
print(f"[{peer_username} reply] {resp.decode().strip()}") |
|
w.close() |
|
await w.wait_closed() |
|
except Exception as e: |
|
print("Send failed:", e) |
|
|
|
async def main(): |
|
if len(sys.argv) != 2: |
|
print("Usage: python p2p_client.py <your_name>") |
|
return |
|
name = sys.argv[1] |
|
client = P2PClient(name) |
|
await client.start() |
|
|
|
loop = asyncio.get_running_loop() |
|
while True: |
|
line = await loop.run_in_executor(None, input, ">> ") |
|
parts = line.strip().split(" ", 2) |
|
if parts[0] == "connect" and len(parts) == 2: |
|
client.request_peer(parts[1]) |
|
elif parts[0] == "send" and len(parts) == 3: |
|
await client.send_message(parts[1], parts[2]) |
|
elif parts[0] in ("exit", "quit"): |
|
break |
|
|
|
if __name__ == "__main__": |
|
asyncio.run(main()) |