gapmind/services/dht/main.go
2026-09-11 23:59:05 -04:00

320 lines
8.6 KiB
Go

package main
import (
"bytes"
"encoding/binary"
"encoding/hex"
"encoding/json"
"flag"
"fmt"
"io"
"net"
"net/http"
"net/url"
"os"
"strings"
"time"
)
var builtInTrackers = []string{
"udp://tracker.opentrackr.org:1337/announce",
"udp://open.stealth.si:80/announce",
"udp://tracker.torrent.eu.org:451/announce",
}
var builtInRSS = []string{
"https://eztv.re/showlist/",
"https://nyaa.si/?page=rss",
"https://www.torrentdownloads.me/rss.xml",
"https://fosstorrents.com/feed/rss.xml",
}
func registryList(envURL, fallbackEnv string, builtIn []string) []string {
if u := os.Getenv(envURL); u != "" {
if req, err := http.NewRequest(http.MethodGet, u, nil); err == nil {
if resp, err := http.DefaultClient.Do(req); err == nil {
var list []string
if json.NewDecoder(resp.Body).Decode(&list) == nil && len(list) > 0 {
resp.Body.Close()
return list
}
resp.Body.Close()
}
}
}
if raw := os.Getenv(fallbackEnv); raw != "" {
return strings.Split(raw, ",")
}
return builtIn
}
func main() {
central := flag.String("central", envOr("CENTRAL_URL", "http://127.0.0.1:3000"), "central gapmind URL")
interval := flag.Duration("interval", 10*time.Minute, "RSS re-ping interval")
listen := flag.String("listen", envOr("DHT_LISTEN", ":6881"), "UDP listen addr for passive DHT node")
passive := flag.Bool("passive", true, "run passive DHT get_peers/announce_peer sniffer")
flag.Parse()
client := &http.Client{Timeout: 20 * time.Second}
names := map[string]string{}
pingRSS := func() []string {
rss := registryList("RSS_REGISTRY_URL", "RSS_FEEDS", builtInRSS)
found := scrapeRSSNames(client, rss, names)
trackers := registryList("TRACKER_REGISTRY_URL", "TRACKERS", builtInTrackers)
raw := make([][]byte, 0, len(found))
for _, h := range found {
if b, err := hex.DecodeString(h); err == nil {
raw = append(raw, b)
}
}
meta := map[string]TrackerMeta{}
for _, t := range trackers {
if u, err := url.Parse(t); err == nil && u.Scheme == "udp" {
for k, v := range udpScrape(u.Host, raw) {
meta[k] = v
}
break
}
}
for _, h := range found {
m := meta[h]
postJSON(client, *central+"/api/dht/torrents", map[string]any{
"infoHash": h, "name": names[h], "source": "rss-ping+tracker-scrape",
"seeders": m.Seeders, "leechers": m.Leechers,
})
}
return found
}
if *passive {
node, err := ListenDHT(*listen, func(ev SniffEvent) {
trackers := registryList("TRACKER_REGISTRY_URL", "TRACKERS", builtInTrackers)
raw, _ := hex.DecodeString(ev.InfoHash)
counts := TrackerMeta{}
for _, t := range trackers {
if u, err := url.Parse(t); err == nil && u.Scheme == "udp" {
if m, ok := udpScrape(u.Host, [][]byte{raw})[ev.InfoHash]; ok {
counts = m
break
}
}
}
postJSON(client, *central+"/api/dht/torrents", map[string]any{
"infoHash": ev.InfoHash, "name": names[ev.InfoHash],
"source": "dht-sniff:" + ev.Query,
"seeders": counts.Seeders, "leechers": counts.Leechers,
})
postJSON(client, *central+"/api/dht/peers", map[string]any{
"infoHash": ev.InfoHash,
"peers": []map[string]any{{"address": ev.IP, "port": ev.Port, "client": "dht-" + ev.Query}},
})
fmt.Printf("sniff %s %s %s:%d seeders=%d\n", ev.Query, ev.InfoHash, ev.IP, ev.Port, counts.Seeders)
})
if err != nil {
fmt.Println("dht listen failed:", err)
} else {
defer node.Close()
fmt.Println("dht passive listener on", *listen)
}
}
for {
hashes := pingRSS()
_ = hashes
time.Sleep(*interval)
}
}
func scrapeRSSNames(client *http.Client, feeds []string, names map[string]string) []string {
seen := map[string]bool{}
for _, f := range feeds {
resp, err := client.Get(f)
if err != nil {
continue
}
body, _ := io.ReadAll(io.LimitReader(resp.Body, 4<<20))
resp.Body.Close()
text := strings.ReplaceAll(string(body), "\"", " ")
for _, tok := range strings.Fields(text) {
lower := strings.ToLower(tok)
if i := strings.Index(lower, "btih:"); i >= 0 {
h := lower[i+5 : min(i+45, len(lower))]
if len(h) == 40 && isHex(h) && !seen[h] {
seen[h] = true
if j := strings.Index(tok, "dn="); j >= 0 {
if name, err := url.QueryUnescape(strings.TrimRight(tok[j+3:], "&")); err == nil && name != "" && names[h] == "" {
names[h] = name[:min(len(name), 200)]
}
}
}
}
if len(tok) == 40 && isHex(strings.ToLower(tok)) {
seen[strings.ToLower(tok)] = true
}
}
}
out := make([]string, 0, len(seen))
for h := range seen {
out = append(out, h)
if len(out) >= 50 {
break
}
}
return out
}
func isHex(s string) bool { _, err := hex.DecodeString(s); return err == nil }
func min(a, b int) int {
if a < b {
return a
}
return b
}
func envOr(k, d string) string {
if v := os.Getenv(k); v != "" {
return v
}
return d
}
func crawlTrackers(infoHash string, trackers []string) []map[string]any {
var peers []map[string]any
raw, _ := hex.DecodeString(infoHash)
for _, t := range trackers {
u, err := url.Parse(t)
if err != nil || u.Scheme != "udp" {
continue
}
if got := udpAnnounce(u.Host, raw); len(got) > 0 {
peers = append(peers, got...)
}
if len(peers) >= 100 {
break
}
}
return peers
}
func udpAnnounce(host string, infoHash []byte) []map[string]any {
addr, err := net.ResolveUDPAddr("udp", host)
if err != nil {
return nil
}
conn, err := net.DialUDP("udp", nil, addr)
if err != nil {
return nil
}
defer conn.Close()
_ = conn.SetDeadline(time.Now().Add(5 * time.Second))
transID := uint32(time.Now().UnixNano())
var req bytes.Buffer
_ = binary.Write(&req, binary.BigEndian, uint64(0x41727101980))
_ = binary.Write(&req, binary.BigEndian, uint32(0))
_ = binary.Write(&req, binary.BigEndian, transID)
if _, err := conn.Write(req.Bytes()); err != nil {
return nil
}
resp := make([]byte, 16)
n, err := conn.Read(resp)
if err != nil || n < 16 {
return nil
}
connID := resp[8:16]
peerID := make([]byte, 20)
copy(peerID, []byte("-GM0001-123456789012"))
var ann bytes.Buffer
ann.Write(connID)
_ = binary.Write(&ann, binary.BigEndian, uint32(1))
_ = binary.Write(&ann, binary.BigEndian, transID+1)
ann.Write(infoHash)
ann.Write(peerID)
_ = binary.Write(&ann, binary.BigEndian, uint64(0))
_ = binary.Write(&ann, binary.BigEndian, uint64(0))
_ = binary.Write(&ann, binary.BigEndian, uint64(0))
_ = binary.Write(&ann, binary.BigEndian, int32(-1))
_ = binary.Write(&ann, binary.BigEndian, uint32(0))
_ = binary.Write(&ann, binary.BigEndian, uint32(0))
_ = binary.Write(&ann, binary.BigEndian, int32(-1))
_ = binary.Write(&ann, binary.BigEndian, uint16(6881))
if _, err := conn.Write(ann.Bytes()); err != nil {
return nil
}
buf := make([]byte, 4096)
_ = conn.SetDeadline(time.Now().Add(5 * time.Second))
n, err = conn.Read(buf)
if err != nil || n < 20 {
return nil
}
var out []map[string]any
for i := 20; i+6 <= n; i += 6 {
ip := net.IP(buf[i : i+4])
port := int(binary.BigEndian.Uint16(buf[i+4 : i+6]))
out = append(out, map[string]any{"address": ip.String(), "port": port})
}
return out
}
func handshakeAll(infoHash string, peers []map[string]any) []map[string]any {
raw, _ := hex.DecodeString(infoHash)
var out []map[string]any
for _, p := range peers {
host, _ := p["address"].(string)
var port int
switch v := p["port"].(type) {
case int:
port = v
case float64:
port = int(v)
}
client := handshakeOne(host, port, raw)
entry := map[string]any{"address": host, "port": port}
if client != "" {
entry["client"] = client
}
out = append(out, entry)
if len(out) >= 100 {
break
}
}
return out
}
func handshakeOne(host string, port int, infoHash []byte) string {
conn, err := net.DialTimeout("tcp", fmt.Sprintf("%s:%d", host, port), 4*time.Second)
if err != nil {
return ""
}
defer conn.Close()
_ = conn.SetDeadline(time.Now().Add(5 * time.Second))
peerID := []byte("-GM0001-123456789012")
hs := append([]byte{19}, []byte("BitTorrent protocol")...)
hs = append(hs, []byte{0, 0, 0, 0, 0, 0x10, 0, 0}...)
hs = append(hs, infoHash...)
hs = append(hs, peerID...)
if _, err := conn.Write(hs); err != nil {
return ""
}
resp := make([]byte, 68+256)
n, err := conn.Read(resp)
if err != nil || n < 68 {
return ""
}
peerStr := string(resp[48:68])
return strings.Trim(peerStr, "\x00 ")
}
func postJSON(client *http.Client, endpoint string, body any) {
raw, _ := json.Marshal(body)
req, _ := http.NewRequest(http.MethodPost, endpoint, bytes.NewReader(raw))
req.Header.Set("Content-Type", "application/json")
if tok := os.Getenv("DHT_API_TOKEN"); tok != "" {
req.Header.Set("Authorization", "Bearer "+tok)
}
resp, err := client.Do(req)
if err == nil {
_, _ = io.Copy(io.Discard, resp.Body)
_ = resp.Body.Close()
}
}