-- coding: utf-8 --

import socket import hashlib import random import time import struct import threading import urllib.parse

'这是一个Python脚本,用于连接到DHT网络并提供查询磁力链接下载人数的接口。' '它不使用btdht库,但可以使用其他库。' '这个脚本需要长期运行。'

DHT协议常量

DHT_PORT = 6881 DHT_NODE_LENGTH = 26 DHT_MAX_NODES = 500 DHT_MAX_RESPONSE_TIME = 10 # 秒 DHT_MAX_RETRY_COUNT = 3 DHT_BOOTSTRAP_NODES = [ ('router.bittorrent.com', 6881), ('dht.transmissionbt.com', 6881), ('router.utorrent.com', 6881) ] DHT_QUERY_NODES = [ b'get_peers', b'announce_peer', b'find_node' ]

其他常量

MAGNET_PREFIX = 'magnet:?xt=urn:btih:' BUFFER_SIZE = 1024

class DHTNode: def init(self, node_id, ip, port): self.node_id = node_id self.ip = ip self.port = port

def __repr__(self):
    return f'<DHTNode: {self.ip}:{self.port}, node_id={self.node_id.hex()}>'

def __eq__(self, other):
    return self.node_id == other.node_id and self.ip == other.ip and self.port == other.port

def __hash__(self):
    return hash((self.node_id, self.ip, self.port))

class DHT: def init(self, node_id=None): if node_id is None: node_id = self.generate_random_node_id() self.node_id = node_id self.socket = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) self.socket.bind(('0.0.0.0', DHT_PORT)) self.socket.settimeout(1) self.nodes = set()

def bootstrap(self):
    '''
    启动DHT节点,加入到DHT网络中
    '''
    for bootstrap_node in DHT_BOOTSTRAP_NODES:
        self.send_find_node(bootstrap_node)

def generate_random_node_id(self):
    '''
    生成一个随机的20字节的node id
    '''
    node_id = b''
    for i in range(20):
        node_id += struct.pack('>B', random.randint(0, 255))
    return node_id

def send_message(self, message, address):
    '''
    发送DHT协议消息
    '''
    try:
        self.socket.sendto(message, address)
    except Exception as e:
        print(f'[send_message error]: {e}')

def send_ping(self, address):
    '''
    发送ping消息
    '''
    message = self.generate_message(b'ping', {})
    self.send_message(message, address)

def send_find_node(self, address):
    '''
    发送find_node消息
    '''
    message = self.generate_message(b'find_node', {'target': self.generate_random_node_id()})
    self.send_message(message, address)

def send_get_peers(self, address, info_hash):
    '''
    发送get_peers消息
    '''
    message = self.generate_message(b'get_peers', {'info_hash': info_hash})
    self.send_message(message, address)

def generate_message(self, message_type, arguments):
    '''
    生成DHT协议消息
    '''
    message = b'd1:ad2:id20:' + self.node_id
    message += message_type
    for key, value in arguments.items():
        message += len(key).to_bytes(1, byteorder='big') + key.encode('utf-8')
        message += len(value).to_bytes(1, byteorder='big') + value
    message += b'e'
    return message

def process_message(self, message, address):
    '''
    处理DHT协议消息
    '''
    try:
        message_type = message[b'y'].decode('utf-8')
        if message_type == 'r':
            self.process_response(message, address)
        elif message_type == 'q':
            self.process_query(message, address)
    except Exception as e:
        print(f'[process_message error]: {e}')

def process_response(self, message, address):
    '''
    处理response消息
    '''
    try:
        nodes = message[b'r'][b'nodes']
        for i in range(0, len(nodes), DHT_NODE_LENGTH):
            node_id = nodes[i:i+20]
            ip = socket.inet_ntoa(nodes[i+20:i+24])
            port = int.from_bytes(nodes[i+24:i+26], byteorder='big')
            node = DHTNode(node_id, ip, port)
            self.nodes.add(node)
    except Exception as e:
        print(f'[process_response error]: {e}')

def process_query(self, message, address):
    '''
    处理query消息
    '''
    try:
        query_type = message[b'q']
        if query_type == b'get_peers':
            info_hash = message[b'a'][b'info_hash']
            self.send_response(address, self.generate_response(b'get_peers', info_hash))
        elif query_type == b'announce_peer':
            pass
        elif query_type == b'find_node':
            self.send_response(address, self.generate_response(b'find_node'))
    except Exception as e:
        print(f'[process_query error]: {e}')

def send_response(self, address, response):
    '''
    发送response消息
    '''
    self.send_message(response, address)

def generate_response(self, message_type, info_hash=None):
    '''
    生成response消息
    '''
    if message_type == b'find_node':
        nodes = b''
        for node in self.nodes:
            nodes += node.node_id + socket.inet_aton(node.ip) + node.port.to_bytes(2, byteorder='big')
        return self.generate_message(b'r', {'id': self.node_id, 'nodes': nodes})
    elif message_type == b'get_peers':
        return self.generate_message(b'r', {'id': self.node_id, 'token': self.generate_random_node_id()})
    elif message_type == b'announce_peer':
        pass

def search_info_hash(self, info_hash):
    '''
    搜索指定info_hash的下载人数
    '''
    nodes = set(DHT_BOOTSTRAP_NODES)
    nodes.update(self.nodes)
    for i in range(DHT_MAX_RETRY_COUNT):
        for node in nodes:
            try:
                self.send_get_peers(node, info_hash)
            except Exception as e:
                print(f'[search_info_hash error]: {e}')
        time.sleep(DHT_MAX_RESPONSE_TIME)
    return len(self.nodes)

def run(self):
    '''
    启动DHT节点
    '''
    self.bootstrap()
    while True:
        try:
            message, address = self.socket.recvfrom(BUFFER_SIZE)
            threading.Thread(target=self.process_message, args=(message, address)).start()
        except Exception as e:
            pass

class DHTServer: def init(self): self.dht = DHT() self.dht_thread = threading.Thread(target=self.dht.run, args=()) self.dht_thread.start()

def search_magnet(self, magnet_link):
    '''
    查询磁力链接下载人数
    '''
    info_hash = hashlib.sha1(urllib.parse.unquote_plus(magnet_link[len(MAGNET_PREFIX):]).encode('utf-8')).digest()
    return self.dht.search_info_hash(info_hash)

def run(self):
    '''
    启动DHT服务
    '''
    server_socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
    server_socket.bind(('0.0.0.0', 8000))
    while True:
        server_socket.listen(1)
        client_socket, client_address = server_socket.accept()
        request_data = client_socket.recv(BUFFER_SIZE).decode('utf-8')
        if request_data.startswith('GET /'):
            magnet_link = request_data.split(' ')[1][1:]
            if magnet_link.startswith(MAGNET_PREFIX):
                count = self.search_magnet(magnet_link)
                response_data = f'HTTP/1.1 200 OK\r\nContent-Type: text/plain\r\nContent-Length: {len(str(count))}\r\n\r\n{count}'
            else:
                response_data = 'HTTP/1.1 400 Bad Request\r\n\r\n'
        else:
            response_data = 'HTTP/1.1 400 Bad Request\r\n\r\n'
        client_socket.send(response_data.encode('utf-8'))
        client_socket.close()

if name == 'main': dht_server = DHTServer() dht_server.run()

Python DHT 网络脚本:查询磁力链接下载人数

原文地址: https://www.cveoy.top/t/topic/nEgu 著作权归作者所有。请勿转载和采集!

免费AI点我,无需注册和登录