"""
NTR Protocol Client for 3DS - v23 NUCLEAR
Based on actual NTR CFW source code + real 3DS testing
v23 NUCLEAR — THE REAL BOTTLENECK ELIMINATION:
- CRITICAL: Replaced timeout-based recv with select.select().
Before: sock.settimeout(3s) → RECV TIMEOUT → stream desync → 198 KB/s.
Now: select() waits efficiently, NO socket timeout, NO stream desync.
- CRITICAL: Lock-free pending reads in recv hot path.
Before: _pending_lock acquired for EVERY response packet (64× per batch).
Now: dict.pop(seq, None) — atomic in CPython, ZERO lock overhead.
- CRITICAL: Batch callback registration.
Before: _pending_lock acquired 64× in read_memory_batch_async.
Now: Single dict.update() — ONE lock acquisition per batch.
- CRITICAL: Process multiple packets per recv() syscall.
Before: One recv() per packet → syscall overhead × 64.
Now: Read up to 4MB, process ALL complete packets in buffer.
- CHUNK_SIZE = 1MB (0x100000). With select-based recv, no more sync loss.
- TCP_NODELAY + TCP_QUICKACK always on.
- SO_RCVBUF = 32MB, SO_SNDBUF = 8MB.
- RECV BUFFER = 8MB initial, auto-grows.
v22 MEGA FIX changes (kept where still used):
- Offset-based recv buffer (no del [:n] memmove).
- Dynamic timeout for _recv_exactly (still used for sync reads).
- CHUNK_SYNC = 0x40000 (256KB) for sync reads.
v21 ULTRA changes (kept):
- ZERO LOGGING IN HOT PATH during scans.
- SINGLE EVENT BatchResult: 1 Condition + counter.
- FAST HEADER PARSING: Single struct.unpack('<21I', header).
- SKIP ASCII DECODE FOR MEMORY READS.
- SEQUENCE INCREMENT = 1.
- FIXED: Pipeline batch tracking.
"""
import socket
import select
import struct
import threading
import time
import logging
from typing import Optional, Callable, List, Tuple
logging.basicConfig(
level=logging.WARNING,
format='[%(name)s] [%(threadName)s] %(levelname)s: %(message)s'
)
logger = logging.getLogger("NTR")
# ============ CONSTANTS ============
NTR_MAGIC = 0x12345678
NTR_HEADER_SIZE = 84
# Packet types
PKT_GENERAL = 0
PKT_GENERAL_MEMORY = 1
# Commands (from PacketCommand.cs)
CMD_HEARTBEAT = 0
CMD_HELLO = 3
CMD_RELOAD = 4
CMD_LISTPROCESSES = 5
CMD_LISTADDRESSES = 8
CMD_READMEMORY = 9
CMD_WRITEMEMORY = 10
# Data types for scanning
DATA_TYPE_U8 = 0
DATA_TYPE_U16 = 1
DATA_TYPE_U32 = 2
DATA_TYPE_U64 = 3
DATA_TYPE_F32 = 4
DATA_TYPE_F64 = 5
DATA_TYPE_ALL_INT = 6 # "Tout" - scan all integer sizes (1/2/4/8 bytes)
DATA_TYPE_SIZES = {0: 1, 1: 2, 2: 4, 3: 8, 4: 4, 5: 8, 6: 1}
DATA_TYPE_NAMES = {0: "u8", 1: "u16", 2: "u32", 3: "u64", 4: "f32", 5: "f64", 6: "Tout"}
# Pre-compiled struct format for fast header parsing
# Header: magic(4) + seq(4) + pkt_type(4) + cmd(4) + args[16](64) + data_len(4) = 84 bytes
_HEADER_STRUCT = struct.Struct('<21I')
# ============ DATA CLASSES ============
class ProcessInfo:
"""Represents a 3DS process from ListProcesses response."""
__slots__ = ('pid', 'name', 'tid')
def __init__(self, pid: int, name: str, tid: str = ""):
self.pid = pid
self.name = name.strip()
self.tid = tid.strip()
def display(self):
return f"[0x{self.pid:08X}] {self.name}"
def __repr__(self):
return f"ProcessInfo(pid=0x{self.pid:08X}, name='{self.name}', tid='{self.tid}')"
class MemoryRegion:
"""A memory region from ListAddresses response."""
__slots__ = ('start', 'size', 'end')
def __init__(self, start: int, size: int):
self.start = start
self.size = size
self.end = start + size
def display(self):
return f"0x{self.start:08X} - 0x{self.end:08X} (0x{self.size:08X})"
class BatchResult:
"""Future-like object for async batch read results - SINGLE CONDITION."""
__slots__ = ('n', 'results', '_cond', '_received', '_seqs')
def __init__(self, n: int):
self.n = n
self.results = [None] * n
self._cond = threading.Condition(threading.Lock())
self._received = 0
self._seqs = [0] * n
def set_result(self, idx: int, data):
"""Set result for index and signal completion."""
with self._cond:
self.results[idx] = data
self._received += 1
if self._received >= self.n:
self._cond.notify_all()
def wait(self, timeout=None):
"""Wait for ALL results. Returns list of results."""
deadline = time.time() + (timeout or 120.0)
with self._cond:
while self._received < self.n:
remaining = deadline - time.time()
if remaining <= 0:
break
self._cond.wait(timeout=min(remaining, 5.0))
return self.results
def wait_one(self, idx: int, timeout=None):
"""Wait for a single result by index."""
with self._cond:
while self.results[idx] is None and self._received < self.n:
self._cond.wait(timeout=5.0)
return self.results[idx]
def is_complete(self):
"""Check if all results received."""
return self._received >= self.n
def count_received(self):
"""Return number of results received so far."""
return self._received
# ============ NTR CLIENT ============
class NTRClient:
def __init__(self):
self.sock = None
self.connected = False
self.ip = ""
self.port = 8000
self.sequence = 0
self.current_pid = None
self.current_pid_hex = ""
self.memory_regions: List[MemoryRegion] = []
self._send_lock = threading.Lock()
self._running = False
self._can_send_heartbeat = True
self._last_heartbeat_time = 0
self.scan_active = False
self._last_listprocess_seq = 0
self._last_listaddr_seq = 0
self._last_readmem_seq = 0
self._listprocess_sent = False
self._listaddr_sent = False
self._waiting_listaddr = False
self._pending_reads = {}
self._pending_lock = threading.Lock()
self._active_batches = []
# v22 FIX: Offset-based recv buffer — NO MORE del [:n] memmove!
self._recv_buf = bytearray(8 * 1024 * 1024) # 8MB pre-allocated
self._recv_off = 0 # Current read offset
self._recv_len = 0 # Total valid data length
self._recv_capacity = 8 * 1024 * 1024
self.on_connected = None
self.on_disconnected = None
self.on_error = None
self.on_process_list = None
self.on_memory_regions = None
# ============ PACKET BUILDING ============
def _build_packet(self, pkt_type: int, cmd: int, args=None, data=b'') -> bytes:
if args is None:
args = []
while len(args) < 16:
args.append(0)
args = args[:16]
self.sequence += 1
self._track_sequence(pkt_type, cmd)
header = struct.pack('<I', NTR_MAGIC)
header += struct.pack('<I', self.sequence)
header += struct.pack('<I', pkt_type)
header += struct.pack('<I', cmd)
for arg in args:
header += struct.pack('<I', arg & 0xFFFFFFFF)
header += struct.pack('<I', len(data))
return header + data
def _build_read_packet(self, pid: int, address: int, size: int, seq: int = None) -> Tuple[bytes, int]:
if seq is None:
seq = self.sequence
args = [pid & 0xFFFFFFFF, address & 0xFFFFFFFF, size]
while len(args) < 16:
args.append(0)
header = struct.pack('<I', NTR_MAGIC)
header += struct.pack('<I', seq)
header += struct.pack('<I', PKT_GENERAL)
header += struct.pack('<I', CMD_READMEMORY)
for arg in args:
header += struct.pack('<I', arg & 0xFFFFFFFF)
header += struct.pack('<I', 0)
return header, seq
def _track_sequence(self, pkt_type: int, cmd: int):
if pkt_type == PKT_GENERAL and cmd == CMD_LISTPROCESSES:
self._last_listprocess_seq = self.sequence
self._listprocess_sent = True
elif pkt_type == PKT_GENERAL and cmd == CMD_LISTADDRESSES:
self._last_listaddr_seq = self.sequence
self._listaddr_sent = True
self._waiting_listaddr = True
elif pkt_type == PKT_GENERAL and cmd == CMD_READMEMORY:
self._last_readmem_seq = self.sequence
def _send_packet(self, pkt_type: int, cmd: int, args=None, data=b'') -> bool:
if not self.connected or not self.sock:
return False
with self._send_lock:
try:
packet = self._build_packet(pkt_type, cmd, args, data)
self.sock.sendall(packet)
return True
except Exception as e:
logger.error(f"SEND ERROR: {e}")
return False
# ============ RECEIVE LOOP — v23 NUCLEAR: SELECT + LOCK-FREE ============
def _recv_compact(self):
"""Compact the recv buffer by moving data to the front."""
if self._recv_off > 0:
data_len = self._recv_len - self._recv_off
if data_len > 0:
self._recv_buf[:data_len] = self._recv_buf[self._recv_off:self._recv_len]
self._recv_off = 0
self._recv_len = data_len
def _recv_exactly(self, n: int, timeout=5.0) -> Optional[bytes]:
"""Read exactly n bytes from socket using offset-based buffer.
Only used for sync reads (hex viewer, etc.) — NOT in the hot path.
v23: Uses select instead of socket timeout to prevent stream desync."""
available = self._recv_len - self._recv_off
if available >= n:
result = bytes(self._recv_buf[self._recv_off:self._recv_off + n])
self._recv_off += n
if self._recv_off > 4 * 1024 * 1024:
self._recv_compact()
return result
# Dynamic timeout for large reads
if n > 100_000:
dynamic_timeout = max(30.0, n / 50000.0 * 2.0)
else:
dynamic_timeout = timeout
deadline = time.time() + dynamic_timeout
while self._recv_len - self._recv_off < n:
remaining = deadline - time.time()
if remaining <= 0:
return None
try:
# Use select to wait for data — NO socket timeout!
try:
readable, _, _ = select.select([self.sock], [], [], min(remaining, 5.0))
except (OSError, ValueError):
return None
if not readable:
continue # select timed out, check deadline
write_start = self._recv_len
space = self._recv_capacity - write_start
if space < 65536:
self._recv_compact()
write_start = self._recv_len
space = self._recv_capacity - write_start
if space < 65536:
new_size = self._recv_capacity + 4 * 1024 * 1024
new_buf = bytearray(new_size)
data_len = self._recv_len - self._recv_off
if data_len > 0:
new_buf[:data_len] = self._recv_buf[self._recv_off:self._recv_len]
self._recv_buf = new_buf
self._recv_off = 0
self._recv_len = data_len
self._recv_capacity = new_size
write_start = data_len
space = self._recv_capacity - write_start
to_read = min(space, max(n - (self._recv_len - self._recv_off), 262144))
chunk = self.sock.recv(to_read)
if not chunk:
return None
self._recv_buf[write_start:write_start + len(chunk)] = chunk
self._recv_len += len(chunk)
except OSError:
return None
result = bytes(self._recv_buf[self._recv_off:self._recv_off + n])
self._recv_off += n
if self._recv_off > 4 * 1024 * 1024:
self._recv_compact()
return result
def _recv_loop(self):
"""v23 NUCLEAR: select-based recv loop.
NO socket timeout → NO RECV TIMEOUT → NO stream desync.
Reads up to 4MB per syscall, processes ALL complete packets in buffer.
Lock-free pending read lookup via dict.pop() — atomic in CPython."""
logger.info("RECV THREAD: started (v23 NUCLEAR)")
while self._running and self.connected:
try:
# Wait for data with select — NO socket timeout!
try:
readable, _, _ = select.select([self.sock], [], [], 30.0)
except (OSError, ValueError):
if self._running and self.connected:
continue
break
if not readable:
# 30s with no data — just continue waiting
if not self._running or not self.connected:
break
continue
# Read as much data as possible into the buffer
self._recv_compact()
write_start = self._recv_len
space = self._recv_capacity - write_start
if space < 65536:
new_size = self._recv_capacity + 4 * 1024 * 1024
new_buf = bytearray(new_size)
data_len = self._recv_len - self._recv_off
if data_len > 0:
new_buf[:data_len] = self._recv_buf[self._recv_off:self._recv_len]
self._recv_buf = new_buf
self._recv_off = 0
self._recv_len = data_len
self._recv_capacity = new_size
write_start = data_len
space = self._recv_capacity - write_start
try:
chunk = self.sock.recv(min(space, 4 * 1024 * 1024)) # Read up to 4MB!
if not chunk:
# Connection closed by remote
break
self._recv_buf[write_start:write_start + len(chunk)] = chunk
self._recv_len += len(chunk)
except OSError:
break
# Process ALL complete packets in the buffer
self._process_buffer()
except Exception as e:
if self._running and self.connected:
logger.error(f"RECV ERROR: {e}")
self._handle_disconnect()
break
logger.info("RECV THREAD: stopped")
def _process_buffer(self):
"""Process all complete NTR packets in the recv buffer.
v23: Processes multiple packets per recv() call — huge throughput gain."""
while self._recv_len - self._recv_off >= NTR_HEADER_SIZE:
# Peek at header without consuming
hdr_view = self._recv_buf[self._recv_off:self._recv_off + NTR_HEADER_SIZE]
try:
fields = _HEADER_STRUCT.unpack(hdr_view)
except struct.error:
self._recv_off += 1 # Skip bad byte
continue
magic = fields[0]
seq = fields[1]
pkt_type = fields[2]
cmd = fields[3]
args = list(fields[4:20])
data_len = fields[20]
if magic != NTR_MAGIC:
# Try to find next valid magic in buffer (fast resync)
search_end = min(self._recv_len - 3, self._recv_off + 65536)
found = -1
for i in range(self._recv_off + 1, search_end):
if (self._recv_buf[i] == 0x78 and
self._recv_buf[i+1] == 0x56 and
self._recv_buf[i+2] == 0x34 and
self._recv_buf[i+3] == 0x12):
found = i
break
if found > 0:
# Try to handle text before the found magic
text_data = self._recv_buf[self._recv_off:found]
try:
text = bytes(text_data).decode('ascii', errors='replace')
if 'finished' in text.lower():
if not self.scan_active:
logger.info("RECV: write finished")
except:
pass
self._recv_off = found
else:
self._recv_off += 1 # Skip one byte, try again
continue
total_len = NTR_HEADER_SIZE + data_len
# Check if we have the full packet
if self._recv_len - self._recv_off < total_len:
break # Need more data from socket
# Extract data payload
if data_len > 0 and data_len < 64 * 1024 * 1024:
data = bytes(self._recv_buf[self._recv_off + NTR_HEADER_SIZE:self._recv_off + total_len])
else:
data = b''
# Advance offset past this packet
self._recv_off += total_len
# Handle the response — LOCK-FREE pending read lookup!
self._handle_response(cmd, seq, pkt_type, args, data)
# Compact buffer if offset is large
if self._recv_off > 2 * 1024 * 1024:
self._recv_compact()
def _handle_response(self, cmd: int, seq: int, pkt_type: int, args: list, data: bytes):
"""v23 NUCLEAR: Lock-free pending read handling.
Uses dict.pop() which is atomic in CPython — NO _pending_lock needed!"""
# Check pending reads FIRST (hot path for scans) — LOCK-FREE!
cb = self._pending_reads.pop(seq, None)
if cb is not None:
try:
cb(cmd, pkt_type, args, data)
except Exception:
pass
return
# Non-read responses (heartbeat, process list, etc.)
if not data:
if cmd == CMD_HEARTBEAT:
self._can_send_heartbeat = True
return
# Only decode ASCII for non-read commands
try:
text = data.decode('ascii', errors='replace')
except:
text = ""
text_cleaned = text
for garbage in ['rtRecvSocket failed: 00000000', 'broken protocol:']:
text_cleaned = text_cleaned.replace(garbage, '')
if 'openProcess failed' in text or 'openFile failed' in text:
if not self.scan_active:
logger.warning(f"ListAddresses FAILED: {text.strip()}")
self._waiting_listaddr = False
if self.on_error:
self.on_error(f"Impossible d'ouvrir ce processus (erreur 3DS). "
f"Essaie un processus de jeu.")
if 'pid:' in text:
self._parse_process_list(data)
return
has_pid_lines = 'pid:' in text_cleaned
has_size_lines = 'size:' in text_cleaned
if has_pid_lines and not has_size_lines:
if not self.scan_active:
logger.info("RECV: Process list detected")
self._parse_process_list(data)
self._listprocess_sent = False
return
if has_size_lines and ' - ' in text_cleaned:
if not self.scan_active:
logger.info("RECV: Address list detected")
self._parse_address_list(data)
self._waiting_listaddr = False
return
if cmd == CMD_HEARTBEAT:
self._can_send_heartbeat = True
if has_pid_lines:
self._parse_process_list(data)
return
if self._last_readmem_seq > 0 and seq == self._last_readmem_seq:
return
# ============ RESPONSE PARSERS ============
def _parse_process_list(self, data: bytes):
processes = []
seen_pids = set()
try:
text = data.decode('ascii', errors='replace')
for line in text.split('\n'):
line = line.strip()
if line.startswith('end of process list'):
break
if not line.startswith('pid'):
continue
parts = line.split(',')
if len(parts) < 2:
continue
try:
pid_str = parts[0].split(': ')
if len(pid_str) < 2:
continue
pid = int(pid_str[1].strip(), 16)
pname_parts = parts[1].split(': ')
if len(pname_parts) < 2:
continue
pname = pname_parts[1].strip()
tid = ""
if len(parts) >= 3:
tid_parts = parts[2].split(': ')
if len(tid_parts) >= 2:
tid = tid_parts[1].strip()
if pid not in seen_pids:
seen_pids.add(pid)
processes.append(ProcessInfo(pid, pname, tid))
except (IndexError, ValueError):
continue
except Exception as e:
logger.error(f"Process list parse error: {e}")
if not self.scan_active:
logger.info(f"Parsed {len(processes)} unique processes")
if processes and self.on_process_list:
self.on_process_list(processes)
def _parse_address_list(self, data: bytes):
regions = []
try:
text = data.decode('ascii', errors='replace')
for line in text.split('\n'):
if 'size' not in line:
continue
try:
parts = line.split(' , ')
if len(parts) < 2:
continue
addr_parts = parts[0].split(' - ')
if len(addr_parts) < 2:
continue
start_str = addr_parts[0].strip().zfill(8)
size_parts = parts[1].split(': ')
if len(size_parts) < 2:
continue
size_str = size_parts[1].strip().zfill(8)
start = int(start_str, 16)
size = int(size_str, 16)
if start > 0 and size > 0:
regions.append(MemoryRegion(start, size))
except (IndexError, ValueError):
continue
except Exception as e:
logger.error(f"Address list parse error: {e}")
if not self.scan_active:
logger.info(f"Parsed {len(regions)} memory regions")
self.memory_regions = regions
if self.on_memory_regions:
self.on_memory_regions(regions)
# ============ PUBLIC API ============
@staticmethod
def test_connection(ip, port=8000, timeout=3.0):
try:
s = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
s.settimeout(timeout)
s.connect((ip, port))
s.close()
return (True, f"Connexion TCP OK sur {ip}:{port}")
except socket.timeout:
return (False, f"Timeout: {ip} ne repond pas.\nVerifie que BootNTR tourne !")
except ConnectionRefusedError:
return (False, f"Connexion refusee.\nBootNTR pas lance.")
except Exception as e:
return (False, f"Erreur: {e}")
def connect(self, ip, port=8000):
logger.info(f"Connecting to {ip}:{port}...")
try:
self.sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
self.sock.settimeout(5) # Only for connect, recv uses select
self.sock.connect((ip, port))
self.ip = ip
self.port = port
self.connected = True
self._running = True
self.sequence = 0
self._can_send_heartbeat = True
self.scan_active = False
# Reset offset-based buffer
self._recv_off = 0
self._recv_len = 0
# Maximum socket buffers
self.sock.setsockopt(socket.IPPROTO_TCP, socket.TCP_NODELAY, 1)
self.sock.setsockopt(socket.SOL_SOCKET, socket.SO_RCVBUF, 32 * 1024 * 1024)
self.sock.setsockopt(socket.SOL_SOCKET, socket.SO_SNDBUF, 8 * 1024 * 1024)
# v23: Remove socket timeout — recv loop uses select.select() instead!
self.sock.settimeout(None) # Blocking mode, select handles timeouts
# TCP_QUICKACK on Linux
try:
self.sock.setsockopt(socket.IPPROTO_TCP, socket.TCP_QUICKACK, 1)
except (OSError, AttributeError):
pass
self._recv_thread = threading.Thread(target=self._recv_loop, daemon=True, name="NTR-Recv")
self._recv_thread.start()
time.sleep(0.1)
self.heartbeat()
time.sleep(1.5)
if self.on_connected:
self.on_connected()
print(f"[NTR] Connecte a {ip}:{port} — v23 NUCLEAR")
return True
except Exception as e:
self.connected = False
logger.error(f"Connect error: {e}")
if self.on_error:
self.on_error(str(e))
return False
def disconnect(self):
was = self.connected
self._running = False
self.connected = False
self.scan_active = False
try:
if self.sock:
self.sock.close()
except:
pass
self.sock = None
self.current_pid = None
self.current_pid_hex = ""
self.memory_regions = []
self._recv_off = 0
self._recv_len = 0
if was and self.on_disconnected:
self.on_disconnected()
def _handle_disconnect(self):
was = self.connected
self.connected = False
self._running = False
self.scan_active = False
try:
if self.sock:
self.sock.close()
except:
pass
self.sock = None
self.current_pid = None
self.current_pid_hex = ""
self.memory_regions = []
self._recv_off = 0
self._recv_len = 0
for batch in self._active_batches:
with batch._cond:
batch._cond.notify_all()
self._active_batches.clear()
if was and self.on_disconnected:
self.on_disconnected()
def heartbeat(self) -> bool:
if not self.connected:
return False
if self.scan_active:
return False
now = time.time() * 1000
if not self._can_send_heartbeat and (now - self._last_heartbeat_time) < 30000:
return False
self._can_send_heartbeat = False
self._last_heartbeat_time = now
return self._send_packet(PKT_GENERAL, CMD_HEARTBEAT)
def send_hello(self) -> bool:
return self._send_packet(PKT_GENERAL, CMD_HELLO)
def send_reload(self) -> bool:
return self._send_packet(PKT_GENERAL, CMD_RELOAD)
def list_processes(self) -> bool:
logger.info("Sending ListProcesses...")
return self._send_packet(PKT_GENERAL, CMD_LISTPROCESSES)
def list_memory_addresses(self, pid: int) -> bool:
logger.info(f"Sending ListAddresses for PID 0x{pid:08X}")
self.current_pid = pid
self.current_pid_hex = f"{pid:08X}"
return self._send_packet(PKT_GENERAL, CMD_LISTADDRESSES, args=[pid & 0xFFFFFFFF])
def read_memory(self, pid: int, address: int, size: int) -> bool:
return self._send_packet(PKT_GENERAL, CMD_READMEMORY,
args=[pid & 0xFFFFFFFF, address & 0xFFFFFFFF, size])
def read_memory_sync(self, address: int, size: int, timeout=5.0) -> Optional[bytes]:
if self.current_pid is None:
return None
CHUNK = 0x40000 # 256KB chunks for sync reads
result_data = bytearray()
remaining = size
addr = address
while remaining > 0:
chunk_size = min(remaining, CHUNK)
data = self._read_chunk_sync(self.current_pid, addr, chunk_size, timeout)
if data is None:
return bytes(result_data) if result_data else None
result_data.extend(data)
addr += chunk_size
remaining -= chunk_size
return bytes(result_data)
def read_memory_batch(self, reads: List[Tuple[int, int, int]], timeout=None) -> List[Optional[bytes]]:
"""Backward-compatible sync batch read."""
batch = self.read_memory_batch_async(reads)
if batch is None:
return []
result = batch.wait()
self._cleanup_batch(batch)
return result
def read_memory_batch_async(self, reads: List[Tuple[int, int, int]]) -> Optional[BatchResult]:
"""v23 NUCLEAR: Batch callback registration — ONE lock per batch, not N.
Registers all callbacks with a single dict.update() call, then sends."""
if not reads:
return None
n = len(reads)
batch = BatchResult(n)
with self._send_lock:
all_packets = bytearray(n * NTR_HEADER_SIZE)
pos = 0
callbacks = {} # Build all callbacks FIRST, then register in ONE operation
for i, (pid, address, size) in enumerate(reads):
self.sequence += 1
seq = self.sequence
batch._seqs[i] = seq
# Create callback — no lock needed during creation
def make_cb(idx):
def cb(cmd, pkt_type, args, data):
batch.set_result(idx, data if data else None)
return cb
callbacks[seq] = make_cb(i)
# Build read packet directly
packet, _ = self._build_read_packet(pid, address, size, seq=seq)
all_packets[pos:pos + len(packet)] = packet
pos += len(packet)
# Register ALL callbacks in ONE operation — single lock acquisition!
with self._pending_lock:
self._pending_reads.update(callbacks)
# Send all packets
try:
self.sock.sendall(bytes(all_packets[:pos]))
except Exception as e:
logger.error(f"BATCH SEND ERROR: {e}")
# Unregister on failure
for seq in batch._seqs:
self._pending_reads.pop(seq, None)
with batch._cond:
batch._cond.notify_all()
return None
self._active_batches.append(batch)
return batch
def _cleanup_batch(self, batch: BatchResult):
if batch in self._active_batches:
self._active_batches.remove(batch)
# Lock-free cleanup — just pop each seq, doesn't matter if already gone
for seq in batch._seqs:
self._pending_reads.pop(seq, None)
def _read_chunk_sync(self, pid: int, address: int, size: int, timeout=5.0) -> Optional[bytes]:
result = {'data': None, 'event': threading.Event()}
def cb(cmd, pkt_type, args, data):
if data:
result['data'] = data
result['event'].set()
self.sequence += 1
expected_seq = self.sequence
with self._pending_lock:
self._pending_reads[expected_seq] = cb
if not self.read_memory(pid, address, size):
with self._pending_lock:
self._pending_reads.pop(expected_seq, None)
return None
if not result['event'].wait(timeout):
with self._pending_lock:
self._pending_reads.pop(expected_seq, None)
return None
return result['data']
def write_memory(self, pid: int, address: int, data: bytes) -> bool:
return self._send_packet(PKT_GENERAL_MEMORY, CMD_WRITEMEMORY,
args=[pid & 0xFFFFFFFF, address & 0xFFFFFFFF, len(data)],
data=data)
def write_u32(self, address: int, value: int):
if self.current_pid is None: return False
return self.write_memory(self.current_pid, address, struct.pack('<I', value & 0xFFFFFFFF))
def write_u16(self, address: int, value: int):
if self.current_pid is None: return False
return self.write_memory(self.current_pid, address, struct.pack('<H', value & 0xFFFF))
def write_u8(self, address: int, value: int):
if self.current_pid is None: return False
return self.write_memory(self.current_pid, address, struct.pack('<B', value & 0xFF))
def write_f32(self, address: int, value: float):
if self.current_pid is None: return False
return self.write_memory(self.current_pid, address, struct.pack('<f', value))
def is_process_attached(self) -> bool:
return self.current_pid is not None
def is_valid_address(self, address: int) -> bool:
for region in self.memory_regions:
if region.start <= address < region.end:
return True
return False
@staticmethod
def parse_value(data, data_type, offset=0):
try:
fmts = {0: '<B', 1: '<H', 2: '<I', 3: '<Q', 4: '<f', 5: '<d'}
if data_type in fmts:
return struct.unpack_from(fmts[data_type], data, offset)[0]
return None
except:
return None
@staticmethod
def pack_value(value, data_type):
try:
fmts = {0: '<B', 1: '<H', 2: '<I', 3: '<Q', 4: '<f', 5: '<d'}
masks = {0: 0xFF, 1: 0xFFFF, 2: 0xFFFFFFFF, 3: 0xFFFFFFFFFFFFFFFF}
if data_type in masks:
return struct.pack(fmts[data_type], int(value) & masks[data_type])
else:
return struct.pack(fmts[data_type], float(value))
except:
return b''5 views