Files
LNXSDK/leenkx/Sources/leenkx/network/torrent/PieceManager.hx
2026-10-03 17:40:22 -07:00

330 lines
9.0 KiB
Haxe

package leenkx.network.torrent;
import haxe.io.Bytes;
import leenkx.network.torrent.storage.Storage.IStorage;
class PieceManager {
public static inline var BLOCK_SIZE = 16384;
var meta:MetaInfo;
var storage:IStorage;
public var have:Array<Bool>;
public var numHave:Int = 0;
public var bitfield:Bytes;
var pending:Map<Int, Map<Int, Bytes>> = [];
var pendingCount:Map<Int, Int> = [];
var requested:Map<String, Map<Int, Map<Int, Bool>>> = [];
var globalRequested:Map<Int, Map<Int, Bool>> = [];
var pendingV2:Map<Int, Bytes> = [];
public var onPieceComplete:Int -> Void;
public var onPieceFailed:Int -> Void;
public var onComplete:Void -> Void;
@:allow(leenkx.network.torrent.Torrent)
@:allow(leenkx.network.torrent.WebSeed)
var selections:Array<{start:Int, end:Int, priority:Int,
notify:Void -> Void}> = [];
public function select(startByte:Int, endByte:Int, ?priority:Int, ?notify:Void -> Void):Void {
selections.push({start: pieceIndex(startByte),
end: pieceIndex(endByte),
priority: priority != null ? priority : 0,
notify: notify});
selections.sort(function(a, b) return b.priority
- a.priority);
}
public function deselect(startByte:Int, endByte:Int, ?priority:Int):Void {
var s = pieceIndex(startByte);
var e = pieceIndex(endByte);
var p = priority != null ? priority : 0;
for (i in 0...selections.length) {
var sel = selections[i];
if (sel.start == s && sel.end == e
&& sel.priority == p) {
selections.splice(i, 1);
return;
}
}
}
function pieceIndex(byte:Int):Int {
var i = Std.int(byte / meta.pieceLength);
if (i < 0) i = 0;
if (i >= meta.numPieces) i = meta.numPieces - 1;
return i;
}
function updateSelections(index:Int):Void {
var i = 0;
while (i < selections.length) {
var sel = selections[i];
if (index < sel.start || index > sel.end
|| sel.notify == null) {
i++;
continue;
}
var done = true;
for (p in sel.start...sel.end + 1) {
if (!have[p]) {
done = false;
break;
}
}
if (done) {
selections.splice(i, 1);
sel.notify();
} else i++;
}
}
public var complete(get, never):Bool;
function get_complete():Bool return numHave == meta.numPieces;
public function new(meta:MetaInfo, storage:IStorage) {
this.meta = meta;
this.storage = storage;
have = [for (i in 0...meta.numPieces) false];
bitfield = Bytes.alloc(Std.int(Math.ceil(meta.numPieces / 8)));
}
public static function forSeeder(meta:MetaInfo, storage:IStorage):PieceManager {
var pm = new PieceManager(meta, storage);
for (i in 0...meta.numPieces) pm.markHave(i);
return pm;
}
public function markHave(index:Int):Void {
if (have[index]) return;
have[index] = true;
numHave++;
bitfield.set(index >> 3, bitfield.get(index >> 3)
| (0x80 >> (index & 7)));
}
public function addBlock(index:Int, begin:Int, data:Bytes):Bool {
if (index < 0 || index >= meta.numPieces || begin < 0
|| data == null
|| begin + data.length > meta.pieceSize(index)) {
return false;
}
if (have[index]) return true;
var blocks = pending.get(index);
if (blocks == null) {
blocks = [];
pending.set(index, blocks);
pendingCount.set(index, 0);
}
if (!blocks.exists(begin)) {
blocks.set(begin, data);
pendingCount.set(index, pendingCount.get(index) + data.length);
var gReq = globalRequested.get(index);
if (gReq != null) {
gReq.remove(begin);
if (!gReq.keys().hasNext()) globalRequested.remove(index);
}
}
if (pendingCount.get(index) < meta.pieceSize(index)) return false;
var piece = Bytes.alloc(meta.pieceSize(index));
var off = 0;
var keys = [for (k in blocks.keys()) k];
keys.sort(Reflect.compare);
for (b in keys) {
var d = blocks.get(b);
piece.blit(b, d, 0, d.length);
}
var ok = true;
if (meta.pieces == null) {
var fi = meta.fileForPiece(index);
var layer = fi >= 0 && meta.fileHashes != null
&& meta.fileHashes.length > fi
? meta.fileHashes[fi] : null;
pending.remove(index);
pendingCount.remove(index);
if (layer == null) {
pendingV2.set(index, piece);
return false;
}
ok = meta.verifyPieceV2(index, piece);
} else {
var hash = haxe.crypto.Sha1.make(piece);
ok = MetaInfo.memcmp(hash, meta.pieceHash(index));
pending.remove(index);
pendingCount.remove(index);
}
if (!ok) {
trace('[PieceManager] piece ' + index
+ ' verification FAILED');
releaseRequested(index);
if (onPieceFailed != null) onPieceFailed(index);
return false;
}
storage.write(index * meta.pieceLength, piece);
markHave(index);
updateSelections(index);
if (onPieceComplete != null) onPieceComplete(index);
if (complete && onComplete != null) onComplete();
return true;
}
public function flushPendingV2(fi:Int):Bool {
if (meta.pieces != null) return false;
var f = meta.files[fi];
var start = Std.int(f.offset / meta.pieceLength);
var np = Std.int(Math.ceil(f.length / meta.pieceLength));
var wrote = false;
for (i in start...start + np) {
var piece = pendingV2.get(i);
if (piece == null) continue;
pendingV2.remove(i);
if (!meta.verifyPieceV2(i, piece)) {
trace('[PieceManager] flush piece ' + i
+ ' verification FAILED');
releaseRequested(i);
if (onPieceFailed != null) onPieceFailed(i);
continue;
}
storage.write(i * meta.pieceLength, piece);
markHave(i);
updateSelections(i);
wrote = true;
if (onPieceComplete != null) onPieceComplete(i);
}
if (complete && onComplete != null) onComplete();
return wrote;
}
public function untouched(index:Int):Bool {
return !have[index] && !pending.exists(index)
&& !pendingV2.exists(index)
&& !globalRequested.exists(index);
}
public function readBlock(index:Int, begin:Int, length:Int):Bytes {
return storage.read(index * meta.pieceLength + begin, length);
}
public function neededBlocks(peerKey:String, index:Int):Array<{begin:Int, length:Int}> {
if (pendingV2.exists(index)) return [];
var size = meta.pieceSize(index);
var blocks = pending.get(index);
var peerReq = requested.get(peerKey);
var req = peerReq != null ? peerReq.get(index) : null;
var globalReq = globalRequested.get(index);
var out:Array<{begin:Int, length:Int}> = [];
var off = 0;
while (off < size) {
var len = size - off;
if (len > BLOCK_SIZE) len = BLOCK_SIZE;
if ((blocks == null || !blocks.exists(off))
&& (req == null || !req.exists(off))
&& (globalReq == null || !globalReq.exists(off))) {
out.push({begin: off, length: len});
}
off += len;
}
return out;
}
public function markRequested(key:String, index:Int, begin:Int):Void {
var r = requested.get(key);
if (r == null) requested.set(key,
r = new Map<Int, Map<Int, Bool>>());
var m = r.get(index);
if (m == null) r.set(index, m = new Map<Int, Bool>());
m.set(begin, true);
var g = globalRequested.get(index);
if (g == null) globalRequested.set(index,
g = new Map<Int, Bool>());
g.set(begin, true);
reqTimes.set(key + "|" + index + "|" + begin,
haxe.Timer.stamp());
}
var reqTimes:Map<String, Float> = [];
public function expireRequests(ttlSec:Float):Map<String, Int> {
var now = haxe.Timer.stamp();
var dead:Array<String> = [];
for (k in reqTimes.keys()) {
if (now - reqTimes.get(k) > ttlSec) dead.push(k);
}
var expired:Map<String, Int> = [];
for (k in dead) {
var p = k.split("|");
var idx = Std.parseInt(p[1]);
var beg = Std.parseInt(p[2]);
var peerReq = requested.get(p[0]);
if (peerReq != null) {
var r = peerReq.get(idx);
if (r != null && r.exists(beg)) {
r.remove(beg);
var n = expired.exists(p[0]) ? expired.get(p[0]) : 0;
expired.set(p[0], n + 1);
}
}
var g = globalRequested.get(idx);
if (g != null) {
g.remove(beg);
var empty = true;
for (b in g.keys()) {
empty = false;
break;
}
if (empty) globalRequested.remove(idx);
}
reqTimes.remove(k);
}
return expired;
}
public function releaseRequested(index:Int):Void {
globalRequested.remove(index);
var keys = [for (k in requested.keys()) k];
for (key in keys) {
var r = requested.get(key);
var m = r.get(index);
if (m != null) {
for (b in m.keys()) {
reqTimes.remove(key + "|" + index + "|" + b);
}
r.remove(index);
}
}
}
public function clearRequested(peerKey:String, index:Int):Void {
var peerReq = requested.get(peerKey);
if (peerReq != null) {
var req = peerReq.get(index);
if (req != null) {
var gReq = globalRequested.get(index);
if (gReq != null) {
for (off in req.keys()) gReq.remove(off);
if (!gReq.keys().hasNext()) globalRequested.remove(index);
}
}
peerReq.remove(index);
}
}
public function clearAllRequested(peerKey:String):Void {
var peerReq = requested.get(peerKey);
if (peerReq != null) {
for (index in peerReq.keys()) {
var req = peerReq.get(index);
var gReq = globalRequested.get(index);
if (gReq != null && req != null) {
for (off in req.keys()) gReq.remove(off);
if (!gReq.keys().hasNext()) globalRequested.remove(index);
}
}
requested.remove(peerKey);
}
}
}