Centaur_Quant_Architecture/Python/router/main.py

408 lines
16 KiB
Python
Raw Permalink Normal View History

#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""
router/main.py Centaur Quant Architecture | Central Gateway (Python Router)
Role
----
Central TCP/IP gateway for the Universal Communication Bridge. The MQL5
Executor (CentaurQuant.mq5) streams SDP JSON frames into this process on
127.0.0.1:5555. This router:
1. Accepts persistent client connections (one thread per connection).
2. Reassembles frames terminated by '\n' or '\r\n' from raw TCP bytes.
3. Parses each frame STRICTLY with the stdlib `json` module.
4. Dispatches on the SDP "action_type" key:
Heartbeat / Tick_Harvest -> lightweight console log
Setup_Detected -> LIVE LLM evaluation (MarketAnalyzer) + reply
Trade_Opened / Closed -> prominent log + SQLite telemetry persist
5. Replies to Setup_Detected with the AI Advisory envelope:
{"algorithmic_confidence_score": <float>}\n
The reply is advisory ONLY (Anti-Veto Principle): the MT5 side may still
execute regardless of the score. MarketAnalyzer returns a neutral 40.0 on
any LLM failure/timeout, which triggers the MT5 quarter-risk path.
Trade_Opened/Trade_Closed are persisted by TelemetryDB (Track B), closing
the score -> execution -> outcome feedback loop for future calibration.
Robustness
----------
- Per-connection threads are daemonized; the server never dies with a client.
- ConnectionResetError / BrokenPipeError / empty reads are handled as a
clean disconnect; a reconnect is simply a new accept().
- Frames and the reassembly buffer are size-bounded to resist floods.
Standard library only: socket, threading, json, logging, argparse.
"""
import argparse
import json
import logging
import os
import socket
import sys
import threading
from typing import Any, Dict, List, Optional
# --- Track A: live LLM evaluation engine (Python/models/analyzer.py) -------
# Make the sibling packages importable whether run as a script or as a
# module: python router/main.py | python -m router.main
sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
from models.analyzer import MarketAnalyzer # noqa: E402
from database.telemetry import TelemetryDB # noqa: E402
# --------------------------------------------------------------------------
# Configuration
# --------------------------------------------------------------------------
HOST = "127.0.0.1" # default listen address (MT5 side connects here)
PORT = 5555 # default listen port (matches CentaurQuant.mq5 input)
BACKLOG = 5 # pending-connection queue depth
RECV_BUFFER_SIZE = 4096 # bytes pulled per recv() call
MAX_FRAME_SIZE = 1_048_576 # 1 MB hard cap per logical frame (flood guard)
ENCODING = "utf-8" # SDP frames are UTF-8 encoded
# --------------------------------------------------------------------------
# AI evaluation engine (Track A)
# Provider/model from env: PROVIDER=auto|openai|google, MODEL, *_API_KEY.
# Anti-Veto: analyze() returns 40.0 (quarter-risk trigger) on any failure.
# --------------------------------------------------------------------------
ANALYZER = MarketAnalyzer()
# --------------------------------------------------------------------------
# Telemetry persistence (Track B) — SQLite feedback-loop store.
# DB path overridable via env: CENTAUR_DB_PATH (default: ./centaur_telemetry.db)
# --------------------------------------------------------------------------
TELEMETRY = TelemetryDB(os.environ.get("CENTAUR_DB_PATH", "centaur_telemetry.db"))
# --------------------------------------------------------------------------
# Logging
# --------------------------------------------------------------------------
def configure_logging() -> None:
"""Console logging with timestamps; INFO shows the full SDP telemetry."""
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s | %(levelname)-8s | %(message)s",
datefmt="%Y-%m-%d %H:%M:%S",
)
# --------------------------------------------------------------------------
# Frame reassembly (newline / CRLF framed protocol)
# --------------------------------------------------------------------------
class FrameReader:
"""
Incrementally reassembles '\n'-terminated text frames from raw bytes.
The MT5 transport (CSocketClient) writes frames terminated by CRLF
("\r\n") and the receiver tolerates bare "\n" so we split on "\n"
and strip one trailing "\r" when present. Multiple frames may arrive
inside a single recv(); the remainder is kept buffered across reads.
"""
def __init__(self) -> None:
self._buffer = b""
def feed(self, data: bytes) -> List[str]:
"""Append raw bytes and return every COMPLETE frame they contain."""
self._buffer += data
frames: List[str] = []
while True:
idx = self._buffer.find(b"\n")
if idx == -1:
break # no complete frame yet — keep buffering
raw = self._buffer[:idx] # frame body (terminator removed)
self._buffer = self._buffer[idx + 1:] # consume terminator too
if raw.endswith(b"\r"):
raw = raw[:-1] # CRLF tolerance
if len(raw) > MAX_FRAME_SIZE:
logging.warning("Frame exceeds %d bytes — dropped.", MAX_FRAME_SIZE)
continue
try:
frames.append(raw.decode(ENCODING, errors="replace"))
except UnicodeDecodeError as exc: # pragma: no cover - defensive
logging.warning("Undecodable frame dropped: %s", exc)
# Flood guard: never let an unframed stream grow the buffer unbounded.
if len(self._buffer) > MAX_FRAME_SIZE:
logging.warning(
"Reassembly buffer exceeded %d bytes without a terminator — cleared.",
MAX_FRAME_SIZE,
)
self._buffer = b""
return frames
# --------------------------------------------------------------------------
# Strict SDP parsing
# --------------------------------------------------------------------------
def parse_sdp_frame(frame: str) -> Optional[Dict[str, Any]]:
"""
Strictly parse one SDP frame.
Returns a dict only when the frame is valid JSON *and* an object *and*
carries a string "action_type"; otherwise logs and returns None.
"""
try:
data = json.loads(frame)
except json.JSONDecodeError as exc:
logging.warning("Malformed JSON dropped: %s | frame=%.120r", exc, frame)
return None
if not isinstance(data, dict):
logging.warning("Non-object JSON dropped: %r", data)
return None
if not isinstance(data.get("action_type"), str):
logging.warning("Frame missing string 'action_type': %r", data)
return None
return data
# --------------------------------------------------------------------------
# Safe transport write (client may vanish between read and reply)
# --------------------------------------------------------------------------
def safe_send(conn: socket.socket, text: str) -> bool:
"""Best-effort frame write; never raises on a dead client."""
try:
conn.sendall(text.encode(ENCODING))
return True
except (ConnectionResetError, BrokenPipeError, OSError) as exc:
logging.warning("Send failed (client gone?): %s", exc)
return False
# --------------------------------------------------------------------------
# Action handlers
# --------------------------------------------------------------------------
def handle_heartbeat(msg: Dict[str, Any]) -> None:
"""Liveness signal — lightweight one-line log."""
logging.info(
"[Heartbeat] symbol=%s timeframe=%s",
msg.get("symbol", "?"),
msg.get("timeframe", "?"),
)
def handle_tick_harvest(msg: Dict[str, Any]) -> None:
"""One harvested tick — compact telemetry line."""
p = msg.get("payload") or {}
logging.info(
"[Tick_Harvest] %s %s bid=%s ask=%s spread_pts=%s tick_vol=%s",
msg.get("symbol", "?"),
msg.get("timeframe", "?"),
p.get("bid", "?"),
p.get("ask", "?"),
p.get("spread_points", "?"),
p.get("tick_volume", "?"),
)
def handle_setup_detected(msg: Dict[str, Any], conn: socket.socket) -> None:
"""
Setup proposal from the Executor.
Extracts symbol / setup_type / historical_context, runs the LIVE LLM
evaluation (MarketAnalyzer), and INSTANTLY replies with the AI Advisory
envelope:
{"algorithmic_confidence_score": <float>}\n
"""
p = msg.get("payload") or {}
symbol = msg.get("symbol", "?")
setup_type = p.get("setup_type", "?")
hist = msg.get("historical_context")
swing_count = len(hist) if isinstance(hist, list) else 0
logging.info(
"[Setup_Detected] %s | setup=%s | entry=%s sl=%s tp=%s | ctx_swings=%d",
symbol,
setup_type,
p.get("entry", "?"),
p.get("sl", "?"),
p.get("tp", "?"),
swing_count,
)
# --- live LLM evaluation (Anti-Veto: neutral 40.0 on failure/timeout) ---
result = ANALYZER.analyze(msg)
score = result["score"]
if result.get("used_fallback"):
logging.warning("[AI_Advisory] %s | fallback score %.2f used — %s",
symbol, score, result.get("fallback_reason", ""))
else:
logging.info("[AI_Advisory] %s | score=%.2f | reason='%s' | provider=%s | %d ms",
symbol, score, result.get("reason", ""),
result.get("provider", "?"), result.get("latency_ms", 0))
# --- AI Advisory envelope, newline-terminated (MT5 frames on '\n') ---
reply = json.dumps({"algorithmic_confidence_score": score}) + "\n"
if safe_send(conn, reply):
logging.info("[AI_Advisory] %s | score=%.2f -> sent to MT5", symbol, score)
else:
logging.warning(
"[AI_Advisory] %s | score=%.2f NOT delivered (client unreachable)",
symbol,
score,
)
def handle_trade_opened(msg: Dict[str, Any]) -> None:
"""Position opened — prominent execution record + SQLite persist."""
p = msg.get("payload") or {}
logging.info(
">>> [Trade_Opened] %s %s | ticket=%s %s lot=%s entry=%s sl=%s tp=%s | ai_score=%s",
msg.get("symbol", "?"),
msg.get("timeframe", "?"),
p.get("ticket", "?"),
p.get("direction", "?"),
p.get("lot", "?"),
p.get("entry_price", "?"),
p.get("sl", "?"),
p.get("tp", "?"),
msg.get("algorithmic_confidence_score", "?"),
)
# --- persist the execution record (Track B) ---
TELEMETRY.log_trade_opened(msg)
def handle_trade_closed(msg: Dict[str, Any]) -> None:
"""Position closed — the feedback-loop training signal + SQLite update."""
p = msg.get("payload") or {}
logging.info(
">>> [Trade_Closed] %s | ticket=%s profit=%s r_multiple=%s initial_ai_score=%s",
msg.get("symbol", "?"),
p.get("ticket", "?"),
p.get("profit", "?"),
p.get("r_multiple", "?"),
p.get("initial_ai_score", "?"),
)
# --- persist the closed outcome (Track B) ---
TELEMETRY.log_trade_closed(msg)
# --------------------------------------------------------------------------
# Dispatcher
# --------------------------------------------------------------------------
# Setup_Detected is handled separately (it needs the socket to reply).
_ACTIONS: Dict[str, Any] = {
"Heartbeat": handle_heartbeat,
"Tick_Harvest": handle_tick_harvest,
"Trade_Opened": handle_trade_opened,
"Trade_Closed": handle_trade_closed,
}
def dispatch(msg: Dict[str, Any], conn: socket.socket) -> None:
"""Route one parsed SDP frame by its action_type."""
action = msg.get("action_type")
if action == "Setup_Detected":
handle_setup_detected(msg, conn)
return
handler = _ACTIONS.get(action)
if handler is None:
logging.warning("[Router] unknown action_type=%r — frame ignored", action)
return
try:
handler(msg)
except Exception: # a broken handler must never kill the connection thread
logging.exception("[Router] handler '%s' raised — frame ignored", action)
# --------------------------------------------------------------------------
# Per-connection worker
# --------------------------------------------------------------------------
def handle_client(conn: socket.socket, addr: tuple) -> None:
"""
Serve one persistent MT5 connection until it disconnects.
Reads raw bytes, reassembles frames, parses and dispatches each one.
A reconnect from the MT5 side simply spawns a fresh thread via accept().
"""
host, port = addr[0], addr[1]
logging.info("[Gateway] client connected: %s:%d", host, port)
reader = FrameReader()
try:
while True:
data = conn.recv(RECV_BUFFER_SIZE)
if not data:
# Peer performed an orderly close — empty byte read.
logging.info("[Gateway] %s:%d closed connection (empty read)", host, port)
break
for frame in reader.feed(data):
msg = parse_sdp_frame(frame)
if msg is not None:
dispatch(msg, conn)
except ConnectionResetError as exc:
# MT5 terminal killed the socket (e.g., chart closed mid-session).
logging.warning("[Gateway] %s:%d connection reset: %s", host, port, exc)
except (BrokenPipeError, OSError) as exc:
logging.warning("[Gateway] %s:%d socket error: %s", host, port, exc)
except Exception:
# Last-resort guard: log and drop the connection, keep the server up.
logging.exception("[Gateway] %s:%d unexpected error — connection closed", host, port)
finally:
try:
conn.shutdown(socket.SHUT_RDWR)
except OSError:
pass # socket may already be closed by the peer
conn.close()
logging.info("[Gateway] %s:%d disconnected", host, port)
# --------------------------------------------------------------------------
# Server entry point
# --------------------------------------------------------------------------
def run_server(host: str, port: int) -> None:
"""Bind, listen, and accept client connections forever (Ctrl+C to stop)."""
server = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
server.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) # fast rebind
server.bind((host, port))
server.listen(BACKLOG)
logging.info("=" * 72)
logging.info("Centaur Gateway listening on %s:%d", host, port)
logging.info("LLM backend: %s",
ANALYZER._backend.name if ANALYZER._backend else "NONE (fallback 40.0)")
logging.info("Telemetry DB: %s", TELEMETRY.path)
logging.info("Press Ctrl+C to stop.")
logging.info("=" * 72)
try:
while True:
conn, addr = server.accept()
# daemon thread: dies with the process; never blocks shutdown
thread = threading.Thread(target=handle_client, args=(conn, addr), daemon=True)
thread.start()
except KeyboardInterrupt:
logging.info("[Gateway] Ctrl+C received — shutting down.")
finally:
server.close()
TELEMETRY.close()
logging.info("[Gateway] server socket closed, telemetry flushed. Bye.")
def main() -> None:
parser = argparse.ArgumentParser(description="Centaur Quant — Central Gateway")
parser.add_argument("--host", default=HOST, help="listen address (default: %(default)s)")
parser.add_argument("--port", type=int, default=PORT, help="listen port (default: %(default)s)")
args = parser.parse_args()
configure_logging()
run_server(args.host, args.port)
if __name__ == "__main__":
main()