Python DHT 网络脚本:查询磁力链接下载人数
-- 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()
原文地址: https://www.cveoy.top/t/topic/nEgu 著作权归作者所有。请勿转载和采集!