XAU_system/bot.py

3570 lines
128 KiB
Python

#!/usr/bin/env python3
"""
XAUUSD five-strategy OCO trailing portfolio for MetaTrader 5.
Required user-maintained files:
1. xauusd_five_strategy_oco_bot_risk_vector_v2.py
2. timetable.csv (UTC macro-news schedule, next to this script)
3. traded_codes.csv (event-code allowlist, next to this script)
Runtime folders/files are created automatically:
logs/YYYY.txt balance changes only (UTC; initial snapshot once)
logs/YYYY_trades.log yearly closed-trade ledger (UTC)
logs/runtime.log rotating operational/error log
state/strategy_state.json restart-safe internal state
Five independent subsystems
---------------------------
1. NEWS
Scheduled events are read from timetable.csv in UTC and filtered by traded_codes.csv.
offset 2.00; initial SL 30.00; activation 5.00;
adaptive trail clip(0.25 * RMS(close changes, 5), 0.50, 2.25);
release-minute expiry 60 s; maximum hold 120 min;
risk-based volume: 6.5% of balance at the initial stop.
2. NY0830
Every weekday at 08:30 America/New_York.
offset 6.00; initial SL 60.00; activation 5.00;
adaptive trail clip(2.00 * RMS(close changes, 20), 3.00, 12.00);
at +8.00 lock +5.00; pending expiry 15 min; maximum hold 300 min;
risk-based volume: 10.5% of balance at the initial stop.
3. TOKYO0845
Every weekday at 08:45 Asia/Tokyo.
offset 3.00; initial SL 30.00; activation 3.00;
adaptive trail clip(30 * Wilder ATR(30), 12.50, 25.00);
at +14.00 lock +6.00; at +32.00 switch permanently to trail 12.00;
pending expiry 15 min; maximum hold 480 min;
risk-based volume: 3.5% of balance at the initial stop.
4. TOKYO1500
Every weekday at 15:00 Asia/Tokyo.
offset 3.00; initial SL 40.00; activation 3.00;
adaptive trail min(15 * median true range(10), 12.50);
at +10.00 lock +6.00; at +20.00 switch permanently to trail 11.00;
pending expiry 15 min; maximum hold 180 min;
risk-based volume: 0.5% of balance at the initial stop.
5. NYSE1545
Every weekday at 15:45 America/New_York.
offset 2.00; initial SL 40.00;
frozen pre-event activation clip(3 * SMA(true range, 10), 2.50, 5.00);
adaptive trail clip(0.30 * SMA(true range, 10), 0.40, 0.85);
pending expiry 5 min; maximum hold 480 min;
risk-based volume: 2.0% of balance at the initial stop.
Concurrency
-----------
Each session subsystem may have at most one pending OCO pair or one open position.
NEWS is keyed by event minute: every distinct enabled-code timetable minute may have
its own pending OCO pair or open position, even while earlier NEWS positions remain
open. Rows without an enabled code in traded_codes.csv are loaded but not traded.
Simultaneous macro releases must be combined into one timetable row, so one UTC
minute always produces at most one NEWS OCO pair. New setups are also governed by
a 23.5% aggregate portfolio initial-stop-risk cap; each intact OCO pair counts once.
Time-zone and scheduling guarantees
-----------------------------------
* All internal timestamps and timetable timestamps are UTC-aware.
* The host Windows/VPS local time zone is never used for strategy scheduling.
* America/New_York ZoneInfo handles US daylight-saving changes automatically.
* Asia/Tokyo ZoneInfo is used for the two Japanese session events.
* MT5 broker-clock timestamps are converted dynamically to and from true UTC
using the live quote clock; no fixed UTC+2/UTC+3 offset is hard-coded.
* timetable.csv and traded_codes.csv are re-read after edits; temporarily
incomplete/invalid replacements do not replace the last valid in-memory data.
Execution and recovery
----------------------
* Orders are submitted approximately 10 seconds before each event.
* There is no live-spread rejection filter.
* Every pending leg carries a visible broker-side initial stop-loss.
* OCO is client-enforced: after one leg fills, the sibling is removed.
* If both legs fill during a disconnection, the later position is closed after
reconnection; every filled position remains protected by its attached SL.
* Trailing uses the final actual bid/ask tick of each fully completed M1 candle.
Distance is volatility-adaptive under the deployment-safe formulas; the entry
candle is excluded, fixed profit floors/late trails remain, and SL never loosens.
* After restart, trailing state is rebuilt from MT5 tick history. Existing
version-1 state files remain valid; no trailing-state migration is required.
* MT5 network status is checked explicitly. A connection retcode aborts
filling/expiry fallbacks immediately and causes a clean reconnect attempt.
* A hedging account is mandatory.
Installation:
pip install MetaTrader5 tzdata
Run while the intended MT5 terminal/account is open and Algo Trading is enabled:
python xauusd_five_strategy_oco_bot_risk_vector_v2.py
"""
from __future__ import annotations
import csv
import json
import logging
from logging.handlers import RotatingFileHandler
import math
import os
from pathlib import Path
import signal
import sys
import time
from dataclasses import dataclass
from datetime import datetime, timedelta, timezone
from typing import Any, Iterable, Optional
from types import SimpleNamespace
from zoneinfo import ZoneInfo, ZoneInfoNotFoundError
try:
import MetaTrader5 as mt5
except ImportError as exc: # pragma: no cover - runtime environment specific
raise SystemExit(
"MetaTrader5 package is not installed. Run: pip install MetaTrader5"
) from exc
# =============================================================================
# USER CONFIGURATION
# =============================================================================
# Exact broker symbol is tried first; common suffix/prefix variants are detected.
PREFERRED_SYMBOL = "XAUUSD"
# -----------------------------------------------------------------------------
# RISK SIZING — MAIN VARIABLES TO ADJUST
# -----------------------------------------------------------------------------
# Values are percentages of the current realized account balance. Volume is
# floored to the broker's lot step so the target risk is never exceeded.
NEWS_RISK_PERCENT = 6.5
NY0830_RISK_PERCENT = 10.5
TOKYO0845_RISK_PERCENT = 3.5
TOKYO1500_RISK_PERCENT = 0.5
NYSE1545_RISK_PERCENT = 2.0
# Aggregate nominal initial-stop risk of all this bot's live positions and OCO
# groups. An intact two-leg OCO pair counts once because only one leg is intended
# to fill. A position plus an unremoved sibling is counted conservatively as both.
PORTFOLIO_RISK_CAP_PERCENT = 23.5
# Conservative allowance used only for position sizing and cap accounting. The
# backtest used $0.15 adverse entry plus $0.15 adverse exit slippage per ounce.
RISK_COST_BUFFER_PER_OUNCE = 0.30
# "SKIP" preserves the requested percentage when the calculated volume is below
# the broker minimum. "FORCE_MIN" trades the broker minimum and may exceed it.
MINIMUM_LOT_POLICY = "SKIP"
# While a setup is blocked by risk, the bot keeps checking until shortly before
# the event. This allows an existing position to close and free risk capacity.
RISK_SKIP_FINALIZE_SECONDS = 0.75
# Optional terminal/account locks. Set these before live deployment when multiple
# terminals or accounts exist on the VPS. None accepts the currently active one.
MT5_TERMINAL_PATH: Optional[str] = None
EXPECTED_ACCOUNT_LOGIN: Optional[int] = None
EXPECTED_SERVER: Optional[str] = None
# Unique IDs. Do not reuse these magic numbers in another EA/script on the account.
NEWS_MAGIC = 271800
NY0830_MAGIC = 271801
TOKYO0845_MAGIC = 271802
TOKYO1500_MAGIC = 271803
NYSE1545_MAGIC = 271804
OUR_MAGICS = {
NEWS_MAGIC,
NY0830_MAGIC,
TOKYO0845_MAGIC,
TOKYO1500_MAGIC,
NYSE1545_MAGIC,
}
# Price both pending legs from the latest live quote approximately ten seconds
# before an event. This leaves operational time for both acknowledgements while
# remaining close to the final pre-event quote used in the M1 backtests.
PLACEMENT_LEAD_SECONDS = 10.0
# Maximum market-order deviation accepted by the client request, in symbol points.
# For a 0.01 XAUUSD point, 15 points = 0.15 USD.
MAX_DEVIATION_POINTS = 15
# Main-loop cadence. It automatically sleeps less when an event is near.
NORMAL_POLL_SECONDS = 0.25
NEAR_EVENT_POLL_SECONDS = 0.05
# Reject stale quotes instead of creating orders from an old market price.
MAX_TICK_AGE_SECONDS = 10.0
# Trailing is allowed to defer locally when the quote feed pauses. This prevents
# an expected short quote gap from aborting reconciliation of the remaining
# subsystems or producing a full top-level traceback on every loop iteration.
TRAIL_TICK_WARNING_INTERVAL_SECONDS = 30.0
# Timetable is reloaded automatically after an edit. The mtime check is cheap
# and a five-second interval keeps last-minute calendar corrections practical.
TIMETABLE_RELOAD_SECONDS = 5
# Runtime connection and history behaviour.
RECONNECT_RETRY_SECONDS = 5
TRADE_LOG_LOOKBACK_DAYS = 45
STATE_RETENTION_DAYS = 120
# Optional hard volume cap. Set to None to use the broker's maximum.
MAX_LOT: Optional[float] = None
# Strategy parameters, in XAUUSD price units (USD per ounce).
NEWS_OFFSET = 2.00
NEWS_INITIAL_SL = 30.00
NEWS_TRAIL_ACTIVATION = 5.00
NEWS_TRAIL_DISTANCE = 0.75
NEWS_TRAIL_ESTIMATOR = "CLOSE_RMS"
NEWS_TRAIL_WINDOW = 5
NEWS_TRAIL_MULTIPLIER = 0.25
NEWS_TRAIL_FLOOR = 0.50
NEWS_TRAIL_CAP = 2.25
NEWS_EXPIRY_SECONDS = 60
NEWS_MAX_HOLD_MINUTES = 120
NEWS_MIN_LOT = 0.01
# NEWS event types are controlled by traded_codes.csv beside this script.
# The bot reloads that file automatically. A timetable minute is tradable when at
# least one of its component event codes is enabled there. The full timetable is
# still loaded; The5ers High Stakes continues to use every enabled timetable row
# for blackout protection, regardless of traded_codes.csv.
TRADED_CODES_RELOAD_SECONDS = 5
SESSION_MIN_LOT = 0.01
NY0830_OFFSET = 6.00
NY0830_INITIAL_SL = 60.00
NY0830_TRAIL_ACTIVATION = 5.00
NY0830_TRAIL_DISTANCE = 6.00
NY0830_TRAIL_ESTIMATOR = "CLOSE_RMS"
NY0830_TRAIL_WINDOW = 20
NY0830_TRAIL_MULTIPLIER = 2.00
NY0830_TRAIL_FLOOR = 3.00
NY0830_TRAIL_CAP = 12.00
NY0830_PROFIT_FLOOR_THRESHOLD = 8.00
NY0830_PROFIT_FLOOR = 5.00
NY0830_EXPIRY_SECONDS = 15 * 60
NY0830_MAX_HOLD_MINUTES = 300
TOKYO0845_OFFSET = 3.00
TOKYO0845_INITIAL_SL = 30.00
TOKYO0845_TRAIL_ACTIVATION = 3.00
TOKYO0845_TRAIL_DISTANCE = 25.00
TOKYO0845_TRAIL_ESTIMATOR = "ATR_WILDER"
TOKYO0845_TRAIL_WINDOW = 30
TOKYO0845_TRAIL_MULTIPLIER = 30.00
TOKYO0845_TRAIL_FLOOR = 12.50
TOKYO0845_TRAIL_CAP = 25.00
TOKYO0845_PROFIT_FLOOR_THRESHOLD = 14.00
TOKYO0845_PROFIT_FLOOR = 6.00
TOKYO0845_LATE_TRAIL_THRESHOLD = 32.00
TOKYO0845_LATE_TRAIL_DISTANCE = 12.00
TOKYO0845_EXPIRY_SECONDS = 15 * 60
TOKYO0845_MAX_HOLD_MINUTES = 480
TOKYO1500_OFFSET = 3.00
TOKYO1500_INITIAL_SL = 40.00
TOKYO1500_TRAIL_ACTIVATION = 3.00
TOKYO1500_TRAIL_DISTANCE = 12.50
TOKYO1500_TRAIL_ESTIMATOR = "TR_MEDIAN"
TOKYO1500_TRAIL_WINDOW = 10
TOKYO1500_TRAIL_MULTIPLIER = 15.00
TOKYO1500_TRAIL_FLOOR = 0.00
TOKYO1500_TRAIL_CAP = 12.50
TOKYO1500_PROFIT_FLOOR_THRESHOLD = 10.00
TOKYO1500_PROFIT_FLOOR = 6.00
TOKYO1500_LATE_TRAIL_THRESHOLD = 20.00
TOKYO1500_LATE_TRAIL_DISTANCE = 11.00
TOKYO1500_EXPIRY_SECONDS = 15 * 60
TOKYO1500_MAX_HOLD_MINUTES = 180
NYSE1545_OFFSET = 2.00
NYSE1545_INITIAL_SL = 40.00
NYSE1545_TRAIL_ACTIVATION = 5.00
NYSE1545_TRAIL_DISTANCE = 0.85
NYSE1545_TRAIL_ESTIMATOR = "TR_SMA"
NYSE1545_TRAIL_WINDOW = 10
NYSE1545_TRAIL_MULTIPLIER = 0.30
NYSE1545_TRAIL_FLOOR = 0.40
NYSE1545_TRAIL_CAP = 0.85
NYSE1545_ACTIVATION_ESTIMATOR = "TR_SMA"
NYSE1545_ACTIVATION_WINDOW = 10
NYSE1545_ACTIVATION_MULTIPLIER = 3.00
NYSE1545_ACTIVATION_FLOOR = 2.50
NYSE1545_ACTIVATION_CAP = 5.00
NYSE1545_EXPIRY_SECONDS = 5 * 60
NYSE1545_MAX_HOLD_MINUTES = 480
# No trade is opened if account currency is not USD, because risk sizing and the
# cost buffer are defined in USD account terms.
REQUIRE_USD_ACCOUNT = True
# =============================================================================
# PATHS AND GLOBALS
# =============================================================================
UTC = timezone.utc
try:
NY_TZ = ZoneInfo("America/New_York")
TOKYO_TZ = ZoneInfo("Asia/Tokyo")
except ZoneInfoNotFoundError as exc: # common on a fresh Windows Python install
raise SystemExit(
"IANA timezone data is missing. Install it with: pip install tzdata"
) from exc
BASE_DIR = Path(__file__).resolve().parent
TIMETABLE_PATH = BASE_DIR / "timetable.csv"
TRADED_CODES_PATH = BASE_DIR / "traded_codes.csv"
LOG_DIR = BASE_DIR / "logs"
STATE_DIR = BASE_DIR / "state"
STATE_PATH = STATE_DIR / "strategy_state.json"
LOCK_PATH = STATE_DIR / "xauusd_combined_oco_bot.lock"
LOG_DIR.mkdir(parents=True, exist_ok=True)
STATE_DIR.mkdir(parents=True, exist_ok=True)
LOGGER = logging.getLogger("xauusd_oco")
STOP_REQUESTED = False
SYMBOL = PREFERRED_SYMBOL
SYMBOL_INFO: Any = None
LOCK_HANDLE: Any = None
LAST_TIMETABLE_RELOAD_MONOTONIC = 0.0
TIMETABLE_SIGNATURE: Optional[tuple[int, int]] = None
TIMETABLE_READY = False
LAST_TRADED_CODES_RELOAD_MONOTONIC = 0.0
TRADED_CODES_SIGNATURE: Optional[tuple[int, int]] = None
TRADED_CODES_READY = False
TRADED_NEWS_CODES: frozenset[str] = frozenset()
NEWS_EVENTS: list["NewsEvent"] = []
LAST_DEAL_SCAN_MONOTONIC = 0.0
LAST_LOGGED_BALANCE_TEXT: Optional[str] = None
LAST_BALANCE_LOG_YEAR: Optional[int] = None
LAST_HORIZON_WARNING_DATE: Optional[str] = None
BROKER_CLOCK_BASIS = "broker_epoch_v2_offset_aware"
BROKER_UTC_OFFSET_SECONDS: Optional[int] = None
MAX_BROKER_UTC_OFFSET_SECONDS = 14 * 60 * 60
COMPLETED_CLOSE_CACHE: dict[int, Optional[tuple[datetime, float, float]]] = {}
M1_RATE_CACHE: dict[int, Optional[list["CompletedM1Bar"]]] = {}
M1_RATE_RETRY_AFTER: dict[int, float] = {}
M1_RATE_SOURCE_WARNINGS: set[int] = set()
VOLATILITY_WARNING_KEYS: set[tuple[str, int]] = set()
RISK_BLOCK_WARNING_KEYS: set[tuple[str, str]] = set()
TRAIL_TICK_WARNING_LAST_MONOTONIC: dict[int, float] = {}
# =============================================================================
# DATA CLASSES
# =============================================================================
@dataclass(frozen=True)
class StrategyConfig:
name: str
magic: int
group_prefix: str
offset: float
initial_sl: float
trail_activation: float
trail_distance: float
expiry_seconds: int
max_hold_minutes: int
risk_percent: float
min_lot: float
profit_floor_threshold: Optional[float] = None
profit_floor: Optional[float] = None
late_trail_threshold: Optional[float] = None
late_trail_distance: Optional[float] = None
trail_estimator: Optional[str] = None
trail_window: int = 0
trail_multiplier: float = 0.0
trail_floor: float = 0.0
trail_cap: float = 0.0
activation_estimator: Optional[str] = None
activation_window: int = 0
activation_multiplier: float = 0.0
activation_floor: float = 0.0
activation_cap: float = 0.0
@dataclass(frozen=True)
class LotSizingResult:
sizing_balance: float
target_risk_usd: float
risk_per_lot_usd: float
raw_lot: float
lot: float
actual_risk_usd: float
NEWS_CONFIG = StrategyConfig(
name="NEWS",
magic=NEWS_MAGIC,
group_prefix="N",
offset=NEWS_OFFSET,
initial_sl=NEWS_INITIAL_SL,
trail_activation=NEWS_TRAIL_ACTIVATION,
trail_distance=NEWS_TRAIL_DISTANCE,
expiry_seconds=NEWS_EXPIRY_SECONDS,
max_hold_minutes=NEWS_MAX_HOLD_MINUTES,
risk_percent=NEWS_RISK_PERCENT,
min_lot=NEWS_MIN_LOT,
trail_estimator=NEWS_TRAIL_ESTIMATOR,
trail_window=NEWS_TRAIL_WINDOW,
trail_multiplier=NEWS_TRAIL_MULTIPLIER,
trail_floor=NEWS_TRAIL_FLOOR,
trail_cap=NEWS_TRAIL_CAP,
)
NY0830_CONFIG = StrategyConfig(
name="NY0830",
magic=NY0830_MAGIC,
group_prefix="Y",
offset=NY0830_OFFSET,
initial_sl=NY0830_INITIAL_SL,
trail_activation=NY0830_TRAIL_ACTIVATION,
trail_distance=NY0830_TRAIL_DISTANCE,
expiry_seconds=NY0830_EXPIRY_SECONDS,
max_hold_minutes=NY0830_MAX_HOLD_MINUTES,
risk_percent=NY0830_RISK_PERCENT,
min_lot=SESSION_MIN_LOT,
profit_floor_threshold=NY0830_PROFIT_FLOOR_THRESHOLD,
profit_floor=NY0830_PROFIT_FLOOR,
trail_estimator=NY0830_TRAIL_ESTIMATOR,
trail_window=NY0830_TRAIL_WINDOW,
trail_multiplier=NY0830_TRAIL_MULTIPLIER,
trail_floor=NY0830_TRAIL_FLOOR,
trail_cap=NY0830_TRAIL_CAP,
)
TOKYO0845_CONFIG = StrategyConfig(
name="TOKYO0845",
magic=TOKYO0845_MAGIC,
group_prefix="A",
offset=TOKYO0845_OFFSET,
initial_sl=TOKYO0845_INITIAL_SL,
trail_activation=TOKYO0845_TRAIL_ACTIVATION,
trail_distance=TOKYO0845_TRAIL_DISTANCE,
expiry_seconds=TOKYO0845_EXPIRY_SECONDS,
max_hold_minutes=TOKYO0845_MAX_HOLD_MINUTES,
risk_percent=TOKYO0845_RISK_PERCENT,
min_lot=SESSION_MIN_LOT,
profit_floor_threshold=TOKYO0845_PROFIT_FLOOR_THRESHOLD,
profit_floor=TOKYO0845_PROFIT_FLOOR,
late_trail_threshold=TOKYO0845_LATE_TRAIL_THRESHOLD,
late_trail_distance=TOKYO0845_LATE_TRAIL_DISTANCE,
trail_estimator=TOKYO0845_TRAIL_ESTIMATOR,
trail_window=TOKYO0845_TRAIL_WINDOW,
trail_multiplier=TOKYO0845_TRAIL_MULTIPLIER,
trail_floor=TOKYO0845_TRAIL_FLOOR,
trail_cap=TOKYO0845_TRAIL_CAP,
)
TOKYO1500_CONFIG = StrategyConfig(
name="TOKYO1500",
magic=TOKYO1500_MAGIC,
group_prefix="T",
offset=TOKYO1500_OFFSET,
initial_sl=TOKYO1500_INITIAL_SL,
trail_activation=TOKYO1500_TRAIL_ACTIVATION,
trail_distance=TOKYO1500_TRAIL_DISTANCE,
expiry_seconds=TOKYO1500_EXPIRY_SECONDS,
max_hold_minutes=TOKYO1500_MAX_HOLD_MINUTES,
risk_percent=TOKYO1500_RISK_PERCENT,
min_lot=SESSION_MIN_LOT,
profit_floor_threshold=TOKYO1500_PROFIT_FLOOR_THRESHOLD,
profit_floor=TOKYO1500_PROFIT_FLOOR,
late_trail_threshold=TOKYO1500_LATE_TRAIL_THRESHOLD,
late_trail_distance=TOKYO1500_LATE_TRAIL_DISTANCE,
trail_estimator=TOKYO1500_TRAIL_ESTIMATOR,
trail_window=TOKYO1500_TRAIL_WINDOW,
trail_multiplier=TOKYO1500_TRAIL_MULTIPLIER,
trail_floor=TOKYO1500_TRAIL_FLOOR,
trail_cap=TOKYO1500_TRAIL_CAP,
)
NYSE1545_CONFIG = StrategyConfig(
name="NYSE1545",
magic=NYSE1545_MAGIC,
group_prefix="C",
offset=NYSE1545_OFFSET,
initial_sl=NYSE1545_INITIAL_SL,
trail_activation=NYSE1545_TRAIL_ACTIVATION,
trail_distance=NYSE1545_TRAIL_DISTANCE,
expiry_seconds=NYSE1545_EXPIRY_SECONDS,
max_hold_minutes=NYSE1545_MAX_HOLD_MINUTES,
risk_percent=NYSE1545_RISK_PERCENT,
min_lot=SESSION_MIN_LOT,
trail_estimator=NYSE1545_TRAIL_ESTIMATOR,
trail_window=NYSE1545_TRAIL_WINDOW,
trail_multiplier=NYSE1545_TRAIL_MULTIPLIER,
trail_floor=NYSE1545_TRAIL_FLOOR,
trail_cap=NYSE1545_TRAIL_CAP,
activation_estimator=NYSE1545_ACTIVATION_ESTIMATOR,
activation_window=NYSE1545_ACTIVATION_WINDOW,
activation_multiplier=NYSE1545_ACTIVATION_MULTIPLIER,
activation_floor=NYSE1545_ACTIVATION_FLOOR,
activation_cap=NYSE1545_ACTIVATION_CAP,
)
ALL_CONFIGS = (
NEWS_CONFIG,
NY0830_CONFIG,
TOKYO0845_CONFIG,
TOKYO1500_CONFIG,
NYSE1545_CONFIG,
)
CONFIG_BY_MAGIC = {config.magic: config for config in ALL_CONFIGS}
SESSION_CONFIGS = (
NY0830_CONFIG,
TOKYO0845_CONFIG,
TOKYO1500_CONFIG,
NYSE1545_CONFIG,
)
CONFIG_BY_MAGIC = {config.magic: config for config in ALL_CONFIGS}
@dataclass(frozen=True)
class DailySessionSpec:
config: StrategyConfig
timezone: ZoneInfo
hour: int
minute: int
event_id_prefix: str
display_name: str
DAILY_SESSION_SPECS = (
DailySessionSpec(
NY0830_CONFIG,
NY_TZ,
8,
30,
"NY0830",
"NY 08:30 ET",
),
DailySessionSpec(
TOKYO0845_CONFIG,
TOKYO_TZ,
8,
45,
"TOKYO0845",
"Tokyo 08:45 JST",
),
DailySessionSpec(
TOKYO1500_CONFIG,
TOKYO_TZ,
15,
0,
"TOKYO1500",
"Tokyo 15:00 JST",
),
DailySessionSpec(
NYSE1545_CONFIG,
NY_TZ,
15,
45,
"NYSE1545",
"NYSE 15:45 ET",
),
)
@dataclass(frozen=True)
class NewsEvent:
event_id: str
event_name: str
when_utc: datetime
enabled: bool
event_codes: tuple[str, ...]
tradable: bool
# =============================================================================
# LOGGING, LOCKING, STATE
# =============================================================================
def setup_logging() -> None:
LOGGER.setLevel(logging.INFO)
LOGGER.handlers.clear()
formatter = logging.Formatter(
"%(asctime)sZ | %(levelname)-8s | %(message)s",
datefmt="%Y-%m-%d %H:%M:%S",
)
formatter.converter = time.gmtime
file_handler = RotatingFileHandler(
LOG_DIR / "runtime.log",
maxBytes=5_000_000,
backupCount=5,
encoding="utf-8",
)
file_handler.setFormatter(formatter)
LOGGER.addHandler(file_handler)
console_handler = logging.StreamHandler(sys.stdout)
console_handler.setFormatter(formatter)
LOGGER.addHandler(console_handler)
def acquire_single_instance_lock() -> None:
"""Hold an OS file lock for the lifetime of the process."""
global LOCK_HANDLE
LOCK_HANDLE = open(LOCK_PATH, "a+b")
try:
if os.name == "nt":
import msvcrt
LOCK_HANDLE.seek(0)
if LOCK_HANDLE.tell() == 0:
LOCK_HANDLE.write(b"0")
LOCK_HANDLE.flush()
LOCK_HANDLE.seek(0)
msvcrt.locking(LOCK_HANDLE.fileno(), msvcrt.LK_NBLCK, 1)
else: # pragma: no cover - production target is Windows
import fcntl
fcntl.flock(LOCK_HANDLE.fileno(), fcntl.LOCK_EX | fcntl.LOCK_NB)
except (OSError, IOError) as exc:
raise SystemExit("Another instance of this bot is already running.") from exc
LOCK_HANDLE.seek(0)
LOCK_HANDLE.truncate()
LOCK_HANDLE.write(str(os.getpid()).encode("ascii"))
LOCK_HANDLE.flush()
def default_state() -> dict[str, Any]:
return {
"version": 1,
"processed_events": {},
"active_groups": {},
"trail": {},
"logged_position_ids": [],
}
def load_state() -> dict[str, Any]:
if not STATE_PATH.exists():
return default_state()
try:
data = json.loads(STATE_PATH.read_text(encoding="utf-8"))
if not isinstance(data, dict):
raise ValueError("state root is not an object")
base = default_state()
for key, value in base.items():
data.setdefault(key, value)
return data
except Exception:
LOGGER.exception("State file is unreadable; backing it up and rebuilding.")
backup = STATE_PATH.with_suffix(f".corrupt.{int(time.time())}.json")
try:
STATE_PATH.replace(backup)
except OSError:
pass
return default_state()
def save_state(state: dict[str, Any]) -> None:
tmp = STATE_PATH.with_suffix(".tmp")
payload = json.dumps(state, indent=2, sort_keys=True, ensure_ascii=False)
tmp.write_text(payload, encoding="utf-8")
os.replace(tmp, STATE_PATH)
def prune_state(state: dict[str, Any], now_utc: datetime) -> None:
cutoff = now_utc - timedelta(days=STATE_RETENTION_DAYS)
processed = state.get("processed_events", {})
for key in list(processed):
stamp = processed[key].get("updated_utc") if isinstance(processed[key], dict) else None
try:
dt = parse_utc(stamp) if stamp else now_utc
except Exception:
dt = now_utc
if dt < cutoff:
del processed[key]
# Trail state for positions that no longer exist is cleared elsewhere.
logged = state.get("logged_position_ids", [])
if len(logged) > 10_000:
state["logged_position_ids"] = logged[-10_000:]
# =============================================================================
# TIME AND CSV
# =============================================================================
def utc_now() -> datetime:
return datetime.now(UTC)
def format_utc(dt: datetime) -> str:
return dt.astimezone(UTC).strftime("%Y-%m-%dT%H:%M:%SZ")
def parse_utc(value: str) -> datetime:
value = value.strip()
if value.endswith("Z"):
value = value[:-1] + "+00:00"
dt = datetime.fromisoformat(value)
if dt.tzinfo is None:
raise ValueError(f"UTC datetime lacks timezone: {value!r}")
return dt.astimezone(UTC)
def bool_from_csv(value: str) -> bool:
return value.strip().lower() in {"1", "true", "yes", "y", "on"}
def refresh_news_tradability() -> None:
"""Re-evaluate loaded timetable rows after traded_codes.csv changes."""
global NEWS_EVENTS
if not NEWS_EVENTS:
return
NEWS_EVENTS = [
NewsEvent(
event_id=event.event_id,
event_name=event.event_name,
when_utc=event.when_utc,
enabled=event.enabled,
event_codes=event.event_codes,
tradable=event.enabled and traded_news_minute(event.event_codes),
)
for event in NEWS_EVENTS
]
def load_traded_codes(force: bool = False) -> None:
"""Atomically load the user-maintained traded_codes.csv file.
Required columns are ``event_code`` and ``enabled``. Extra columns such as
``family`` and ``notes`` are ignored. Missing or invalid replacements never
overwrite the last valid in-memory code set. To intentionally disable all
NEWS entries, keep at least one row and set every row's enabled field to 0.
"""
global TRADED_NEWS_CODES, TRADED_CODES_SIGNATURE
global LAST_TRADED_CODES_RELOAD_MONOTONIC, TRADED_CODES_READY
now_mono = time.monotonic()
if (
not force
and now_mono - LAST_TRADED_CODES_RELOAD_MONOTONIC
< TRADED_CODES_RELOAD_SECONDS
):
return
LAST_TRADED_CODES_RELOAD_MONOTONIC = now_mono
if not TRADED_CODES_PATH.exists():
raise FileNotFoundError(f"Missing traded-code file: {TRADED_CODES_PATH}")
last_change_error: Optional[Exception] = None
for attempt in range(1, 4):
before = TRADED_CODES_PATH.stat()
signature = (before.st_mtime_ns, before.st_size)
if not force and TRADED_CODES_SIGNATURE == signature:
return
enabled_codes: set[str] = set()
seen_codes: set[str] = set()
data_rows = 0
try:
with TRADED_CODES_PATH.open(
"r", newline="", encoding="utf-8-sig"
) as handle:
reader = csv.DictReader(handle)
required = {"event_code", "enabled"}
if reader.fieldnames is None or not required.issubset(
set(reader.fieldnames)
):
raise ValueError(
"traded_codes.csv must contain columns: event_code,enabled"
)
for line_no, row in enumerate(reader, start=2):
if not any((value or "").strip() for value in row.values()):
continue
data_rows += 1
event_code = (row.get("event_code") or "").strip().lower()
enabled_text = (row.get("enabled") or "").strip()
if not event_code:
raise ValueError(
f"Blank event_code at traded_codes.csv line {line_no}"
)
if event_code in seen_codes:
raise ValueError(
f"Duplicate event_code {event_code!r} at "
f"traded_codes.csv line {line_no}"
)
normalized_enabled = enabled_text.lower()
valid_true = {"1", "true", "yes", "y", "on"}
valid_false = {"0", "false", "no", "n", "off"}
if normalized_enabled not in valid_true | valid_false:
raise ValueError(
f"Invalid enabled value {enabled_text!r} at "
f"traded_codes.csv line {line_no}"
)
seen_codes.add(event_code)
if normalized_enabled in valid_true:
enabled_codes.add(event_code)
if data_rows == 0:
raise ValueError(
"traded_codes.csv contains no data rows; retain rows with "
"enabled=0 to intentionally disable NEWS trading"
)
after = TRADED_CODES_PATH.stat()
after_signature = (after.st_mtime_ns, after.st_size)
if signature != after_signature:
last_change_error = RuntimeError(
"traded_codes.csv changed while it was being read"
)
time.sleep(0.10)
continue
TRADED_NEWS_CODES = frozenset(enabled_codes)
TRADED_CODES_SIGNATURE = after_signature
TRADED_CODES_READY = True
refresh_news_tradability()
future_tradable = sum(
event.tradable
and event.when_utc > utc_now() - timedelta(minutes=1)
for event in NEWS_EVENTS
)
LOGGER.info(
"Traded-code file loaded atomically: %d rows, %d enabled codes, "
"%d current/future tradable timetable minute(s); %s",
data_rows,
len(TRADED_NEWS_CODES),
future_tradable,
TRADED_CODES_PATH,
)
if not TRADED_NEWS_CODES:
LOGGER.warning(
"traded_codes.csv currently enables no event codes; NEWS "
"entries are disabled while session strategies continue."
)
return
except (FileNotFoundError, OSError) as exc:
last_change_error = exc
if attempt < 3:
time.sleep(0.10)
continue
raise
except Exception:
raise
if last_change_error is not None:
raise last_change_error
raise RuntimeError("Could not load a stable traded_codes.csv snapshot")
def event_codes_from_csv_row(row: dict[str, str]) -> tuple[str, ...]:
"""Extract normalized MT5 event codes from a timetable row.
Newer exporters may provide a dedicated event_codes column. The current EA
stores the same value inside notes as ``event_codes=a|b|c``. Both layouts are
accepted so the filter remains compatible with future exporter revisions.
"""
raw = (row.get("event_codes") or "").strip()
if not raw:
notes = row.get("notes") or ""
marker = "event_codes="
marker_at = notes.lower().find(marker)
if marker_at >= 0:
raw = notes[marker_at + len(marker):].split(";", 1)[0].strip()
normalized = {
code.strip().lower()
for code in raw.replace(",", "|").split("|")
if code.strip()
}
return tuple(sorted(normalized))
def traded_news_minute(event_codes: tuple[str, ...]) -> bool:
"""Return True when any component code is enabled in traded_codes.csv."""
return bool(TRADED_NEWS_CODES.intersection(event_codes))
def load_timetable(force: bool = False) -> None:
"""Load a complete, stable timetable.csv from the bot's own folder.
The MT5 exporter publishes the file atomically, but this loader still verifies
that size and modification time did not change while it was being parsed.
A failed reload raises without replacing NEWS_EVENTS, so the main loop keeps
the last valid in-memory calendar.
"""
global NEWS_EVENTS, TIMETABLE_SIGNATURE, LAST_TIMETABLE_RELOAD_MONOTONIC
global TIMETABLE_READY
now_mono = time.monotonic()
if not force and now_mono - LAST_TIMETABLE_RELOAD_MONOTONIC < TIMETABLE_RELOAD_SECONDS:
return
LAST_TIMETABLE_RELOAD_MONOTONIC = now_mono
if not TIMETABLE_PATH.exists():
raise FileNotFoundError(f"Missing timetable: {TIMETABLE_PATH}")
last_change_error: Optional[Exception] = None
for attempt in range(1, 4):
before = TIMETABLE_PATH.stat()
signature = (before.st_mtime_ns, before.st_size)
if not force and TIMETABLE_SIGNATURE == signature:
return
rows: list[NewsEvent] = []
seen_ids: set[str] = set()
seen_times: set[datetime] = set()
try:
with TIMETABLE_PATH.open("r", newline="", encoding="utf-8-sig") as handle:
reader = csv.DictReader(handle)
required = {"event_id", "event_name", "utc_datetime", "enabled"}
if reader.fieldnames is None or not required.issubset(set(reader.fieldnames)):
raise ValueError(
"timetable.csv must contain columns: "
"event_id,event_name,utc_datetime,enabled"
)
for line_no, row in enumerate(reader, start=2):
# Permit harmless completely blank lines at the end of a CSV.
if not any((value or "").strip() for value in row.values()):
continue
event_id = (row.get("event_id") or "").strip()
event_name = (row.get("event_name") or "").strip()
utc_text = (row.get("utc_datetime") or "").strip()
enabled_text = (row.get("enabled") or "").strip()
event_codes = event_codes_from_csv_row(row)
if not event_id or not event_name or not utc_text:
raise ValueError(
f"Blank event_id/event_name/utc_datetime at timetable line {line_no}"
)
if event_id in seen_ids:
raise ValueError(f"Duplicate event_id {event_id!r} at line {line_no}")
when = parse_utc(utc_text)
enabled = bool_from_csv(enabled_text)
if enabled and when in seen_times:
raise ValueError(
f"Duplicate enabled UTC minute {format_utc(when)}. "
"Combine simultaneous releases into one row."
)
seen_ids.add(event_id)
if enabled:
seen_times.add(when)
rows.append(
NewsEvent(
event_id=event_id,
event_name=event_name,
when_utc=when,
enabled=enabled,
event_codes=event_codes,
tradable=enabled and traded_news_minute(event_codes),
)
)
after = TIMETABLE_PATH.stat()
after_signature = (after.st_mtime_ns, after.st_size)
if signature != after_signature:
last_change_error = RuntimeError(
"timetable.csv changed while it was being read"
)
time.sleep(0.10)
continue
enabled_count = sum(event.enabled for event in rows)
future_enabled = sum(
event.enabled and event.when_utc > utc_now() - timedelta(minutes=1)
for event in rows
)
tradable_count = sum(event.tradable for event in rows)
future_tradable = sum(
event.tradable and event.when_utc > utc_now() - timedelta(minutes=1)
for event in rows
)
missing_code_count = sum(
event.enabled and not event.event_codes for event in rows
)
if not rows:
raise ValueError("timetable.csv contains no event rows")
if enabled_count == 0:
raise ValueError("timetable.csv contains no enabled events")
if future_enabled == 0:
raise ValueError("timetable.csv contains no current or future enabled events")
NEWS_EVENTS = sorted(rows, key=lambda event: event.when_utc)
TIMETABLE_SIGNATURE = after_signature
TIMETABLE_READY = True
LOGGER.info(
"Timetable loaded atomically: %d rows, %d enabled, %d current/future, "
"%d traded-code rows, %d traded-code current/future; %s",
len(NEWS_EVENTS),
enabled_count,
future_enabled,
tradable_count,
future_tradable,
TIMETABLE_PATH,
)
if missing_code_count:
LOGGER.warning(
"%d enabled timetable row(s) have no event_codes and are fail-closed "
"for NEWS trading.",
missing_code_count,
)
if future_tradable == 0:
LOGGER.warning(
"No current or future timetable minute matches the enabled codes in traded_codes.csv; "
"NEWS trading is disabled until a matching row appears."
)
return
except (FileNotFoundError, OSError) as exc:
last_change_error = exc
if attempt < 3:
time.sleep(0.10)
continue
raise
raise RuntimeError(
"Could not obtain a stable timetable.csv after three read attempts"
) from last_change_error
def today_session_event(
now_utc: datetime,
spec: DailySessionSpec,
) -> Optional[tuple[str, datetime, str]]:
"""Build today's event in its market-local calendar, then convert to UTC.
Host/VPS local time is never consulted. New York DST is handled by ZoneInfo;
Tokyo remains fixed at JST through the same timezone-aware mechanism.
"""
local_now = now_utc.astimezone(spec.timezone)
local_day = local_now.date()
if local_day.weekday() >= 5:
return None
local_event = datetime(
local_day.year,
local_day.month,
local_day.day,
spec.hour,
spec.minute,
tzinfo=spec.timezone,
)
event_utc = local_event.astimezone(UTC)
event_id = f"{spec.event_id_prefix}_{local_day:%Y%m%d}"
return event_id, event_utc, f"{spec.display_name} {local_day.isoformat()}"
def all_daily_session_events(
now_utc: datetime,
) -> list[tuple[DailySessionSpec, str, datetime, str]]:
events: list[tuple[DailySessionSpec, str, datetime, str]] = []
for spec in DAILY_SESSION_SPECS:
event = today_session_event(now_utc, spec)
if event is not None:
event_id, event_utc, event_name = event
events.append((spec, event_id, event_utc, event_name))
return events
def seconds_to_nearest_event(now_utc: datetime) -> float:
candidates: list[float] = []
for event in NEWS_EVENTS:
if event.tradable and event.when_utc >= now_utc - timedelta(seconds=5):
candidates.append(abs((event.when_utc - now_utc).total_seconds()))
break
for _spec, _event_id, event_utc, _event_name in all_daily_session_events(now_utc):
candidates.append(abs((event_utc - now_utc).total_seconds()))
return min(candidates) if candidates else 9_999.0
# =============================================================================
# MT5 CONNECTION AND SYMBOL HELPERS
# =============================================================================
def mt5_error_text() -> str:
try:
return repr(mt5.last_error())
except Exception:
return "unknown MT5 error"
def connect_mt5() -> None:
global SYMBOL, SYMBOL_INFO
initialized = (
mt5.initialize(path=MT5_TERMINAL_PATH)
if MT5_TERMINAL_PATH
else mt5.initialize()
)
if not initialized:
raise RuntimeError(f"mt5.initialize() failed: {mt5_error_text()}")
account = mt5.account_info()
terminal = mt5.terminal_info()
if account is None or terminal is None:
raise RuntimeError(f"MT5 account/terminal unavailable: {mt5_error_text()}")
if hasattr(terminal, "connected") and not bool(terminal.connected):
raise RuntimeError("MT5 terminal is open but not connected to the broker server.")
if EXPECTED_ACCOUNT_LOGIN is not None and int(account.login) != int(EXPECTED_ACCOUNT_LOGIN):
raise RuntimeError(
f"Wrong MT5 account: connected to {account.login}, expected {EXPECTED_ACCOUNT_LOGIN}."
)
if EXPECTED_SERVER is not None and str(account.server) != str(EXPECTED_SERVER):
raise RuntimeError(
f"Wrong MT5 server: connected to {account.server!r}, expected {EXPECTED_SERVER!r}."
)
hedging_mode = getattr(mt5, "ACCOUNT_MARGIN_MODE_RETAIL_HEDGING", 2)
if account.margin_mode != hedging_mode:
raise RuntimeError(
"A hedging MT5 account is mandatory. This account is netting/exchange mode, "
"so the five XAUUSD strategy positions cannot be managed independently."
)
if REQUIRE_USD_ACCOUNT and str(account.currency).upper() != "USD":
raise RuntimeError(
f"Account currency is {account.currency}, but lot formulas are defined for USD balance."
)
if not account.trade_allowed or not account.trade_expert:
raise RuntimeError("Trading or expert trading is disabled on the account.")
if hasattr(terminal, "trade_allowed") and not terminal.trade_allowed:
raise RuntimeError("Algo Trading is disabled in the MT5 terminal.")
SYMBOL = resolve_symbol(PREFERRED_SYMBOL)
if not mt5.symbol_select(SYMBOL, True):
raise RuntimeError(f"Could not select {SYMBOL} in Market Watch: {mt5_error_text()}")
SYMBOL_INFO = mt5.symbol_info(SYMBOL)
if SYMBOL_INFO is None:
raise RuntimeError(f"symbol_info({SYMBOL}) failed: {mt5_error_text()}")
LOGGER.info(
"Connected: account=%s server=%s balance=%.2f currency=%s symbol=%s digits=%d point=%g",
account.login,
account.server,
account.balance,
account.currency,
SYMBOL,
SYMBOL_INFO.digits,
SYMBOL_INFO.point,
)
try:
current_tick(require_fresh=False)
except Exception as exc:
LOGGER.warning(
"MT5 broker clock offset is not yet available; it will be detected "
"from the next usable live quote: %s",
exc,
)
def ensure_connection() -> bool:
"""Verify actual broker connectivity, not merely cached terminal objects."""
global SYMBOL_INFO
account = mt5.account_info()
terminal = mt5.terminal_info()
connected = (
account is not None
and terminal is not None
and (not hasattr(terminal, "connected") or bool(terminal.connected))
)
if connected:
info = mt5.symbol_info(SYMBOL)
if info is not None:
SYMBOL_INFO = info
return True
LOGGER.error("MT5 broker connection unavailable; attempting reconnect: %s", mt5_error_text())
try:
mt5.shutdown()
except Exception:
pass
try:
connect_mt5()
return True
except Exception:
LOGGER.exception("MT5 reconnect failed.")
return False
def resolve_symbol(preferred: str) -> str:
exact = mt5.symbol_info(preferred)
if exact is not None:
return preferred
symbols = mt5.symbols_get()
if not symbols:
raise RuntimeError(f"No MT5 symbols available: {mt5_error_text()}")
preferred_upper = preferred.upper()
candidates = [s.name for s in symbols if preferred_upper in s.name.upper()]
if not candidates:
raise RuntimeError(f"Could not find a symbol containing {preferred!r}.")
# Prefer names beginning with XAUUSD and then the shortest broker suffix.
candidates.sort(key=lambda name: (not name.upper().startswith(preferred_upper), len(name), name))
chosen = candidates[0]
LOGGER.warning("Preferred symbol %s not found; using broker symbol %s.", preferred, chosen)
return chosen
def _mt5_object_time_msc(value: Any) -> int:
"""Return an MT5 object's timestamp in the broker's own epoch space.
RoboForex may expose tick/position/deal timestamps shifted by the broker-server
UTC offset. Arithmetic between two MT5 timestamps is exact because the same
offset is present in both values.
"""
raw_msc = int(getattr(value, "time_msc", 0) or 0)
if raw_msc <= 0:
raw_msc = int(float(getattr(value, "time", 0) or 0) * 1000)
return raw_msc
def _broker_clock_datetime(value: Any) -> datetime:
"""Represent an MT5 timestamp as an aware broker-clock datetime.
This is only an internal coordinate for MT5 history queries and completed-M1
comparisons. It is not used for strategy scheduling or presented as true UTC.
"""
return datetime.fromtimestamp(_mt5_object_time_msc(value) / 1000.0, UTC)
def _broker_elapsed_seconds(earlier: Any, later: Any) -> float:
"""Elapsed real seconds between two MT5 objects, independent of timezone."""
return (_mt5_object_time_msc(later) - _mt5_object_time_msc(earlier)) / 1000.0
def _broker_offset_candidate(clock_tick: Any) -> tuple[int, float]:
"""Return the live whole-hour broker offset and quote age in seconds.
Some RoboForex terminals expose MT5 epoch fields in broker-server wall-clock
space rather than true UTC. A live quote supplies both the broker epoch and a
current wall-clock reference, so their difference identifies the active server
offset without hard-coding summer/winter time.
"""
raw_msc = _mt5_object_time_msc(clock_tick)
if raw_msc <= 0:
raise RuntimeError(f"Tick for {SYMBOL} has no usable timestamp")
wall_msc = int(time.time() * 1000)
offset_seconds = int(round((raw_msc - wall_msc) / 3_600_000)) * 3600
quote_age = (wall_msc + offset_seconds * 1000 - raw_msc) / 1000.0
return offset_seconds, quote_age
def broker_utc_offset_seconds(clock_tick: Any = None) -> int:
"""Return the dynamically detected broker-clock offset from true UTC.
The cached value is refreshed from every usable live tick. It is deliberately
constrained to a plausible civil-time range so a weekend-old quote cannot be
mistaken for a multi-day timezone offset.
"""
global BROKER_UTC_OFFSET_SECONDS
tick = clock_tick if clock_tick is not None else mt5.symbol_info_tick(SYMBOL)
if tick is None:
if BROKER_UTC_OFFSET_SECONDS is not None:
return BROKER_UTC_OFFSET_SECONDS
raise RuntimeError(f"Cannot detect {SYMBOL} broker clock: {mt5_error_text()}")
candidate, quote_age = _broker_offset_candidate(tick)
if abs(candidate) > MAX_BROKER_UTC_OFFSET_SECONDS:
if BROKER_UTC_OFFSET_SECONDS is not None:
return BROKER_UTC_OFFSET_SECONDS
raise RuntimeError(
f"Implausible broker UTC offset {candidate / 3600:+.1f}h; "
f"the latest {SYMBOL} quote is probably too old"
)
if BROKER_UTC_OFFSET_SECONDS != candidate:
previous = BROKER_UTC_OFFSET_SECONDS
BROKER_UTC_OFFSET_SECONDS = candidate
COMPLETED_CLOSE_CACHE.clear()
M1_RATE_CACHE.clear()
M1_RATE_RETRY_AFTER.clear()
M1_RATE_SOURCE_WARNINGS.clear()
if previous is None:
LOGGER.info(
"Detected MT5 broker clock offset UTC%+.0f (quote age %.2fs)",
candidate / 3600.0,
quote_age,
)
else:
LOGGER.warning(
"MT5 broker clock offset changed UTC%+.0f -> UTC%+.0f; "
"cleared time-dependent caches.",
previous / 3600.0,
candidate / 3600.0,
)
return candidate
def utc_to_broker_clock(value_utc: datetime, clock_tick: Any = None) -> datetime:
"""Convert a true UTC datetime into the MT5 broker-clock coordinate."""
if value_utc.tzinfo is None:
raise ValueError("UTC-to-broker conversion requires an aware datetime")
offset = broker_utc_offset_seconds(clock_tick)
return value_utc.astimezone(UTC) + timedelta(seconds=offset)
def broker_clock_to_utc(value_broker: datetime, clock_tick: Any = None) -> datetime:
"""Convert an MT5 broker-clock coordinate into true UTC."""
if value_broker.tzinfo is None:
raise ValueError("Broker-to-UTC conversion requires an aware datetime")
offset = broker_utc_offset_seconds(clock_tick)
return (value_broker.astimezone(UTC) - timedelta(seconds=offset)).astimezone(UTC)
def _mt5_object_utc_datetime(value: Any, clock_tick: Any = None) -> datetime:
"""Return an MT5 object's broker-shifted epoch as a true UTC datetime."""
return broker_clock_to_utc(_broker_clock_datetime(value), clock_tick)
class TransientTickError(RuntimeError):
"""Temporary missing/stale quote condition that is safe to retry."""
def current_tick(require_fresh: bool = True) -> Any:
tick = mt5.symbol_info_tick(SYMBOL)
if tick is None or tick.bid <= 0 or tick.ask <= 0:
raise TransientTickError(f"No valid tick for {SYMBOL}: {mt5_error_text()}")
raw_msc = _mt5_object_time_msc(tick)
offset = broker_utc_offset_seconds(tick)
wall_msc = int(time.time() * 1000)
age = (wall_msc + offset * 1000 - raw_msc) / 1000.0
if require_fresh and abs(age) > MAX_TICK_AGE_SECONDS:
direction = "old" if age >= 0 else "ahead of wall clock"
raise TransientTickError(
f"Stale {SYMBOL} tick: {abs(age):.1f} seconds {direction}"
)
return tick
def normalize_price(price: float) -> float:
return round(float(price), int(SYMBOL_INFO.digits))
def volume_decimals(step: float) -> int:
text = f"{step:.10f}".rstrip("0")
return len(text.split(".")[1]) if "." in text else 0
def normalize_volume(raw: float, minimum: float) -> float:
info = SYMBOL_INFO
step = float(info.volume_step)
broker_min = float(info.volume_min)
broker_max = float(info.volume_max)
effective_min = max(minimum, broker_min)
effective_max = min(broker_max, MAX_LOT) if MAX_LOT is not None else broker_max
normalized_min = math.ceil((effective_min - 1e-12) / step) * step
normalized_max = math.floor((effective_max + 1e-12) / step) * step
if normalized_max + 1e-12 < normalized_min:
raise RuntimeError(
f"Volume constraints are impossible: min={normalized_min}, max={normalized_max}, step={step}"
)
policy = str(MINIMUM_LOT_POLICY).strip().upper()
if policy not in {"SKIP", "FORCE_MIN"}:
raise RuntimeError(
f"MINIMUM_LOT_POLICY must be 'SKIP' or 'FORCE_MIN', got {MINIMUM_LOT_POLICY!r}"
)
# Floor to avoid exceeding the requested percentage risk.
floored = math.floor((raw + 1e-12) / step) * step
if floored + 1e-12 < normalized_min:
if policy == "SKIP":
return 0.0
floored = normalized_min
volume = min(floored, normalized_max)
return round(volume, volume_decimals(step))
def symbol_contract_size() -> float:
value = float(getattr(SYMBOL_INFO, "trade_contract_size", 0.0) or 0.0)
if not math.isfinite(value) or value <= 0:
raise RuntimeError(f"Invalid {SYMBOL} contract size: {value!r}")
return value
def nominal_risk_per_lot(config: StrategyConfig) -> float:
return (float(config.initial_sl) + float(RISK_COST_BUFFER_PER_OUNCE)) * symbol_contract_size()
def nominal_risk_for_volume(config: StrategyConfig, volume: float) -> float:
return nominal_risk_per_lot(config) * float(volume)
def lot_for_config(config: StrategyConfig) -> LotSizingResult:
account = mt5.account_info()
if account is None:
raise RuntimeError(f"account_info() failed: {mt5_error_text()}")
sizing_balance = float(account.balance)
if not math.isfinite(sizing_balance) or sizing_balance <= 0:
raise RuntimeError(f"Invalid account balance for sizing: {sizing_balance!r}")
if not math.isfinite(config.risk_percent) or config.risk_percent <= 0:
raise RuntimeError(f"Invalid {config.name} risk percentage: {config.risk_percent!r}")
risk_per_lot = nominal_risk_per_lot(config)
target_risk = sizing_balance * float(config.risk_percent) / 100.0
raw_lot = target_risk / risk_per_lot
lot = normalize_volume(raw_lot, config.min_lot)
actual_risk = nominal_risk_for_volume(config, lot) if lot > 0 else 0.0
return LotSizingResult(
sizing_balance=sizing_balance,
target_risk_usd=target_risk,
risk_per_lot_usd=risk_per_lot,
raw_lot=raw_lot,
lot=lot,
actual_risk_usd=actual_risk,
)
def estimated_sl_loss(
order_type: int,
volume: float,
entry_price: float,
stop_price: float,
) -> Optional[float]:
"""Return broker-calculated loss at the initial SL, when available."""
try:
value = mt5.order_calc_profit(
int(order_type),
SYMBOL,
float(volume),
float(entry_price),
float(stop_price),
)
except Exception:
return None
if value is None:
return None
return abs(min(0.0, float(value)))
def orders_for_magic(magic: int) -> list[Any]:
orders = mt5.orders_get(symbol=SYMBOL)
if orders is None:
return []
return [order for order in orders if int(order.magic) == magic]
def positions_for_magic(magic: int) -> list[Any]:
positions = mt5.positions_get(symbol=SYMBOL)
if positions is None:
return []
return [position for position in positions if int(position.magic) == magic]
def subsystem_busy(config: StrategyConfig) -> bool:
return bool(orders_for_magic(config.magic) or positions_for_magic(config.magic))
def group_id_from_comment(comment: str, config: StrategyConfig) -> Optional[str]:
"""Extract the stable prefix+YYMMDDHHMM event group from an MT5 comment."""
text = str(comment or "").strip()
expected_length = 1 + 10
if len(text) < expected_length:
return None
candidate = text[:expected_length]
if candidate[0] != config.group_prefix or not candidate[1:].isdigit():
return None
return candidate
def orders_for_group(config: StrategyConfig, group_id: str) -> list[Any]:
return [
order
for order in orders_for_magic(config.magic)
if str(getattr(order, "comment", "")).strip().startswith(group_id)
]
def positions_for_group(config: StrategyConfig, group_id: str) -> list[Any]:
return [
position
for position in positions_for_magic(config.magic)
if str(getattr(position, "comment", "")).strip().startswith(group_id)
]
def event_setup_busy(config: StrategyConfig, event_utc: datetime) -> bool:
"""NEWS is busy only for this event minute; sessions remain subsystem-wide."""
if config.name != "NEWS":
return subsystem_busy(config)
group_id = make_group_id(config, event_utc)
return bool(orders_for_group(config, group_id) or positions_for_group(config, group_id))
def portfolio_initial_risk_snapshot() -> tuple[float, list[str]]:
"""Return aggregate nominal initial-stop risk for this bot's live groups.
An intact pending OCO pair counts once. Duplicate positions are summed. If a
position and an unremoved sibling order coexist, both are counted because a
second fill remains possible until reconciliation removes the sibling.
"""
total = 0.0
details: list[str] = []
for config in ALL_CONFIGS:
position_risk: dict[str, float] = {}
order_risk: dict[str, float] = {}
for position in positions_for_magic(config.magic):
group_id = group_id_from_comment(getattr(position, "comment", ""), config)
key = group_id or f"P{config.magic}:{int(position.ticket)}"
risk = nominal_risk_for_volume(config, float(position.volume))
position_risk[key] = position_risk.get(key, 0.0) + risk
for order in orders_for_magic(config.magic):
group_id = group_id_from_comment(getattr(order, "comment", ""), config)
key = group_id or f"O{config.magic}:{int(order.ticket)}"
risk = nominal_risk_for_volume(config, float(order.volume_current))
# Two intact OCO legs represent one intended position, so use max.
order_risk[key] = max(order_risk.get(key, 0.0), risk)
for key in set(position_risk) | set(order_risk):
risk = position_risk.get(key, 0.0) + order_risk.get(key, 0.0)
total += risk
details.append(f"{config.name}:{key}={risk:.2f}")
return total, sorted(details)
def risk_block_log_once(
event_key: str,
reason: str,
message: str,
*args: Any,
) -> None:
key = (event_key, reason)
if key in RISK_BLOCK_WARNING_KEYS:
return
RISK_BLOCK_WARNING_KEYS.add(key)
LOGGER.warning(message, *args)
def finalize_risk_skip_if_due(
state: dict[str, Any],
event_key: str,
event_name: str,
event_utc: datetime,
status: str,
detail: str,
) -> None:
remaining = (event_utc - utc_now()).total_seconds()
if remaining <= RISK_SKIP_FINALIZE_SECONDS:
mark_event(state, event_key, status, event_utc, detail)
LOGGER.warning("Risk sizing skipped %s event %s: %s", status, event_name, detail)
# =============================================================================
# ORDER SENDING
# =============================================================================
def accepted_retcode(retcode: int) -> bool:
accepted = {
getattr(mt5, "TRADE_RETCODE_DONE", 10009),
getattr(mt5, "TRADE_RETCODE_PLACED", 10008),
getattr(mt5, "TRADE_RETCODE_DONE_PARTIAL", 10010),
}
return int(retcode) in accepted
def connection_retcode(retcode: int) -> bool:
return int(retcode) == int(getattr(mt5, "TRADE_RETCODE_CONNECTION", 10031))
def transient_retcode(retcode: int) -> bool:
names = (
"TRADE_RETCODE_REQUOTE",
"TRADE_RETCODE_PRICE_CHANGED",
"TRADE_RETCODE_PRICE_OFF",
"TRADE_RETCODE_TIMEOUT",
"TRADE_RETCODE_CONNECTION",
"TRADE_RETCODE_TOO_MANY_REQUESTS",
"TRADE_RETCODE_LOCKED",
)
values = {getattr(mt5, name, -999999) for name in names}
return int(retcode) in values
def order_result_text(result: Any) -> str:
if result is None:
return f"None, last_error={mt5_error_text()}"
return (
f"retcode={result.retcode} comment={getattr(result, 'comment', '')!r} "
f"order={getattr(result, 'order', 0)} deal={getattr(result, 'deal', 0)} "
f"price={getattr(result, 'price', 0)}"
)
def recover_matching_request(request: dict[str, Any]) -> Any:
"""Recover an order accepted by the server after an ambiguous client timeout.
Retrying an order_send blindly after a timeout can create a duplicate leg. Every
pending leg has a unique magic/comment pair, so inspect live orders and positions
before any retry. A lightweight result object is returned when the request is
already present on the server.
"""
comment = str(request.get("comment", "")).strip()
magic = int(request.get("magic", 0) or 0)
symbol = str(request.get("symbol", SYMBOL))
if not comment or magic <= 0:
return None
request_type = int(request.get("type", -1))
orders = mt5.orders_get(symbol=symbol)
if orders is not None:
for order in orders:
if int(getattr(order, "magic", 0)) != magic:
continue
if str(getattr(order, "comment", "")).strip() != comment:
continue
if request_type >= 0 and int(getattr(order, "type", -2)) != request_type:
continue
LOGGER.warning(
"Recovered already accepted pending order after ambiguous send: "
"ticket=%s comment=%s",
order.ticket,
comment,
)
return SimpleNamespace(
retcode=getattr(mt5, "TRADE_RETCODE_PLACED", 10008),
comment="recovered existing pending order",
order=int(order.ticket),
deal=0,
price=float(getattr(order, "price_open", 0.0)),
)
positions = mt5.positions_get(symbol=symbol)
if positions is not None:
buy_types = {
int(getattr(mt5, "ORDER_TYPE_BUY", 0)),
int(getattr(mt5, "ORDER_TYPE_BUY_LIMIT", 2)),
int(getattr(mt5, "ORDER_TYPE_BUY_STOP", 4)),
int(getattr(mt5, "ORDER_TYPE_BUY_STOP_LIMIT", 6)),
}
expected_position_type = (
int(getattr(mt5, "POSITION_TYPE_BUY", 0))
if request_type in buy_types
else int(getattr(mt5, "POSITION_TYPE_SELL", 1))
)
for position in positions:
if int(getattr(position, "magic", 0)) != magic:
continue
if str(getattr(position, "comment", "")).strip() != comment:
continue
if int(getattr(position, "type", -2)) != expected_position_type:
continue
LOGGER.warning(
"Recovered position filled after ambiguous order send: "
"ticket=%s comment=%s",
position.ticket,
comment,
)
return SimpleNamespace(
retcode=getattr(mt5, "TRADE_RETCODE_DONE", 10009),
comment="recovered filled position",
order=int(position.ticket),
deal=0,
price=float(getattr(position, "price_open", 0.0)),
)
return None
def check_and_send(request: dict[str, Any], retries: int = 1) -> Any:
"""Validate and send one request without duplicating ambiguous fills.
A broker-connection failure is not a filling-policy problem. It aborts this
request immediately, resets the Python/terminal bridge, and lets the main loop
reconnect and retry the complete pair while the placement window remains open.
"""
last_result = None
for attempt in range(retries + 1):
recovered = recover_matching_request(request)
if recovered is not None:
return recovered
check = mt5.order_check(request)
if check is None:
LOGGER.error("order_check returned None: %s", mt5_error_text())
terminal = mt5.terminal_info()
if terminal is None or (
hasattr(terminal, "connected") and not bool(terminal.connected)
):
try:
mt5.shutdown()
except Exception:
pass
return SimpleNamespace(
retcode=getattr(mt5, "TRADE_RETCODE_CONNECTION", 10031),
comment="order_check unavailable: broker connection lost",
order=0,
deal=0,
price=0.0,
)
return None
if int(check.retcode) != 0:
LOGGER.error(
"order_check rejected request: retcode=%s comment=%r request=%s",
check.retcode,
check.comment,
request,
)
if connection_retcode(int(check.retcode)):
try:
mt5.shutdown()
except Exception:
pass
return check
return None
last_result = mt5.order_send(request)
if last_result is not None and accepted_retcode(last_result.retcode):
return last_result
LOGGER.error("order_send failed: %s request=%s", order_result_text(last_result), request)
if last_result is not None and connection_retcode(last_result.retcode):
LOGGER.error(
"Broker connection failure during order submission; "
"aborting policy fallbacks and forcing reconnect."
)
try:
mt5.shutdown()
except Exception:
pass
return last_result
# A timeout can mean acceptance with a lost acknowledgement. Inspect the
# server before retrying; comments/magic numbers make recovery idempotent.
time.sleep(0.10 * (attempt + 1))
recovered = recover_matching_request(request)
if recovered is not None:
return recovered
if (
last_result is not None
and not transient_retcode(last_result.retcode)
) or attempt >= retries:
break
return last_result
def send_pending_request(request: dict[str, Any]) -> Any:
"""Send a pending leg with broker-compatible expiry/filling fallbacks.
Specified server expiry is preferred. If unsupported, the bot uses GTC and
enforces the same expiry locally. Filling policies vary by broker/symbol, so
RETURN, IOC and FOK are tried without ever duplicating a successfully accepted
magic/comment pair.
"""
time_variants: list[dict[str, Any]] = [dict(request)]
if request.get("type_time") == mt5.ORDER_TIME_SPECIFIED:
fallback = dict(request)
fallback["type_time"] = mt5.ORDER_TIME_GTC
fallback.pop("expiration", None)
time_variants.append(fallback)
filling_candidates = [
request.get("type_filling", getattr(mt5, "ORDER_FILLING_RETURN", 2)),
getattr(mt5, "ORDER_FILLING_RETURN", 2),
getattr(mt5, "ORDER_FILLING_IOC", 1),
getattr(mt5, "ORDER_FILLING_FOK", 0),
]
last_result = None
warned_gtc = False
tried: set[tuple[int, int]] = set()
for variant in time_variants:
if variant.get("type_time") == mt5.ORDER_TIME_GTC and not warned_gtc:
LOGGER.warning(
"Specified pending-order expiration is unavailable; trying GTC. "
"The bot will enforce expiry locally."
)
warned_gtc = True
for filling in filling_candidates:
key = (int(variant.get("type_time", -1)), int(filling))
if key in tried:
continue
tried.add(key)
candidate = dict(variant)
candidate["type_filling"] = int(filling)
last_result = check_and_send(candidate, retries=1)
if last_result is not None and accepted_retcode(last_result.retcode):
return last_result
if last_result is not None and connection_retcode(last_result.retcode):
return last_result
return last_result
def remove_pending_order(order: Any, reason: str) -> bool:
request = {
"action": mt5.TRADE_ACTION_REMOVE,
"order": int(order.ticket),
"symbol": SYMBOL,
"magic": int(order.magic),
"comment": reason[:31],
}
result = mt5.order_send(request)
ok = result is not None and accepted_retcode(result.retcode)
if ok:
LOGGER.info("Removed pending order ticket=%s reason=%s", order.ticket, reason)
else:
LOGGER.error(
"Failed to remove pending order ticket=%s: %s",
order.ticket,
order_result_text(result),
)
return ok
def cancel_all_pending(config: StrategyConfig, reason: str) -> None:
for order in orders_for_magic(config.magic):
remove_pending_order(order, reason)
def pair_comment(group_id: str, side: str) -> str:
# Keep below the typical MT5 31-character comment limit.
return f"{group_id}{side}"[:31]
def make_group_id(config: StrategyConfig, event_utc: datetime) -> str:
# N2608071230 / Y2608071230, plus B/S in the order comment.
return f"{config.group_prefix}{event_utc:%y%m%d%H%M}"
def place_oco_pair(
config: StrategyConfig,
event_id: str,
event_name: str,
event_utc: datetime,
state: dict[str, Any],
) -> bool:
group_id = make_group_id(config, event_utc)
if event_setup_busy(config, event_utc):
return False
tick = current_tick(require_fresh=True)
expiry_utc = event_utc + timedelta(seconds=config.expiry_seconds)
buy_price = normalize_price(float(tick.ask) + config.offset)
sell_price = normalize_price(float(tick.bid) - config.offset)
buy_sl = normalize_price(buy_price - config.initial_sl)
sell_sl = normalize_price(sell_price + config.initial_sl)
sizing = lot_for_config(config)
lot = sizing.lot
balance = sizing.sizing_balance
if lot <= 0:
detail = (
f"{config.name} raw lot {sizing.raw_lot:.4f} is below broker/configured "
f"minimum {max(config.min_lot, float(SYMBOL_INFO.volume_min)):.2f}; "
f"policy={MINIMUM_LOT_POLICY}"
)
risk_block_log_once(
event_id,
"minimum_lot",
"%s event %s waiting/skipping: %s",
config.name,
event_name,
detail,
)
finalize_risk_skip_if_due(
state, event_id, event_name, event_utc, "skipped_minimum_lot", detail
)
return False
current_portfolio_risk, risk_details = portfolio_initial_risk_snapshot()
cap_usd = balance * PORTFOLIO_RISK_CAP_PERCENT / 100.0
projected_risk = current_portfolio_risk + sizing.actual_risk_usd
if projected_risk > cap_usd + 1e-9:
detail = (
f"current={current_portfolio_risk:.2f}, proposed={sizing.actual_risk_usd:.2f}, "
f"projected={projected_risk:.2f}, cap={cap_usd:.2f} "
f"({PORTFOLIO_RISK_CAP_PERCENT:.2f}% of balance {balance:.2f})"
)
risk_block_log_once(
event_id,
"portfolio_cap",
"%s event %s blocked by portfolio risk cap: %s | live=%s",
config.name,
event_name,
detail,
"; ".join(risk_details) if risk_details else "none",
)
finalize_risk_skip_if_due(
state, event_id, event_name, event_utc, "skipped_portfolio_risk", detail
)
return False
buy_risk = estimated_sl_loss(mt5.ORDER_TYPE_BUY, lot, buy_price, buy_sl)
sell_risk = estimated_sl_loss(mt5.ORDER_TYPE_SELL, lot, sell_price, sell_sl)
broker_estimated_risk = (
max(value for value in (buy_risk, sell_risk) if value is not None)
if buy_risk is not None or sell_risk is not None
else None
)
LOGGER.info(
"%s risk sizing | target=%.2f%% | balance=%.2f | raw=%.4f | lot=%.2f | "
"nominal initial risk=%.2f (%.2f%%) | portfolio %.2f -> %.2f / %.2f",
config.name,
config.risk_percent,
balance,
sizing.raw_lot,
lot,
sizing.actual_risk_usd,
100.0 * sizing.actual_risk_usd / balance,
current_portfolio_risk,
projected_risk,
cap_usd,
)
if broker_estimated_risk is not None:
LOGGER.info(
"%s broker order_calc_profit initial-SL estimate=%.2f; sizing model with "
"$%.2f/oz cost buffer=%.2f",
config.name,
broker_estimated_risk,
RISK_COST_BUFFER_PER_OUNCE,
sizing.actual_risk_usd,
)
min_distance = max(
int(getattr(SYMBOL_INFO, "trade_stops_level", 0)),
int(getattr(SYMBOL_INFO, "trade_freeze_level", 0)),
) * float(SYMBOL_INFO.point)
if buy_price - float(tick.ask) < min_distance or float(tick.bid) - sell_price < min_distance:
LOGGER.error(
"%s pair invalid: broker minimum pending distance %.5f exceeds configured offset %.5f",
config.name,
min_distance,
config.offset,
)
return False
expiry_broker = utc_to_broker_clock(expiry_utc, tick)
base = {
"action": mt5.TRADE_ACTION_PENDING,
"symbol": SYMBOL,
"volume": lot,
"deviation": MAX_DEVIATION_POINTS,
"magic": config.magic,
"type_time": mt5.ORDER_TIME_SPECIFIED,
"expiration": int(expiry_broker.timestamp()),
"type_filling": mt5.ORDER_FILLING_RETURN,
}
buy_request = {
**base,
"type": mt5.ORDER_TYPE_BUY_STOP,
"price": buy_price,
"sl": buy_sl,
"tp": 0.0,
"comment": pair_comment(group_id, "B"),
}
sell_request = {
**base,
"type": mt5.ORDER_TYPE_SELL_STOP,
"price": sell_price,
"sl": sell_sl,
"tp": 0.0,
"comment": pair_comment(group_id, "S"),
}
buy_result = send_pending_request(buy_request)
if buy_result is None or not accepted_retcode(buy_result.retcode):
return False
sell_result = send_pending_request(sell_request)
if sell_result is None or not accepted_retcode(sell_result.retcode):
LOGGER.critical(
"%s OCO second leg failed after first leg accepted. group=%s first=%s second=%s",
config.name,
group_id,
order_result_text(buy_result),
order_result_text(sell_result),
)
# Cancel the first leg if it is still pending. If it filled during the tiny
# submission interval, it remains protected by its server-side initial SL.
for order in orders_for_group(config, group_id):
remove_pending_order(order, "PAIR_FAIL")
return bool(positions_for_group(config, group_id))
group_record = {
"strategy": config.name,
"event_id": event_id,
"event_name": event_name,
"event_utc": format_utc(event_utc),
"expiry_utc": format_utc(expiry_utc),
"lot": lot,
"risk_percent": config.risk_percent,
"sizing_balance": balance,
"raw_lot": sizing.raw_lot,
"initial_risk_usd": sizing.actual_risk_usd,
"portfolio_risk_before_usd": current_portfolio_risk,
"portfolio_risk_after_usd": projected_risk,
"portfolio_risk_cap_usd": cap_usd,
"reference_bid": float(tick.bid),
"reference_ask": float(tick.ask),
"buy_order": int(getattr(buy_result, "order", 0)),
"sell_order": int(getattr(sell_result, "order", 0)),
}
state["active_groups"][group_id] = group_record
save_state(state)
LOGGER.info(
"%s OCO placed | event=%s | group=%s | lot=%.2f | ref=%.2f/%.2f | "
"buy_stop=%.2f SL=%.2f | sell_stop=%.2f SL=%.2f | expiry=%s",
config.name,
event_name,
group_id,
lot,
tick.bid,
tick.ask,
buy_price,
buy_sl,
sell_price,
sell_sl,
format_utc(expiry_utc),
)
return True
def close_position(position: Any, comment: str) -> bool:
tick = current_tick(require_fresh=True)
position_type_buy = int(position.type) == int(mt5.POSITION_TYPE_BUY)
order_type = mt5.ORDER_TYPE_SELL if position_type_buy else mt5.ORDER_TYPE_BUY
price = float(tick.bid) if position_type_buy else float(tick.ask)
# RoboForex market execution usually accepts IOC; retry supported policies.
filling_candidates = [
getattr(mt5, "ORDER_FILLING_IOC", 1),
getattr(mt5, "ORDER_FILLING_FOK", 0),
getattr(mt5, "ORDER_FILLING_RETURN", 2),
]
tried: set[int] = set()
for filling in filling_candidates:
if filling in tried:
continue
tried.add(filling)
request = {
"action": mt5.TRADE_ACTION_DEAL,
"symbol": SYMBOL,
"position": int(position.ticket),
"volume": float(position.volume),
"type": order_type,
"price": normalize_price(price),
"deviation": MAX_DEVIATION_POINTS,
"magic": int(position.magic),
"comment": comment[:31],
"type_time": mt5.ORDER_TIME_GTC,
"type_filling": filling,
}
result = mt5.order_send(request)
if result is not None and accepted_retcode(result.retcode):
LOGGER.info(
"Closed position ticket=%s strategy=%s reason=%s result=%s",
position.ticket,
CONFIG_BY_MAGIC.get(int(position.magic), NEWS_CONFIG).name,
comment,
order_result_text(result),
)
return True
LOGGER.error(
"Close attempt failed ticket=%s filling=%s: %s",
position.ticket,
filling,
order_result_text(result),
)
time.sleep(0.1)
tick = current_tick(require_fresh=True)
price = float(tick.bid) if position_type_buy else float(tick.ask)
return False
def modify_position_sl(position: Any, new_sl: float) -> bool:
request = {
"action": mt5.TRADE_ACTION_SLTP,
"symbol": SYMBOL,
"position": int(position.ticket),
"sl": normalize_price(new_sl),
"tp": float(position.tp) if float(position.tp) > 0 else 0.0,
"magic": int(position.magic),
"comment": "TRAIL",
}
result = mt5.order_send(request)
ok = result is not None and accepted_retcode(result.retcode)
if ok:
LOGGER.info(
"Position SL modified ticket=%s old_sl=%.2f new_sl=%.2f",
position.ticket,
float(position.sl),
new_sl,
)
else:
LOGGER.error(
"Position SL modification failed ticket=%s new_sl=%.2f: %s",
position.ticket,
new_sl,
order_result_text(result),
)
return ok
def ensure_initial_protection(position: Any, config: StrategyConfig) -> bool:
"""Ensure a recovered/live position never remains without its initial SL.
Pending orders normally install the SL server-side. This guard repairs a
missing or manually widened SL after a restart. If the market is already at
or beyond the intended stop, the position is closed instead of widening risk.
"""
buy = int(position.type) == int(mt5.POSITION_TYPE_BUY)
target = normalize_price(
float(position.price_open) - config.initial_sl
if buy
else float(position.price_open) + config.initial_sl
)
old_sl = float(getattr(position, "sl", 0.0) or 0.0)
point = float(SYMBOL_INFO.point)
protected = (old_sl >= target - point / 2) if buy else (0 < old_sl <= target + point / 2)
if protected:
return True
tick = current_tick(require_fresh=True)
market = float(tick.bid) if buy else float(tick.ask)
if (buy and market <= target) or ((not buy) and market >= target):
LOGGER.critical(
"%s position ticket=%s lacks its required initial SL and market has "
"already crossed it; closing immediately.",
config.name,
position.ticket,
)
close_position(position, f"{config.name}_INITIAL_GAP")
return False
min_distance = max(
int(getattr(SYMBOL_INFO, "trade_stops_level", 0)),
int(getattr(SYMBOL_INFO, "trade_freeze_level", 0)),
) * point
if (buy and target > market - min_distance) or ((not buy) and target < market + min_distance):
LOGGER.critical(
"%s position ticket=%s cannot restore the intended initial SL without "
"widening risk; closing immediately.",
config.name,
position.ticket,
)
close_position(position, f"{config.name}_SL_REPAIR")
return False
LOGGER.warning(
"%s position ticket=%s has missing/wider SL %.2f; restoring %.2f.",
config.name,
position.ticket,
old_sl,
target,
)
return modify_position_sl(position, target)
# =============================================================================
# COMPLETED-M1 TICK CLOSES AND TRAILING
# =============================================================================
def minute_floor(dt: datetime) -> datetime:
return dt.astimezone(UTC).replace(second=0, microsecond=0)
def completed_minute_close(boundary_broker: datetime) -> Optional[tuple[datetime, float, float]]:
"""Return the final bid/ask of the completed broker-clock M1 candle."""
boundary = minute_floor(boundary_broker)
key = int(boundary.timestamp())
if key in COMPLETED_CLOSE_CACHE:
return COMPLETED_CLOSE_CACHE[key]
start = boundary - timedelta(seconds=90)
end = boundary - timedelta(microseconds=1)
ticks = mt5.copy_ticks_range(SYMBOL, start, end, mt5.COPY_TICKS_ALL)
if ticks is None or len(ticks) == 0:
LOGGER.warning(
"No ticks available for completed broker minute ending %s",
boundary.isoformat(),
)
COMPLETED_CLOSE_CACHE.clear()
COMPLETED_CLOSE_CACHE[key] = None
return None
last = ticks[-1]
last_msc = int(last["time_msc"])
completed_start_msc = int((boundary - timedelta(minutes=1)).timestamp() * 1000)
if last_msc < completed_start_msc:
LOGGER.warning(
"No tick belongs to completed broker minute %s; refusing to carry an older quote forward.",
(boundary - timedelta(minutes=1)).isoformat(),
)
COMPLETED_CLOSE_CACHE.clear()
COMPLETED_CLOSE_CACHE[key] = None
return None
bid = float(last["bid"])
ask = float(last["ask"])
if bid <= 0 or ask <= 0:
COMPLETED_CLOSE_CACHE.clear()
COMPLETED_CLOSE_CACHE[key] = None
return None
result = (boundary - timedelta(minutes=1), bid, ask)
COMPLETED_CLOSE_CACHE.clear() # only the current boundary is useful live
COMPLETED_CLOSE_CACHE[key] = result
return result
@dataclass(frozen=True)
class CompletedM1Bar:
minute_ts: int
open: float
high: float
low: float
close: float
def _clamp(value: float, lower: float, upper: float) -> float:
value = max(float(lower), float(value))
if float(upper) > 0:
value = min(float(upper), value)
return value
def completed_m1_rates(boundary_broker: datetime) -> Optional[list[CompletedM1Bar]]:
"""Return completed M1 bid bars ending before a broker-clock boundary.
This RoboForex terminal exposes rate epochs in the same broker-clock coordinate
as live ticks and positions. Querying in that coordinate prevents a UTC+2/+3
server offset from making the newest completed bars appear unavailable.
"""
boundary = minute_floor(boundary_broker)
key = int(boundary.timestamp())
if key in M1_RATE_CACHE:
return M1_RATE_CACHE[key]
if time.monotonic() < M1_RATE_RETRY_AFTER.get(key, 0.0):
return None
start = boundary - timedelta(days=1)
end = boundary - timedelta(microseconds=1)
rates = mt5.copy_rates_range(SYMBOL, mt5.TIMEFRAME_M1, start, end)
if rates is None or len(rates) == 0:
if key not in M1_RATE_SOURCE_WARNINGS:
M1_RATE_SOURCE_WARNINGS.add(key)
LOGGER.warning(
"No completed M1 rates available through broker boundary %s: %s; retrying.",
boundary.isoformat(),
mt5_error_text(),
)
M1_RATE_RETRY_AFTER[key] = time.monotonic() + 1.0
return None
by_minute: dict[int, CompletedM1Bar] = {}
for rate in rates:
minute_ts = int(rate["time"])
if minute_ts >= key:
continue
open_price = float(rate["open"])
high_price = float(rate["high"])
low_price = float(rate["low"])
close_price = float(rate["close"])
if min(open_price, high_price, low_price, close_price) <= 0:
continue
by_minute[minute_ts] = CompletedM1Bar(
minute_ts=minute_ts,
open=open_price,
high=high_price,
low=low_price,
close=close_price,
)
bars = [by_minute[minute] for minute in sorted(by_minute)]
expected_last = key - 60
if not bars or bars[-1].minute_ts < expected_last:
if key not in M1_RATE_SOURCE_WARNINGS:
M1_RATE_SOURCE_WARNINGS.add(key)
LOGGER.warning(
"Latest completed M1 rate is stale at broker boundary %s; latest=%s; retrying.",
boundary.isoformat(),
datetime.fromtimestamp(bars[-1].minute_ts, UTC).isoformat() if bars else None,
)
M1_RATE_RETRY_AFTER[key] = time.monotonic() + 1.0
return None
M1_RATE_CACHE.clear() # live calculations need only the current/frozen boundary
M1_RATE_RETRY_AFTER.clear()
M1_RATE_SOURCE_WARNINGS.discard(key)
M1_RATE_CACHE[key] = bars
return bars
def volatility_from_bars(
bars: list[CompletedM1Bar],
estimator: str,
window: int,
) -> Optional[float]:
"""Calculate the exact estimator definitions used in the brute-force study."""
if window <= 0 or len(bars) < window + 1:
return None
estimator = str(estimator).upper()
closes = [bar.close for bar in bars]
if estimator == "CLOSE_RMS":
changes = [
closes[index] - closes[index - 1]
for index in range(len(closes) - window, len(closes))
]
return math.sqrt(sum(change * change for change in changes) / window)
true_ranges: list[float] = []
for index in range(1, len(bars)):
bar = bars[index]
previous_close = bars[index - 1].close
true_ranges.append(
max(
bar.high - bar.low,
abs(bar.high - previous_close),
abs(bar.low - previous_close),
)
)
if len(true_ranges) < window:
return None
if estimator == "TR_SMA":
values = true_ranges[-window:]
return sum(values) / window
if estimator == "TR_MEDIAN":
values = sorted(true_ranges[-window:])
middle = window // 2
if window % 2:
return values[middle]
return (values[middle - 1] + values[middle]) / 2.0
if estimator == "ATR_WILDER":
# Conventional Wilder ATR seed followed by recursive smoothing. The
# four-day history above makes the result effectively independent of the
# seed for the selected 30-minute window.
atr = sum(true_ranges[:window]) / window
for true_range in true_ranges[window:]:
atr += (true_range - atr) / window
return atr
raise ValueError(f"Unsupported volatility estimator: {estimator!r}")
def volatility_at_boundary(
boundary_broker: datetime,
estimator: str,
window: int,
) -> Optional[float]:
bars = completed_m1_rates(boundary_broker)
if bars is None:
return None
return volatility_from_bars(bars, estimator, window)
def frozen_activation_for_event(
config: StrategyConfig,
event_utc: datetime,
clock_tick: Any = None,
) -> Optional[float]:
if not config.activation_estimator or config.activation_window <= 0:
return None
event_broker = utc_to_broker_clock(event_utc, clock_tick)
volatility = volatility_at_boundary(
event_broker,
config.activation_estimator,
config.activation_window,
)
if volatility is None or not math.isfinite(volatility):
return None
return _clamp(
config.activation_multiplier * volatility,
config.activation_floor,
config.activation_cap,
)
def _group_event_utc(group_id: str, group: Any) -> Optional[datetime]:
if isinstance(group, dict) and group.get("event_utc"):
try:
return parse_utc(str(group["event_utc"]))
except Exception:
pass
try:
return datetime.strptime(group_id[1:11], "%y%m%d%H%M").replace(tzinfo=UTC)
except Exception:
return None
def effective_trail_activation(
position: Any,
config: StrategyConfig,
state: dict[str, Any],
clock_tick: Any,
) -> float:
"""Return fixed activation, or the frozen pre-event NYSE activation."""
if not config.activation_estimator or config.activation_window <= 0:
return config.trail_activation
group_id = group_id_from_comment(getattr(position, "comment", ""), config)
groups = state.setdefault("active_groups", {})
group = groups.get(group_id) if group_id else None
if isinstance(group, dict):
stored = group.get("frozen_trail_activation")
if stored is not None:
try:
value = float(stored)
if math.isfinite(value) and value > 0:
return value
except (TypeError, ValueError):
pass
event_utc = _group_event_utc(group_id, group) if group_id else None
if event_utc is None:
return config.trail_activation
calculated = frozen_activation_for_event(config, event_utc, clock_tick)
if calculated is None:
warning_key = (f"{config.name}_ACTIVATION", int(event_utc.timestamp()))
if warning_key not in VOLATILITY_WARNING_KEYS:
VOLATILITY_WARNING_KEYS.add(warning_key)
LOGGER.warning(
"%s frozen activation unavailable; using fixed fallback %.2f",
config.name,
config.trail_activation,
)
return config.trail_activation
if isinstance(group, dict):
group["frozen_trail_activation"] = calculated
save_state(state)
return calculated
def effective_trail_distance(
config: StrategyConfig,
boundary_broker: datetime,
broker_min_distance: float,
) -> float:
"""Return deployment-safe adaptive distance, never below broker constraints."""
distance = config.trail_distance
if config.trail_estimator and config.trail_window > 0:
volatility = volatility_at_boundary(
boundary_broker,
config.trail_estimator,
config.trail_window,
)
if volatility is not None and math.isfinite(volatility):
distance = _clamp(
config.trail_multiplier * volatility,
config.trail_floor,
config.trail_cap,
)
else:
warning_key = (config.name, int(minute_floor(boundary_broker).timestamp()))
if warning_key not in VOLATILITY_WARNING_KEYS:
VOLATILITY_WARNING_KEYS.add(warning_key)
LOGGER.warning(
"%s volatility unavailable at %s; using fixed fallback trail %.2f",
config.name,
minute_floor(boundary_broker).isoformat(),
config.trail_distance,
)
return max(float(distance), float(broker_min_distance))
def rebuild_trail_state(
position: Any,
config: StrategyConfig,
clock_tick: Any = None,
) -> dict[str, Any]:
"""Reconstruct completed-M1 executable closes after the entry candle.
Position and tick timestamps remain in MT5 broker-clock space. Their common
offset cancels, so reconstruction is independent of Windows/VPS timezone.
A caller-supplied tick keeps one trailing pass on a single clock snapshot.
"""
entry_dt = _broker_clock_datetime(position)
entry_minute = minute_floor(entry_dt)
if clock_tick is None:
clock_tick = current_tick(require_fresh=False)
end_boundary = minute_floor(_broker_clock_datetime(clock_tick))
start = entry_minute
end = end_boundary - timedelta(microseconds=1)
best: Optional[float] = None
last_minute: Optional[datetime] = None
history_query_succeeded = end <= start
if end > start:
ticks = mt5.copy_ticks_range(SYMBOL, start, end, mt5.COPY_TICKS_ALL)
history_query_succeeded = ticks is not None
if ticks is None:
LOGGER.error(
"Could not rebuild trailing history ticket=%s: %s",
position.ticket,
mt5_error_text(),
)
elif len(ticks):
last_by_minute: dict[int, tuple[float, float]] = {}
for tick in ticks:
minute_ts = int(tick["time_msc"] // 60_000) * 60
last_by_minute[minute_ts] = (float(tick["bid"]), float(tick["ask"]))
buy = int(position.type) == int(mt5.POSITION_TYPE_BUY)
# The entry candle is excluded: only minute starts strictly greater
# than the entry candle start are eligible.
for minute_ts in sorted(last_by_minute):
minute_dt = datetime.fromtimestamp(minute_ts, UTC)
if minute_dt <= entry_minute:
continue
bid, ask = last_by_minute[minute_ts]
executable = bid if buy else ask
if executable <= 0:
continue
if best is None:
best = executable
elif buy:
best = max(best, executable)
else:
best = min(best, executable)
last_minute = minute_dt
LOGGER.info(
"Rebuilt trailing state ticket=%s best=%s last_broker_minute=%s",
position.ticket,
best,
last_minute.isoformat() if last_minute else None,
)
rebuilt_through = (
end_boundary - timedelta(minutes=1)
if history_query_succeeded
else None
)
return {
"strategy": config.name,
"clock_basis": BROKER_CLOCK_BASIS,
"best_executable_close": best,
"last_completed_minute": format_utc(last_minute) if last_minute else None,
"rebuilt_through": format_utc(rebuilt_through) if rebuilt_through else None,
}
def theoretical_trailing_sl(
config: StrategyConfig,
buy: bool,
entry_price: float,
best_close: float,
trail_activation: float,
trail_distance: float,
broker_min_distance: float,
) -> Optional[float]:
"""Return the completed-M1 stop under the adaptive deployment-safe rules.
Volatility is measured only from fully completed M1 bars. The current stop is
evaluated before this candidate is applied, and callers still enforce the
never-loosen rule. Existing profit floors and permanent late trails remain
fixed exactly as in the prior specification.
"""
favourable = best_close - entry_price if buy else entry_price - best_close
if favourable + 1e-9 < trail_activation:
return None
active_distance = max(float(trail_distance), float(broker_min_distance))
if (
config.late_trail_threshold is not None
and config.late_trail_distance is not None
and favourable + 1e-9 >= config.late_trail_threshold
):
active_distance = max(
float(config.late_trail_distance),
float(broker_min_distance),
)
candidate = best_close - active_distance if buy else best_close + active_distance
if (
config.profit_floor_threshold is not None
and config.profit_floor is not None
and favourable + 1e-9 >= config.profit_floor_threshold
):
floor_price = entry_price + config.profit_floor if buy else entry_price - config.profit_floor
candidate = max(candidate, floor_price) if buy else min(candidate, floor_price)
return normalize_price(candidate)
def manage_trailing(position: Any, config: StrategyConfig, state: dict[str, Any]) -> None:
ticket = int(position.ticket)
ticket_key = str(ticket)
# A short broker quote pause is expected occasionally. Defer only this
# position's trailing pass; do not abort the rest of the main reconciliation
# loop. Entry placement still uses the same strict freshness check and is not
# relaxed by this handling.
try:
clock_tick = current_tick(require_fresh=True)
except TransientTickError as exc:
now_mono = time.monotonic()
last_warning = TRAIL_TICK_WARNING_LAST_MONOTONIC.get(ticket, -math.inf)
if now_mono - last_warning >= TRAIL_TICK_WARNING_INTERVAL_SECONDS:
TRAIL_TICK_WARNING_LAST_MONOTONIC[ticket] = now_mono
LOGGER.warning(
"%s trailing deferred ticket=%s current_sl=%.2f: %s",
config.name,
ticket,
float(getattr(position, "sl", 0.0) or 0.0),
exc,
)
return
trail_map = state.setdefault("trail", {})
trail_state = trail_map.get(ticket_key)
if (
not isinstance(trail_state, dict)
or trail_state.get("clock_basis") != BROKER_CLOCK_BASIS
):
# Existing version-1 state files remain valid. Only an open position's old
# trail subrecord is rebuilt once under the corrected broker-clock basis.
trail_state = rebuild_trail_state(position, config, clock_tick)
trail_map[ticket_key] = trail_state
save_state(state)
entry_dt = _broker_clock_datetime(position)
now_broker = _broker_clock_datetime(clock_tick)
latest_completed_minute = minute_floor(now_broker) - timedelta(minutes=1)
last_text = trail_state.get("last_completed_minute")
rebuilt_text = trail_state.get("rebuilt_through")
last_minute = parse_utc(last_text) if last_text else None
rebuilt_through = parse_utc(rebuilt_text) if rebuilt_text else None
covered_candidates = [dt for dt in (last_minute, rebuilt_through) if dt is not None]
covered_minute = max(covered_candidates) if covered_candidates else None
if covered_minute is None or latest_completed_minute > covered_minute + timedelta(minutes=1):
trail_state = rebuild_trail_state(position, config, clock_tick)
trail_map[ticket_key] = trail_state
save_state(state)
completed = completed_minute_close(now_broker)
buy = int(position.type) == int(mt5.POSITION_TYPE_BUY)
if completed is not None:
candle_minute, bid_close, ask_close = completed
entry_minute = minute_floor(entry_dt)
last_text = trail_state.get("last_completed_minute")
last_minute = parse_utc(last_text) if last_text else None
# Entry candle is not eligible, and every minute is processed once.
if candle_minute > entry_minute and (last_minute is None or candle_minute > last_minute):
executable_close = bid_close if buy else ask_close
best = trail_state.get("best_executable_close")
best_value = float(best) if best is not None else None
if best_value is None:
best_value = executable_close
elif buy:
best_value = max(best_value, executable_close)
else:
best_value = min(best_value, executable_close)
trail_state["best_executable_close"] = best_value
trail_state["last_completed_minute"] = format_utc(candle_minute)
trail_state["rebuilt_through"] = format_utc(candle_minute)
trail_state["clock_basis"] = BROKER_CLOCK_BASIS
save_state(state)
best_raw = trail_state.get("best_executable_close")
if best_raw is None:
return
best_close = float(best_raw)
entry_price = float(position.price_open)
# Use the same fresh snapshot obtained at the beginning of this trailing pass.
# A second quote request here previously created another avoidable failure point.
tick = clock_tick
current_bid = float(tick.bid)
current_ask = float(tick.ask)
broker_min_distance = max(
int(getattr(SYMBOL_INFO, "trade_stops_level", 0)),
int(getattr(SYMBOL_INFO, "trade_freeze_level", 0)),
) * float(SYMBOL_INFO.point)
trail_activation = effective_trail_activation(
position,
config,
state,
clock_tick,
)
trail_distance = effective_trail_distance(
config,
minute_floor(now_broker),
broker_min_distance,
)
theoretical_sl = theoretical_trailing_sl(
config,
buy,
entry_price,
best_close,
trail_activation,
trail_distance,
broker_min_distance,
)
if theoretical_sl is None:
return
# Conservative restart/gap handling: if the market is already through the
# recovered stop, close at the current market price instead of pretending the
# old stop could be installed retroactively.
if (buy and current_bid <= theoretical_sl) or (not buy and current_ask >= theoretical_sl):
LOGGER.warning(
"Price is through theoretical trailing stop; closing at market. "
"ticket=%s theoretical_sl=%.2f market=%.2f/%.2f",
position.ticket,
theoretical_sl,
current_bid,
current_ask,
)
close_position(position, f"{config.name}_TRAIL_GAP")
return
min_distance = broker_min_distance
if buy:
candidate = min(theoretical_sl, current_bid - min_distance)
candidate = normalize_price(candidate)
old_sl = float(position.sl)
if old_sl > 0 and candidate <= old_sl + float(SYMBOL_INFO.point) / 2:
return
if candidate >= current_bid:
return
else:
candidate = max(theoretical_sl, current_ask + min_distance)
candidate = normalize_price(candidate)
old_sl = float(position.sl)
if old_sl > 0 and candidate >= old_sl - float(SYMBOL_INFO.point) / 2:
return
if candidate <= current_ask:
return
modify_position_sl(position, candidate)
# =============================================================================
# OCO, EXPIRY, POSITION MANAGEMENT
# =============================================================================
def order_expiry_utc(order: Any, state: dict[str, Any]) -> Optional[datetime]:
comment = str(getattr(order, "comment", ""))
group_id = comment[:-1] if comment.endswith(("B", "S")) else comment
group = state.get("active_groups", {}).get(group_id)
if isinstance(group, dict) and group.get("expiry_utc"):
try:
return parse_utc(group["expiry_utc"])
except Exception:
pass
raw = int(getattr(order, "time_expiration", 0) or 0)
if raw > 0:
try:
clock_tick = current_tick(require_fresh=False)
return broker_clock_to_utc(datetime.fromtimestamp(raw, UTC), clock_tick)
except Exception:
LOGGER.exception(
"Could not convert broker expiration for order ticket=%s",
getattr(order, "ticket", 0),
)
return None
def pending_oco_pair_is_complete(orders: list[Any]) -> bool:
"""Return True only for exactly one BUY leg and one SELL leg of one group."""
if len(orders) != 2:
return False
comments = [str(getattr(order, "comment", "")).strip() for order in orders]
if not all(comment.endswith(("B", "S")) for comment in comments):
return False
group_ids = {comment[:-1] for comment in comments}
sides = {comment[-1] for comment in comments}
if len(group_ids) != 1 or sides != {"B", "S"}:
return False
order_types = {int(getattr(order, "type", -1)) for order in orders}
required_types = {
int(getattr(mt5, "ORDER_TYPE_BUY_STOP", 4)),
int(getattr(mt5, "ORDER_TYPE_SELL_STOP", 5)),
}
return order_types == required_types
def manage_open_position(
position: Any,
config: StrategyConfig,
state: dict[str, Any],
) -> None:
if not ensure_initial_protection(position, config):
return
# Maximum hold is elapsed time, not a wall-clock comparison. The live MT5
# tick and the position opening timestamp share the same broker clock, so the
# timezone offset cancels exactly.
clock_tick = current_tick(require_fresh=False)
elapsed_seconds = _broker_elapsed_seconds(position, clock_tick)
if elapsed_seconds < -5.0:
LOGGER.error(
"%s position ticket=%s has a future broker timestamp; deferring time exit. "
"position_msc=%s tick_msc=%s",
config.name,
position.ticket,
_mt5_object_time_msc(position),
_mt5_object_time_msc(clock_tick),
)
elif elapsed_seconds >= config.max_hold_minutes * 60.0:
LOGGER.info(
"%s maximum hold reached ticket=%s elapsed=%.1fs limit=%ss",
config.name,
position.ticket,
elapsed_seconds,
config.max_hold_minutes * 60,
)
close_position(position, f"{config.name}_TIME")
return
manage_trailing(position, config, state)
def reconcile_news_events(config: StrategyConfig, state: dict[str, Any]) -> None:
"""Reconcile each NEWS event minute independently.
Distinct event minutes may have simultaneous positions. OCO sibling removal,
duplicate-fill repair, expiry, time exit and trailing are all scoped to one
group_id, never to every position sharing NEWS_MAGIC.
"""
all_orders = sorted(
orders_for_magic(config.magic),
key=lambda item: (int(getattr(item, "time_setup_msc", 0)), int(item.ticket)),
)
all_positions = sorted(
positions_for_magic(config.magic),
key=lambda item: (int(getattr(item, "time_msc", 0)), int(item.ticket)),
)
grouped_orders: dict[str, list[Any]] = {}
grouped_positions: dict[str, list[Any]] = {}
malformed_orders: list[Any] = []
malformed_positions: list[Any] = []
for order in all_orders:
group_id = group_id_from_comment(getattr(order, "comment", ""), config)
if group_id is None:
malformed_orders.append(order)
else:
grouped_orders.setdefault(group_id, []).append(order)
for position in all_positions:
group_id = group_id_from_comment(getattr(position, "comment", ""), config)
if group_id is None:
malformed_positions.append(position)
else:
grouped_positions.setdefault(group_id, []).append(position)
for order in malformed_orders:
LOGGER.critical(
"NEWS pending order ticket=%s has an unrecognized event comment %r; removing it.",
order.ticket,
str(getattr(order, "comment", "")),
)
remove_pending_order(order, "NEWS_COMMENT_REPAIR")
now = utc_now()
group_ids = set(grouped_orders) | set(grouped_positions)
for group_id in sorted(group_ids):
orders = grouped_orders.get(group_id, [])
positions = grouped_positions.get(group_id, [])
if not positions and orders and not pending_oco_pair_is_complete(orders):
LOGGER.critical(
"NEWS group %s has an incomplete/duplicate OCO set (%d order(s)); "
"cancelling only this event's legs.",
group_id,
len(orders),
)
for order in orders:
remove_pending_order(order, "NEWS_PAIR_REPAIR")
continue
# Local OCO: a fill cancels only the sibling leg from the same event minute.
if positions and orders:
for order in orders:
remove_pending_order(order, "NEWS_OCO")
orders = []
# Both directions may fill during a gap, but only within one event group.
if len(positions) > 1:
keeper = positions[0]
LOGGER.critical(
"NEWS OCO breach for group %s: %d positions. Keeping earliest ticket=%s.",
group_id,
len(positions),
keeper.ticket,
)
for extra in positions[1:]:
close_position(extra, "NEWS_OCO_BREACH")
positions = [keeper]
if not positions:
for order in orders:
expiry = order_expiry_utc(order, state)
if expiry and now >= expiry + timedelta(seconds=1):
remove_pending_order(order, "NEWS_EXPIRED")
continue
manage_open_position(positions[0], config, state)
# A broker may alter a position comment. Do not close such a protected position
# merely because its group cannot be reconstructed; manage it independently.
for position in malformed_positions:
LOGGER.warning(
"NEWS position ticket=%s has unrecognized comment %r; managing it as a "
"standalone NEWS position.",
position.ticket,
str(getattr(position, "comment", "")),
)
manage_open_position(position, config, state)
def reconcile_subsystem(config: StrategyConfig, state: dict[str, Any]) -> None:
if config.name == "NEWS":
reconcile_news_events(config, state)
return
positions = sorted(
positions_for_magic(config.magic),
key=lambda p: (int(getattr(p, "time_msc", 0)), int(p.ticket)),
)
orders = orders_for_magic(config.magic)
# Session subsystems retain the original one-OCO/one-position rule.
if not positions and orders and not pending_oco_pair_is_complete(orders):
LOGGER.critical(
"%s has an incomplete/duplicate pending OCO set (%d order(s)); "
"cancelling all legs before any new placement.",
config.name,
len(orders),
)
for order in orders:
remove_pending_order(order, f"{config.name}_PAIR_REPAIR")
orders = orders_for_magic(config.magic)
if positions and orders:
for order in orders:
remove_pending_order(order, f"{config.name}_OCO")
orders = []
if len(positions) > 1:
keeper = positions[0]
LOGGER.critical(
"%s OCO breach: %d simultaneous positions. Keeping earliest ticket=%s.",
config.name,
len(positions),
keeper.ticket,
)
for extra in positions[1:]:
close_position(extra, f"{config.name}_OCO_BREACH")
positions = [keeper]
now = utc_now()
if not positions:
for order in orders:
expiry = order_expiry_utc(order, state)
if expiry and now >= expiry + timedelta(seconds=1):
remove_pending_order(order, f"{config.name}_EXPIRED")
return
manage_open_position(positions[0], config, state)
def cleanup_trail_state(state: dict[str, Any]) -> None:
live_tickets = {
str(int(position.ticket))
for magic in OUR_MAGICS
for position in positions_for_magic(magic)
}
trail = state.get("trail", {})
changed = False
for ticket in list(trail):
if ticket not in live_tickets:
del trail[ticket]
TRAIL_TICK_WARNING_LAST_MONOTONIC.pop(int(ticket), None)
changed = True
if changed:
save_state(state)
def cleanup_active_groups(state: dict[str, Any]) -> None:
live_comments = {
str(order.comment)
for magic in OUR_MAGICS
for order in orders_for_magic(magic)
} | {
str(position.comment)
for magic in OUR_MAGICS
for position in positions_for_magic(magic)
}
now = utc_now()
groups = state.get("active_groups", {})
changed = False
for group_id, group in list(groups.items()):
still_live = any(comment.startswith(group_id) for comment in live_comments)
if still_live:
continue
try:
expiry = parse_utc(group.get("expiry_utc", ""))
except Exception:
expiry = now - timedelta(days=1)
if now > expiry + timedelta(minutes=10):
del groups[group_id]
changed = True
if changed:
save_state(state)
# =============================================================================
# EVENT PROCESSING
# =============================================================================
def mark_event(
state: dict[str, Any],
event_key: str,
status: str,
event_utc: datetime,
detail: str = "",
) -> None:
state.setdefault("processed_events", {})[event_key] = {
"status": status,
"event_utc": format_utc(event_utc),
"updated_utc": format_utc(utc_now()),
"detail": detail,
}
save_state(state)
def process_scheduled_event(
config: StrategyConfig,
event_key: str,
event_name: str,
event_utc: datetime,
state: dict[str, Any],
) -> None:
if event_key in state.get("processed_events", {}):
return
now = utc_now()
placement_start = event_utc - timedelta(seconds=PLACEMENT_LEAD_SECONDS)
if now < placement_start:
return
# Never create a new pair after the scheduled event time. A late bot start is
# recorded as missed rather than silently changing the tested entry timing.
if now >= event_utc:
status = "skipped_busy" if event_setup_busy(config, event_utc) else "missed_late"
mark_event(state, event_key, status, event_utc, event_name)
LOGGER.warning("%s event %s: %s", config.name, event_name, status)
return
if event_setup_busy(config, event_utc):
# For NEWS this checks only the same UTC event minute. Earlier NEWS event
# positions do not block this event; session subsystems remain single-setup.
return
try:
placed = place_oco_pair(config, event_key, event_name, event_utc, state)
except Exception:
LOGGER.exception("Failed to place %s pair for %s", config.name, event_name)
return
if placed:
RISK_BLOCK_WARNING_KEYS.discard((event_key, "minimum_lot"))
RISK_BLOCK_WARNING_KEYS.discard((event_key, "portfolio_cap"))
mark_event(state, event_key, "placed", event_utc, event_name)
def process_events(state: dict[str, Any]) -> None:
now = utc_now()
# Macro-news timetable. These rows are explicitly UTC and independent of the
# machine, VPS, terminal, broker-server, Moldova or European local timezone.
for event in NEWS_EVENTS:
if not event.tradable:
continue
if event.when_utc < now - timedelta(minutes=1):
if event.event_id not in state.get("processed_events", {}):
mark_event(
state,
event.event_id,
"expired_before_start",
event.when_utc,
event.event_name,
)
continue
if event.when_utc > now + timedelta(minutes=2):
break
process_scheduled_event(
NEWS_CONFIG,
event.event_id,
event.event_name,
event.when_utc,
state,
)
# Four daily market-session events. Each event is constructed in its own
# exchange timezone and converted to UTC before scheduling.
for spec, event_id, event_utc, event_name in all_daily_session_events(now):
process_scheduled_event(
spec.config,
event_id,
event_name,
event_utc,
state,
)
def timetable_horizon_warning(now: datetime) -> None:
global LAST_HORIZON_WARNING_DATE
today_key = now.date().isoformat()
if LAST_HORIZON_WARNING_DATE == today_key:
return
LAST_HORIZON_WARNING_DATE = today_key
future = [event for event in NEWS_EVENTS if event.tradable and event.when_utc > now]
if not future:
LOGGER.critical("No future traded-code NEWS events are available; check traded_codes.csv and timetable.csv.")
return
latest = max(event.when_utc for event in future)
if latest < now + timedelta(days=10):
LOGGER.warning(
"Strong-news timetable horizon is short: last tradable event is %s; check the calendar EA.",
format_utc(latest),
)
# =============================================================================
# BALANCE AND TRADE LOGS
# =============================================================================
def read_last_logged_balance_text(path: Path) -> Optional[str]:
"""Return the last valid balance recorded in *path*, normalized to 2 decimals."""
if not path.exists():
return None
last_balance_text: Optional[str] = None
try:
with path.open("r", encoding="utf-8") as handle:
for raw_line in handle:
line = raw_line.strip()
if not line or "\t" not in line:
continue
_, raw_balance = line.rsplit("\t", 1)
try:
balance = float(raw_balance.strip())
except ValueError:
continue
if math.isfinite(balance):
last_balance_text = f"{balance:.2f}"
except OSError:
LOGGER.exception("Cannot read existing balance log %s", path)
return None
return last_balance_text
def append_balance(force: bool = False) -> None:
"""Append the current balance only when its logged 2-decimal value changes."""
global LAST_LOGGED_BALANCE_TEXT, LAST_BALANCE_LOG_YEAR
now = utc_now()
path = LOG_DIR / f"{now.year}.txt"
# On startup/reconnect calls and at a UTC year rollover, recover the last
# persisted value so an unchanged balance is not duplicated after restart.
if force or LAST_BALANCE_LOG_YEAR != now.year:
LAST_LOGGED_BALANCE_TEXT = read_last_logged_balance_text(path)
LAST_BALANCE_LOG_YEAR = now.year
account = mt5.account_info()
if account is None:
LOGGER.error("Cannot log balance: %s", mt5_error_text())
return
balance_text = f"{float(account.balance):.2f}"
if LAST_LOGGED_BALANCE_TEXT == balance_text:
return
with path.open("a", encoding="utf-8") as handle:
# User requested only date/time and balance, nothing else.
handle.write(f"{now:%Y-%m-%d %H:%M:%S} UTC\t{balance_text}\n")
LAST_LOGGED_BALANCE_TEXT = balance_text
LAST_BALANCE_LOG_YEAR = now.year
def deal_reason_name(value: int) -> str:
mapping = {}
for name in (
"DEAL_REASON_CLIENT",
"DEAL_REASON_MOBILE",
"DEAL_REASON_WEB",
"DEAL_REASON_EXPERT",
"DEAL_REASON_SL",
"DEAL_REASON_TP",
"DEAL_REASON_SO",
"DEAL_REASON_ROLLOVER",
"DEAL_REASON_VMARGIN",
"DEAL_REASON_SPLIT",
"DEAL_REASON_CORPORATE_ACTION",
):
if hasattr(mt5, name):
mapping[int(getattr(mt5, name))] = name.replace("DEAL_REASON_", "")
return mapping.get(int(value), str(value))
def weighted_price(deals: Iterable[Any]) -> float:
deals = list(deals)
total_volume = sum(float(deal.volume) for deal in deals)
if total_volume <= 0:
return 0.0
return sum(float(deal.price) * float(deal.volume) for deal in deals) / total_volume
def append_trade_row(
exit_dt: datetime,
entry_dt: datetime,
strategy: str,
direction: str,
volume: float,
entry_price: float,
exit_price: float,
net: float,
reason: str,
position_id: int,
comments: str,
) -> None:
# Closed trades are grouped by their UTC exit year. Entry and exit fields
# below retain the full UTC date and time for every position.
path = LOG_DIR / f"{exit_dt:%Y}_trades.log"
header = (
"ENTRY DATE/TIME UTC | EXIT DATE/TIME UTC | STRATEGY | DIR | LOT | "
"ENTRY | EXIT | NET USD | OUTCOME | REASON | POSITION | COMMENTS\n"
"--------------------+---------------------+-------------+------+--------+"
"------------+------------+------------+---------+--------------+----------+---------\n"
)
if not path.exists() or path.stat().st_size == 0:
path.write_text(header, encoding="utf-8")
outcome = "PROFIT" if net > 1e-9 else "LOSS" if net < -1e-9 else "FLAT"
line = (
f"{entry_dt:%Y-%m-%d %H:%M:%S} | {exit_dt:%Y-%m-%d %H:%M:%S} | "
f"{strategy:<11} | {direction:<4} | {volume:>6.2f} | "
f"{entry_price:>10.2f} | {exit_price:>10.2f} | {net:>10.2f} | "
f"{outcome:<7} | {reason:<12} | {position_id:<8} | {comments}\n"
)
with path.open("a", encoding="utf-8") as handle:
handle.write(line)
def reconcile_trade_logs(state: dict[str, Any], force: bool = False) -> None:
global LAST_DEAL_SCAN_MONOTONIC
now_mono = time.monotonic()
if not force and now_mono - LAST_DEAL_SCAN_MONOTONIC < 30:
return
LAST_DEAL_SCAN_MONOTONIC = now_mono
end_utc = utc_now() + timedelta(minutes=1)
start_utc = end_utc - timedelta(days=TRADE_LOG_LOOKBACK_DAYS)
try:
clock_tick = current_tick(require_fresh=False)
start_broker = utc_to_broker_clock(start_utc, clock_tick)
end_broker = utc_to_broker_clock(end_utc, clock_tick)
except Exception:
LOGGER.exception("Cannot establish broker-clock range for deal-history reconciliation.")
return
deals = mt5.history_deals_get(start_broker, end_broker)
if deals is None:
LOGGER.error("history_deals_get failed: %s", mt5_error_text())
return
grouped: dict[int, list[Any]] = {}
for deal in deals:
position_id = int(getattr(deal, "position_id", 0) or 0)
if position_id <= 0:
continue
grouped.setdefault(position_id, []).append(deal)
logged = {int(value) for value in state.get("logged_position_ids", [])}
changed = False
entry_in = int(getattr(mt5, "DEAL_ENTRY_IN", 0))
entry_out = int(getattr(mt5, "DEAL_ENTRY_OUT", 1))
entry_out_by = int(getattr(mt5, "DEAL_ENTRY_OUT_BY", 3))
entry_inout = int(getattr(mt5, "DEAL_ENTRY_INOUT", 2))
for position_id, group in sorted(grouped.items()):
if position_id in logged:
continue
if not any(int(getattr(deal, "magic", 0)) in OUR_MAGICS for deal in group):
continue
group.sort(key=lambda deal: (int(getattr(deal, "time_msc", 0)), int(deal.ticket)))
entries = [deal for deal in group if int(deal.entry) in {entry_in, entry_inout}]
exits = [deal for deal in group if int(deal.entry) in {entry_out, entry_out_by, entry_inout}]
if not entries or not exits:
continue
entry_volume = sum(float(deal.volume) for deal in entries)
exit_volume = sum(float(deal.volume) for deal in exits)
if exit_volume + 1e-9 < entry_volume:
continue # still partially open
entry = entries[0]
magic = next(
(int(deal.magic) for deal in entries if int(deal.magic) in OUR_MAGICS),
int(entry.magic),
)
config = CONFIG_BY_MAGIC.get(magic)
if config is None:
continue
direction = (
"BUY"
if int(entry.type) == int(getattr(mt5, "DEAL_TYPE_BUY", 0))
else "SELL"
)
entry_dt = _mt5_object_utc_datetime(entry, clock_tick)
exit_dt = _mt5_object_utc_datetime(exits[-1], clock_tick)
entry_price = weighted_price(entries)
exit_price = weighted_price(exits)
net = sum(
float(getattr(deal, "profit", 0.0))
+ float(getattr(deal, "commission", 0.0))
+ float(getattr(deal, "swap", 0.0))
+ float(getattr(deal, "fee", 0.0))
for deal in group
)
reasons = sorted({deal_reason_name(int(deal.reason)) for deal in exits})
comments = " / ".join(
dict.fromkeys(str(deal.comment).strip() for deal in group if str(deal.comment).strip())
)
append_trade_row(
exit_dt=exit_dt,
entry_dt=entry_dt,
strategy=config.name,
direction=direction,
volume=entry_volume,
entry_price=entry_price,
exit_price=exit_price,
net=net,
reason=",".join(reasons),
position_id=position_id,
comments=comments,
)
LOGGER.info(
"Trade logged position_id=%s strategy=%s direction=%s net=%.2f",
position_id,
config.name,
direction,
net,
)
logged.add(position_id)
changed = True
if changed:
state["logged_position_ids"] = sorted(logged)
save_state(state)
# =============================================================================
# STARTUP VALIDATION AND MAIN LOOP
# =============================================================================
def describe_configuration() -> None:
LOGGER.info(
"Scheduling is UTC-aware and host-timezone independent | "
"news=timetable UTC | New York=%s | Tokyo=%s | placement lead=%.1fs",
getattr(NY_TZ, "key", "America/New_York"),
getattr(TOKYO_TZ, "key", "Asia/Tokyo"),
PLACEMENT_LEAD_SECONDS,
)
LOGGER.info(
"Risk sizing: portfolio cap=%.2f%% | cost buffer=$%.2f/oz | minimum-lot policy=%s",
PORTFOLIO_RISK_CAP_PERCENT,
RISK_COST_BUFFER_PER_OUNCE,
MINIMUM_LOT_POLICY,
)
for config in ALL_CONFIGS:
trail_rule = (
f"clip({config.trail_multiplier:g}*{config.trail_estimator}{config.trail_window},"
f"{config.trail_floor:g},{config.trail_cap:g})"
if config.trail_estimator and config.trail_window > 0
else f"fixed {config.trail_distance:g}"
)
activation_rule = (
f"frozen clip({config.activation_multiplier:g}*{config.activation_estimator}"
f"{config.activation_window},{config.activation_floor:g},{config.activation_cap:g})"
if config.activation_estimator and config.activation_window > 0
else f"fixed {config.trail_activation:g}"
)
LOGGER.info(
"%s parameters: offset=%.2f SL=%.2f activation=%s trail=%s "
"floor_at=%s floor=%s late_at=%s late_trail=%s expiry=%ss hold=%sm "
"risk=%.2f%% min_lot=%.2f",
config.name,
config.offset,
config.initial_sl,
activation_rule,
trail_rule,
config.profit_floor_threshold,
config.profit_floor,
config.late_trail_threshold,
config.late_trail_distance,
config.expiry_seconds,
config.max_hold_minutes,
config.risk_percent,
config.min_lot,
)
def signal_handler(signum: int, _frame: Any) -> None:
global STOP_REQUESTED
STOP_REQUESTED = True
LOGGER.warning(
"Shutdown requested by signal %s. Existing pending orders/positions for all strategies are left "
"server-side with their initial/current SL and pending expiration.",
signum,
)
def main() -> int:
setup_logging()
acquire_single_instance_lock()
signal.signal(signal.SIGINT, signal_handler)
if hasattr(signal, "SIGTERM"):
signal.signal(signal.SIGTERM, signal_handler)
state = load_state()
try:
try:
load_traded_codes(force=True)
except FileNotFoundError:
LOGGER.error(
"traded_codes.csv is not present; NEWS remains disabled until a "
"valid file appears. Session strategies continue."
)
except Exception:
LOGGER.exception(
"Initial traded-code load failed; NEWS remains disabled until a "
"valid replacement is loaded. Session strategies continue."
)
try:
load_timetable(force=True)
except FileNotFoundError:
LOGGER.error(
"timetable.csv is not present yet; NEWS remains disabled until the "
"EA publishes a valid file. Session strategies continue."
)
except Exception:
LOGGER.exception(
"Initial timetable load failed; NEWS remains disabled until a valid "
"replacement is loaded. Session strategies continue."
)
connect_mt5()
describe_configuration()
append_balance(force=True)
reconcile_trade_logs(state, force=True)
prune_state(state, utc_now())
save_state(state)
while not STOP_REQUESTED:
loop_started = time.monotonic()
try:
if not ensure_connection():
time.sleep(RECONNECT_RETRY_SECONDS)
continue
try:
load_traded_codes(force=False)
except FileNotFoundError:
LOGGER.error(
"traded_codes.csv is temporarily missing; retaining the last "
"valid in-memory code set. Set enabled=0 in a valid file to "
"intentionally disable NEWS trading."
)
except Exception:
LOGGER.exception(
"Traded-code reload failed; retaining the last valid in-memory "
"code set."
)
try:
load_timetable(force=False)
except FileNotFoundError:
LOGGER.error(
"timetable.csv is temporarily missing; retaining the last "
"valid in-memory NEWS calendar while session strategies continue."
)
except Exception:
LOGGER.exception(
"Timetable reload failed; retaining the last valid in-memory "
"NEWS calendar while session strategies continue."
)
now = utc_now()
timetable_horizon_warning(now)
# OCO cancellation and position protection run before new entries.
for config in ALL_CONFIGS:
reconcile_subsystem(config, state)
cleanup_trail_state(state)
cleanup_active_groups(state)
process_events(state)
append_balance(force=False)
reconcile_trade_logs(state, force=False)
except Exception:
LOGGER.exception("Unhandled loop error; bot will continue.")
near_window = max(15.0, PLACEMENT_LEAD_SECONDS + 2.0)
near = seconds_to_nearest_event(utc_now()) <= near_window
target_sleep = NEAR_EVENT_POLL_SECONDS if near else NORMAL_POLL_SECONDS
elapsed = time.monotonic() - loop_started
time.sleep(max(0.01, target_sleep - elapsed))
finally:
try:
reconcile_trade_logs(state, force=True)
except Exception:
LOGGER.exception("Final trade-log reconciliation failed.")
try:
mt5.shutdown()
except Exception:
pass
LOGGER.info("Bot stopped.")
return 0
if __name__ == "__main__":
raise SystemExit(main())