""" swarm_download.py - Download a real torrent from a real swarm using the engine, to measure how it performs against actual peers. Division of labour: for .torrent files, this harness parses metainfo locally and uses the sibling torrent-tracker library for DHT get_peers and HTTP/UDP tracker announces. Magnet metadata resolution still falls back to libtorrent. Every discovered peer endpoint is fed to the engine, which does all the data transfer. Pieces are reassembled and SHA-1-verified here; the engine never hashes or persists anything. Reality check (important for interpreting the numbers): a real public swarm contains unreachable peers, peers behind NAT, peers that only accept a different transport/encryption combination, and peers that do not actually have useful pieces. Those show up as "failed". The headline metric this prints is therefore how many discovered peers were actually usable, and the sustained rate across them. That is the honest "how well does it work today" answer. Usage ----- python harness/swarm_download.py path/to/file.torrent python harness/swarm_download.py 'magnet:?xt=urn:btih:...' python harness/swarm_download.py file.torrent --max-peers 200 --output /tmp/out --timeout 600 python harness/swarm_download.py file.torrent --output /tmp/out --resume Pick a well-seeded torrent (e.g. a current Linux distro ISO) for a meaningful test; obscure or dead torrents will show few usable peers regardless. """ from __future__ import annotations import argparse import hashlib import os import secrets import sys import tempfile import time ROOT = os.path.dirname(os.path.dirname(os.path.abspath(__file__))) sys.path.insert(0, os.path.join(ROOT, "harness")) import libtorrent as lt # noqa: E402 # only used for magnet metadata fallback from harness import load_metadata # noqa: E402 from engine_ffi import (BLOCK_SIZE, Engine, EngineConfig, STATE_NAMES, # noqa: E402 ERROR_NAMES) from torrent_meta import load_torrent # noqa: E402 from tracker_ffi import DHTClient, TrackerClient # noqa: E402 LIB = os.path.join(ROOT, "build", "libtorrentpeer.so") def _discovery_session() -> lt.session: s = lt.session({ "listen_interfaces": "0.0.0.0:0,[::]:0", "enable_dht": True, "enable_lsd": True, "enable_upnp": True, "enable_natpmp": True, "alert_mask": 0, }) for host, port in (("router.bittorrent.com", 6881), ("dht.transmissionbt.com", 6881), ("router.utorrent.com", 6881)): try: s.add_dht_node((host, port)) except Exception: pass return s class LibtorrentDiscovery: def __init__(self, handle, session): self.handle = handle self.session = session def collect(self, _max_peers: int) -> set[tuple[str, int]]: eps = set() try: for pi in self.handle.get_peer_info(): ip = pi.ip if isinstance(ip, tuple) and len(ip) == 2 and ip[1]: eps.add((ip[0], int(ip[1]))) except Exception: pass return eps def close(self): try: self.session.remove_torrent(self.handle) except Exception: pass class TrackerDiscovery: def __init__(self, trackers: list[str], meta, *, timeout: float, max_trackers: int, use_dht: bool, dht_timeout: float, dht_queries: int): self.trackers = trackers[:max_trackers] if max_trackers > 0 else trackers self.meta = meta self.timeout = timeout self.client = TrackerClient() self.dht = DHTClient() if use_dht else None self.dht_timeout = dht_timeout self.dht_queries = dht_queries self.dht_done = False self.peer_id = b"-PC0001-" + secrets.token_bytes(12) self.key = secrets.randbits(32) self.endpoints: set[tuple[str, int]] = set() self.index = 0 self.next_cycle_at = 0.0 def collect(self, max_peers: int) -> set[tuple[str, int]]: now = time.time() if len(self.endpoints) >= max_peers: return set(self.endpoints) if self.dht and not self.dht_done: self.dht_done = True result = self.dht.lookup( self.meta.info_hash, timeout=self.dht_timeout, max_queries=self.dht_queries, max_peers=max_peers - len(self.endpoints), ) self.endpoints.update(result.peers) msg = (f"dht: {len(result.peers)} peers, " f"{result.nodes_queried} queried/{result.nodes_discovered} learned " f"({result.elapsed_ms:.0f} ms)") if result.error: msg += f": {result.error}" print(msg, flush=True) if len(self.endpoints) >= max_peers: return set(self.endpoints) if self.index >= len(self.trackers): if now < self.next_cycle_at: return set(self.endpoints) self.index = 0 if not self.trackers: return set() url = self.trackers[self.index] self.index += 1 if self.index >= len(self.trackers): self.next_cycle_at = now + 60.0 result = self.client.announce( url, self.meta, self.peer_id, port=6881, key=self.key, numwant=max(1, min(200, max_peers - len(self.endpoints))), event="started", timeout=self.timeout) if result.ok: self.endpoints.update(result.peers) print(f"tracker: {url} returned {len(result.peers)} peers " f"({result.elapsed_ms:.0f} ms)", flush=True) else: print(f"tracker: {url} failed: {result.error}", flush=True) return set(self.endpoints) def close(self): pass def _resolve_magnet(arg: str, scratch: str): ses = _discovery_session() params = lt.parse_magnet_uri(arg) params.save_path = scratch params.flags |= lt.torrent_flags.upload_mode h = ses.add_torrent(params) print("resolving magnet metadata from the swarm using libtorrent...", flush=True) deadline = time.time() + 120 while time.time() < deadline and not h.status().has_metadata: time.sleep(0.5) if not h.status().has_metadata: raise TimeoutError("could not fetch metadata for magnet within 120s") tpath = os.path.join(scratch, "resolved.torrent") with open(tpath, "wb") as f: f.write(lt.bencode(lt.create_torrent(h.torrent_file()).generate())) return tpath, LibtorrentDiscovery(h, ses) def _completed_bytes(meta, done: bytearray) -> int: return sum(meta.piece_len(i) for i, is_done in enumerate(done) if is_done) def verify_existing_output(out_fh, meta) -> tuple[bytearray, int, int]: """Hash pieces already present in the output file. Returns (done_bitfield, done_count, verified_bytes). Only pieces whose full range exists and whose SHA-1 matches the torrent metadata are marked done. This intentionally ignores partial pieces because the engine currently only persists data after a whole piece has passed verification. """ done = bytearray(meta.num_pieces) done_count = 0 file_size = os.fstat(out_fh.fileno()).st_size for piece in range(meta.num_pieces): piece_len = meta.piece_len(piece) offset = piece * meta.piece_length if offset + piece_len > file_size: continue out_fh.seek(offset) data = out_fh.read(piece_len) if len(data) != piece_len: continue if hashlib.sha1(data).digest() != meta.piece_hashes[piece]: continue done[piece] = 1 done_count += 1 return done, done_count, _completed_bytes(meta, done) class PieceAssembler: """Reassemble pieces while ignoring duplicate blocks. Endgame mode deliberately re-requests unfinished pieces from multiple peers. That means the same block may arrive more than once. Counting raw bytes would mark a piece complete too early, so completion is based on unique block offsets within each piece. """ def __init__(self, meta, done: bytearray): self.meta = meta self.done = done self.buffers: dict[int, bytearray] = {} self.received = [0] * meta.num_pieces self.seen: dict[int, bytearray] = {} def _block_count(self, piece: int) -> int: piece_len = self.meta.piece_len(piece) return (piece_len + BLOCK_SIZE - 1) // BLOCK_SIZE def reset_piece(self, piece: int) -> int: previous = self.received[piece] self.received[piece] = 0 self.buffers.pop(piece, None) self.seen.pop(piece, None) return previous def add_block(self, piece: int, begin: int, data) -> tuple[bool, bool]: if self.done[piece]: return False, False if begin % BLOCK_SIZE != 0: raise ValueError(f"unaligned block for piece {piece}: begin={begin}") piece_len = self.meta.piece_len(piece) length = len(data) if begin + length > piece_len: raise ValueError( f"block overruns piece {piece}: begin={begin} len={length} " f"piece_len={piece_len}") buf = self.buffers.get(piece) if buf is None: buf = bytearray(piece_len) self.buffers[piece] = buf seen = self.seen.get(piece) if seen is None: seen = bytearray(self._block_count(piece)) self.seen[piece] = seen block = begin // BLOCK_SIZE buf[begin:begin + length] = data if seen[block]: return False, self.received[piece] >= piece_len seen[block] = 1 self.received[piece] += length return True, self.received[piece] >= piece_len def piece_bytes(self, piece: int) -> bytes: return bytes(self.buffers[piece]) def finish_piece(self, piece: int) -> None: self.done[piece] = 1 self.buffers.pop(piece, None) self.seen.pop(piece, None) class EndgameController: """Conservative piece-level endgame for slow tail pieces.""" def __init__(self, meta, *, min_pieces: int, peer_factor: float, interval: float): self.meta = meta self.min_pieces = min_pieces self.peer_factor = peer_factor self.interval = interval self.active = False self.last_rearm = 0.0 def maybe_rearm(self, eng: Engine, tid: int, done: bytearray, done_count: int, connected: int, now: float) -> None: remaining = self.meta.num_pieces - done_count if remaining <= 0: return threshold = max(self.min_pieces, int(max(1, connected) * self.peer_factor)) if remaining > threshold: return if not self.active: self.active = True print(f"endgame: {remaining} pieces left; duplicating tail requests", flush=True) for piece, is_done in enumerate(done): if not is_done: eng.set_priority(tid, piece, 255) self.last_rearm = 0.0 if now - self.last_rearm < self.interval: return for piece, is_done in enumerate(done): if not is_done: eng.request_piece(tid, piece) self.last_rearm = now def open_output(path: str, meta, resume: bool): if resume: existed = os.path.exists(path) out_fh = open(path, "r+b" if existed else "w+b") if existed: print("resume: verifying existing output pieces...", flush=True) done, done_count, verified_bytes = verify_existing_output(out_fh, meta) pct = 100.0 * done_count / meta.num_pieces if meta.num_pieces else 100.0 print(f"resume: found {done_count}/{meta.num_pieces} verified pieces " f"({pct:4.1f}%, {verified_bytes/1e6:.1f} MB)", flush=True) else: done = bytearray(meta.num_pieces) done_count = 0 print("resume: output file does not exist yet; starting fresh", flush=True) else: out_fh = open(path, "w+b") done = bytearray(meta.num_pieces) done_count = 0 out_fh.truncate(meta.total_size) return out_fh, done, done_count def main() -> int: ap = argparse.ArgumentParser( description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter) ap.add_argument("torrent", help="path to a .torrent file or a magnet: URI") ap.add_argument("--max-peers", type=int, default=100, help="cap on peers fed to the engine (default 100)") ap.add_argument("--output", "-o", help="write verified data here (else discard)") ap.add_argument("--resume", action="store_true", help="resume from an existing --output file by hashing " "verified pieces and skipping them") ap.add_argument("--timeout", type=float, default=600.0, help="overall stall timeout in seconds") ap.add_argument("--tracker-timeout", type=float, default=4.0, help="per-tracker announce timeout for direct tracker " "discovery (default 4)") ap.add_argument("--max-trackers", type=int, default=16, help="maximum trackers to announce to from a .torrent " "(0 = all, default 16)") ap.add_argument("--no-dht", action="store_true", help="disable DHT get_peers discovery for .torrent files") ap.add_argument("--dht-timeout", type=float, default=6.0, help="total DHT lookup budget in seconds (default 6)") ap.add_argument("--dht-queries", type=int, default=32, help="maximum DHT nodes to query per torrent (default 32)") ap.add_argument("--lib", default=LIB) ap.add_argument("--loops", type=int, default=0, help="engine loops (0=auto)") ap.add_argument("--encryption", type=int, default=1, choices=[0, 1, 2], help="0=plaintext, 1=MSE offer RC4+plaintext (default), " "2=MSE require RC4. Encryption reaches far more of a " "real swarm.") ap.add_argument("--utp", type=int, default=0, choices=[0, 1], help="0=TCP (default), 1=µTP/UDP. A swarm has a mix; this " "selects which transport the engine dials with.") ap.add_argument("--no-endgame", action="store_true", help="disable duplicate tail-piece requests") ap.add_argument("--endgame-min-pieces", type=int, default=16, help="enter endgame when remaining pieces are at or below " "this count, also scaled by connected peers " "(default 16)") ap.add_argument("--endgame-peer-factor", type=float, default=2.0, help="also enter endgame below peers*factor remaining " "pieces (default 2.0)") ap.add_argument("--endgame-interval", type=float, default=3.0, help="seconds between tail-piece re-arms in endgame " "(default 3.0)") args = ap.parse_args() if args.resume and not args.output: ap.error("--resume requires --output") scratch = tempfile.mkdtemp(prefix="swarm_dl_") if args.torrent.startswith("magnet:"): tpath, discovery = _resolve_magnet(args.torrent, scratch) else: torrent_file = load_torrent(args.torrent) tpath = args.torrent discovery = TrackerDiscovery( torrent_file.trackers, torrent_file.metadata, timeout=args.tracker_timeout, max_trackers=args.max_trackers, use_dht=not args.no_dht, dht_timeout=args.dht_timeout, dht_queries=args.dht_queries, ) meta = load_metadata(tpath) print(f"torrent: {meta.name} {meta.total_size/1e6:.1f} MB " f"{meta.num_pieces} pieces x {meta.piece_length}", flush=True) out_fh = None done = bytearray(meta.num_pieces) done_count = 0 if args.output: out_fh, done, done_count = open_output(args.output, meta, args.resume) cfg = EngineConfig(loop_count=args.loops, slots_per_loop=4096, max_pipeline=1024, encryption=args.encryption, utp=args.utp) print(f"transport: {'µTP' if args.utp else 'TCP'}, " f"encryption={['off','offer','require'][args.encryption]}", flush=True) eng = Engine(cfg, lib_path=args.lib) tid = eng.add_torrent(meta.info_hash, b"-PC0001-" + secrets.token_bytes(12), meta.piece_length, meta.total_size, meta.num_pieces) eng.set_priorities(tid, [0 if done[i] else 1 for i in range(meta.num_pieces)]) assembler = PieceAssembler(meta, done) endgame = None if args.no_endgame else EndgameController( meta, min_pieces=max(1, args.endgame_min_pieces), peer_factor=max(1.0, args.endgame_peer_factor), interval=max(0.5, args.endgame_interval), ) added: set = set() peak_connected = 0 useful_bytes = _completed_bytes(meta, done) t0 = time.time() deadline = t0 + args.timeout last_print = 0.0 last_useful_bytes = useful_bytes last_progress_t = t0 try: while done_count < meta.num_pieces: # Feed any newly-discovered peers to the engine, up to the cap. if len(added) < args.max_peers: for ip, port in discovery.collect(args.max_peers): if (ip, port) in added: continue added.add((ip, port)) try: eng.add_peer(tid, ip, port) except Exception: pass if len(added) >= args.max_peers: break st = eng.status(tid) peak_connected = max(peak_connected, st.peers_connected) now = time.time() if endgame: endgame.maybe_rearm( eng, tid, done, done_count, st.peers_connected, now) descs = eng.poll_ready() if not descs: eng.wait(200) now = time.time() if useful_bytes != last_useful_bytes: last_useful_bytes = useful_bytes last_progress_t = now if now - last_print >= 1.0: last_print = now pct = 100.0 * done_count / meta.num_pieces print(f" {done_count}/{meta.num_pieces} pieces ({pct:4.1f}%) " f"{st.rate_bps/1e6:6.1f} MB/s " f"peers {st.peers_connected} up / {st.peers_failed} failed " f"/ {len(added)} tried outstanding={st.outstanding}", flush=True) if now > deadline or (now - last_progress_t) > args.timeout: print("stalled; giving up", flush=True) break continue for x in descs: if done[x.piece]: # A block can arrive after the piece was completed and # de-prioritized because it was already in flight. Drop it. eng.release(x.loop, x.slot) continue block = bytes(eng.block_data(x.loop, x.slot, x.len)) eng.release(x.loop, x.slot) added_unique, complete = assembler.add_block( x.piece, x.begin, block) if added_unique: useful_bytes += x.len if complete: piece_data = assembler.piece_bytes(x.piece) if hashlib.sha1(piece_data).digest() != meta.piece_hashes[x.piece]: # Corrupt/garbage block from a misbehaving peer: re-arm. useful_bytes -= assembler.reset_piece(x.piece) eng.request_piece(tid, x.piece) continue done[x.piece] = 1 done_count += 1 eng.set_priority(tid, x.piece, 0) if out_fh: out_fh.seek(x.piece * meta.piece_length) out_fh.write(piece_data) out_fh.flush() assembler.finish_piece(x.piece) # free as we go dt = time.time() - t0 st = eng.status(tid) mb = _completed_bytes(meta, done) / 1e6 print("-" * 70) print(f"downloaded {done_count}/{meta.num_pieces} pieces " f"({mb:.1f} MB) in {dt:.1f}s = {mb/dt:.1f} MB/s" if dt > 0 else "") print(f"peers: {len(added)} discovered+tried, " f"{peak_connected} usable at peak, " f"{st.peers_failed} failed (likely incompatible transport/" f"encryption or unreachable)") if done_count < meta.num_pieces: print(f"engine state: {STATE_NAMES[st.state]} " f"err={ERROR_NAMES[st.error]}") return 0 if done_count == meta.num_pieces else 2 finally: if out_fh: out_fh.close() eng.close() discovery.close() if __name__ == "__main__": raise SystemExit(main())