nigig-org/crates/apps/streem/mitm/mitm_proxy.py
Arena Bot 33dbece7d8
Some checks failed
nigig-build (CAD) / supply-chain (push) Has been cancelled
nigig-build (CAD) / cad-module (push) Has been cancelled
nigig-build (CAD) / full-crate-check (push) Has been cancelled
Payment domain and storage / isolated-payment-tests (push) Has been cancelled
PDF engine / engine (push) Has been cancelled
PDF engine / makepad-integration (push) Has been cancelled
PDF engine / fuzz (push) Has been cancelled
fix(ui): use valid syntax for assigning ids to metric values
2026-07-29 01:56:41 +00:00

317 lines
9.1 KiB
Python

#!/usr/bin/env python3
"""MITM Proxy - captures HTTP/HTTPS traffic with logging."""
import asyncio
import ssl
import socket
import os
import subprocess
import logging
import sys
from urllib.parse import urlparse
import select
PROXY_HOST = "0.0.0.0"
PROXY_PORT = 8080
CA_CERT = "/tmp/mitm_ca.crt"
CA_KEY = "/tmp/mitm_ca.key"
CERT_DIR = "/tmp/mitm_certs"
LOG_FILE = "/tmp/mitm_capture.txt"
os.makedirs(CERT_DIR, exist_ok=True)
logging.basicConfig(level=logging.DEBUG,
format="%(asctime)s|%(message)s",
handlers=[logging.FileHandler(LOG_FILE)])
log = logging.getLogger("mitm")
def gen_cert(hostname):
safe = hostname.replace("*", "_").replace("/", "_").replace(":", "_")
cp = os.path.join(CERT_DIR, f"{safe}.pem")
kp = os.path.join(CERT_DIR, f"{safe}.key")
if os.path.exists(cp) and os.path.exists(kp):
return cp, kp
subprocess.run(["openssl", "genrsa", "-out", kp, "2048"],
capture_output=True, check=True)
csr = os.path.join(CERT_DIR, f"{safe}.csr")
subprocess.run(["openssl", "req", "-new", "-key", kp, "-out", csr,
"-subj", f"/CN={hostname}"],
capture_output=True, check=True)
ext = os.path.join(CERT_DIR, f"{safe}.ext")
with open(ext, "w") as f:
f.write(f"subjectAltName=DNS:{hostname}\n")
subprocess.run(["openssl", "x509", "-req", "-in", csr,
"-CA", CA_CERT, "-CAkey", CA_KEY, "-CAcreateserial",
"-out", cp, "-days", "365", "-sha256", "-extfile", ext],
capture_output=True, check=True)
for p in [csr, ext, os.path.join(CERT_DIR, f"{safe}.srl")]:
try: os.remove(p)
except: pass
return cp, kp
def recv_headers(sock):
data = b""
while b"\r\n\r\n" not in data:
try:
c = sock.recv(65536)
if not c: return None, b""
data += c
except (socket.timeout, ConnectionError):
return None, b""
idx = data.find(b"\r\n\r\n")
return data[:idx+4], data[idx+4:]
def get_cl(head):
for l in head.split(b"\r\n"):
if l.lower().startswith(b"content-length:"):
return int(l.split(b":")[1].strip())
return 0
def is_chunked(head):
return b"transfer-encoding: chunked" in head.lower()
def recv_exact(sock, n):
d = b""
while len(d) < n:
c = sock.recv(min(65536, n - len(d)))
if not c: raise ConnectionError("closed")
d += c
return d
def recv_chunked(sock):
body = b""
while True:
sl = b""
while not sl.endswith(b"\n"):
c = sock.recv(1)
if not c: return body
sl += c
sz = int(sl.strip(), 16)
if sz == 0:
sock.recv(2)
break
c = recv_exact(sock, sz)
sock.recv(2)
body += c
return body
def relay_http(cli, srv, host):
try:
while True:
h, leftover = recv_headers(cli)
if h is None: break
fl = h.split(b"\r\n")[0].decode(errors="replace")
log.info(f"REQ {host} {fl}")
srv.sendall(h)
cl = get_cl(h)
if cl > 0:
body = leftover + recv_exact(cli, cl - len(leftover))
srv.sendall(body)
elif is_chunked(h):
body = leftover + recv_chunked(cli)
srv.sendall(body)
elif leftover:
srv.sendall(leftover)
rsph, leftover = recv_headers(srv)
if rsph is None: break
st = rsph.split(b"\r\n")[0].decode(errors="replace")
log.info(f"RESP {host} {st}")
is_101 = b" 101 " in rsph.split(b"\r\n")[0]
cli.sendall(rsph)
if is_101:
log.info(f"WS UPGRADE {host}")
if leftover:
cli.sendall(leftover)
fwd_raw(cli, srv)
break
cl = get_cl(rsph)
if cl > 0:
body = leftover + recv_exact(srv, cl - len(leftover))
cli.sendall(body)
log.info(f"RESP body: {cl} bytes | {body[:200]}")
elif is_chunked(rsph):
body = leftover + recv_chunked(srv)
cli.sendall(body)
elif leftover:
cli.sendall(leftover)
except (ConnectionError, socket.timeout, ssl.SSLError) as e:
log.debug(f"Relay end {host}: {e}")
except Exception as e:
log.error(f"Relay {host}: {e}")
def fwd_raw(a, b):
socks = {a: b, b: a}
while True:
try:
r, _, _ = select.select([a, b], [], [], 60)
if not r: break
for s in r:
d = s.recv(65536)
if not d: return
socks[s].sendall(d)
except: return
def mitm_thread(sock, host, port, cp, kp):
try:
ctx = ssl.SSLContext(ssl.PROTOCOL_TLS_SERVER)
ctx.load_cert_chain(cp, kp)
cli = ctx.wrap_socket(sock, server_side=True)
cli.settimeout(60)
srv = socket.create_connection((host, port), timeout=30)
srv_ctx = ssl.create_default_context()
srv = srv_ctx.wrap_socket(srv, server_hostname=host)
srv.settimeout(60)
log.info(f"Tunnel {host}:{port}")
relay_http(cli, srv, host)
except Exception as e:
log.error(f"MITM {host}:{port}: {e}")
finally:
try: sock.close()
except: pass
async def handle_connect(reader, writer, host, port):
cp, kp = gen_cert(host)
writer.write(b"HTTP/1.1 200 Connection Established\r\n\r\n")
await writer.drain()
tr = writer.transport
sock = tr.get_extra_info("socket")
fd = os.dup(sock.fileno())
tr.close()
new_sock = socket.socket(fileno=fd)
new_sock.settimeout(60)
loop = asyncio.get_event_loop()
await loop.run_in_executor(None, mitm_thread, new_sock, host, int(port), cp, kp)
async def handle_http(reader, writer, data):
try:
while b"\r\n\r\n" not in data:
c = await reader.read(4096)
if not c: return
data += c
fl = data.split(b"\r\n")[0].decode()
parts = fl.split(" ")
method, url = parts[0], parts[1]
up = urlparse(url)
host, port = up.hostname, up.port or 80
path = up.path or "/"
if up.query: path += "?" + up.query
nfl = f"{method} {path} HTTP/1.1\r\n".encode()
req = data.replace(fl.encode() + b"\r\n", nfl, 1)
cl = 0
if b"\r\n\r\n" in req:
hd = req.split(b"\r\n\r\n")[0]
for l in hd.split(b"\r\n"):
if l.lower().startswith(b"content-length:"):
cl = int(l.split(b":")[1].strip())
break
bd = req.split(b"\r\n\r\n", 1)[1]
rem = cl - len(bd)
while rem > 0:
c = await reader.read(min(65536, rem))
if not c: break
req += c
rem -= len(c)
log.info(f"HTTP {host} {method} {path}")
sr, sw = await asyncio.open_connection(host, port)
sw.write(req)
await sw.drain()
rsph = await sr.readuntil(b"\r\n\r\n")
cl = 0
for l in rsph.split(b"\r\n"):
if l.lower().startswith(b"content-length:"):
cl = int(l.split(b":")[1].strip())
break
body = b""
if cl > 0:
body = await sr.readexactly(cl)
writer.write(rsph + body)
await writer.drain()
sw.close()
_resp_status = rsph.split(b'\r\n')[0].decode(errors='replace')
log.info(f"HTTP RESP {_resp_status} | body={cl}")
except Exception as e:
log.error(f"HTTP proxy: {e}")
try:
writer.write(b"HTTP/1.1 502 Bad Gateway\r\n\r\n")
await writer.drain()
except: pass
finally:
try: writer.close()
except: pass
async def handle_client(reader, writer):
try:
data = await reader.readuntil(b"\r\n")
fl = data.decode(errors="replace").strip()
if fl.startswith("CONNECT"):
while b"\r\n\r\n" not in data:
try:
data += await asyncio.wait_for(reader.readuntil(b"\r\n"), 10)
except (asyncio.TimeoutError, asyncio.IncompleteReadError):
break
hp = fl.split(" ")[1]
host, port = hp.split(":") if ":" in hp else (hp, "443")
await handle_connect(reader, writer, host, port)
else:
await handle_http(reader, writer, data)
except Exception as e:
log.error(f"Client handler: {e}")
finally:
try: writer.close()
except: pass
async def main():
if not os.path.exists(CA_CERT):
print("ERROR: Generate CA cert first", file=sys.stderr)
sys.exit(1)
srv = await asyncio.start_server(handle_client, PROXY_HOST, PROXY_PORT)
addr = srv.sockets[0].getsockname()
print(f"MITM Proxy: {addr[0]}:{addr[1]}", file=sys.stderr)
print(f"Logs: {LOG_FILE}", file=sys.stderr)
async with srv:
await srv.serve_forever()
if __name__ == "__main__":
try:
asyncio.run(main())
except KeyboardInterrupt:
print("\nStopped", file=sys.stderr)