"""ZT 启动与串行调度:成交同步、买回、卖出、建仓。""" import logging as log import time from contextlib import closing from datetime import datetime from pathlib import Path from tempfile import TemporaryDirectory from concurrent.futures import Future, ThreadPoolExecutor import config from libs.calc import trading_time from libs.collector import collector_push from libs.grid_take_profit import GridTrailingTracker from libs.market import market_allow_open from libs.order import OrderBook from libs.overview import Overview from libs.runtime import Runtime from libs.signal import SignalItem, init_signals from libs.state import State from libs.watch import DipWatch from sdk import Client, DealItem, PositionItem from .open import open_signal from .positions import manage_positions, t_rounds from libs.snapshot import cache_portfolio def StartZT() -> None: client = Client( config.global_config.qmt_base_url, config.global_config.qmt_token, config.HTTP_TIMEOUT, ) executor = None try: portfolio = client.portfolio() assets = portfolio.assets positions = list(portfolio.positions.values()) state = State(Path(config.global_config.qmt_data_dir) / f'zt_{config.account_config.account_id}_state.db') run = Runtime( client=client, global_cfg=config.global_config, account_cfg=config.account_config, orders=OrderBook('zt'), open_watch=DipWatch(), add_watch=DipWatch(), profit_tracker=GridTrailingTracker(config.account_config.grid_step_pct), ) # 获取本策略的信号开仓数据 signals = init_signals( config.global_config, config.account_config.signal_allow, ) log.info( "[启动] Trend策略已启动,账户=%s,信号=%d,持仓=%d", config.account_config.account_id, len(signals), len(positions), ) deals = client.deals() cache_portfolio(config.account_config.account_id, assets, positions, deals) state.load() state.sync_deals(deals) state.sync_state(positions) state.archiving() run.orders.refresh(client, portfolio.orders) Overview(assets, positions, config.account_config) executor = ThreadPoolExecutor(max_workers=2, thread_name_prefix="zt") DEFAULT_TICK_INTERVAL = 30 while True: lt = time.localtime() if (lt.tm_hour, lt.tm_min, lt.tm_sec) >= (15, 0, 0): log.info("[ZT] 已到 15:00,结束趋势策略") return current_sec = lt.tm_sec # 计算距离下一个目标时间点(0秒或30秒)的等待时间 if current_sec < DEFAULT_TICK_INTERVAL: wait_seconds = DEFAULT_TICK_INTERVAL - current_sec elif current_sec < 60: wait_seconds = 60 - current_sec else: wait_seconds = DEFAULT_TICK_INTERVAL # 等待到目标时间点 time.sleep(wait_seconds) # 单轮失败不能杀死唯一的交易定时线程。 try: RunOnce(run, state, signals) except Exception as e: log.error( f"[ZT] 本 tick 执行失败,下一 tick 继续: {e}", exc_info=True ) finally: try: if executor is not None: executor.shutdown(wait=True) finally: client.close() def RunOnce(run: Runtime, state: State, signals: list[SignalItem]) -> None: now = datetime.now() if not trading_time(now): return print( "=" * 40 + f" Ticker {datetime.now().strftime('%Y-%m-%d %H:%M:%S')} " + "=" * 40 ) started_at = time.monotonic() futures: list[tuple[str, Future]] = [] try: deals = run.client.deals() portfolio = run.client.portfolio() assets = portfolio.assets positions = list(portfolio.positions.values()) position_codes = list(portfolio.positions) state.sync_deals(deals) state.archiving() run.orders.refresh(run.client, portfolio.orders) except Exception: log.exception("[Portfolio] 刷新账户快照失败") return # 2. 验证可用资金;低于资金安全线时禁止开新仓。 allow_open_by_cash = ( assets.available >= assets.total * run.account_cfg.min_cash_ratio ) if not allow_open_by_cash: log.info( "[Status] 禁止开仓:可用资金不足,可用=%.2f,总资产=%.2f", assets.available, assets.total, ) # 3. 获取大盘状态,只有大盘信号允许时才执行开仓。 market_ok = market_allow_open() # 4. 验证有效开仓信号:排除已有持仓。 allow_open: list[SignalItem] = [] allow_codes: list[str] = [] for signal in signals: if signal.code not in portfolio.positions: allow_open.append(signal) allow_codes.append(signal.code) if allow_open and not market_ok: log.info("[开仓] 禁止开仓:大盘信号不允许,候选=%d", len(allow_open)) # 5. 获取持仓和待开仓证券的实时行情 tick。 all_codes = list(dict.fromkeys(position_codes + allow_codes)) try: ticks = run.client.full_tick(all_codes) except Exception: log.exception("[行情] 获取行情失败,代码数量=%d", len(all_codes)) return log.info( "[RunOnce] 本轮就绪,持仓=%d,候选=%d,大盘允许=%s,资金允许=%s", len(positions), len(allow_open), market_ok, allow_open_by_cash, ) # 启动线程,开始计算 # 7. 持仓计算。 futures.append( ( "持仓计算", run.executor.submit( manage_positions, run, ticks, positions, market_ok, assets.available, state ), ) ) # 8. 开仓计算:必须同时存在有效信号且大盘允许开仓。 if allow_open and market_ok and allow_open_by_cash: futures.append( ("开仓计算", run.executor.submit(open_signal, run, ticks, allow_open)) ) # 9. 开始执行 for name, future in futures: _wait_worker(name, future) log.info( "[RunOnce] 本轮完成,耗时=%d毫秒", int((time.monotonic() - started_at) * 1000) ) def _wait_worker(name: str, future: Future) -> None: """保留单轮继续运行的语义,分别记录工作线程异常。""" try: future.result() except Exception: log.exception("[运行] %s线程失败", name)