Alpaca One-Shot 异步价格监控系统

本卡片记录专为日内标的(如 QQQ 等美股/ETF)价格监控设计的高可用、异步(Async/Await)One-Shot 自动销毁监控系统的 app.py 源码、环境参数、部署命令及 API 调用限制。


一、 核心功能设计 (One-Shot & 自毁机制)

  1. 单次触发即停 (One-Shot):满足条件(低于跌破价或高于涨破价)后,瞬间向飞书发送警报通知,并立即将任务激活状态重置为 is_active = False 退出监控,防止价格在边界线反复摩擦带来消息轰炸。
  2. 当日开盘有效 (单日自毁):通过美东时间(America/New_York)时区感知,检测到已收盘(常规 16:00,或含盘前盘后 20:00)或跨天变动时,自动将监控挂起并写入磁盘,收市后自动释放资源,对后续无交易日静默。
  3. IO 并发锁与原子写入:在修改状态写入磁盘时,引入 asyncio.Lock() 协程锁,且写入过程通过临时文件写入再调用 os.replace 原子替换,规避并发写盘冲突和突然断电导致的数据损耗崩溃。
  4. 接口安全与只读设计:系统不具备任何交易、开仓、平仓函数,并配置 X-Hermes-Token 进行安全头鉴权,防范未授权的第三方越权篡改。

二、 核心监控服务源码 (app.py)

import os
import json
import logging
import asyncio
from datetime import datetime
from zoneinfo import ZoneInfo
import httpx
import aiofiles
from fastapi import FastAPI, HTTPException, Depends
from fastapi.security.api_key import APIKeyHeader
from pydantic import BaseModel
from contextlib import asynccontextmanager
 
# ================= 1. 系统初始化与并发锁 =================
logging.basicConfig(level=logging.INFO, format='%(asctime)s [%(levelname)s] %(message)s', datefmt='%Y-%m-%d %H:%M:%S')
logger = logging.getLogger(__name__)
 
STATE_FILE = "monitor_state.json"
STATE_FILE_TMP = f"{STATE_FILE}.tmp"
NY_TZ = ZoneInfo("America/New_York")
API_KEY = os.getenv("HERMES_MONITOR_TOKEN", "DefaultSecretToken123")
api_key_header = APIKeyHeader(name="X-Hermes-Token", auto_error=False)
 
state_lock = asyncio.Lock()
 
monitor_state = {
    "symbol": "QQQ",
    "target_price": None,
    "direction": "BELOW",
    "extended_hours": False,
    "is_active": False,
    "created_date": None
}
 
async def save_state_to_disk():
    async with state_lock:
        try:
            state_json = json.dumps(monitor_state, ensure_ascii=False, indent=4)
            async with aiofiles.open(STATE_FILE_TMP, "w", encoding="utf-8") as f:
                await f.write(state_json)
                await f.flush()
                try:
                    os.fsync(f.fileno())
                except Exception:
                    pass
            os.replace(STATE_FILE_TMP, STATE_FILE)
        except Exception as e:
            logger.error(f"💾 状态持久化失败: {e}")
 
async def load_state_from_disk():
    global monitor_state
    if os.path.exists(STATE_FILE):
        try:
            async with state_lock:
                async with aiofiles.open(STATE_FILE, "r", encoding="utf-8") as f:
                    content = await f.read()
                    if content.strip():
                        monitor_state.update(json.loads(content))
            logger.info(f"💾 成功恢复状态: {monitor_state}")
        except json.JSONDecodeError:
            logger.error("💾 状态文件已损坏,将使用默认状态并尝试覆盖修复。")
        except Exception as e:
            logger.error(f"💾 读取状态失败: {e}")
 
# ================= 2. 异步网络组件 =================
APCA_API_KEY_ID = os.getenv("APCA_API_KEY_ID", "你的APCA_API_KEY_ID")
APCA_API_SECRET_KEY = os.getenv("APCA_API_SECRET_KEY", "你的APCA_API_SECRET_KEY")
PROXY_URL = os.getenv("SOCKS_PROXY", None)
 
HERMES_API_URL = os.getenv("HERMES_API_URL", "http://127.0.0.1:8317/api/send_message")
FEISHU_USER_ID = os.getenv("FEISHU_USER_ID", "ou_5a1339ced5fdcdfdbea4b784c6d92c9e")
 
client = httpx.AsyncClient(
    limits=httpx.Limits(max_keepalive_connections=5, max_connections=50),
    transport=httpx.AsyncHTTPTransport(retries=3),
    proxy=PROXY_URL if PROXY_URL else None,
    headers={
        "APCA-API-KEY-ID": APCA_API_KEY_ID,
        "APCA-API-SECRET-KEY": APCA_API_SECRET_KEY,
        "Accept": "application/json"
    }
)
 
async def get_market_price(symbol: str) -> float | None:
    try:
        url = f"https://data.alpaca.markets/v2/stocks/{symbol}/trades/latest"
        resp = await client.get(url, timeout=5.0)
        resp.raise_for_status()
        price = resp.json().get("trade", {}).get("p")
        return float(price) if price is not None else None
    except Exception as e:
        logger.error(f"❌ 行情解析故障 (Alpaca Data API): {e}")
        return None
 
async def send_feishu_notification(message: str):
    try:
        payload = {
            "chat_id": FEISHU_USER_ID,
            "text": message
        }
        await client.post(HERMES_API_URL, json=payload, timeout=5.0)
        logger.info("🔔 飞书本地通知推送成功。")
    except Exception as e:
        logger.error(f"❌ 飞书本地通知推送失败: {e}")
 
# ================= 3. 核心监控逻辑 =================
def check_time_window(ny_now: datetime, allow_extended: bool):
    if ny_now.weekday() >= 5: return False, "非交易日"
    if monitor_state["created_date"] and monitor_state["created_date"] != ny_now.strftime("%Y-%m-%d"):
        return False, "单日自毁"
        
    open_hr, close_hr = (4, 20) if allow_extended else (9, 16)
    market_open = ny_now.replace(hour=open_hr, minute=30 if not allow_extended else 0, second=0, microsecond=0)
    market_close = ny_now.replace(hour=close_hr, minute=0, second=0, microsecond=0)
    
    if ny_now < market_open: return False, "盘前等待中..."
    if ny_now >= market_close: return False, "今日已收盘"
    return True, "交易窗口内"
 
async def monitoring_loop():
    logger.info("🚀 异步极简版监控引擎启动 (One-Shot 模式)...")
    while True:
        try:
            if not monitor_state["is_active"]:
                await asyncio.sleep(2)
                continue
 
            ny_now = datetime.now(NY_TZ)
            is_in_window, reason = check_time_window(ny_now, monitor_state.get("extended_hours", False))
 
            if not is_in_window:
                if "单日自毁" in reason or "今日已收盘" in reason:
                    logger.info(f"🔒 {reason},当日监控任务结束,自动退出监控。")
                    monitor_state["is_active"] = False
                    await save_state_to_disk()
                await asyncio.sleep(5)
                continue
 
            sym, target = monitor_state["symbol"], monitor_state["target_price"]
            direction = monitor_state["direction"]
            price = await get_market_price(sym)
            
            if price and price > 1.0:
                trigger = (direction == "BELOW" and price <= target) or (direction == "ABOVE" and price >= target)
                if trigger:
                    dir_txt = "跌破" if direction == "BELOW" else "涨破"
                    msg = f"🚨【Hermes 信号提示】\n标的: {sym}\n现价: {price}{dir_txt}目标价 {target}\n\n✅ 监控任务已完成并自动退出。"
                    await send_feishu_notification(msg)
                    
                    monitor_state["is_active"] = False
                    logger.info("✅ 目标已触发,One-Shot 任务终止。")
                    await save_state_to_disk()
 
            await asyncio.sleep(5)
            
        except asyncio.CancelledError:
            logger.info("🛑 收到取消信号,监控循环安全退出。")
            raise
        except Exception as e:
            logger.error(f"💥 监控协程异常: {e}")
            await asyncio.sleep(10)
 
# ================= 4. FastAPI 接口 =================
@asynccontextmanager
async def lifespan(app: FastAPI):
    await load_state_from_disk()
    task = asyncio.create_task(monitoring_loop())
    
    def handle_task_result(t):
        try:
            t.result()
        except asyncio.CancelledError:
            pass
        except Exception as err:
            logger.critical(f"💥 后台监控任务异常崩溃: {err}", exc_info=True)
            
    task.add_done_callback(handle_task_result)
    yield
    task.cancel()
    try:
        await task
    except asyncio.CancelledError:
        pass
    await client.aclose()
 
app = FastAPI(title="Hermes Agent API - OneShot", lifespan=lifespan)
 
class MonitorTask(BaseModel):
    symbol: str = "QQQ"
    target_price: float
    direction: str = "BELOW"
    extended_hours: bool = False
 
def verify_token(header_token: str = Depends(api_key_header)):
    if header_token != API_KEY: raise HTTPException(status_code=403, detail="Forbidden")
    return header_token
 
@app.post("/set_task", dependencies=[Depends(verify_token)])
async def set_task(task: MonitorTask):
    if task.direction.upper() not in ["BELOW", "ABOVE"]:
        raise HTTPException(status_code=400, detail="Direction must be BELOW or ABOVE")
 
    monitor_state.update({
        "symbol": task.symbol.upper(),
        "target_price": task.target_price,
        "direction": task.direction.upper(),
        "extended_hours": task.extended_hours,
        "is_active": True,
        "created_date": datetime.now(NY_TZ).strftime("%Y-%m-%d")
    })
    
    await save_state_to_disk()
    logger.info(f"📥 Agent 下发任务 -> {task.symbol} {task.direction} {task.target_price}")
    return {"status": "success", "data": monitor_state}
 
@app.post("/stop_task", dependencies=[Depends(verify_token)])
async def stop_task():
    monitor_state["is_active"] = False
    await save_state_to_disk()
    return {"status": "success", "message": "Monitoring stopped"}
 
@app.get("/status", dependencies=[Depends(verify_token)])
async def get_status():
    return monitor_state
 
if __name__ == "__main__":
    import uvicorn
    uvicorn.run("app:app", host="127.0.0.1", port=8000, reload=False, workers=1)

三、 Alpaca API 限流与数据量说明

  • 频率上限 (Rate Limit)
    • 免费版账户:限制为 每分钟 200 次 API 请求
    • 付费版账户:限制提升至 每分钟 1000 - 10000 次
  • 本脚本消耗计算
    • 采用 5 秒轮询间隔。每分钟发送 12 次请求(仅占用免费版限制的 6%),安全空间充裕,绝不影响并行交易程序的 API 额度。
  • 免费版数据源局限
    • 免费版的实时数据源仅包含 IEX 交易所 的成交,当成交发生在其他主流交易所时,数据获取可能产生几秒的细微滞后或些许价差,但不妨碍大波动的宏观监控。

四、 宿主机 Systemd 守护运行配置

1. 服务配置文件 /etc/systemd/system/hermes-monitor.service

[Unit]
Description=Hermes Agent One-Shot Price Monitor Service
After=network.target
 
[Service]
Type=simple
User=root
WorkingDirectory=/root/HermesAgent/price-monitor
Environment="HERMES_MONITOR_TOKEN=DefaultSecretToken123"
Environment="APCA_API_KEY_ID=[REDACTED]"
Environment="APCA_API_SECRET_KEY=[REDACTED]"
Environment="HERMES_API_URL=http://127.0.0.1:8317/api/send_message"
Environment="FEISHU_USER_ID=[REDACTED]"
ExecStart=/usr/bin/python3 app.py
Restart=always
RestartSec=5
 
[Install]
WantedBy=multi-user.target

2. 状态控制常用命令

# 重载系统守护进程并启动服务
systemctl daemon-reload
systemctl enable hermes-monitor --now
 
# 状态与日志监控
systemctl status hermes-monitor
journalctl -u hermes-monitor -n 50 --no-pager

相关链接