Files
LNXSDK/leenkx/Sources/leenkx/network/torrent/Torrent.hx
2026-10-06 02:10:08 -07:00

970 lines
27 KiB
Haxe

package leenkx.network.torrent;
import haxe.io.Bytes;
import leenkx.network.torrent.peer.IPeerChannel;
import leenkx.network.torrent.peer.WireProtocol;
import leenkx.network.torrent.peer.RtcPeerPool;
import leenkx.network.torrent.storage.Storage;
import leenkx.network.torrent.storage.Storage.IStorage;
class Torrent {
static inline var MAX_OUTSTANDING = 8;
static inline var META_BLOCK = 16384;
static inline var MAX_META_SIZE = 10000000;
static inline var KEEPALIVE_MS = 60000;
public var meta:MetaInfo = null;
public var infoHash:Bytes;
public var infoHashHex:String;
public var infoHashV2:Bytes = null;
public var infoHashV2Hex:String = null;
public var pendingWebSeeds:Array<String> = null;
public var magnetURI:String = null;
public var files:Array<TorrentFile> = [];
public var ready = false;
public var done = false;
public var numPeers(get, never):Int;
function get_numPeers():Int return Lambda.count(wires);
var listeners:Map<String, Array<Dynamic>> = [];
public var onPeerExtension:String -> String -> Bytes -> Void;
public var uploadedBytes:Float = 0;
public var stopped = false;
public var client:TorrentClient = null;
var webSeeds:Array<WebSeed> = [];
var lastStatT:Float = 0;
var lastStatDown:Float = 0;
var lastStatUp:Float = 0;
public var announcePort:Int = 0;
public var downloadDir:String = null;
public function getStats():{downloaded:Float, uploaded:Float, left:Float, downRate:Float, upRate:Float} {
var total = meta != null ? meta.totalLength : 0;
var downloaded = pieces != null ? Math.min(pieces.numHave * (meta.pieceLength:Float), total) : 0;
var now = haxe.Timer.stamp();
var dt = now - lastStatT;
var downRate:Float = 0;
var upRate:Float = 0;
if (lastStatT > 0 && dt > 0) {
downRate = (downloaded - lastStatDown) / dt;
upRate = (uploadedBytes - lastStatUp) / dt;
if (downRate < 0) downRate = 0;
if (upRate < 0) upRate = 0;
}
lastStatT = now;
lastStatDown = downloaded;
lastStatUp = uploadedBytes;
return {
downloaded: downloaded,
uploaded: uploadedBytes,
left: total - downloaded,
downRate: downRate,
upRate: upRate
};
}
var peerId:Bytes;
var trackers:Array<String>;
var storage:IStorage;
public var trackerObjects:Array<Dynamic> = [];
@:allow(leenkx.network.torrent.TorrentFile)
@:allow(leenkx.network.torrent.WebSeed)
var pieces:PieceManager = null;
var wires:Map<String, WireProtocol> = [];
var peerBitfields:Map<String, Array<Bool>> = [];
var outstanding:Map<String, Int> = [];
var pieceCursor:Map<String, Int> = [];
var rtcPool:RtcPeerPool = null;
var httpTrackers:Array<Dynamic> = [];
var udpTrackers:Array<Dynamic> = [];
var keepaliveTimer:haxe.Timer = null;
var destroyed = false;
var metaSize:Int = 0;
var metaReceived:Int = 0;
var metaRequested:Array<Bool> = [];
var metaRejects:Int = 0;
var metaWire:WireProtocol = null;
public function new(infoHash:Bytes, peerId:Bytes, trackers:Array<String>, ?meta:MetaInfo, ?storage:IStorage, ?downloadDir:String) {
TorrentClient.initMain();
this.infoHash = infoHash;
this.infoHashHex = infoHash != null ? Crypto.toHex(infoHash) : null;
this.peerId = peerId;
this.trackers = trackers;
this.storage = storage;
this.downloadDir = downloadDir;
if (meta != null) setMeta(meta);
}
public function on(event:String, cb:Dynamic):Void {
var l = listeners.get(event);
if (l == null) listeners.set(event, l = []);
l.push(cb);
}
function emit(event:String):Void {
var l = listeners.get(event);
if (l != null) for (cb in l) cb();
}
function emit1(event:String, v:Dynamic):Void {
var l = listeners.get(event);
if (l != null) for (cb in l) cb(v);
}
function setMeta(m:MetaInfo):Void {
meta = m;
if (storage == null) {
var sz = m.storageLength > 0 ? m.storageLength : m.totalLength;
#if (sys || kha_krom)
storage = downloadDir != null ? new FileTreeStorage(downloadDir, m.files) : new MemoryStorage(sz);
#else
storage = new MemoryStorage(sz);
#end
}
if (pieces == null) {
pieces = new PieceManager(m, storage);
pieces.onPieceComplete = onPieceComplete;
pieces.onComplete = onTorrentComplete;
}
if (m.fileHashes == null && m.files != null) {
m.fileHashes = [for (f in m.files) null];
}
files = [for (f in m.files) new TorrentFile(f, storage, this)];
infoHashV2 = m.infoHashV2;
infoHashV2Hex = m.infoHashV2Hex;
if (infoHash == null && m.infoHash != null) {
infoHash = m.infoHash;
infoHashHex = m.infoHashHex;
}
if (infoHashHex == null) infoHashHex = m.infoHashV2Hex;
if (m.urlList == null && pendingWebSeeds != null) {
m.urlList = pendingWebSeeds;
}
magnetURI = Magnet.build(m.infoHash, m.name, trackers, m.urlList, m.infoHashV2);
ready = true;
initWebSeeds();
emit1("metadata", m);
emit("ready");
if (!pieces.complete) {
for (w in wires) w.sendInterested();
}
}
public function markComplete():Void {
if (pieces == null) return;
for (i in 0...meta.numPieces) pieces.markHave(i);
done = true;
}
public function start():Void {
if (destroyed || keepaliveTimer != null) return;
stopped = false;
startTrackers();
keepaliveTimer = new haxe.Timer(KEEPALIVE_MS);
keepaliveTimer.run = function() {
for (w in wires) w.sendKeepAlive();
workWebSeeds();
};
workWebSeeds();
}
public function stop():Void {
if (destroyed || stopped) return;
stopped = true;
if (keepaliveTimer != null) {
keepaliveTimer.stop();
keepaliveTimer = null;
}
for (w in [for (x in wires) x]) w.close();
wires.clear();
outstanding.clear();
if (rtcPool != null) {
rtcPool.close();
rtcPool = null;
}
for (t in httpTrackers) t.stop();
for (t in udpTrackers) t.stop();
httpTrackers = [];
udpTrackers = [];
trackerObjects = [];
}
function startTrackers():Void {
for (url in trackers) startTracker(url);
}
public function addTracker(url:String):Dynamic {
if (url == null || url == "") return null;
for (t in trackers) {
if (t == url) return null;
}
trackers.push(url);
return startTracker(url);
}
public function setTrackers(urls:Array<String>):Void {
if (stopped) {
trackers = urls != null ? urls.copy() : [];
return;
}
if (rtcPool != null) {
rtcPool.close();
rtcPool = null;
}
for (t in httpTrackers) t.stop();
for (t in udpTrackers) t.stop();
httpTrackers = [];
udpTrackers = [];
trackerObjects = [];
trackers = urls != null ? urls.copy() : [];
startTrackers();
}
function startTracker(url:String):Dynamic {
var infoHashBin = TorrentInfo.toBinaryString(infoHash);
var peerIdBin = TorrentInfo.toBinaryString(peerId);
var obj:Dynamic = null;
if (url.indexOf("ws") == 0) {
if (rtcPool == null) {
rtcPool = new RtcPeerPool();
rtcPool.onPeerChannel = onRtcPeer;
rtcPool.statsProvider = getStats;
rtcPool.onError = function(e) {
emit1("error", e);
};
}
obj = rtcPool.connect(url, infoHashBin, peerIdBin);
}
#if (sys || kha_krom)
else if (url.indexOf("http") == 0) {
var t = new leenkx.network.torrent.tracker.HttpTracker(url, infoHash, peerId, announcePort > 0 ? announcePort : 6881);
t.statsProvider = getStats;
t.onPeers = onClassicPeers;
t.onError = function(e) {
emit1("error", e);
};
httpTrackers.push(t);
t.start();
obj = t;
}
#end
#if (sys || kha_krom)
else if (url.indexOf("udp") == 0) {
var t = new leenkx.network.torrent.tracker.UdpTracker(url, infoHash, peerId, announcePort > 0 ? announcePort : 6881);
t.statsProvider = getStats;
t.onPeers = onClassicPeers;
t.onError = function(e) {
emit1("error", e);
};
udpTrackers.push(t);
t.start();
obj = t;
}
#end
if (obj != null) trackerObjects.push(obj);
return obj;
}
function onRtcPeer(peerIdHex:String, ch:IPeerChannel):Void {
addWire("rtc:" + peerIdHex, ch);
}
#if (sys || kha_krom)
function onClassicPeers(list:Array<{host:String, port:Int}>):Void {
for (p in list) {
var key = "tcp:" + p.host + ":" + p.port;
if (wires.exists(key)) continue;
var ch = new leenkx.network.torrent.peer.TcpPeer(p.host, p.port);
addWire(key, ch);
}
}
#end
public function addIncomingPeer(key:String, ch:IPeerChannel):Void {
addWire(key, ch);
}
function addWire(key:String, ch:IPeerChannel):Void {
if (destroyed || stopped || wires.exists(key)) {
ch.close();
return;
}
var wire = new WireProtocol(ch, infoHash, peerId);
wires.set(key, wire);
outstanding.set(key, 0);
wire.metaSize = meta != null ? meta.rawInfo.length : 0;
wire.onHandshake = function(remotePeerId) {
if (ready) wire.sendBitfield(pieces.bitfield);
};
wire.onExtendedHandshake = function(hs) {
if (meta == null && wire.extIds.exists("ut_metadata")) {
requestMetadata(wire);
}
if (ready && !pieces.complete) wire.sendInterested();
};
wire.onUnchoke = function() requestBlocks(key);
wire.onBitfield = function(bf) {
peerBitfields.set(key, bitfieldToArray(bf));
if (ready && !pieces.complete) wire.sendInterested();
requestBlocks(key);
};
wire.onHave = function(index) {
if (meta == null || index < 0 || index >= meta.numPieces) return;
var bf = peerBitfields.get(key);
if (bf != null && index < bf.length) bf[index] = true;
requestBlocks(key);
};
wire.onPiece = function(index, begin, data) {
var n = outstanding.get(key);
if (n != null && n > 0) outstanding.set(key, n - 1);
if (pieces != null && index >= 0 && index < meta.numPieces && pieces.addBlock(index, begin, data)) {
broadcastHave(index);
}
requestBlocks(key);
};
wire.onRequest = function(index, begin, length) {
if (wire.amChoking || pieces == null || index < 0 || index >= meta.numPieces || begin < 0 || length <= 0 || begin + length > meta.pieceSize(index)) return;
if (!pieces.have[index]) return;
var data = pieces.readBlock(index, begin, length);
uploadedBytes += data.length;
wire.sendPiece(index, begin, data);
};
wire.onChoke = function() {
outstanding.set(key, 0);
if (pieces != null) pieces.clearAllRequested(key);
};
wire.onInterested = function() wire.sendUnchoke();
wire.onHashRequest = function(p) onHashRequest(wire, p);
wire.onHashes = onHashesMsg;
wire.onHashReject = onHashRejectMsg;
wire.onExtended = function(name, payload) {
if (name == "ut_metadata") {
onUtMetadata(wire, payload);
} else if (onPeerExtension != null) {
onPeerExtension(key, name, payload);
}
};
wire.onClose = function() {
wires.remove(key);
peerBitfields.remove(key);
outstanding.remove(key);
if (pieces != null) pieces.clearAllRequested(key);
if (meta == null && metaSize > 0 && metaWire == wire) {
resetMetaRequests();
metaWire = null;
for (k in wires.keys()) {
var w = wires.get(k);
if (w != wire && w.extIds.exists("ut_metadata")) {
metaWire = w;
sendMetaRequest(w, nextMetaPiece());
break;
}
}
}
};
wire.onError = function(e) {
emit1("error", "wire " + key + ": " + e);
};
wire.start(["ut_metadata", "lx_channel"]);
emit1("wire", wire);
}
static function bytesEq(a:Bytes, b:Bytes):Bool {
if (a == null || b == null || a.length != b.length) return false;
for (i in 0...a.length) if (a.get(i) != b.get(i)) return false;
return true;
}
function bitfieldToArray(bf:Bytes):Array<Bool> {
var n = meta != null ? meta.numPieces : bf.length * 8;
var out:Array<Bool> = [];
for (i in 0...n) {
var byte = i >> 3;
out.push(byte < bf.length && (bf.get(byte) & (0x80 >> (i & 7))) != 0);
}
return out;
}
function tryRequest(wire:WireProtocol, key:String, bf:Array<Bool>, i:Int):Bool {
if (pieces.have[i] || i >= bf.length || !bf[i]) return false;
if (meta.pieces == null && meta.fileHashes != null) {
var fi = meta.fileForPiece(i);
if (fi >= 0 && meta.fileHashes[fi] == null) {
requestPieceLayer(wire, fi);
return false;
}
}
var needed = pieces.neededBlocks(key, i);
if (needed.length == 0) return false;
var b = needed[0];
wire.sendRequest(i, b.begin, b.length);
pieces.markRequested(key, i, b.begin);
return true;
}
function requestBlocks(key:String):Void {
if (pieces == null || pieces.complete) return;
var expired = pieces.expireRequests(60);
for (pk in expired.keys()) {
var c = outstanding.get(pk);
if (c != null) {
c -= expired.get(pk);
outstanding.set(pk, c < 0 ? 0 : c);
}
}
var wire = wires.get(key);
var bf = peerBitfields.get(key);
if (wire == null || bf == null || wire.peerChoking) return;
var n = outstanding.get(key);
var cursor = pieceCursor.get(key);
if (cursor == null) cursor = 0;
var numPieces = meta.numPieces;
while (n < MAX_OUTSTANDING) {
var sent = false;
for (sel in pieces.selections) {
for (i in sel.start...sel.end + 1) {
if (tryRequest(wire, key, bf, i)) {
n++;
sent = true;
break;
}
}
if (sent) break;
}
if (!sent) {
for (j in 0...numPieces) {
var i = (cursor + j) % numPieces;
if (tryRequest(wire, key, bf, i)) {
n++;
sent = true;
pieceCursor.set(key, (i + 1) % numPieces);
break;
}
}
}
if (!sent) break;
}
outstanding.set(key, n);
workWebSeeds();
}
public function select(start:Int, end:Int, ?priority:Int, ?notify:Void -> Void):Void {
if (pieces == null) return;
pieces.select(start, end, priority, notify);
for (key in wires.keys()) requestBlocks(key);
workWebSeeds();
}
public function deselect(start:Int, end:Int, ?priority:Int):Void {
if (pieces == null) return;
pieces.deselect(start, end, priority);
}
public function critical(start:Int, end:Int):Void {
select(start, end, 0x7fffffff);
}
public function addPeer(peer:String):Bool {
#if (sys || kha_krom)
if (peer == null) return false;
var i = peer.lastIndexOf(":");
if (i < 0) return false;
var host = peer.substr(0, i);
var port = Std.parseInt(peer.substr(i + 1));
if (host == "" || port == null) return false;
var key = "tcp:" + host + ":" + port;
if (wires.exists(key)) return false;
addWire(key, new leenkx.network.torrent.peer.TcpPeer(host, port));
return true;
#else
return false;
#end
}
function broadcastHave(index:Int):Void {
for (w in wires) w.sendHave(index);
}
public function sendExtension(key:String, name:String, payload:Bytes):Void {
if (key != null) {
var w = wires.get(key);
if (w != null) w.sendExtended(name, payload);
return;
}
for (w in wires) w.sendExtended(name, payload);
}
function onPieceComplete(index:Int):Void {
emit1("piece", index);
}
var hashReqNext:Map<Int, Int> = [];
var hashReqTime:Map<Int, Float> = [];
var pendingLayers:Map<Int, {arr:Array<Bytes>, got:Int}> = [];
function pieceLayerBase():Int {
return Std.int(Math.log(meta.pieceLength / 16384) / Math.log(2));
}
function requestPieceLayer(wire:WireProtocol, fi:Int):Void {
var f = meta.files[fi];
var root = f.piecesRoot;
if (root == null || f.length <= 0) return;
var np = Std.int(Math.ceil(f.length / meta.pieceLength));
var sent = hashReqNext.exists(fi) ? hashReqNext.get(fi) : 0;
if (sent >= np) {
var t = hashReqTime.exists(fi) ? hashReqTime.get(fi) : 0;
if (haxe.Timer.stamp() - t < 60) return;
sent = 0;
}
var chunk = np - sent;
if (chunk > 512) chunk = 512;
var len = 2;
while (len < chunk) len <<= 1;
wire.sendHashRequest(root, pieceLayerBase(), sent, len, 0);
hashReqNext.set(fi, sent + len);
hashReqTime.set(fi, haxe.Timer.stamp());
}
function onHashesMsg(payload:Bytes):Void {
if (meta == null || meta.fileHashes == null) return;
if (payload.length < 80) return;
var root = payload.sub(0, 32);
var base = WireProtocol.readU32(payload, 32);
var index = WireProtocol.readU32(payload, 36);
var length = WireProtocol.readU32(payload, 40);
var proofs = WireProtocol.readU32(payload, 44);
if (base != pieceLayerBase()) return;
var count = Std.int((payload.length - 48) / 32);
if (count <= 0) return;
for (fi in 0...meta.files.length) {
var f = meta.files[fi];
if (f.piecesRoot == null || !bytesEq(root, f.piecesRoot)) continue;
if (meta.fileHashes[fi] != null) return;
var np = Std.int(Math.ceil(f.length / meta.pieceLength));
var pl = pendingLayers.get(fi);
if (pl == null) {
pl = {arr: [for (i in 0...np) null], got: 0};
pendingLayers.set(fi, pl);
}
var n = count < length ? count : length;
for (j in 0...n) {
var slot = index + j;
if (slot < np && pl.arr[slot] == null) {
pl.arr[slot] = payload.sub(48 + j * 32, 32);
pl.got++;
}
}
if (pl.got >= np) {
pendingLayers.remove(fi);
if (meta.verifyLayer(fi, pl.arr)) {
meta.fileHashes[fi] = pl.arr;
pieces.flushPendingV2(fi);
for (k in wires.keys()) requestBlocks(k);
} else {
trace('[Torrent] v2 layer REJECTED fi=' + fi);
hashReqNext.remove(fi);
}
}
return;
}
}
function onHashRejectMsg(payload:Bytes):Void {
if (meta == null) return;
var root = payload.sub(0, 32);
for (fi in 0...meta.files.length) {
var f = meta.files[fi];
if (f.piecesRoot != null && bytesEq(root, f.piecesRoot)) {
hashReqNext.remove(fi);
pendingLayers.remove(fi);
return;
}
}
}
var layerCache:Map<Int, Array<Bytes>> = [];
var treeCache:Map<Int, Array<Array<Bytes>>> = [];
function hashLayerForFile(fi:Int):Array<Bytes> {
var l = layerCache.get(fi);
if (l != null) return l;
var f = meta.files[fi];
var np = Std.int(Math.ceil(f.length / meta.pieceLength));
l = [];
for (p in 0...np) {
var size = Std.int(Math.min(meta.pieceLength, f.length - p * meta.pieceLength));
l.push(meta.pieceRootV2(storage.read(f.offset + p * meta.pieceLength, size)));
}
layerCache.set(fi, l);
return l;
}
function treeLevelsFor(fi:Int):Array<Array<Bytes>> {
var t = treeCache.get(fi);
if (t != null) return t;
var level = hashLayerForFile(fi);
t = [level];
var depth = 0;
while (level.length > 1) {
var zero = MetaInfo.layerZero(pieceLayerBase() + depth + 1);
var next:Array<Bytes> = [];
var i = 0;
while (i < level.length) {
var b = i + 1 < level.length ? level[i + 1] : zero;
next.push(MetaInfo.hashPair(level[i], b));
i += 2;
}
t.push(next);
level = next;
depth++;
}
if (t.length > 1) t.pop();
treeCache.set(fi, t);
return t;
}
function onHashRequest(wire:WireProtocol, payload:Bytes):Void {
var root = payload.sub(0, 32);
var base = WireProtocol.readU32(payload, 32);
var index = WireProtocol.readU32(payload, 36);
var length = WireProtocol.readU32(payload, 40);
var proofs = WireProtocol.readU32(payload, 44);
var reject = function() wire.sendHashReject(root, base,
index, length, proofs);
if (meta == null || meta.pieces != null || storage == null
|| length < 2 || length > 512
|| (length & (length - 1)) != 0
|| (length > 0 && index % length != 0)) {
reject();
return;
}
var pieceBase = pieceLayerBase();
if (base != pieceBase) {
reject();
return;
}
for (fi in 0...meta.files.length) {
var f = meta.files[fi];
if (f.piecesRoot == null || f.length <= 0
|| !bytesEq(root, f.piecesRoot)) continue;
var layer = hashLayerForFile(fi);
if (index >= layer.length) {
reject();
return;
}
var n = layer.length - index;
if (n > length) n = length;
var buf = new haxe.io.BytesBuffer();
for (i in index...index + n) {
buf.addBytes(layer[i], 0, 32);
}
var lg = 0;
while ((1 << lg) < length) lg++;
var tree = treeLevelsFor(fi);
for (p in 0...(proofs - lg + 1)) {
var lvl = lg + p;
if (lvl >= tree.length) break;
var sib = (index >> lvl) ^ 1;
var h = sib < tree[lvl].length
? tree[lvl][sib]
: MetaInfo.layerZero(pieceBase + lvl);
buf.addBytes(h, 0, 32);
}
wire.sendHashes(root, base, index, n,
proofs, buf.getBytes());
return;
}
reject();
}
function onTorrentComplete():Void {
done = true;
storage.flush();
#if sys
if (downloadDir != null
&& Std.isOfType(storage, MemoryStorage)) {
writeToDisk(downloadDir);
}
#end
for (t in trackerObjects) {
var f:Dynamic = Reflect.field(t, "announceEvent");
if (f != null) {
try {
Reflect.callMethod(t, f, ["completed"]);
} catch(e:Dynamic) {}
}
}
emit("done");
if (client != null && client.onTorrentDone != null) {
client.onTorrentDone(this);
}
}
#if sys
function writeToDisk(dir:String):Void {
for (f in files) {
var p = dir + "/" + MetaInfo.sanitizePath(f.path);
try {
var i = Std.int(Math.max(p.lastIndexOf("/"),
p.lastIndexOf("\\")));
if (i > 0) {
var d = p.substr(0, i);
if (!sys.FileSystem.exists(d)) {
sys.FileSystem.createDirectory(d);
}
}
var out = sys.io.File.write(p, true);
out.write(f.getBytes());
out.close();
} catch(e:Dynamic) {
emit1("error", "write " + p + ": " + Std.string(e));
}
}
}
#end
function initWebSeeds():Void {
if (meta == null || meta.urlList == null) return;
for (u in meta.urlList) addWebSeed(u);
}
public function addWebSeed(url:String):Void {
if (url == null || url == "") return;
for (w in webSeeds) if (w.url == url) return;
var ws = new WebSeed(url, this);
var self = this;
ws.onError = function(e) {
self.emit1("error", "webseed " + url + ": " + e);
haxe.Timer.delay(self.workWebSeeds, 10000);
};
ws.onDone = workWebSeeds;
webSeeds.push(ws);
ws.work();
}
function workWebSeeds():Void {
if (destroyed || stopped || done) return;
for (w in webSeeds) w.work();
}
function requestMetadata(wire:WireProtocol):Void {
if (metaSize > 0) {
var alive = false;
for (w in wires) {
if (w == metaWire) alive = true;
}
if (!alive) {
resetMetaRequests();
metaWire = wire;
sendMetaRequest(wire, nextMetaPiece());
}
return;
}
var hs = wire.peerExtendedHandshake;
metaWire = wire;
var size:Dynamic = Reflect.field(hs, "metadata_size");
if (size == null) return;
var sizeF = Std.parseFloat(Std.string(size));
if (Math.isNaN(sizeF) || sizeF <= 0 || sizeF > MAX_META_SIZE) {
emit1("error", "invalid metadata_size " + size);
return;
}
metaSize = Std.int(sizeF);
metaRejects = 0;
var count = Std.int(Math.ceil(metaSize / META_BLOCK));
metaChunks.clear();
metaReceived = 0;
metaRequested = [for (i in 0...count) false];
sendMetaRequest(wire, 0);
}
function resetMetaRequests():Void {
for (i in 0...metaRequested.length) {
if (metaRequested[i] && !metaChunks.exists(i)) {
metaRequested[i] = false;
}
}
}
function sendMetaRequest(wire:WireProtocol, piece:Int):Void {
if (piece < 0 || piece >= metaRequested.length
|| metaRequested[piece]) return;
metaRequested[piece] = true;
var req = Bencode.encode({msg_type: 0, piece: piece});
wire.sendExtended("ut_metadata", req);
}
function onUtMetadata(wire:WireProtocol, payload:Bytes):Void {
var headerEnd = findDictEnd(payload);
if (headerEnd < 0) return;
var header:Dynamic;
try {
header = Bencode.decode(payload.sub(0, headerEnd));
} catch(e:Dynamic) {
return;
}
var msgType = Std.int(metaField(header, "msg_type"));
var piece = Std.int(metaField(header, "piece"));
if (msgType == 0) {
if (meta == null || piece < 0) {
var rej = Bencode.encode({msg_type: 2, piece: piece});
wire.sendExtended("ut_metadata", rej);
return;
}
var startF = piece * (META_BLOCK : Float);
if (startF >= meta.rawInfo.length) return;
var start = Std.int(startF);
var len = meta.rawInfo.length - start;
if (len > META_BLOCK) len = META_BLOCK;
var head = Bencode.encode({
msg_type: 1, piece: piece,
total_size: meta.rawInfo.length
});
var out = new haxe.io.BytesBuffer();
out.addBytes(head, 0, head.length);
out.addBytes(meta.rawInfo, start, len);
wire.sendExtended("ut_metadata", out.getBytes());
} else if (msgType == 1) {
var total = metaField(header, "total_size");
if (Math.isNaN(total) || total <= 0 || total > MAX_META_SIZE) {
return;
}
var totalI = Std.int(total);
if (metaSize == 0 || totalI != metaSize) {
metaSize = totalI;
metaRejects = 0;
metaChunks.clear();
metaReceived = 0;
metaRequested = [
for (i in 0...Std.int(Math.ceil(totalI / META_BLOCK)))
false
];
}
if (piece < 0 || piece >= metaRequested.length) return;
var expected = metaSize - piece * META_BLOCK;
if (expected > META_BLOCK) expected = META_BLOCK;
var chunk = payload.sub(headerEnd, payload.length - headerEnd);
if (chunk.length > expected) {
chunk = chunk.sub(0, expected);
}
appendMetaChunk(piece, chunk);
if (metaReceived >= metaSize) finishMetadata(wire);
else sendMetaRequest(wire, nextMetaPiece());
} else if (msgType == 2) {
if (piece >= 0 && piece < metaRequested.length) {
metaRequested[piece] = false;
metaRejects++;
if (metaRejects > 2 * metaRequested.length) {
metaSize = 0;
metaChunks.clear();
emit1("error", "ut_metadata rejected by peers");
} else {
for (k in wires.keys()) {
var w = wires.get(k);
if (w != wire
&& w.extIds.exists("ut_metadata")) {
sendMetaRequest(w, piece);
break;
}
}
}
}
}
}
var metaChunks:Map<Int, Bytes> = [];
function appendMetaChunk(piece:Int, chunk:Bytes):Void {
if (!metaChunks.exists(piece)) {
metaChunks.set(piece, chunk);
metaReceived += chunk.length;
}
}
function nextMetaPiece():Int {
for (i in 0...metaRequested.length) {
if (!metaRequested[i]) return i;
}
return -1;
}
function finishMetadata(wire:WireProtocol):Void {
var buf = new haxe.io.BytesBuffer();
var count = Std.int(Math.ceil(metaSize / META_BLOCK));
for (i in 0...count) {
var c = metaChunks.get(i);
if (c != null) buf.addBytes(c, 0, c.length);
}
var raw = buf.getBytes().sub(0, metaSize);
var okMeta = false;
var v1Hash:Bytes = null;
var sha1 = haxe.crypto.Sha1.make(raw);
if (infoHash != null && bytesEq(sha1, infoHash)) {
okMeta = true;
v1Hash = infoHash;
}
if (!okMeta) {
var sha2 = haxe.crypto.Sha256.make(raw);
if (infoHashV2 != null && bytesEq(sha2, infoHashV2)) {
okMeta = true;
} else if (infoHash != null
&& bytesEq(sha2.sub(0, 20), infoHash)) {
okMeta = true;
}
}
if (!okMeta) {
emit1("error", "ut_metadata hash mismatch");
metaRejects += metaRequested.length;
metaChunks.clear();
metaReceived = 0;
if (metaRejects > 2 * metaRequested.length) {
metaSize = 0;
return;
}
for (i in 0...metaRequested.length) metaRequested[i] = false;
sendMetaRequest(wire, nextMetaPiece());
return;
}
setMeta(MetaInfo.fromInfoDict(raw, v1Hash, trackers));
metaChunks.clear();
for (key in wires.keys()) requestBlocks(key);
}
static function metaField(d:Dynamic, name:String):Float {
var v = Reflect.field(d, name);
if (v == null) return 0;
return Std.parseFloat(Std.string(v));
}
static function findDictEnd(data:Bytes):Int {
if (data.length == 0 || data.get(0) != 'd'.code) return -1;
var end = MetaInfo.skipValue(data, 0);
return end <= data.length ? end : -1;
}
public function destroy():Void {
if (destroyed) return;
destroyed = true;
if (keepaliveTimer != null) {
keepaliveTimer.stop();
keepaliveTimer = null;
}
if (rtcPool != null) rtcPool.close();
for (t in httpTrackers) t.stop();
for (t in udpTrackers) t.stop();
for (w in [for (x in wires) x]) w.close();
wires.clear();
if (storage != null) storage.close();
}
public function deleteData():Void {
if (storage != null) storage.deleteData();
}
}