Naut-Peer/harness/harness.py
ookami125 d8208685a2 Initial commit: multi-peer torrent download engine
Reactor/loop-pool engine with TCP/µTP/MSE transports, per-connection
pipelining, priority-driven piece selection with endgame, and the Python
FFI test harness.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-21 23:12:32 -04:00

202 lines
7.7 KiB
Python

"""
Test/driver harness for the C peer.
The Python harness parses .torrent metadata and can discover peers through the
sibling torrent-tracker C library's DHT and HTTP/UDP tracker helpers. The C peer
does the fast part: pull the requested blocks from one peer. This harness owns
what to download, reassembles pieces, and verifies SHA-1 hashes -- the peer
never hashes or persists anything.
"""
from __future__ import annotations
import argparse
import hashlib
import os
import secrets
import time
from peer_ffi import Peer, PeerConfig, STATE_ERROR, STATE_NAMES, ERROR_NAMES
from torrent_meta import Metadata, load_metadata, load_torrent
from tracker_ffi import DHTClient, TrackerClient
def make_peer_id() -> bytes:
return b"-PC0001-" + secrets.token_bytes(12)
def discover_peers(torrent_path: str, max_wait: float = 20.0) -> list[tuple[str, int]]:
"""Discover peer endpoints through DHT first, then torrent trackers."""
tf = load_torrent(torrent_path)
peer_id = make_peer_id()
key = secrets.randbits(32)
deadline = time.time() + max_wait
seen: set[tuple[str, int]] = set()
dht_budget = max(0.5, min(6.0, max_wait / 2.0))
dht_result = DHTClient().lookup(tf.metadata.info_hash, timeout=dht_budget)
seen.update(dht_result.peers)
if seen:
return sorted(seen)
client = TrackerClient()
for url in tf.trackers:
if time.time() >= deadline:
break
result = client.announce(
url, tf.metadata, peer_id, port=6881, key=key, numwant=50,
event="started", timeout=max(0.5, deadline - time.time()))
if result.ok:
seen.update(result.peers)
if seen:
break
return sorted(seen)
class Downloader:
def __init__(self, meta: Metadata, *, peer_id: bytes | None = None,
num_slots: int = 0, max_pipeline: int = 0,
request_timeout_ms: int = 0, recv_buffer_bytes: int = 0,
lib_path: str | None = None):
self.meta = meta
cfg = PeerConfig()
cfg.info_hash[:] = meta.info_hash
cfg.peer_id[:] = peer_id or make_peer_id()
cfg.piece_length = meta.piece_length
cfg.total_size = meta.total_size
cfg.num_pieces = meta.num_pieces
cfg.num_slots = num_slots
cfg.max_pipeline = max_pipeline
cfg.request_timeout_ms = request_timeout_ms
cfg.recv_buffer_bytes = recv_buffer_bytes
self.peer = Peer(cfg, lib_path)
self.buffers: list[bytearray | None] = [None] * meta.num_pieces
self.received = [0] * meta.num_pieces
self.done = bytearray(meta.num_pieces)
self.done_count = 0
def download(self, ip: str, port: int, pieces=None, priorities=None,
timeout: float = 60.0, progress_every: float = 1.0,
output: str | None = None) -> bytes | None:
meta = self.meta
pieces = list(pieces) if pieces is not None else list(range(meta.num_pieces))
out_fh = open(output, "wb") if output else None
if out_fh:
out_fh.truncate(meta.total_size)
# Build the priority vector: caller-supplied scheme, or uniform "1" over
# the wanted pieces (0 = not wanted). The peer masks this with what the
# remote actually has, so a peer missing a piece is simply skipped.
prio = bytearray(priorities) if priorities is not None \
else bytearray(meta.num_pieces)
for i in pieces:
self.buffers[i] = bytearray(meta.piece_len(i))
if priorities is None:
prio[i] = 1
self.peer.start(ip, port)
self.peer.set_priorities(prio)
want = len(pieces)
deadline = time.time() + timeout
last_print = 0.0
last_progress_bytes = 0
last_progress_time = time.time()
while self.done_count < want:
st = self.peer.status()
if st.state == STATE_ERROR:
raise RuntimeError(f"peer error: {ERROR_NAMES[st.error]}")
descs = self.peer.poll_ready()
if not descs:
self.peer.wait(100)
now = time.time()
if st.bytes_received != last_progress_bytes:
last_progress_bytes = st.bytes_received
last_progress_time = now
if now > deadline and now - last_progress_time > timeout:
raise TimeoutError(
f"stalled: {self.done_count}/{want} pieces, "
f"state={STATE_NAMES[st.state]}")
if now - last_print >= progress_every:
last_print = now
print(f" {self.done_count}/{want} pieces "
f"{st.rate_bps/1e6:.1f} MB/s "
f"outstanding={st.outstanding} free={st.free_slots}")
continue
for d in descs:
buf = self.buffers[d.piece]
buf[d.begin:d.begin + d.len] = self.peer.block_data(d.slot, d.len)
self.peer.release(d.slot)
self.received[d.piece] += d.len
if (not self.done[d.piece]
and self.received[d.piece] >= meta.piece_len(d.piece)):
digest = hashlib.sha1(bytes(buf)).digest()
if digest != meta.piece_hashes[d.piece]:
raise ValueError(f"piece {d.piece} hash mismatch")
self.done[d.piece] = 1
self.done_count += 1
# Drop priority so a verified piece is no longer a
# selection candidate (the peer also won't re-request it).
self.peer.set_priority(d.piece, 0)
if out_fh:
out_fh.seek(d.piece * meta.piece_length)
out_fh.write(buf)
if not output:
pass # keep in memory for return
else:
self.buffers[d.piece] = None # free once flushed
if out_fh:
out_fh.close()
return None
return b"".join(bytes(self.buffers[i]) for i in pieces)
def close(self):
self.peer.stop()
self.peer.close()
def main() -> int:
ap = argparse.ArgumentParser(description="Drive the C peer to download a torrent.")
ap.add_argument("torrent", help="path to .torrent file")
ap.add_argument("--peer", help="explicit peer ip:port (skip tracker)")
ap.add_argument("--output", "-o", help="write downloaded data here")
ap.add_argument("--slots", type=int, default=0, help="arena slots (16 KiB each)")
ap.add_argument("--pipeline", type=int, default=0, help="max outstanding requests")
ap.add_argument("--timeout", type=float, default=120.0)
args = ap.parse_args()
meta = load_metadata(args.torrent)
print(f"torrent: {meta.name} {meta.total_size} bytes "
f"{meta.num_pieces} pieces x {meta.piece_length}")
if args.peer:
host, port = args.peer.rsplit(":", 1)
endpoints = [(host, int(port))]
else:
print("discovering peers via tracker/DHT...")
endpoints = discover_peers(args.torrent)
if not endpoints:
print("no peers found")
return 1
print(f"found {len(endpoints)} peer(s); using {endpoints[0]}")
dl = Downloader(meta, num_slots=args.slots, max_pipeline=args.pipeline)
try:
ip, port = endpoints[0]
t0 = time.time()
dl.download(ip, port, timeout=args.timeout, output=args.output)
dt = time.time() - t0
mb = meta.total_size / 1e6
print(f"done: {mb:.1f} MB in {dt:.2f}s = {mb/dt:.1f} MB/s")
finally:
dl.close()
return 0
if __name__ == "__main__":
raise SystemExit(main())