mirror of
https://ghfast.top/https://github.com/aeroxw/easy_tdx_max.git
synced 2026-09-12 16:54:20 +08:00
Merge pull request #2 from awayings/feature/fix0910
fix: 适配 2026-09 主站新式握手,修复 K线/市场统计返回空数据
This commit is contained in:
@@ -1,10 +1,7 @@
|
||||
from .base import BaseCommand
|
||||
from .setup import SETUP_CMD1, SETUP_CMD2, SETUP_CMD3, SETUP_COMMANDS
|
||||
from .setup import build_handshake_command
|
||||
|
||||
__all__ = [
|
||||
"BaseCommand",
|
||||
"SETUP_CMD1",
|
||||
"SETUP_CMD2",
|
||||
"SETUP_CMD3",
|
||||
"SETUP_COMMANDS",
|
||||
"build_handshake_command",
|
||||
]
|
||||
|
||||
@@ -1,15 +1,31 @@
|
||||
"""握手命令原始字节(从 pytdx/parser/setup_commands.py 移植,已在真实服务器验证)。
|
||||
"""新式握手命令(2026-09 起主站强制要求,勿回退旧三条命令)。
|
||||
|
||||
连接建立后必须按序发送三条握手命令,每条均需读取并丢弃响应。
|
||||
历史(2026-09-10 与 eltdx 逐字节比对 + 对照实测定案):
|
||||
|
||||
2026-09 起行情主站拒绝"旧版客户端握手"建立的连接:旧握手 = pytdx 三条
|
||||
固定 msg_id 的 0x000d 命令(0x1893/0x1894/0x1899,payload 01/02/签名串)。
|
||||
握手本身有响应,但连接上所有 K 线请求一律返回 2 字节空包(0x0320,声称
|
||||
800 条)、880xxx 统计指数快照返回空——服务器不报错,只是不给数据
|
||||
(部分主站 setup2 响应明示"客户端与行情主站不匹配")。web /market/stat
|
||||
500、/bars/index 空列表均源于此。
|
||||
|
||||
对照实验(三组):
|
||||
1. 旧三条命令 + 随机 msg_id → 仍被拒(拒绝标记是三条命令序列本身);
|
||||
2. 新式单条握手 + 固定 msg_id 的业务请求 → 全部正常;
|
||||
3. 新式握手连接上连续 8 个固定 msg_id 业务请求 → 全部正常。
|
||||
|
||||
新式握手 = 单条 0x000d 命令、payload 0x01、msg_id 随机(每连接新生成)。
|
||||
业务请求格式(28 字节 K 线 / 0x053e 快照等)无需任何改动。
|
||||
"""
|
||||
|
||||
from typing import Final
|
||||
import random
|
||||
import struct
|
||||
|
||||
# 从 pytdx 源码原文复制,去除空格
|
||||
SETUP_CMD1: Final[bytes] = bytes.fromhex("0c0218930001030003000d0001")
|
||||
SETUP_CMD2: Final[bytes] = bytes.fromhex("0c0218940001030003000d0002")
|
||||
SETUP_CMD3: Final[bytes] = bytes.fromhex(
|
||||
"0c031899000120002000db0fd5d0c9ccd6a4a8af0000008fc22540130000d500c9ccbdf0d7ea00000002"
|
||||
)
|
||||
|
||||
SETUP_COMMANDS: Final[tuple[bytes, ...]] = (SETUP_CMD1, SETUP_CMD2, SETUP_CMD3)
|
||||
def build_handshake_command() -> bytes:
|
||||
"""生成一条新式握手命令(msg_id 每次调用随机生成)。
|
||||
|
||||
旧三条 setup 常量(SETUP_CMD1/2/3)已删除;发送握手一律走本函数。
|
||||
"""
|
||||
msg_id = random.randint(1, 0xFFFFFFFE)
|
||||
return struct.pack("<HIHHH", 0x010C, msg_id, 0x0003, 0x0003, 0x000D) + b"\x01"
|
||||
|
||||
@@ -5,7 +5,7 @@ from types import TracebackType
|
||||
from typing import TYPE_CHECKING, TypeVar
|
||||
|
||||
from ..codec.frame import HEADER_SIZE, decompress_body, parse_header
|
||||
from ..commands.setup import SETUP_COMMANDS
|
||||
from ..commands.setup import build_handshake_command
|
||||
from ..config import get_best_host, get_port, get_timeout
|
||||
from ..exceptions import TdxConnectionError
|
||||
|
||||
@@ -122,19 +122,23 @@ class AsyncTdxConnection:
|
||||
# ------------------------------------------------------------------ #
|
||||
|
||||
async def _send_setup(self) -> None:
|
||||
"""按序发送三条握手命令并丢弃响应。"""
|
||||
"""发送一条新式握手命令并丢弃响应。
|
||||
|
||||
2026-09 起主站拒绝旧三条握手(pytdx 固定 msg_id)建立的连接:握手
|
||||
有响应但后续 K 线/快照一律返回空包。新式握手 = 单条 0x000d 命令、
|
||||
payload 0x01、随机 msg_id(见 commands/setup.py 模块注释)。
|
||||
"""
|
||||
assert self._writer is not None
|
||||
assert self._reader is not None
|
||||
for cmd_bytes in SETUP_COMMANDS:
|
||||
self._writer.write(cmd_bytes)
|
||||
await asyncio.wait_for(self._writer.drain(), timeout=self.timeout)
|
||||
try:
|
||||
hdr_buf = await self._recv_exact(HEADER_SIZE)
|
||||
hdr = parse_header(hdr_buf)
|
||||
if hdr.zipsize > 0:
|
||||
await self._recv_exact(hdr.zipsize)
|
||||
except (OSError, asyncio.TimeoutError, asyncio.IncompleteReadError):
|
||||
pass
|
||||
self._writer.write(build_handshake_command())
|
||||
await asyncio.wait_for(self._writer.drain(), timeout=self.timeout)
|
||||
try:
|
||||
hdr_buf = await self._recv_exact(HEADER_SIZE)
|
||||
hdr = parse_header(hdr_buf)
|
||||
if hdr.zipsize > 0:
|
||||
await self._recv_exact(hdr.zipsize)
|
||||
except (OSError, asyncio.TimeoutError, asyncio.IncompleteReadError):
|
||||
pass
|
||||
|
||||
async def _recv_exact(self, n: int) -> bytes:
|
||||
"""读满 n 字节。"""
|
||||
|
||||
@@ -7,7 +7,7 @@ from types import TracebackType
|
||||
from typing import TYPE_CHECKING, TypeVar
|
||||
|
||||
from ..codec.frame import HEADER_SIZE, decompress_body, parse_header
|
||||
from ..commands.setup import SETUP_COMMANDS
|
||||
from ..commands.setup import build_handshake_command
|
||||
from ..config import (
|
||||
get_best_host,
|
||||
get_calc_hosts,
|
||||
@@ -49,8 +49,8 @@ def ping_host(
|
||||
sock.settimeout(timeout)
|
||||
try:
|
||||
sock.connect((host, port))
|
||||
# 发送第一条握手命令并等待响应作为可用性验证
|
||||
sock.sendall(SETUP_COMMANDS[0])
|
||||
# 发送握手命令并等待响应作为可用性验证(新式单条握手,随机 msg_id)
|
||||
sock.sendall(build_handshake_command())
|
||||
hdr_buf = _recv_exact_sock(sock, HEADER_SIZE)
|
||||
hdr = parse_header(hdr_buf)
|
||||
if hdr.zipsize > 0:
|
||||
@@ -140,6 +140,7 @@ class TdxConnection:
|
||||
self.port = port if port is not None else get_port()
|
||||
self.timeout = timeout if timeout is not None else get_timeout()
|
||||
self._sock: socket.socket | None = None
|
||||
self._handshake_cmd: bytes = b"" # connect 时生成,心跳复用
|
||||
self._lock = threading.Lock()
|
||||
self._heartbeat_interval: float = 0 # 0 = disabled
|
||||
self._stop_event: threading.Event | None = None
|
||||
@@ -258,7 +259,7 @@ class TdxConnection:
|
||||
self._sock = None
|
||||
return
|
||||
try:
|
||||
self._sock.sendall(SETUP_COMMANDS[0])
|
||||
self._sock.sendall(self._handshake_cmd)
|
||||
hdr_buf = _recv_exact_sock(self._sock, HEADER_SIZE)
|
||||
hdr = parse_header(hdr_buf)
|
||||
if hdr.zipsize > 0:
|
||||
@@ -276,19 +277,26 @@ class TdxConnection:
|
||||
# ------------------------------------------------------------------ #
|
||||
|
||||
def _send_setup(self) -> None:
|
||||
"""按序发送三条握手命令并丢弃响应。"""
|
||||
"""发送一条新式握手命令并丢弃响应。
|
||||
|
||||
2026-09 起主站拒绝旧三条握手(pytdx 固定 msg_id)建立的连接:握手
|
||||
有响应但后续 K 线/快照一律返回空包。新式握手 = 单条 0x000d 命令、
|
||||
payload 0x01、随机 msg_id(见 commands/setup.py 模块注释)。
|
||||
"""
|
||||
assert self._sock is not None
|
||||
for cmd_bytes in SETUP_COMMANDS:
|
||||
self._sock.sendall(cmd_bytes)
|
||||
# 读取并丢弃握手响应
|
||||
try:
|
||||
hdr_buf = self._recv_exact(HEADER_SIZE)
|
||||
hdr = parse_header(hdr_buf)
|
||||
if hdr.zipsize > 0:
|
||||
self._recv_exact(hdr.zipsize)
|
||||
except OSError:
|
||||
# 部分服务器的握手无响应,忽略错误
|
||||
pass
|
||||
# 每连接生成一次握手字节:随机 msg_id 是本连接"新式客户端"的标记,
|
||||
# 心跳复用它(与旧行为对称——旧代码心跳重发同一条固定 setup 命令)。
|
||||
self._handshake_cmd = build_handshake_command()
|
||||
self._sock.sendall(self._handshake_cmd)
|
||||
# 读取并丢弃握手响应
|
||||
try:
|
||||
hdr_buf = self._recv_exact(HEADER_SIZE)
|
||||
hdr = parse_header(hdr_buf)
|
||||
if hdr.zipsize > 0:
|
||||
self._recv_exact(hdr.zipsize)
|
||||
except OSError:
|
||||
# 部分服务器的握手无响应,忽略错误
|
||||
pass
|
||||
|
||||
def _recv_exact(self, n: int) -> bytes:
|
||||
"""循环 recv 直到读满 n 字节。"""
|
||||
|
||||
@@ -8,7 +8,7 @@ import time
|
||||
|
||||
from easy_tdx import AsyncTdxClient, Market
|
||||
from easy_tdx.commands.security_count import GetSecurityCountCmd
|
||||
from easy_tdx.commands.setup import SETUP_COMMANDS
|
||||
from easy_tdx.commands.setup import build_handshake_command
|
||||
from easy_tdx.exceptions import TdxConnectionError
|
||||
|
||||
|
||||
@@ -16,15 +16,24 @@ def _pack_frame(body: bytes) -> bytes:
|
||||
return struct.pack("<IIIHH", 0, 0, 0, len(body), len(body)) + body
|
||||
|
||||
|
||||
_HANDSHAKE_LEN = len(build_handshake_command())
|
||||
|
||||
|
||||
async def _read_and_ack_handshake(
|
||||
reader: asyncio.StreamReader, writer: asyncio.StreamWriter
|
||||
) -> None:
|
||||
"""假服务器:读取一条握手命令并回一个空响应帧(新式握手,2026-09)。"""
|
||||
await reader.readexactly(_HANDSHAKE_LEN)
|
||||
writer.write(_pack_frame(b""))
|
||||
await writer.drain()
|
||||
|
||||
|
||||
def test_async_client_serializes_concurrent_calls() -> None:
|
||||
request_len = len(GetSecurityCountCmd(Market.SH).build_request())
|
||||
|
||||
async def handle(reader: asyncio.StreamReader, writer: asyncio.StreamWriter) -> None:
|
||||
try:
|
||||
for setup_cmd in SETUP_COMMANDS:
|
||||
await reader.readexactly(len(setup_cmd))
|
||||
writer.write(_pack_frame(b""))
|
||||
await writer.drain()
|
||||
await _read_and_ack_handshake(reader, writer)
|
||||
|
||||
await reader.readexactly(request_len)
|
||||
writer.write(_pack_frame(struct.pack("<H", 5)))
|
||||
@@ -68,10 +77,7 @@ def test_async_client_auto_reconnect() -> None:
|
||||
connection_ids.append(len(connection_ids) + 1)
|
||||
connection_id = connection_ids[-1]
|
||||
try:
|
||||
for setup_cmd in SETUP_COMMANDS:
|
||||
await reader.readexactly(len(setup_cmd))
|
||||
writer.write(_pack_frame(b""))
|
||||
await writer.drain()
|
||||
await _read_and_ack_handshake(reader, writer)
|
||||
|
||||
await reader.readexactly(request_len)
|
||||
writer.write(_pack_frame(struct.pack("<H", 10 + connection_id)))
|
||||
@@ -104,10 +110,7 @@ def test_async_client_request_timeout() -> None:
|
||||
|
||||
async def handle(reader: asyncio.StreamReader, writer: asyncio.StreamWriter) -> None:
|
||||
try:
|
||||
for setup_cmd in SETUP_COMMANDS:
|
||||
await reader.readexactly(len(setup_cmd))
|
||||
writer.write(_pack_frame(b""))
|
||||
await writer.drain()
|
||||
await _read_and_ack_handshake(reader, writer)
|
||||
|
||||
await reader.readexactly(request_len)
|
||||
await asyncio.sleep(1.0)
|
||||
|
||||
Reference in New Issue
Block a user