TriloopTem_App/tools/mock_device.py
2026-06-19 22:08:10 +08:00

371 lines
14 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

#!/usr/bin/env python3
"""
TEM 下位机模拟器 — 用于 App 调试
用法:
python tools/mock_device.py
python tools/mock_device.py --host 0.0.0.0 --port 4321
python tools/mock_device.py --channels 3 --freq 2 # 3通道, 2Hz
在 App 连接界面将 IP 改为本机局域网 IP或 10.0.2.2 若在 Android 模拟器中)。
"""
import socket
import struct
import threading
import time
import math
import random
import argparse
import sys
from datetime import datetime
# ── 协议常量 ──────────────────────────────────────────────────────────────────
MAGIC = b'\x68\x68\xff\xff'
FRAME_FLAG = 0xFE
class FC:
SETUP_REQ = 0x01
CONTINUOUS_REQ = 0x02
SINGLE_REQ = 0x03
STOP_REQ = 0x04
ACTIVE_REQ = 0x05
SPLITFRAME_REQ = 0x08
SETUP_ACK = 0x81
CONTINUOUS_ACK = 0x82
SINGLE_ACK = 0x83
STOP_ACK = 0x84
DATA_ACK = 0x85
_FC_NAME = {
0x01: 'SETUP_REQ', 0x02: 'CONTINUOUS_REQ', 0x03: 'SINGLE_REQ',
0x04: 'STOP_REQ', 0x05: 'ACTIVE_REQ', 0x08: 'SPLITFRAME_REQ',
0x81: 'SETUP_ACK', 0x82: 'CONTINUOUS_ACK', 0x83: 'SINGLE_ACK',
0x84: 'STOP_ACK', 0x85: 'DATA_ACK',
}
# sendFreq code → Hz
SEND_FREQ_HZ = {
0x00: 0.5, 0x01: 1.0, 0x02: 2.0, 0x03: 4.0,
0x04: 8.0, 0x05: 12.5, 0x06: 16.0, 0x07: 25.0,
0x08: 32.0, 0x09: 50.0, 0x0a: 64.0,
}
AMP_GAIN = [0.125, 0.25, 0.5, 1.0, 2.0, 4.0, 8.0, 16.0, 32.0, 64.0, 128.0, 1.0]
# ── 帧构建 ──────────────────────────────────────────────────────────────────
def build_frame(func: int, payload: bytes) -> bytes:
header = MAGIC + struct.pack('<BBI', FRAME_FLAG, func, len(payload))
return header + payload
def build_ack(func_ack: int, ok: bool = True, reason: str = '') -> bytes:
payload = bytearray(54)
payload[0] = 0x01 if ok else 0x02
if reason:
encoded = reason.encode('ascii', errors='replace')[:31]
payload[22:22 + len(encoded)] = encoded
return build_frame(func_ack, bytes(payload))
def build_data_frame(cfg: dict, phase: float) -> bytes:
"""
构建 DATA_ACK 帧,包含 54 字节元数据 + channelNum*sampleDepth 个 int32 ADC 采样。
phase: 当前帧的相位偏移(秒),用于生成连续正弦波形。
"""
ch_num = cfg['channelNum']
depth = cfg['sampleDepth']
amp_idx = cfg['ampRatio']
acc_num = max(cfg['accNum'], 1)
gain = AMP_GAIN[amp_idx] if amp_idx < len(AMP_GAIN) else 1.0
# 54 字节元数据(与 parser.ts parseMetadata 对齐)
meta = bytearray(54)
meta[0] = 1 # devId
struct.pack_into('<I', meta, 1, int(time.time())) # utc
struct.pack_into('<d', meta, 5, 116.3972) # longitude北京
struct.pack_into('<d', meta, 13, 39.9086) # latitude
struct.pack_into('<f', meta, 21, 50.0) # altitude
struct.pack_into('<f', meta, 25, 1.0) # height
meta[29] = 0x11 # sdStatus=1, gpsStatus=1
meta[30] = amp_idx & 0xFF # ampRatio
struct.pack_into('<f', meta, 31, 0.0) # roll
struct.pack_into('<f', meta, 35, 2.5) # pitch模拟轻微倾斜
struct.pack_into('<f', meta, 39, 0.0) # yaw
meta[43] = ch_num & 0xFF # channelNum
struct.pack_into('<I', meta, 44, 1000) # current
struct.pack_into('<H', meta, 48, 2100) # temperatureraw ADC ~25°C
struct.pack_into('<H', meta, 50, 3760) # batteryVoltraw ~3.7V
meta[52] = cfg.get('sourceMode', 0) & 0xFF
# ADC 采样数据:每通道不同频率的正弦信号
SIGNAL_UV = 5e5 # 信号幅度 0.5V 换算为 μV
NOISE_RATIO = 0.02 # 2% 噪声
# 各通道信号频率(单位 Hz相对于 sampleDepth 归一化)
CH_FREQS = [1.0, 2.0, 3.0, 5.0, 7.0, 11.0]
adc_bytes = bytearray(ch_num * depth * 4)
for ch in range(ch_num):
freq = CH_FREQS[ch % len(CH_FREQS)]
for s in range(depth):
t_norm = s / depth + phase * freq
uv = SIGNAL_UV * math.sin(2 * math.pi * t_norm)
uv += random.gauss(0, SIGNAL_UV * NOISE_RATIO)
raw = int(uv * gain * acc_num)
raw = max(-0x80000000, min(0x7FFFFFFF, raw))
struct.pack_into('<i', adc_bytes, (ch * depth + s) * 4, raw)
return build_frame(FC.DATA_ACK, bytes(meta) + bytes(adc_bytes))
# ── TCP 帧解析器 ────────────────────────────────────────────────────────────
class FrameParser:
def __init__(self):
self.buf = bytearray()
def feed(self, data: bytes):
self.buf += data
frames = []
while len(self.buf) >= 10:
# 搜索帧头
idx = -1
for i in range(len(self.buf) - 3):
if self.buf[i:i+4] == b'\x68\x68\xff\xff':
idx = i
break
if idx == -1:
self.buf = bytearray(self.buf[-3:]) if len(self.buf) >= 3 else bytearray()
break
if idx > 0:
self.buf = self.buf[idx:]
if len(self.buf) < 10:
break
payload_len = struct.unpack_from('<I', self.buf, 6)[0]
if payload_len > 400 * 1024: # 异常长度,跳过此帧头
self.buf = self.buf[4:]
continue
total = 10 + payload_len
if len(self.buf) < total:
break
func = self.buf[5]
payload = bytes(self.buf[10:total])
frames.append((func, payload))
self.buf = self.buf[total:]
return frames
# ── 客户端会话 ──────────────────────────────────────────────────────────────
class Session:
def __init__(self, sock: socket.socket, addr, cli_cfg: dict):
self.sock = sock
self.addr = addr
self.parser = FrameParser()
self.phase = 0.0 # 当前正弦相位(秒)
self.running = False # 连续采集中?
self._lock = threading.Lock()
# 当前设备配置(可被 SETUP_REQ 覆盖)
self.cfg = {
'channelNum': cli_cfg.get('channels', 6),
'sendFreq': cli_cfg.get('freq_code', 0x01),
'sampleFreq': 0x03,
'sampleDepth': cli_cfg.get('depth', 256),
'accNum': 1,
'ampRatio': 3, # 1×
'sourceMode': 0x00,
}
def _ts(self):
return datetime.now().strftime('%H:%M:%S.%f')[:-3]
def log(self, msg: str, direction: str = ' '):
tag = f'{direction} [{self.addr[0]}:{self.addr[1]}]'
print(f'[{self._ts()}] {tag} {msg}')
def send(self, data: bytes) -> bool:
try:
self.sock.sendall(data)
return True
except OSError:
return False
# ── 连续发送线程 ──────────────────────────────────────────────────────
def _data_loop(self):
freq_hz = SEND_FREQ_HZ.get(self.cfg['sendFreq'], 1.0)
interval = 1.0 / freq_hz
self.log(f'开始连续推送 @ {freq_hz} Hz间隔 {interval*1000:.0f} ms', '')
while self.running:
t0 = time.monotonic()
frame = build_data_frame(self.cfg, self.phase)
if not self.send(frame):
break
payload_size = len(frame) - 10
self.log(f'DATA_ACK ch={self.cfg["channelNum"]} '
f'depth={self.cfg["sampleDepth"]} '
f'payload={payload_size}B', '')
self.phase += interval
elapsed = time.monotonic() - t0
sleep_t = interval - elapsed
if sleep_t > 0:
time.sleep(sleep_t)
self.log('连续推送已停止', '')
def start_continuous(self):
if self.running:
return
self.running = True
t = threading.Thread(target=self._data_loop, daemon=True)
t.start()
def stop_continuous(self):
self.running = False
# ── 帧处理 ────────────────────────────────────────────────────────────
def _parse_setup(self, payload: bytes):
if len(payload) < 14:
return
self.cfg['channelNum'] = max(1, min(6, payload[0]))
self.cfg['sendFreq'] = payload[1]
self.cfg['sampleFreq'] = payload[2]
depth = struct.unpack_from('<H', payload, 3)[0]
self.cfg['sampleDepth'] = depth if depth > 0 else 256
acc_flags = struct.unpack_from('<H', payload, 5)[0]
acc_num = (acc_flags >> 2) & 0x3FFF
self.cfg['accNum'] = acc_num if acc_num > 0 else 1
self.cfg['ampRatio'] = payload[7]
self.cfg['sourceMode'] = payload[13] if len(payload) > 13 else 0
prefix = ''
if len(payload) >= 54:
raw = payload[38:54]
prefix = raw.split(b'\x00', 1)[0].decode('ascii', errors='replace')
freq_hz = SEND_FREQ_HZ.get(self.cfg['sendFreq'], 1.0)
self.log(
f'SETUP ch={self.cfg["channelNum"]} '
f'freq={freq_hz}Hz depth={self.cfg["sampleDepth"]} '
f'acc={self.cfg["accNum"]} amp=0x{self.cfg["ampRatio"]:02X} '
f'prefix="{prefix}"', ''
)
def handle(self, func: int, payload: bytes):
name = _FC_NAME.get(func, f'0x{func:02X}')
if func == FC.SETUP_REQ:
self._parse_setup(payload)
self.send(build_ack(FC.SETUP_ACK))
self.log('SETUP_ACK OK', '')
elif func == FC.CONTINUOUS_REQ:
self.log('CONTINUOUS_REQ', '')
self.send(build_ack(FC.CONTINUOUS_ACK))
self.log('CONTINUOUS_ACK OK', '')
self.start_continuous()
elif func == FC.SINGLE_REQ:
self.log('SINGLE_REQ', '')
self.send(build_ack(FC.SINGLE_ACK))
self.log('SINGLE_ACK OK', '')
frame = build_data_frame(self.cfg, self.phase)
self.send(frame)
self.log(f'DATA_ACK (单次) {len(frame)-10}B', '')
self.phase += 1.0
elif func == FC.STOP_REQ:
self.log('STOP_REQ', '')
self.stop_continuous()
self.send(build_ack(FC.STOP_ACK))
self.log('STOP_ACK OK', '')
elif func == FC.ACTIVE_REQ:
pass # keepalive无需回应
elif func == FC.SPLITFRAME_REQ:
self.log(f'SPLITFRAME_REQ未实现忽略', '')
else:
self.log(f'未知帧 0x{func:02X} ({len(payload)}B payload)', '')
# ── 主循环 ────────────────────────────────────────────────────────────
def run(self):
self.log('已连接', '')
try:
while True:
data = self.sock.recv(4096)
if not data:
break
for func, payload in self.parser.feed(data):
self.handle(func, payload)
except OSError:
pass
finally:
self.stop_continuous()
self.sock.close()
self.log('已断开', '')
# ── 服务器主入口 ────────────────────────────────────────────────────────────
def main():
parser = argparse.ArgumentParser(
description='TEM 下位机模拟器',
formatter_class=argparse.RawDescriptionHelpFormatter,
epilog="""
连接方式:
Android 模拟器 → App IP 填 10.0.2.2
真机(同 WiFi → App IP 填本机局域网 IP运行 ipconfig 查看)
Expo Web/桌面 → App IP 填 127.0.0.1
""",
)
parser.add_argument('--host', default='0.0.0.0', help='监听地址 (默认: 0.0.0.0)')
parser.add_argument('--port', type=int, default=4321, help='监听端口 (默认: 4321)')
parser.add_argument('--channels', type=int, default=6, help='默认通道数 1-6 (默认: 6)')
parser.add_argument('--freq', type=float, default=1.0, help='默认发送频率 Hz (默认: 1.0)')
parser.add_argument('--depth', type=int, default=256, help='默认采样深度 (默认: 256)')
args = parser.parse_args()
# 找到最接近的 freq code
freq_code = min(SEND_FREQ_HZ, key=lambda k: abs(SEND_FREQ_HZ[k] - args.freq))
actual_freq = SEND_FREQ_HZ[freq_code]
cli_cfg = {
'channels': min(6, max(1, args.channels)),
'freq_code': freq_code,
'depth': max(16, args.depth),
}
srv = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
srv.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
srv.bind((args.host, args.port))
srv.listen(5)
print(f'\n TEM 下位机模拟器')
print(f' ─────────────────────────────────────────')
print(f' 监听地址 : {args.host}:{args.port}')
print(f' 默认通道数: {cli_cfg["channels"]}')
print(f' 默认频率 : {actual_freq} Hz')
print(f' 默认采样深: {cli_cfg["depth"]} pts')
print(f' ─────────────────────────────────────────')
print(f' 在 App 连接界面将 IP 改为本机局域网 IP')
print(f' Android 模拟器请填 10.0.2.2')
print(f' Ctrl+C 停止\n')
try:
while True:
conn, addr = srv.accept()
conn.setsockopt(socket.IPPROTO_TCP, socket.TCP_NODELAY, 1)
session = Session(conn, addr, cli_cfg)
t = threading.Thread(target=session.run, daemon=True)
t.start()
except KeyboardInterrupt:
print('\n模拟器已停止')
finally:
srv.close()
if __name__ == '__main__':
main()