Compare commits
2 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 18df96872b | |||
| 5b1d54bfe9 |
@@ -136,7 +136,24 @@ update_history_with_config(
|
|||||||
- **Rate view resolution**: use `resolve_rate_view_name()` / `resolve_rate_view_names()` to map symbols and granularities to existing SQLite compatibility views without creating databases. Both accept `None` (or a missing path) and return deterministic default names unless `require_existing=True`.
|
- **Rate view resolution**: use `resolve_rate_view_name()` / `resolve_rate_view_names()` to map symbols and granularities to existing SQLite compatibility views without creating databases. Both accept `None` (or a missing path) and return deterministic default names unless `require_existing=True`.
|
||||||
- **Rate view loading**: use `load_rate_data()` / `load_rate_data_from_connection()` to load a SQLite rate table or view into a `DatetimeIndex` DataFrame.
|
- **Rate view loading**: use `load_rate_data()` / `load_rate_data_from_connection()` to load a SQLite rate table or view into a `DatetimeIndex` DataFrame.
|
||||||
- **Multi-series rate loading**: use `build_rate_targets()` to build neutral `RateTarget(symbol, timeframe)` pairs, `resolve_rate_tables()` to map them to table/view names (pass `require_existing=True` for strict resolution), and `load_rate_series_from_sqlite()` to load them into a mapping keyed by `(symbol, integer timeframe)`. The loader requires existing managed views unless `explicit_tables` is supplied, and rejects duplicate `(symbol, timeframe)` targets.
|
- **Multi-series rate loading**: use `build_rate_targets()` to build neutral `RateTarget(symbol, timeframe)` pairs, `resolve_rate_tables()` to map them to table/view names (pass `require_existing=True` for strict resolution), and `load_rate_series_from_sqlite()` to load them into a mapping keyed by `(symbol, integer timeframe)`. The loader requires existing managed views unless `explicit_tables` is supplied, and rejects duplicate `(symbol, timeframe)` targets.
|
||||||
- **Multi-account latest rates**: use `collect_latest_rates_for_accounts()` with `AccountSpec` to read the latest bars for several account groups, merged into a `(symbol, integer timeframe)` mapping.
|
- **Multi-account latest rates**: use `collect_latest_rates_for_accounts()` with `AccountSpec` to read the latest bars for several account groups, merged into a `(symbol, integer timeframe)` mapping. For long-running pollers, `collect_latest_rates_for_accounts_with_retries()` adds bounded exponential backoff that retries only `pdmt5.Mt5TradingError` / `pdmt5.Mt5RuntimeError` and re-raises once `retry_count` is exhausted.
|
||||||
|
- **Latest closed bars**: use `collect_latest_closed_rates_for_accounts()` when downstream logic must exclude the still-forming current bar. It fetches `count + 1` bars at `start_pos=0`, drops the last row with `drop_forming_rate_bar()`, and validates each series is non-empty. `collect_latest_closed_rates_by_granularity()` returns the same data keyed by `(symbol, granularity_name)` such as `("EURUSD", "M1")`.
|
||||||
|
|
||||||
|
```python
|
||||||
|
from mt5cli import AccountSpec, collect_latest_closed_rates_by_granularity
|
||||||
|
|
||||||
|
rates = collect_latest_closed_rates_by_granularity(
|
||||||
|
[AccountSpec(symbols=["EURUSD", "GBPUSD"], login=12345)],
|
||||||
|
["M1", "H1"],
|
||||||
|
count=500,
|
||||||
|
retry_count=3,
|
||||||
|
)
|
||||||
|
eurusd_m1 = rates["EURUSD", "M1"] # closed bars only
|
||||||
|
```
|
||||||
|
|
||||||
|
- **Credential resolution**: use `resolve_account_spec()` / `resolve_account_specs()` to merge explicit override values over `AccountSpec` fields and expand `${ENV_VAR}` placeholders (via `substitute_env_placeholders()`), raising `ValueError` for missing variables. This keeps secrets out of plan/config files without coupling to any strategy code.
|
||||||
|
- **Throttled history updates**: use `ThrottledHistoryUpdater` to wrap `update_history()` with a minimum `interval_seconds` between successful runs (monotonic clock). Call `should_update()` / `update(client, symbols)` from an application loop; errors propagate by default, or pass `suppress_errors=True` to swallow recoverable `Mt5*Error`/`sqlite3.Error` and let the caller decide logging.
|
||||||
|
- **Granularity-keyed rate loading**: `load_rate_series_by_granularity()` builds targets with `build_rate_targets()`, loads them with `load_rate_series_from_sqlite()`, and returns a mapping keyed by `(symbol | None, granularity_name)` such as `("EURUSD", "M1")` to reduce downstream boilerplate.
|
||||||
- **MT5 session helper**: use the `mt5_session()` context manager to attach to (or, when `Mt5Config.path` is set, launch) an MT5 terminal, log in, and yield a connected `Mt5CliClient` that shuts down on exit.
|
- **MT5 session helper**: use the `mt5_session()` context manager to attach to (or, when `Mt5Config.path` is set, launch) an MT5 terminal, log in, and yield a connected `Mt5CliClient` that shuts down on exit.
|
||||||
- **SQLite export helpers**: use `export_dataframe_to_sqlite()` for append mode, optional index export, and post-write deduplication by key columns.
|
- **SQLite export helpers**: use `export_dataframe_to_sqlite()` for append mode, optional index export, and post-write deduplication by key columns.
|
||||||
- **Recent ticks and margins**: `recent_ticks()` and `minimum_margins()` SDK helpers (and matching CLI commands) cover common downstream read-only queries.
|
- **Recent ticks and margins**: `recent_ticks()` and `minimum_margins()` SDK helpers (and matching CLI commands) cover common downstream read-only queries.
|
||||||
|
|||||||
@@ -215,3 +215,15 @@ frame = series["EURUSD", 1] # keyed by (symbol, integer timeframe)
|
|||||||
requires existing managed `rate_*` compatibility views and raises
|
requires existing managed `rate_*` compatibility views and raises
|
||||||
`ValueError` when they are missing. Duplicate `(symbol, timeframe)` targets
|
`ValueError` when they are missing. Duplicate `(symbol, timeframe)` targets
|
||||||
are rejected.
|
are rejected.
|
||||||
|
- `load_rate_series_by_granularity()` is a thin wrapper that builds the targets,
|
||||||
|
loads the series, and rekeys the result by granularity name to avoid
|
||||||
|
converting integer timeframes downstream:
|
||||||
|
|
||||||
|
```python
|
||||||
|
from mt5cli import load_rate_series_by_granularity
|
||||||
|
|
||||||
|
series = load_rate_series_by_granularity(
|
||||||
|
"history.db", ["EURUSD"], ["M1", "H1"], count=1000
|
||||||
|
)
|
||||||
|
frame = series["EURUSD", "M1"] # keyed by (symbol | None, granularity_name)
|
||||||
|
```
|
||||||
|
|||||||
+100
@@ -1,3 +1,103 @@
|
|||||||
# SDK Module
|
# SDK Module
|
||||||
|
|
||||||
::: mt5cli.sdk
|
::: mt5cli.sdk
|
||||||
|
|
||||||
|
## Resilient multi-account orchestration
|
||||||
|
|
||||||
|
The SDK ships strategy-agnostic helpers for building long-running collectors on
|
||||||
|
top of the read-only client. None of them depend on a particular trading
|
||||||
|
application.
|
||||||
|
|
||||||
|
### Retrying transient rate collection
|
||||||
|
|
||||||
|
`collect_latest_rates_for_accounts_with_retries()` wraps
|
||||||
|
`collect_latest_rates_for_accounts()` with bounded exponential backoff. Only
|
||||||
|
`pdmt5.Mt5TradingError` and `pdmt5.Mt5RuntimeError` are retried; the final
|
||||||
|
failure is re-raised once `retry_count` is exhausted.
|
||||||
|
|
||||||
|
```python
|
||||||
|
from mt5cli import AccountSpec, collect_latest_rates_for_accounts_with_retries
|
||||||
|
|
||||||
|
accounts = [AccountSpec(symbols=["EURUSD"], login=12345)]
|
||||||
|
rates = collect_latest_rates_for_accounts_with_retries(
|
||||||
|
accounts,
|
||||||
|
["M1", "H1"],
|
||||||
|
count=500,
|
||||||
|
retry_count=3,
|
||||||
|
backoff_base=2, # sleeps 2s, 4s, 8s between attempts
|
||||||
|
)
|
||||||
|
```
|
||||||
|
|
||||||
|
### Latest closed rate bars
|
||||||
|
|
||||||
|
MetaTrader 5 `start_pos=0` includes the still-forming current bar as the last
|
||||||
|
row. `collect_latest_closed_rates_for_accounts()` fetches `count + 1` bars,
|
||||||
|
drops that row with `drop_forming_rate_bar()`, and validates each series is
|
||||||
|
non-empty. Use `collect_latest_closed_rates_by_granularity()` when callers
|
||||||
|
prefer keys such as `("EURUSD", "M1")` instead of integer timeframes.
|
||||||
|
|
||||||
|
```python
|
||||||
|
from mt5cli import AccountSpec, collect_latest_closed_rates_by_granularity
|
||||||
|
|
||||||
|
rates = collect_latest_closed_rates_by_granularity(
|
||||||
|
[AccountSpec(symbols=["EURUSD"], login=12345)],
|
||||||
|
["M1", "H1"],
|
||||||
|
count=500,
|
||||||
|
retry_count=3,
|
||||||
|
)
|
||||||
|
closed_m1 = rates["EURUSD", "M1"]
|
||||||
|
```
|
||||||
|
|
||||||
|
### Resolving credentials and `${ENV_VAR}` placeholders
|
||||||
|
|
||||||
|
`resolve_account_spec()` / `resolve_account_specs()` merge explicit override
|
||||||
|
values over `AccountSpec` fields and expand `${ENV_VAR}` placeholders, keeping
|
||||||
|
secrets out of plan/config files. A missing environment variable raises
|
||||||
|
`ValueError`.
|
||||||
|
|
||||||
|
```python
|
||||||
|
import os
|
||||||
|
|
||||||
|
from mt5cli import AccountSpec, resolve_account_specs
|
||||||
|
|
||||||
|
os.environ["MT5_LOGIN"] = "12345"
|
||||||
|
os.environ["MT5_PASSWORD"] = "secret"
|
||||||
|
accounts = [
|
||||||
|
AccountSpec(symbols=["EURUSD"], login="${MT5_LOGIN}", password="${MT5_PASSWORD}")
|
||||||
|
]
|
||||||
|
|
||||||
|
resolved = resolve_account_specs(accounts, server="Broker-Demo")
|
||||||
|
# resolved[0].login == "12345", resolved[0].server == "Broker-Demo"
|
||||||
|
```
|
||||||
|
|
||||||
|
### Throttled incremental history updates
|
||||||
|
|
||||||
|
`ThrottledHistoryUpdater` wraps `update_history()` with a minimum interval
|
||||||
|
between successful runs (using a monotonic clock), so an application loop can
|
||||||
|
call it every iteration without over-fetching.
|
||||||
|
|
||||||
|
```python
|
||||||
|
from pdmt5 import Mt5Config, Mt5DataClient
|
||||||
|
|
||||||
|
from mt5cli import Dataset, ThrottledHistoryUpdater
|
||||||
|
|
||||||
|
updater = ThrottledHistoryUpdater(
|
||||||
|
output="history.db",
|
||||||
|
datasets={Dataset.rates},
|
||||||
|
timeframes=["M1"],
|
||||||
|
interval_seconds=60, # <= 0 updates on every call
|
||||||
|
)
|
||||||
|
|
||||||
|
client = Mt5DataClient(config=Mt5Config(login=12345))
|
||||||
|
client.initialize_and_login_mt5()
|
||||||
|
try:
|
||||||
|
while True:
|
||||||
|
updater.update(client, ["EURUSD", "GBPUSD"]) # no-op until 60s elapse
|
||||||
|
# ... do other work; break when shutting down ...
|
||||||
|
finally:
|
||||||
|
client.shutdown()
|
||||||
|
```
|
||||||
|
|
||||||
|
By default `Mt5TradingError`, `Mt5RuntimeError`, and `sqlite3.Error` propagate so
|
||||||
|
the caller controls logging; pass `suppress_errors=True` to swallow them and
|
||||||
|
return `False` without advancing the throttle.
|
||||||
|
|||||||
@@ -6,8 +6,10 @@ from .history import (
|
|||||||
RateTarget,
|
RateTarget,
|
||||||
build_rate_targets,
|
build_rate_targets,
|
||||||
build_rate_view_name,
|
build_rate_view_name,
|
||||||
|
drop_forming_rate_bar,
|
||||||
load_rate_data,
|
load_rate_data,
|
||||||
load_rate_data_from_connection,
|
load_rate_data_from_connection,
|
||||||
|
load_rate_series_by_granularity,
|
||||||
load_rate_series_from_sqlite,
|
load_rate_series_from_sqlite,
|
||||||
resolve_history_datasets,
|
resolve_history_datasets,
|
||||||
resolve_history_tick_flags,
|
resolve_history_tick_flags,
|
||||||
@@ -19,11 +21,15 @@ from .history import (
|
|||||||
from .sdk import (
|
from .sdk import (
|
||||||
AccountSpec,
|
AccountSpec,
|
||||||
Mt5CliClient,
|
Mt5CliClient,
|
||||||
|
ThrottledHistoryUpdater,
|
||||||
account_info,
|
account_info,
|
||||||
build_config,
|
build_config,
|
||||||
collect_history,
|
collect_history,
|
||||||
|
collect_latest_closed_rates_by_granularity,
|
||||||
|
collect_latest_closed_rates_for_accounts,
|
||||||
collect_latest_rates,
|
collect_latest_rates,
|
||||||
collect_latest_rates_for_accounts,
|
collect_latest_rates_for_accounts,
|
||||||
|
collect_latest_rates_for_accounts_with_retries,
|
||||||
copy_rates_from,
|
copy_rates_from,
|
||||||
copy_rates_from_pos,
|
copy_rates_from_pos,
|
||||||
copy_rates_range,
|
copy_rates_range,
|
||||||
@@ -42,6 +48,9 @@ from .sdk import (
|
|||||||
positions,
|
positions,
|
||||||
recent_history_deals,
|
recent_history_deals,
|
||||||
recent_ticks,
|
recent_ticks,
|
||||||
|
resolve_account_spec,
|
||||||
|
resolve_account_specs,
|
||||||
|
substitute_env_placeholders,
|
||||||
symbol_info,
|
symbol_info,
|
||||||
symbol_info_tick,
|
symbol_info_tick,
|
||||||
symbols,
|
symbols,
|
||||||
@@ -75,19 +84,24 @@ __all__ = [
|
|||||||
"IfExists",
|
"IfExists",
|
||||||
"Mt5CliClient",
|
"Mt5CliClient",
|
||||||
"RateTarget",
|
"RateTarget",
|
||||||
|
"ThrottledHistoryUpdater",
|
||||||
"account_info",
|
"account_info",
|
||||||
"build_config",
|
"build_config",
|
||||||
"build_rate_targets",
|
"build_rate_targets",
|
||||||
"build_rate_view_name",
|
"build_rate_view_name",
|
||||||
"collect_history",
|
"collect_history",
|
||||||
|
"collect_latest_closed_rates_by_granularity",
|
||||||
|
"collect_latest_closed_rates_for_accounts",
|
||||||
"collect_latest_rates",
|
"collect_latest_rates",
|
||||||
"collect_latest_rates_for_accounts",
|
"collect_latest_rates_for_accounts",
|
||||||
|
"collect_latest_rates_for_accounts_with_retries",
|
||||||
"copy_rates_from",
|
"copy_rates_from",
|
||||||
"copy_rates_from_pos",
|
"copy_rates_from_pos",
|
||||||
"copy_rates_range",
|
"copy_rates_range",
|
||||||
"copy_ticks_from",
|
"copy_ticks_from",
|
||||||
"copy_ticks_range",
|
"copy_ticks_range",
|
||||||
"detect_format",
|
"detect_format",
|
||||||
|
"drop_forming_rate_bar",
|
||||||
"export_dataframe",
|
"export_dataframe",
|
||||||
"export_dataframe_to_sqlite",
|
"export_dataframe_to_sqlite",
|
||||||
"history_deals",
|
"history_deals",
|
||||||
@@ -96,6 +110,7 @@ __all__ = [
|
|||||||
"latest_rates",
|
"latest_rates",
|
||||||
"load_rate_data",
|
"load_rate_data",
|
||||||
"load_rate_data_from_connection",
|
"load_rate_data_from_connection",
|
||||||
|
"load_rate_series_by_granularity",
|
||||||
"load_rate_series_from_sqlite",
|
"load_rate_series_from_sqlite",
|
||||||
"market_book",
|
"market_book",
|
||||||
"minimum_margins",
|
"minimum_margins",
|
||||||
@@ -110,12 +125,15 @@ __all__ = [
|
|||||||
"positions",
|
"positions",
|
||||||
"recent_history_deals",
|
"recent_history_deals",
|
||||||
"recent_ticks",
|
"recent_ticks",
|
||||||
|
"resolve_account_spec",
|
||||||
|
"resolve_account_specs",
|
||||||
"resolve_history_datasets",
|
"resolve_history_datasets",
|
||||||
"resolve_history_tick_flags",
|
"resolve_history_tick_flags",
|
||||||
"resolve_history_timeframes",
|
"resolve_history_timeframes",
|
||||||
"resolve_rate_tables",
|
"resolve_rate_tables",
|
||||||
"resolve_rate_view_name",
|
"resolve_rate_view_name",
|
||||||
"resolve_rate_view_names",
|
"resolve_rate_view_names",
|
||||||
|
"substitute_env_placeholders",
|
||||||
"symbol_info",
|
"symbol_info",
|
||||||
"symbol_info_tick",
|
"symbol_info_tick",
|
||||||
"symbols",
|
"symbols",
|
||||||
|
|||||||
@@ -106,6 +106,23 @@ def resolve_granularity_name(timeframe: int) -> str:
|
|||||||
return str(timeframe)
|
return str(timeframe)
|
||||||
|
|
||||||
|
|
||||||
|
def drop_forming_rate_bar(df_rate: pd.DataFrame) -> pd.DataFrame:
|
||||||
|
"""Return closed bars from chronologically ordered MT5 rate data.
|
||||||
|
|
||||||
|
MetaTrader 5 ``copy_rates_from_pos(start_pos=0)`` includes the still-forming
|
||||||
|
current bar as the last row. Slice it off so downstream logic only sees
|
||||||
|
completed bars. Empty frames and single-row frames return empty results.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
df_rate: Rate data ordered oldest-to-newest with the forming bar last.
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
A new DataFrame with all rows except the last. Index and columns are
|
||||||
|
preserved. The input frame is not modified.
|
||||||
|
"""
|
||||||
|
return df_rate.iloc[:-1].copy()
|
||||||
|
|
||||||
|
|
||||||
def build_rate_view_name(
|
def build_rate_view_name(
|
||||||
*,
|
*,
|
||||||
symbol: str,
|
symbol: str,
|
||||||
@@ -706,6 +723,55 @@ def load_rate_series_from_sqlite(
|
|||||||
conn.close()
|
conn.close()
|
||||||
|
|
||||||
|
|
||||||
|
def load_rate_series_by_granularity(
|
||||||
|
conn_or_path: SqliteConnOrPath,
|
||||||
|
symbols: Sequence[str],
|
||||||
|
granularities: Sequence[int | str],
|
||||||
|
count: int,
|
||||||
|
*,
|
||||||
|
explicit_tables: Sequence[str] | None = None,
|
||||||
|
allow_missing_symbol: bool = False,
|
||||||
|
) -> dict[tuple[str | None, str], pd.DataFrame]:
|
||||||
|
"""Load rate series keyed by symbol and string granularity name.
|
||||||
|
|
||||||
|
Builds targets with :func:`build_rate_targets` and loads them with
|
||||||
|
:func:`load_rate_series_from_sqlite`, then rekeys the result by granularity
|
||||||
|
name (for example ``M1``) instead of the integer timeframe to reduce
|
||||||
|
downstream boilerplate.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
conn_or_path: SQLite database path or open connection.
|
||||||
|
symbols: MT5 symbol names. May be empty when ``allow_missing_symbol``.
|
||||||
|
granularities: MT5 timeframes as integers or names (for example ``M1``).
|
||||||
|
count: Number of most recent rows to load per series.
|
||||||
|
explicit_tables: Optional explicit table or view names matching the
|
||||||
|
built targets in row-major order. Required when symbols are omitted.
|
||||||
|
allow_missing_symbol: When True and ``symbols`` is empty, build targets
|
||||||
|
with ``symbol=None`` for each granularity instead of raising.
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
Mapping keyed by ``(symbol | None, granularity_name)`` to each rate
|
||||||
|
DataFrame. Propagates ``ValueError`` (via :func:`build_rate_targets` and
|
||||||
|
:func:`load_rate_series_from_sqlite`) when inputs are empty or invalid,
|
||||||
|
table resolution fails, or duplicate targets are present.
|
||||||
|
"""
|
||||||
|
targets = build_rate_targets(
|
||||||
|
symbols,
|
||||||
|
granularities,
|
||||||
|
allow_missing_symbol=allow_missing_symbol,
|
||||||
|
)
|
||||||
|
series = load_rate_series_from_sqlite(
|
||||||
|
conn_or_path,
|
||||||
|
targets,
|
||||||
|
count,
|
||||||
|
explicit_tables=explicit_tables,
|
||||||
|
)
|
||||||
|
return {
|
||||||
|
(symbol, resolve_granularity_name(timeframe)): frame
|
||||||
|
for (symbol, timeframe), frame in series.items()
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
def get_table_columns(conn: sqlite3.Connection, table: str) -> set[str]:
|
def get_table_columns(conn: sqlite3.Connection, table: str) -> set[str]:
|
||||||
"""Return existing SQLite columns for a table."""
|
"""Return existing SQLite columns for a table."""
|
||||||
quoted_table = quote_sqlite_identifier(table)
|
quoted_table = quote_sqlite_identifier(table)
|
||||||
|
|||||||
+442
-2
@@ -4,7 +4,10 @@ from __future__ import annotations
|
|||||||
|
|
||||||
import json
|
import json
|
||||||
import logging
|
import logging
|
||||||
|
import os
|
||||||
|
import re
|
||||||
import sqlite3
|
import sqlite3
|
||||||
|
import time
|
||||||
from contextlib import contextmanager
|
from contextlib import contextmanager
|
||||||
from dataclasses import dataclass, field
|
from dataclasses import dataclass, field
|
||||||
from datetime import UTC, datetime, timedelta
|
from datetime import UTC, datetime, timedelta
|
||||||
@@ -12,12 +15,14 @@ from pathlib import Path
|
|||||||
from typing import TYPE_CHECKING, Self, TypeVar, cast
|
from typing import TYPE_CHECKING, Self, TypeVar, cast
|
||||||
|
|
||||||
import pandas as pd
|
import pandas as pd
|
||||||
from pdmt5 import Mt5Config, Mt5DataClient
|
from pdmt5 import Mt5Config, Mt5DataClient, Mt5RuntimeError, Mt5TradingError
|
||||||
|
|
||||||
from .history import (
|
from .history import (
|
||||||
create_cash_events_view,
|
create_cash_events_view,
|
||||||
create_history_indexes,
|
create_history_indexes,
|
||||||
create_positions_reconstructed_view,
|
create_positions_reconstructed_view,
|
||||||
|
drop_forming_rate_bar,
|
||||||
|
resolve_granularity_name,
|
||||||
resolve_history_datasets,
|
resolve_history_datasets,
|
||||||
resolve_history_tick_flags,
|
resolve_history_tick_flags,
|
||||||
resolve_history_timeframes,
|
resolve_history_timeframes,
|
||||||
@@ -42,11 +47,15 @@ logger = logging.getLogger(__name__)
|
|||||||
__all__ = [
|
__all__ = [
|
||||||
"AccountSpec",
|
"AccountSpec",
|
||||||
"Mt5CliClient",
|
"Mt5CliClient",
|
||||||
|
"ThrottledHistoryUpdater",
|
||||||
"account_info",
|
"account_info",
|
||||||
"build_config",
|
"build_config",
|
||||||
"collect_history",
|
"collect_history",
|
||||||
|
"collect_latest_closed_rates_by_granularity",
|
||||||
|
"collect_latest_closed_rates_for_accounts",
|
||||||
"collect_latest_rates",
|
"collect_latest_rates",
|
||||||
"collect_latest_rates_for_accounts",
|
"collect_latest_rates_for_accounts",
|
||||||
|
"collect_latest_rates_for_accounts_with_retries",
|
||||||
"copy_rates_from",
|
"copy_rates_from",
|
||||||
"copy_rates_from_pos",
|
"copy_rates_from_pos",
|
||||||
"copy_rates_range",
|
"copy_rates_range",
|
||||||
@@ -65,6 +74,9 @@ __all__ = [
|
|||||||
"positions",
|
"positions",
|
||||||
"recent_history_deals",
|
"recent_history_deals",
|
||||||
"recent_ticks",
|
"recent_ticks",
|
||||||
|
"resolve_account_spec",
|
||||||
|
"resolve_account_specs",
|
||||||
|
"substitute_env_placeholders",
|
||||||
"symbol_info",
|
"symbol_info",
|
||||||
"symbol_info_tick",
|
"symbol_info_tick",
|
||||||
"symbols",
|
"symbols",
|
||||||
@@ -121,6 +133,12 @@ def _require_positive(value: float, name: str) -> None:
|
|||||||
raise ValueError(msg)
|
raise ValueError(msg)
|
||||||
|
|
||||||
|
|
||||||
|
def _require_non_negative(value: int, name: str) -> None:
|
||||||
|
if value < 0:
|
||||||
|
msg = f"{name} must be non-negative."
|
||||||
|
raise ValueError(msg)
|
||||||
|
|
||||||
|
|
||||||
def _call_required_client_method(client: Mt5DataClient, name: str) -> object:
|
def _call_required_client_method(client: Mt5DataClient, name: str) -> object:
|
||||||
try:
|
try:
|
||||||
method = getattr(client, name)
|
method = getattr(client, name)
|
||||||
@@ -960,6 +978,115 @@ def update_history_with_config( # noqa: PLR0913
|
|||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
class ThrottledHistoryUpdater:
|
||||||
|
"""Throttled incremental SQLite history updater for long-running apps.
|
||||||
|
|
||||||
|
Wraps :func:`update_history` with a minimum interval between successful
|
||||||
|
updates, so a tight application loop can call :meth:`update` every
|
||||||
|
iteration without re-fetching MT5 history more often than desired. Timing
|
||||||
|
uses a monotonic clock, so it is unaffected by wall-clock changes.
|
||||||
|
"""
|
||||||
|
|
||||||
|
def __init__(
|
||||||
|
self,
|
||||||
|
*,
|
||||||
|
output: Path | str,
|
||||||
|
datasets: set[Dataset] | None = None,
|
||||||
|
timeframes: Sequence[int | str] | None = None,
|
||||||
|
flags: int | str = "ALL",
|
||||||
|
lookback_hours: float = 24.0,
|
||||||
|
with_views: bool = False,
|
||||||
|
include_account_events: bool = True,
|
||||||
|
interval_seconds: float = 0.0,
|
||||||
|
suppress_errors: bool = False,
|
||||||
|
) -> None:
|
||||||
|
"""Initialize the throttled updater.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
output: SQLite database path.
|
||||||
|
datasets: Datasets to include (defaults to all).
|
||||||
|
timeframes: Rate timeframes to update (defaults to all fixed MT5
|
||||||
|
timeframes).
|
||||||
|
flags: Tick copy flags as integer or name (e.g. ``ALL``).
|
||||||
|
lookback_hours: First-run lookback when a table has no prior rows.
|
||||||
|
with_views: Create ``cash_events`` and ``positions_reconstructed``
|
||||||
|
views.
|
||||||
|
include_account_events: Include account-level cash events.
|
||||||
|
interval_seconds: Minimum seconds between successful updates. Values
|
||||||
|
``<= 0`` update on every call.
|
||||||
|
suppress_errors: When True, ``Mt5TradingError``, ``Mt5RuntimeError``,
|
||||||
|
and ``sqlite3.Error`` raised during an update are swallowed and
|
||||||
|
:meth:`update` returns False without advancing the throttle. When
|
||||||
|
False (default), such errors propagate so callers control logging.
|
||||||
|
"""
|
||||||
|
self.output = output
|
||||||
|
self.datasets = datasets
|
||||||
|
self.timeframes = timeframes
|
||||||
|
self.flags = flags
|
||||||
|
self.lookback_hours = lookback_hours
|
||||||
|
self.with_views = with_views
|
||||||
|
self.include_account_events = include_account_events
|
||||||
|
self.interval_seconds = interval_seconds
|
||||||
|
self.suppress_errors = suppress_errors
|
||||||
|
self._last_update_monotonic: float | None = None
|
||||||
|
|
||||||
|
@property
|
||||||
|
def last_update_monotonic(self) -> float | None:
|
||||||
|
"""Return the monotonic timestamp of the last successful update."""
|
||||||
|
return self._last_update_monotonic
|
||||||
|
|
||||||
|
def should_update(self) -> bool:
|
||||||
|
"""Return whether enough time has elapsed to run another update.
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
True when ``interval_seconds <= 0``, when no update has succeeded
|
||||||
|
yet, or when at least ``interval_seconds`` have elapsed since the
|
||||||
|
last successful update.
|
||||||
|
"""
|
||||||
|
if self.interval_seconds <= 0 or self._last_update_monotonic is None:
|
||||||
|
return True
|
||||||
|
return (time.monotonic() - self._last_update_monotonic) >= self.interval_seconds
|
||||||
|
|
||||||
|
def update(self, client: Mt5DataClient, symbols: Sequence[str]) -> bool:
|
||||||
|
"""Run a throttled incremental history update.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
client: Connected MT5 data client.
|
||||||
|
symbols: Symbols to update.
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
True if an update ran successfully, False if it was throttled or
|
||||||
|
(when ``suppress_errors`` is True) failed with a recoverable error.
|
||||||
|
|
||||||
|
Raises:
|
||||||
|
Mt5TradingError: If the update fails and ``suppress_errors`` is False.
|
||||||
|
Mt5RuntimeError: If the update fails and ``suppress_errors`` is False.
|
||||||
|
sqlite3.Error: If the SQLite write fails and ``suppress_errors`` is
|
||||||
|
False.
|
||||||
|
"""
|
||||||
|
if not self.should_update():
|
||||||
|
return False
|
||||||
|
try:
|
||||||
|
update_history(
|
||||||
|
client=client,
|
||||||
|
output=self.output,
|
||||||
|
symbols=symbols,
|
||||||
|
datasets=self.datasets,
|
||||||
|
timeframes=self.timeframes,
|
||||||
|
flags=self.flags,
|
||||||
|
lookback_hours=self.lookback_hours,
|
||||||
|
with_views=self.with_views,
|
||||||
|
include_account_events=self.include_account_events,
|
||||||
|
)
|
||||||
|
except (Mt5TradingError, Mt5RuntimeError, sqlite3.Error):
|
||||||
|
if self.suppress_errors:
|
||||||
|
logger.warning("Suppressed history update error", exc_info=True)
|
||||||
|
return False
|
||||||
|
raise
|
||||||
|
self._last_update_monotonic = time.monotonic()
|
||||||
|
return True
|
||||||
|
|
||||||
|
|
||||||
def collect_history(
|
def collect_history(
|
||||||
output: Path,
|
output: Path,
|
||||||
symbols: list[str],
|
symbols: list[str],
|
||||||
@@ -1113,13 +1240,155 @@ class AccountSpec:
|
|||||||
"""
|
"""
|
||||||
|
|
||||||
symbols: Sequence[str]
|
symbols: Sequence[str]
|
||||||
login: int | str | None = None
|
login: int | str | None = field(default=None, repr=False)
|
||||||
password: str | None = field(default=None, repr=False)
|
password: str | None = field(default=None, repr=False)
|
||||||
server: str | None = None
|
server: str | None = None
|
||||||
path: str | None = None
|
path: str | None = None
|
||||||
timeout: int | None = None
|
timeout: int | None = None
|
||||||
|
|
||||||
|
|
||||||
|
_ENV_PLACEHOLDER_PATTERN = re.compile(r"\$\{(?P<name>[A-Za-z_][A-Za-z0-9_]*)\}")
|
||||||
|
|
||||||
|
|
||||||
|
def substitute_env_placeholders(value: str) -> str:
|
||||||
|
"""Replace ``${ENV_VAR}`` placeholders in a string with environment values.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
value: String that may contain one or more ``${ENV_VAR}`` placeholders.
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
The string with every placeholder replaced by its environment value.
|
||||||
|
|
||||||
|
Raises:
|
||||||
|
ValueError: If a referenced environment variable is not set.
|
||||||
|
"""
|
||||||
|
parts: list[str] = []
|
||||||
|
last_end = 0
|
||||||
|
for match in _ENV_PLACEHOLDER_PATTERN.finditer(value):
|
||||||
|
parts.append(value[last_end : match.start()])
|
||||||
|
name = match.group("name")
|
||||||
|
if name not in os.environ:
|
||||||
|
msg = f"Environment variable {name!r} is not set."
|
||||||
|
raise ValueError(msg)
|
||||||
|
parts.append(os.environ[name])
|
||||||
|
last_end = match.end()
|
||||||
|
parts.append(value[last_end:])
|
||||||
|
return "".join(parts)
|
||||||
|
|
||||||
|
|
||||||
|
def _resolve_field(override: str | None, account_value: str | None) -> str | None:
|
||||||
|
"""Resolve a string field from an override or account value with env subst.
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
The explicit override when provided, otherwise the account value, with
|
||||||
|
any ``${ENV_VAR}`` placeholders substituted.
|
||||||
|
"""
|
||||||
|
value = override if override is not None else account_value
|
||||||
|
if value is None:
|
||||||
|
return None
|
||||||
|
return substitute_env_placeholders(value)
|
||||||
|
|
||||||
|
|
||||||
|
def _resolve_login(
|
||||||
|
override: int | str | None,
|
||||||
|
account_login: int | str | None,
|
||||||
|
) -> int | str | None:
|
||||||
|
"""Resolve a login from an override or account value with env substitution.
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
The explicit override when provided, otherwise the account login.
|
||||||
|
Integer values are preserved; string values have ``${ENV_VAR}``
|
||||||
|
placeholders substituted.
|
||||||
|
"""
|
||||||
|
if override is not None:
|
||||||
|
if isinstance(override, int):
|
||||||
|
return override
|
||||||
|
return substitute_env_placeholders(override)
|
||||||
|
if account_login is None or isinstance(account_login, int):
|
||||||
|
return account_login
|
||||||
|
return substitute_env_placeholders(account_login)
|
||||||
|
|
||||||
|
|
||||||
|
def resolve_account_spec(
|
||||||
|
account: AccountSpec,
|
||||||
|
*,
|
||||||
|
login: int | str | None = None,
|
||||||
|
password: str | None = None,
|
||||||
|
server: str | None = None,
|
||||||
|
path: str | None = None,
|
||||||
|
timeout: int | None = None,
|
||||||
|
) -> AccountSpec:
|
||||||
|
"""Resolve an account's credentials from overrides and ``${ENV_VAR}`` values.
|
||||||
|
|
||||||
|
Explicit override arguments take precedence over the corresponding
|
||||||
|
:class:`AccountSpec` fields. The resolved string fields (``login``,
|
||||||
|
``password``, ``server``, ``path``) have any ``${ENV_VAR}`` placeholders
|
||||||
|
substituted from the environment.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
account: Source account specification.
|
||||||
|
login: Optional explicit login override.
|
||||||
|
password: Optional explicit password override.
|
||||||
|
server: Optional explicit server override.
|
||||||
|
path: Optional explicit terminal path override.
|
||||||
|
timeout: Optional explicit connection timeout override.
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
A new :class:`AccountSpec` with resolved credentials and the original
|
||||||
|
symbols preserved. Raises ``ValueError`` (via
|
||||||
|
:func:`substitute_env_placeholders`) if a referenced environment
|
||||||
|
variable is not set.
|
||||||
|
"""
|
||||||
|
return AccountSpec(
|
||||||
|
symbols=account.symbols,
|
||||||
|
login=_resolve_login(login, account.login),
|
||||||
|
password=_resolve_field(password, account.password),
|
||||||
|
server=_resolve_field(server, account.server),
|
||||||
|
path=_resolve_field(path, account.path),
|
||||||
|
timeout=timeout if timeout is not None else account.timeout,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def resolve_account_specs(
|
||||||
|
accounts: Sequence[AccountSpec],
|
||||||
|
*,
|
||||||
|
login: int | str | None = None,
|
||||||
|
password: str | None = None,
|
||||||
|
server: str | None = None,
|
||||||
|
path: str | None = None,
|
||||||
|
timeout: int | None = None,
|
||||||
|
) -> list[AccountSpec]:
|
||||||
|
"""Resolve credentials for multiple accounts.
|
||||||
|
|
||||||
|
Applies the same overrides and ``${ENV_VAR}`` substitution as
|
||||||
|
:func:`resolve_account_spec` to every account.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
accounts: Source account specifications.
|
||||||
|
login: Optional explicit login override applied to each account.
|
||||||
|
password: Optional explicit password override applied to each account.
|
||||||
|
server: Optional explicit server override applied to each account.
|
||||||
|
path: Optional explicit terminal path override applied to each account.
|
||||||
|
timeout: Optional explicit timeout override applied to each account.
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
Resolved account specifications in the original order. Raises
|
||||||
|
``ValueError`` (via :func:`substitute_env_placeholders`) if a referenced
|
||||||
|
environment variable is not set.
|
||||||
|
"""
|
||||||
|
return [
|
||||||
|
resolve_account_spec(
|
||||||
|
account,
|
||||||
|
login=login,
|
||||||
|
password=password,
|
||||||
|
server=server,
|
||||||
|
path=path,
|
||||||
|
timeout=timeout,
|
||||||
|
)
|
||||||
|
for account in accounts
|
||||||
|
]
|
||||||
|
|
||||||
|
|
||||||
def _coerce_login(login: int | str | None) -> int | None:
|
def _coerce_login(login: int | str | None) -> int | None:
|
||||||
"""Coerce a login value to int, treating empty strings as unset.
|
"""Coerce a login value to int, treating empty strings as unset.
|
||||||
|
|
||||||
@@ -1212,6 +1481,177 @@ def collect_latest_rates_for_accounts(
|
|||||||
return result
|
return result
|
||||||
|
|
||||||
|
|
||||||
|
def collect_latest_rates_for_accounts_with_retries(
|
||||||
|
accounts: Sequence[AccountSpec],
|
||||||
|
timeframes: Sequence[int | str],
|
||||||
|
count: int,
|
||||||
|
*,
|
||||||
|
start_pos: int = 0,
|
||||||
|
base_config: Mt5Config | None = None,
|
||||||
|
retry_count: int = 0,
|
||||||
|
backoff_base: float = 2.0,
|
||||||
|
) -> dict[tuple[str, int], pd.DataFrame]:
|
||||||
|
"""Collect latest rates across accounts, retrying transient MT5 failures.
|
||||||
|
|
||||||
|
Wraps :func:`collect_latest_rates_for_accounts` with bounded exponential
|
||||||
|
backoff. Only ``pdmt5.Mt5TradingError`` and ``pdmt5.Mt5RuntimeError`` are
|
||||||
|
retried; other exceptions propagate immediately. The final failure is
|
||||||
|
re-raised once retries are exhausted.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
accounts: Account groups to read. Each must define at least one symbol.
|
||||||
|
timeframes: MT5 timeframes as integers or names (for example ``M1``).
|
||||||
|
count: Number of most recent bars to read per symbol/timeframe.
|
||||||
|
start_pos: Initial bar position offset.
|
||||||
|
base_config: Optional base configuration whose fields fill any value not
|
||||||
|
set on an individual account.
|
||||||
|
retry_count: Maximum number of retries after the first attempt. ``0``
|
||||||
|
disables retries.
|
||||||
|
backoff_base: Base for exponential backoff. The delay before retry
|
||||||
|
attempt ``n`` (1-indexed) is ``backoff_base ** n`` seconds.
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
Mapping keyed by ``(symbol, timeframe_int)``. Propagates ``ValueError``
|
||||||
|
for invalid inputs (see :func:`collect_latest_rates_for_accounts`) and
|
||||||
|
re-raises the last ``pdmt5.Mt5TradingError`` or ``pdmt5.Mt5RuntimeError``
|
||||||
|
once retries are exhausted.
|
||||||
|
"""
|
||||||
|
attempts = max(retry_count, 0) + 1
|
||||||
|
|
||||||
|
def _collect() -> dict[tuple[str, int], pd.DataFrame]:
|
||||||
|
return collect_latest_rates_for_accounts(
|
||||||
|
accounts,
|
||||||
|
timeframes,
|
||||||
|
count,
|
||||||
|
start_pos=start_pos,
|
||||||
|
base_config=base_config,
|
||||||
|
)
|
||||||
|
|
||||||
|
for attempt in range(attempts - 1):
|
||||||
|
try:
|
||||||
|
return _collect()
|
||||||
|
except (Mt5TradingError, Mt5RuntimeError) as exc:
|
||||||
|
delay = backoff_base ** (attempt + 1)
|
||||||
|
logger.warning(
|
||||||
|
"Rate collection failed (attempt %d/%d): %s; retrying in %.1fs",
|
||||||
|
attempt + 1,
|
||||||
|
attempts,
|
||||||
|
exc,
|
||||||
|
delay,
|
||||||
|
)
|
||||||
|
time.sleep(delay)
|
||||||
|
return _collect()
|
||||||
|
|
||||||
|
|
||||||
|
def collect_latest_closed_rates_for_accounts(
|
||||||
|
accounts: Sequence[AccountSpec],
|
||||||
|
timeframes: Sequence[int | str],
|
||||||
|
count: int,
|
||||||
|
*,
|
||||||
|
start_pos: int = 0,
|
||||||
|
base_config: Mt5Config | None = None,
|
||||||
|
retry_count: int = 0,
|
||||||
|
backoff_base: float = 2.0,
|
||||||
|
) -> dict[tuple[str, int], pd.DataFrame]:
|
||||||
|
"""Collect latest closed rate bars across multiple MT5 account groups.
|
||||||
|
|
||||||
|
When ``start_pos`` is ``0`` (the default), MetaTrader 5 includes the
|
||||||
|
still-forming current bar as the last row. This helper fetches
|
||||||
|
``count + 1`` bars, drops that bar with :func:`drop_forming_rate_bar`, and
|
||||||
|
validates that each resulting frame is non-empty. When ``start_pos`` is
|
||||||
|
greater than zero the forming bar is not in range, so only ``count`` bars
|
||||||
|
are fetched and no row is dropped.
|
||||||
|
|
||||||
|
Wraps :func:`collect_latest_rates_for_accounts_with_retries` for transient
|
||||||
|
MT5 error handling.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
accounts: Account groups to read. Each must define at least one symbol.
|
||||||
|
timeframes: MT5 timeframes as integers or names (for example ``M1``).
|
||||||
|
count: Number of closed bars to return per symbol/timeframe.
|
||||||
|
start_pos: Initial bar position offset passed to the underlying collector.
|
||||||
|
base_config: Optional base configuration whose fields fill any value not
|
||||||
|
set on an individual account.
|
||||||
|
retry_count: Maximum number of retries after the first attempt. ``0``
|
||||||
|
disables retries.
|
||||||
|
backoff_base: Base for exponential backoff between retry attempts.
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
Mapping keyed by ``(symbol, timeframe_int)``.
|
||||||
|
|
||||||
|
Raises:
|
||||||
|
ValueError: If inputs are invalid, or any series is empty (after
|
||||||
|
dropping the still-forming bar when ``start_pos`` is ``0``).
|
||||||
|
"""
|
||||||
|
_require_positive(count, "count")
|
||||||
|
_require_non_negative(start_pos, "start_pos")
|
||||||
|
fetch_count = count + 1 if start_pos == 0 else count
|
||||||
|
loaded = collect_latest_rates_for_accounts_with_retries(
|
||||||
|
accounts,
|
||||||
|
timeframes,
|
||||||
|
fetch_count,
|
||||||
|
start_pos=start_pos,
|
||||||
|
base_config=base_config,
|
||||||
|
retry_count=retry_count,
|
||||||
|
backoff_base=backoff_base,
|
||||||
|
)
|
||||||
|
result: dict[tuple[str, int], pd.DataFrame] = {}
|
||||||
|
for key, df_rate in loaded.items():
|
||||||
|
closed = drop_forming_rate_bar(df_rate) if start_pos == 0 else df_rate
|
||||||
|
if closed.empty:
|
||||||
|
symbol, timeframe = key
|
||||||
|
msg = f"Rate data is empty for {symbol!r} at timeframe {timeframe}."
|
||||||
|
raise ValueError(msg)
|
||||||
|
result[key] = closed
|
||||||
|
return result
|
||||||
|
|
||||||
|
|
||||||
|
def collect_latest_closed_rates_by_granularity(
|
||||||
|
accounts: Sequence[AccountSpec],
|
||||||
|
granularities: Sequence[int | str],
|
||||||
|
count: int,
|
||||||
|
*,
|
||||||
|
start_pos: int = 0,
|
||||||
|
base_config: Mt5Config | None = None,
|
||||||
|
retry_count: int = 0,
|
||||||
|
backoff_base: float = 2.0,
|
||||||
|
) -> dict[tuple[str, str], pd.DataFrame]:
|
||||||
|
"""Collect latest closed rate bars keyed by symbol and granularity name.
|
||||||
|
|
||||||
|
Thin wrapper around :func:`collect_latest_closed_rates_for_accounts` that
|
||||||
|
rekeys the result by granularity name (for example ``M1``) instead of the
|
||||||
|
integer timeframe.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
accounts: Account groups to read. Each must define at least one symbol.
|
||||||
|
granularities: MT5 timeframes as integers or names (for example ``M1``).
|
||||||
|
count: Number of closed bars to return per symbol/timeframe.
|
||||||
|
start_pos: Initial bar position offset passed to the underlying collector.
|
||||||
|
base_config: Optional base configuration whose fields fill any value not
|
||||||
|
set on an individual account.
|
||||||
|
retry_count: Maximum number of retries after the first attempt. ``0``
|
||||||
|
disables retries.
|
||||||
|
backoff_base: Base for exponential backoff between retry attempts.
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
Mapping keyed by ``(symbol, granularity_name)``. Propagates
|
||||||
|
``ValueError`` from :func:`collect_latest_closed_rates_for_accounts`.
|
||||||
|
"""
|
||||||
|
loaded = collect_latest_closed_rates_for_accounts(
|
||||||
|
accounts,
|
||||||
|
granularities,
|
||||||
|
count,
|
||||||
|
start_pos=start_pos,
|
||||||
|
base_config=base_config,
|
||||||
|
retry_count=retry_count,
|
||||||
|
backoff_base=backoff_base,
|
||||||
|
)
|
||||||
|
return {
|
||||||
|
(symbol, resolve_granularity_name(timeframe)): frame
|
||||||
|
for (symbol, timeframe), frame in loaded.items()
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
def copy_rates_range(
|
def copy_rates_range(
|
||||||
symbol: str,
|
symbol: str,
|
||||||
timeframe: int | str,
|
timeframe: int | str,
|
||||||
|
|||||||
+1
-1
@@ -1,6 +1,6 @@
|
|||||||
[project]
|
[project]
|
||||||
name = "mt5cli"
|
name = "mt5cli"
|
||||||
version = "0.5.2"
|
version = "0.6.0"
|
||||||
description = "Command-line tool for MetaTrader 5"
|
description = "Command-line tool for MetaTrader 5"
|
||||||
authors = [{name = "dceoy", email = "dceoy@users.noreply.github.com"}]
|
authors = [{name = "dceoy", email = "dceoy@users.noreply.github.com"}]
|
||||||
maintainers = [{name = "dceoy", email = "dceoy@users.noreply.github.com"}]
|
maintainers = [{name = "dceoy", email = "dceoy@users.noreply.github.com"}]
|
||||||
|
|||||||
@@ -30,6 +30,7 @@ from mt5cli.history import (
|
|||||||
create_rate_compatibility_views,
|
create_rate_compatibility_views,
|
||||||
deduplicate_history_tables,
|
deduplicate_history_tables,
|
||||||
drop_duplicates_in_table,
|
drop_duplicates_in_table,
|
||||||
|
drop_forming_rate_bar,
|
||||||
filter_incremental_history_deals_frame,
|
filter_incremental_history_deals_frame,
|
||||||
filter_trade_history_frame,
|
filter_trade_history_frame,
|
||||||
get_history_deals_account_event_start_datetime,
|
get_history_deals_account_event_start_datetime,
|
||||||
@@ -38,6 +39,7 @@ from mt5cli.history import (
|
|||||||
load_incremental_start_datetimes,
|
load_incremental_start_datetimes,
|
||||||
load_rate_data,
|
load_rate_data,
|
||||||
load_rate_data_from_connection,
|
load_rate_data_from_connection,
|
||||||
|
load_rate_series_by_granularity,
|
||||||
load_rate_series_from_sqlite,
|
load_rate_series_from_sqlite,
|
||||||
parse_sqlite_timestamp,
|
parse_sqlite_timestamp,
|
||||||
quote_sqlite_identifier,
|
quote_sqlite_identifier,
|
||||||
@@ -533,6 +535,46 @@ class TestResolveHistorySettings:
|
|||||||
assert resolve_granularity_name(1) == "M1"
|
assert resolve_granularity_name(1) == "M1"
|
||||||
|
|
||||||
|
|
||||||
|
class TestDropFormingRateBar:
|
||||||
|
"""Tests for drop_forming_rate_bar."""
|
||||||
|
|
||||||
|
def test_drops_still_forming_last_bar(self) -> None:
|
||||||
|
"""Test the still-forming last bar is removed."""
|
||||||
|
df_rate = pd.DataFrame(
|
||||||
|
{"time": [1, 2, 3], "close": [1.1, 1.2, 1.3]},
|
||||||
|
index=pd.Index(["a", "b", "c"], name="idx"),
|
||||||
|
)
|
||||||
|
|
||||||
|
result = drop_forming_rate_bar(df_rate)
|
||||||
|
|
||||||
|
pd.testing.assert_frame_equal(
|
||||||
|
result,
|
||||||
|
pd.DataFrame(
|
||||||
|
{"time": [1, 2], "close": [1.1, 1.2]},
|
||||||
|
index=pd.Index(["a", "b"], name="idx"),
|
||||||
|
),
|
||||||
|
)
|
||||||
|
assert df_rate.shape == (3, 2)
|
||||||
|
|
||||||
|
def test_returns_empty_frame_when_input_empty(self) -> None:
|
||||||
|
"""Test empty frames stay empty."""
|
||||||
|
df_rate = pd.DataFrame(columns=["time", "close"])
|
||||||
|
|
||||||
|
result = drop_forming_rate_bar(df_rate)
|
||||||
|
|
||||||
|
assert result.empty
|
||||||
|
assert list(result.columns) == ["time", "close"]
|
||||||
|
|
||||||
|
def test_returns_empty_frame_when_only_forming_bar_present(self) -> None:
|
||||||
|
"""Test a single-bar frame becomes empty after dropping the forming bar."""
|
||||||
|
df_rate = pd.DataFrame({"time": [1], "close": [1.1]})
|
||||||
|
|
||||||
|
result = drop_forming_rate_bar(df_rate)
|
||||||
|
|
||||||
|
assert result.empty
|
||||||
|
assert list(result.columns) == ["time", "close"]
|
||||||
|
|
||||||
|
|
||||||
class TestParseSqliteTimestamp:
|
class TestParseSqliteTimestamp:
|
||||||
"""Tests for parse_sqlite_timestamp."""
|
"""Tests for parse_sqlite_timestamp."""
|
||||||
|
|
||||||
@@ -2257,6 +2299,56 @@ class TestRateSourceHelpers:
|
|||||||
assert set(result) == {("EURUSD", 1)}
|
assert set(result) == {("EURUSD", 1)}
|
||||||
assert len(result["EURUSD", 1]) == 2
|
assert len(result["EURUSD", 1]) == 2
|
||||||
|
|
||||||
|
def test_load_rate_series_by_granularity(self, tmp_path: Path) -> None:
|
||||||
|
"""Test loading rate series keyed by symbol and granularity name."""
|
||||||
|
db_path = tmp_path / "granularity.db"
|
||||||
|
with sqlite3.connect(db_path) as conn:
|
||||||
|
conn.execute(
|
||||||
|
"CREATE TABLE rates("
|
||||||
|
" symbol TEXT, timeframe INTEGER, time TEXT, close REAL)",
|
||||||
|
)
|
||||||
|
conn.executemany(
|
||||||
|
"INSERT INTO rates(symbol, timeframe, time, close) VALUES (?, ?, ?, ?)",
|
||||||
|
[
|
||||||
|
("EURUSD", 1, "2024-01-01T00:00:00+00:00", 1.0),
|
||||||
|
("EURUSD", 16385, "2024-01-01T00:00:00+00:00", 1.1),
|
||||||
|
],
|
||||||
|
)
|
||||||
|
create_rate_compatibility_views(conn)
|
||||||
|
|
||||||
|
result = load_rate_series_by_granularity(
|
||||||
|
db_path,
|
||||||
|
["EURUSD"],
|
||||||
|
["M1", "H1"],
|
||||||
|
count=1,
|
||||||
|
)
|
||||||
|
|
||||||
|
assert set(result) == {("EURUSD", "M1"), ("EURUSD", "H1")}
|
||||||
|
|
||||||
|
def test_load_rate_series_by_granularity_explicit_tables(
|
||||||
|
self,
|
||||||
|
tmp_path: Path,
|
||||||
|
) -> None:
|
||||||
|
"""Test explicit tables with None-symbol targets key by granularity."""
|
||||||
|
db_path = tmp_path / "granularity-explicit.db"
|
||||||
|
with sqlite3.connect(db_path) as conn:
|
||||||
|
conn.execute("CREATE TABLE custom_view(time TEXT, close REAL)")
|
||||||
|
conn.execute(
|
||||||
|
"INSERT INTO custom_view(time, close) VALUES (?, ?)",
|
||||||
|
("2024-01-01T00:00:00+00:00", 1.0),
|
||||||
|
)
|
||||||
|
|
||||||
|
result = load_rate_series_by_granularity(
|
||||||
|
db_path,
|
||||||
|
[],
|
||||||
|
["M1"],
|
||||||
|
count=1,
|
||||||
|
explicit_tables=["custom_view"],
|
||||||
|
allow_missing_symbol=True,
|
||||||
|
)
|
||||||
|
|
||||||
|
assert set(result) == {(None, "M1")}
|
||||||
|
|
||||||
def test_load_rate_series_reuses_path_connection(
|
def test_load_rate_series_reuses_path_connection(
|
||||||
self,
|
self,
|
||||||
tmp_path: Path,
|
tmp_path: Path,
|
||||||
|
|||||||
@@ -10,6 +10,7 @@ from unittest.mock import MagicMock, call
|
|||||||
|
|
||||||
import pandas as pd
|
import pandas as pd
|
||||||
import pytest
|
import pytest
|
||||||
|
from pdmt5 import Mt5RuntimeError, Mt5TradingError
|
||||||
from pytest_mock import MockerFixture # noqa: TC002
|
from pytest_mock import MockerFixture # noqa: TC002
|
||||||
|
|
||||||
if TYPE_CHECKING:
|
if TYPE_CHECKING:
|
||||||
@@ -22,11 +23,15 @@ from mt5cli.history import DEFAULT_HISTORY_TIMEFRAMES
|
|||||||
from mt5cli.sdk import (
|
from mt5cli.sdk import (
|
||||||
AccountSpec,
|
AccountSpec,
|
||||||
Mt5CliClient,
|
Mt5CliClient,
|
||||||
|
ThrottledHistoryUpdater,
|
||||||
account_info,
|
account_info,
|
||||||
build_config,
|
build_config,
|
||||||
collect_history,
|
collect_history,
|
||||||
|
collect_latest_closed_rates_by_granularity,
|
||||||
|
collect_latest_closed_rates_for_accounts,
|
||||||
collect_latest_rates,
|
collect_latest_rates,
|
||||||
collect_latest_rates_for_accounts,
|
collect_latest_rates_for_accounts,
|
||||||
|
collect_latest_rates_for_accounts_with_retries,
|
||||||
copy_rates_from,
|
copy_rates_from,
|
||||||
copy_rates_from_pos,
|
copy_rates_from_pos,
|
||||||
copy_rates_range,
|
copy_rates_range,
|
||||||
@@ -45,6 +50,9 @@ from mt5cli.sdk import (
|
|||||||
positions,
|
positions,
|
||||||
recent_history_deals,
|
recent_history_deals,
|
||||||
recent_ticks,
|
recent_ticks,
|
||||||
|
resolve_account_spec,
|
||||||
|
resolve_account_specs,
|
||||||
|
substitute_env_placeholders,
|
||||||
symbol_info,
|
symbol_info,
|
||||||
symbol_info_tick,
|
symbol_info_tick,
|
||||||
symbols,
|
symbols,
|
||||||
@@ -1425,3 +1433,502 @@ class TestCollectLatestRatesForAccounts:
|
|||||||
collect_latest_rates_for_accounts(accounts, ["M1"], count=1)
|
collect_latest_rates_for_accounts(accounts, ["M1"], count=1)
|
||||||
|
|
||||||
mt5_data_client.assert_not_called()
|
mt5_data_client.assert_not_called()
|
||||||
|
|
||||||
|
|
||||||
|
class TestCollectLatestRatesForAccountsWithRetries:
|
||||||
|
"""Tests for collect_latest_rates_for_accounts_with_retries."""
|
||||||
|
|
||||||
|
def test_returns_result_on_first_success(self, mocker: MockerFixture) -> None:
|
||||||
|
"""Test no retry happens when the first attempt succeeds."""
|
||||||
|
expected = {("EURUSD", 1): pd.DataFrame()}
|
||||||
|
wrapped = mocker.patch(
|
||||||
|
"mt5cli.sdk.collect_latest_rates_for_accounts",
|
||||||
|
return_value=expected,
|
||||||
|
)
|
||||||
|
sleep = mocker.patch("mt5cli.sdk.time.sleep")
|
||||||
|
accounts = [AccountSpec(symbols=["EURUSD"])]
|
||||||
|
|
||||||
|
result = collect_latest_rates_for_accounts_with_retries(
|
||||||
|
accounts,
|
||||||
|
["M1"],
|
||||||
|
count=1,
|
||||||
|
retry_count=3,
|
||||||
|
)
|
||||||
|
|
||||||
|
assert result is expected
|
||||||
|
assert wrapped.call_count == 1
|
||||||
|
sleep.assert_not_called()
|
||||||
|
|
||||||
|
def test_retries_then_succeeds(self, mocker: MockerFixture) -> None:
|
||||||
|
"""Test transient MT5 errors are retried with exponential backoff."""
|
||||||
|
expected = {("EURUSD", 1): pd.DataFrame()}
|
||||||
|
wrapped = mocker.patch(
|
||||||
|
"mt5cli.sdk.collect_latest_rates_for_accounts",
|
||||||
|
side_effect=[
|
||||||
|
Mt5TradingError("boom"),
|
||||||
|
Mt5RuntimeError("boom"),
|
||||||
|
expected,
|
||||||
|
],
|
||||||
|
)
|
||||||
|
sleep = mocker.patch("mt5cli.sdk.time.sleep")
|
||||||
|
accounts = [AccountSpec(symbols=["EURUSD"])]
|
||||||
|
|
||||||
|
result = collect_latest_rates_for_accounts_with_retries(
|
||||||
|
accounts,
|
||||||
|
["M1"],
|
||||||
|
count=1,
|
||||||
|
retry_count=2,
|
||||||
|
backoff_base=2,
|
||||||
|
)
|
||||||
|
|
||||||
|
assert result is expected
|
||||||
|
assert wrapped.call_count == 3
|
||||||
|
assert sleep.call_args_list == [call(2), call(4)]
|
||||||
|
|
||||||
|
def test_reraises_after_exhausting_retries(self, mocker: MockerFixture) -> None:
|
||||||
|
"""Test the final error is re-raised once retries are exhausted."""
|
||||||
|
wrapped = mocker.patch(
|
||||||
|
"mt5cli.sdk.collect_latest_rates_for_accounts",
|
||||||
|
side_effect=Mt5RuntimeError("boom"),
|
||||||
|
)
|
||||||
|
sleep = mocker.patch("mt5cli.sdk.time.sleep")
|
||||||
|
accounts = [AccountSpec(symbols=["EURUSD"])]
|
||||||
|
|
||||||
|
with pytest.raises(Mt5RuntimeError, match="boom"):
|
||||||
|
collect_latest_rates_for_accounts_with_retries(
|
||||||
|
accounts,
|
||||||
|
["M1"],
|
||||||
|
count=1,
|
||||||
|
retry_count=2,
|
||||||
|
)
|
||||||
|
|
||||||
|
assert wrapped.call_count == 3
|
||||||
|
assert sleep.call_count == 2
|
||||||
|
|
||||||
|
def test_does_not_retry_unrelated_errors(self, mocker: MockerFixture) -> None:
|
||||||
|
"""Test non-MT5 errors propagate without retrying."""
|
||||||
|
wrapped = mocker.patch(
|
||||||
|
"mt5cli.sdk.collect_latest_rates_for_accounts",
|
||||||
|
side_effect=ValueError("bad input"),
|
||||||
|
)
|
||||||
|
sleep = mocker.patch("mt5cli.sdk.time.sleep")
|
||||||
|
|
||||||
|
with pytest.raises(ValueError, match="bad input"):
|
||||||
|
collect_latest_rates_for_accounts_with_retries(
|
||||||
|
[AccountSpec(symbols=["EURUSD"])],
|
||||||
|
["M1"],
|
||||||
|
count=1,
|
||||||
|
retry_count=3,
|
||||||
|
)
|
||||||
|
|
||||||
|
assert wrapped.call_count == 1
|
||||||
|
sleep.assert_not_called()
|
||||||
|
|
||||||
|
|
||||||
|
class TestCollectLatestClosedRatesForAccounts:
|
||||||
|
"""Tests for collect_latest_closed_rates_for_accounts."""
|
||||||
|
|
||||||
|
def test_fetches_count_plus_one_and_drops_forming_bar(
|
||||||
|
self,
|
||||||
|
mocker: MockerFixture,
|
||||||
|
) -> None:
|
||||||
|
"""Test closed-bar collection requests one extra bar at start_pos=0."""
|
||||||
|
df_rate = pd.DataFrame({"time": [1, 2, 3], "close": [1.1, 1.2, 1.3]})
|
||||||
|
wrapped = mocker.patch(
|
||||||
|
"mt5cli.sdk.collect_latest_rates_for_accounts_with_retries",
|
||||||
|
return_value={("EURUSD", 1): df_rate},
|
||||||
|
)
|
||||||
|
accounts = [AccountSpec(symbols=["EURUSD"])]
|
||||||
|
|
||||||
|
result = collect_latest_closed_rates_for_accounts(
|
||||||
|
accounts,
|
||||||
|
["M1"],
|
||||||
|
count=2,
|
||||||
|
retry_count=1,
|
||||||
|
backoff_base=3,
|
||||||
|
)
|
||||||
|
|
||||||
|
wrapped.assert_called_once_with(
|
||||||
|
accounts,
|
||||||
|
["M1"],
|
||||||
|
3,
|
||||||
|
start_pos=0,
|
||||||
|
base_config=None,
|
||||||
|
retry_count=1,
|
||||||
|
backoff_base=3,
|
||||||
|
)
|
||||||
|
pd.testing.assert_frame_equal(
|
||||||
|
result["EURUSD", 1],
|
||||||
|
pd.DataFrame({"time": [1, 2], "close": [1.1, 1.2]}),
|
||||||
|
)
|
||||||
|
|
||||||
|
def test_rejects_forming_bar_only_frames(self, mocker: MockerFixture) -> None:
|
||||||
|
"""Test empty results after dropping the forming bar raise ValueError."""
|
||||||
|
mocker.patch(
|
||||||
|
"mt5cli.sdk.collect_latest_rates_for_accounts_with_retries",
|
||||||
|
return_value={("EURUSD", 1): pd.DataFrame({"time": [1], "close": [1.1]})},
|
||||||
|
)
|
||||||
|
|
||||||
|
with pytest.raises(ValueError, match="Rate data is empty"):
|
||||||
|
collect_latest_closed_rates_for_accounts(
|
||||||
|
[AccountSpec(symbols=["EURUSD"])],
|
||||||
|
["M1"],
|
||||||
|
count=1,
|
||||||
|
)
|
||||||
|
|
||||||
|
def test_skips_extra_fetch_when_start_pos_nonzero(
|
||||||
|
self,
|
||||||
|
mocker: MockerFixture,
|
||||||
|
) -> None:
|
||||||
|
"""Test start_pos > 0 fetches count bars without dropping the last row."""
|
||||||
|
df_rate = pd.DataFrame({"time": [1, 2], "close": [1.1, 1.2]})
|
||||||
|
wrapped = mocker.patch(
|
||||||
|
"mt5cli.sdk.collect_latest_rates_for_accounts_with_retries",
|
||||||
|
return_value={("EURUSD", 1): df_rate},
|
||||||
|
)
|
||||||
|
|
||||||
|
result = collect_latest_closed_rates_for_accounts(
|
||||||
|
[AccountSpec(symbols=["EURUSD"])],
|
||||||
|
["M1"],
|
||||||
|
count=2,
|
||||||
|
start_pos=1,
|
||||||
|
)
|
||||||
|
|
||||||
|
wrapped.assert_called_once_with(
|
||||||
|
[AccountSpec(symbols=["EURUSD"])],
|
||||||
|
["M1"],
|
||||||
|
2,
|
||||||
|
start_pos=1,
|
||||||
|
base_config=None,
|
||||||
|
retry_count=0,
|
||||||
|
backoff_base=2.0,
|
||||||
|
)
|
||||||
|
pd.testing.assert_frame_equal(result["EURUSD", 1], df_rate)
|
||||||
|
|
||||||
|
def test_rejects_zero_count_before_fetching(self, mocker: MockerFixture) -> None:
|
||||||
|
"""Test count=0 is rejected before any MT5 collection attempt."""
|
||||||
|
wrapped = mocker.patch(
|
||||||
|
"mt5cli.sdk.collect_latest_rates_for_accounts_with_retries",
|
||||||
|
)
|
||||||
|
|
||||||
|
with pytest.raises(ValueError, match="count must be positive"):
|
||||||
|
collect_latest_closed_rates_for_accounts(
|
||||||
|
[AccountSpec(symbols=["EURUSD"])],
|
||||||
|
["M1"],
|
||||||
|
count=0,
|
||||||
|
)
|
||||||
|
|
||||||
|
wrapped.assert_not_called()
|
||||||
|
|
||||||
|
def test_rejects_negative_start_pos(self, mocker: MockerFixture) -> None:
|
||||||
|
"""Test negative start_pos is rejected before any MT5 collection attempt."""
|
||||||
|
wrapped = mocker.patch(
|
||||||
|
"mt5cli.sdk.collect_latest_rates_for_accounts_with_retries",
|
||||||
|
)
|
||||||
|
|
||||||
|
with pytest.raises(ValueError, match="start_pos must be non-negative"):
|
||||||
|
collect_latest_closed_rates_for_accounts(
|
||||||
|
[AccountSpec(symbols=["EURUSD"])],
|
||||||
|
["M1"],
|
||||||
|
count=1,
|
||||||
|
start_pos=-1,
|
||||||
|
)
|
||||||
|
|
||||||
|
wrapped.assert_not_called()
|
||||||
|
|
||||||
|
def test_rejects_empty_frames_with_start_pos_nonzero(
|
||||||
|
self,
|
||||||
|
mocker: MockerFixture,
|
||||||
|
) -> None:
|
||||||
|
"""Test empty upstream frames raise ValueError when start_pos > 0."""
|
||||||
|
mocker.patch(
|
||||||
|
"mt5cli.sdk.collect_latest_rates_for_accounts_with_retries",
|
||||||
|
return_value={("EURUSD", 1): pd.DataFrame(columns=["time", "close"])},
|
||||||
|
)
|
||||||
|
|
||||||
|
with pytest.raises(ValueError, match="Rate data is empty"):
|
||||||
|
collect_latest_closed_rates_for_accounts(
|
||||||
|
[AccountSpec(symbols=["EURUSD"])],
|
||||||
|
["M1"],
|
||||||
|
count=1,
|
||||||
|
start_pos=1,
|
||||||
|
)
|
||||||
|
|
||||||
|
def test_processes_multiple_symbol_timeframe_pairs(
|
||||||
|
self,
|
||||||
|
mocker: MockerFixture,
|
||||||
|
) -> None:
|
||||||
|
"""Test each returned series is trimmed and validated independently."""
|
||||||
|
mocker.patch(
|
||||||
|
"mt5cli.sdk.collect_latest_rates_for_accounts_with_retries",
|
||||||
|
return_value={
|
||||||
|
("EURUSD", 1): pd.DataFrame(
|
||||||
|
{"time": [1, 2, 3], "close": [1.1, 1.2, 1.3]},
|
||||||
|
),
|
||||||
|
("GBPUSD", 16385): pd.DataFrame(
|
||||||
|
{"time": [4, 5, 6], "close": [2.1, 2.2, 2.3]},
|
||||||
|
),
|
||||||
|
},
|
||||||
|
)
|
||||||
|
|
||||||
|
result = collect_latest_closed_rates_for_accounts(
|
||||||
|
[AccountSpec(symbols=["EURUSD", "GBPUSD"])],
|
||||||
|
["M1", "H1"],
|
||||||
|
count=2,
|
||||||
|
)
|
||||||
|
|
||||||
|
assert set(result) == {("EURUSD", 1), ("GBPUSD", 16385)}
|
||||||
|
pd.testing.assert_frame_equal(
|
||||||
|
result["EURUSD", 1],
|
||||||
|
pd.DataFrame({"time": [1, 2], "close": [1.1, 1.2]}),
|
||||||
|
)
|
||||||
|
pd.testing.assert_frame_equal(
|
||||||
|
result["GBPUSD", 16385],
|
||||||
|
pd.DataFrame({"time": [4, 5], "close": [2.1, 2.2]}),
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
class TestCollectLatestClosedRatesByGranularity:
|
||||||
|
"""Tests for collect_latest_closed_rates_by_granularity."""
|
||||||
|
|
||||||
|
def test_rekeys_by_granularity_name(self, mocker: MockerFixture) -> None:
|
||||||
|
"""Test closed rates are keyed by symbol and granularity name."""
|
||||||
|
df_rate = pd.DataFrame({"time": [1, 2], "close": [1.1, 1.2]})
|
||||||
|
wrapped = mocker.patch(
|
||||||
|
"mt5cli.sdk.collect_latest_closed_rates_for_accounts",
|
||||||
|
return_value={("EURUSD", 1): df_rate},
|
||||||
|
)
|
||||||
|
|
||||||
|
result = collect_latest_closed_rates_by_granularity(
|
||||||
|
[AccountSpec(symbols=["EURUSD"])],
|
||||||
|
["M1"],
|
||||||
|
count=2,
|
||||||
|
)
|
||||||
|
|
||||||
|
wrapped.assert_called_once_with(
|
||||||
|
[AccountSpec(symbols=["EURUSD"])],
|
||||||
|
["M1"],
|
||||||
|
2,
|
||||||
|
start_pos=0,
|
||||||
|
base_config=None,
|
||||||
|
retry_count=0,
|
||||||
|
backoff_base=2.0,
|
||||||
|
)
|
||||||
|
assert ("EURUSD", "M1") in result
|
||||||
|
pd.testing.assert_frame_equal(result["EURUSD", "M1"], df_rate)
|
||||||
|
|
||||||
|
|
||||||
|
class TestSubstituteEnvPlaceholders:
|
||||||
|
"""Tests for ${ENV_VAR} substitution."""
|
||||||
|
|
||||||
|
def test_substitutes_known_variables(
|
||||||
|
self,
|
||||||
|
monkeypatch: pytest.MonkeyPatch,
|
||||||
|
) -> None:
|
||||||
|
"""Test placeholders are replaced with environment values."""
|
||||||
|
monkeypatch.setenv("MT5_LOGIN", "12345")
|
||||||
|
monkeypatch.setenv("MT5_SERVER", "Broker-Demo")
|
||||||
|
|
||||||
|
assert substitute_env_placeholders("${MT5_LOGIN}") == "12345"
|
||||||
|
assert substitute_env_placeholders("srv=${MT5_SERVER}!") == "srv=Broker-Demo!"
|
||||||
|
|
||||||
|
def test_returns_plain_strings_unchanged(self) -> None:
|
||||||
|
"""Test strings without placeholders are returned as-is."""
|
||||||
|
assert substitute_env_placeholders("plain") == "plain"
|
||||||
|
|
||||||
|
def test_raises_on_missing_variable(
|
||||||
|
self,
|
||||||
|
monkeypatch: pytest.MonkeyPatch,
|
||||||
|
) -> None:
|
||||||
|
"""Test a missing environment variable raises a clear error."""
|
||||||
|
monkeypatch.delenv("MT5_MISSING", raising=False)
|
||||||
|
|
||||||
|
with pytest.raises(ValueError, match="'MT5_MISSING' is not set"):
|
||||||
|
substitute_env_placeholders("${MT5_MISSING}")
|
||||||
|
|
||||||
|
|
||||||
|
class TestResolveAccountSpec:
|
||||||
|
"""Tests for resolve_account_spec and resolve_account_specs."""
|
||||||
|
|
||||||
|
def test_substitutes_env_placeholders_in_account(
|
||||||
|
self,
|
||||||
|
monkeypatch: pytest.MonkeyPatch,
|
||||||
|
) -> None:
|
||||||
|
"""Test account string fields resolve ${ENV_VAR} placeholders."""
|
||||||
|
monkeypatch.setenv("MT5_PASSWORD", "secret")
|
||||||
|
account = AccountSpec(
|
||||||
|
symbols=["EURUSD"],
|
||||||
|
login="${MT5_LOGIN}",
|
||||||
|
password="${MT5_PASSWORD}",
|
||||||
|
)
|
||||||
|
monkeypatch.setenv("MT5_LOGIN", "999")
|
||||||
|
|
||||||
|
resolved = resolve_account_spec(account)
|
||||||
|
|
||||||
|
assert resolved.login == "999"
|
||||||
|
assert resolved.password == "secret" # noqa: S105
|
||||||
|
assert resolved.symbols == ["EURUSD"]
|
||||||
|
|
||||||
|
def test_explicit_overrides_take_precedence(self) -> None:
|
||||||
|
"""Test explicit override values win over account fields."""
|
||||||
|
account = AccountSpec(symbols=["EURUSD"], login=111, server="Acct")
|
||||||
|
|
||||||
|
resolved = resolve_account_spec(
|
||||||
|
account,
|
||||||
|
login=222,
|
||||||
|
server="Override",
|
||||||
|
timeout=5000,
|
||||||
|
)
|
||||||
|
|
||||||
|
assert resolved.login == 222
|
||||||
|
assert resolved.server == "Override"
|
||||||
|
assert resolved.timeout == 5000
|
||||||
|
|
||||||
|
def test_resolves_string_login_override(
|
||||||
|
self,
|
||||||
|
monkeypatch: pytest.MonkeyPatch,
|
||||||
|
) -> None:
|
||||||
|
"""Test string login overrides expand ${ENV_VAR} placeholders."""
|
||||||
|
monkeypatch.setenv("MT5_LOGIN", "777")
|
||||||
|
account = AccountSpec(symbols=["EURUSD"], login=111)
|
||||||
|
|
||||||
|
resolved = resolve_account_spec(account, login="${MT5_LOGIN}")
|
||||||
|
|
||||||
|
assert resolved.login == "777"
|
||||||
|
|
||||||
|
def test_preserves_integer_login_without_coercion(self) -> None:
|
||||||
|
"""Test integer logins remain integers after resolution."""
|
||||||
|
account = AccountSpec(symbols=["EURUSD"], login=111)
|
||||||
|
|
||||||
|
resolved = resolve_account_spec(account)
|
||||||
|
|
||||||
|
assert resolved.login == 111
|
||||||
|
assert isinstance(resolved.login, int)
|
||||||
|
|
||||||
|
def test_raises_on_missing_env_variable(
|
||||||
|
self,
|
||||||
|
monkeypatch: pytest.MonkeyPatch,
|
||||||
|
) -> None:
|
||||||
|
"""Test missing environment variables raise ValueError."""
|
||||||
|
monkeypatch.delenv("MT5_NOPE", raising=False)
|
||||||
|
account = AccountSpec(symbols=["EURUSD"], server="${MT5_NOPE}")
|
||||||
|
|
||||||
|
with pytest.raises(ValueError, match="'MT5_NOPE' is not set"):
|
||||||
|
resolve_account_spec(account)
|
||||||
|
|
||||||
|
def test_resolve_account_specs_applies_to_all(
|
||||||
|
self,
|
||||||
|
monkeypatch: pytest.MonkeyPatch,
|
||||||
|
) -> None:
|
||||||
|
"""Test resolve_account_specs resolves every account in order."""
|
||||||
|
monkeypatch.setenv("MT5_SERVER", "Shared")
|
||||||
|
accounts = [
|
||||||
|
AccountSpec(symbols=["EURUSD"], server="${MT5_SERVER}"),
|
||||||
|
AccountSpec(symbols=["GBPUSD"], server="Fixed"),
|
||||||
|
]
|
||||||
|
|
||||||
|
resolved = resolve_account_specs(accounts, timeout=1000)
|
||||||
|
|
||||||
|
assert [a.server for a in resolved] == ["Shared", "Fixed"]
|
||||||
|
assert all(a.timeout == 1000 for a in resolved)
|
||||||
|
|
||||||
|
|
||||||
|
class TestThrottledHistoryUpdater:
|
||||||
|
"""Tests for the throttled incremental history updater."""
|
||||||
|
|
||||||
|
def test_updates_every_call_when_interval_non_positive(
|
||||||
|
self,
|
||||||
|
mocker: MockerFixture,
|
||||||
|
) -> None:
|
||||||
|
"""Test interval_seconds <= 0 updates on every call."""
|
||||||
|
update = mocker.patch("mt5cli.sdk.update_history")
|
||||||
|
client = MagicMock()
|
||||||
|
updater = ThrottledHistoryUpdater(output="history.db", interval_seconds=0)
|
||||||
|
|
||||||
|
assert updater.update(client, ["EURUSD"]) is True
|
||||||
|
assert updater.update(client, ["EURUSD"]) is True
|
||||||
|
assert update.call_count == 2
|
||||||
|
|
||||||
|
def test_throttles_within_interval(self, mocker: MockerFixture) -> None:
|
||||||
|
"""Test updates are skipped until the interval elapses."""
|
||||||
|
update = mocker.patch("mt5cli.sdk.update_history")
|
||||||
|
monotonic = mocker.patch("mt5cli.sdk.time.monotonic")
|
||||||
|
# Calls: set(t=100), check(t=105), check(t=200), set(t=200).
|
||||||
|
monotonic.side_effect = [100.0, 105.0, 200.0, 200.0]
|
||||||
|
client = MagicMock()
|
||||||
|
updater = ThrottledHistoryUpdater(output="history.db", interval_seconds=60)
|
||||||
|
|
||||||
|
assert updater.update(client, ["EURUSD"]) is True # first update at t=100
|
||||||
|
assert updater.update(client, ["EURUSD"]) is False # t=105, throttled
|
||||||
|
assert updater.update(client, ["EURUSD"]) is True # t=200, elapsed
|
||||||
|
assert update.call_count == 2
|
||||||
|
|
||||||
|
def test_update_passes_expected_arguments(
|
||||||
|
self,
|
||||||
|
mocker: MockerFixture,
|
||||||
|
) -> None:
|
||||||
|
"""Test update_history is called with the configured arguments."""
|
||||||
|
update = mocker.patch("mt5cli.sdk.update_history")
|
||||||
|
client = MagicMock()
|
||||||
|
updater = ThrottledHistoryUpdater(
|
||||||
|
output="history.db",
|
||||||
|
datasets={Dataset.rates},
|
||||||
|
timeframes=["M1", "H1"],
|
||||||
|
flags="INFO",
|
||||||
|
lookback_hours=12.0,
|
||||||
|
with_views=True,
|
||||||
|
include_account_events=False,
|
||||||
|
)
|
||||||
|
|
||||||
|
updater.update(client, ["EURUSD", "GBPUSD"])
|
||||||
|
|
||||||
|
update.assert_called_once_with(
|
||||||
|
client=client,
|
||||||
|
output="history.db",
|
||||||
|
symbols=["EURUSD", "GBPUSD"],
|
||||||
|
datasets={Dataset.rates},
|
||||||
|
timeframes=["M1", "H1"],
|
||||||
|
flags="INFO",
|
||||||
|
lookback_hours=12.0,
|
||||||
|
with_views=True,
|
||||||
|
include_account_events=False,
|
||||||
|
)
|
||||||
|
|
||||||
|
def test_propagates_errors_by_default(self, mocker: MockerFixture) -> None:
|
||||||
|
"""Test MT5/SQLite errors propagate and do not advance the throttle."""
|
||||||
|
mocker.patch(
|
||||||
|
"mt5cli.sdk.update_history",
|
||||||
|
side_effect=Mt5RuntimeError("boom"),
|
||||||
|
)
|
||||||
|
updater = ThrottledHistoryUpdater(output="history.db")
|
||||||
|
|
||||||
|
with pytest.raises(Mt5RuntimeError, match="boom"):
|
||||||
|
updater.update(MagicMock(), ["EURUSD"])
|
||||||
|
|
||||||
|
assert updater.last_update_monotonic is None
|
||||||
|
|
||||||
|
@pytest.mark.parametrize(
|
||||||
|
"error",
|
||||||
|
[
|
||||||
|
Mt5RuntimeError("boom"),
|
||||||
|
Mt5TradingError("trade failed"),
|
||||||
|
sqlite3.OperationalError("locked"),
|
||||||
|
],
|
||||||
|
)
|
||||||
|
def test_suppresses_errors_when_requested(
|
||||||
|
self,
|
||||||
|
mocker: MockerFixture,
|
||||||
|
error: Exception,
|
||||||
|
) -> None:
|
||||||
|
"""Test suppress_errors swallows recoverable errors and returns False."""
|
||||||
|
mocker.patch(
|
||||||
|
"mt5cli.sdk.update_history",
|
||||||
|
side_effect=error,
|
||||||
|
)
|
||||||
|
updater = ThrottledHistoryUpdater(
|
||||||
|
output="history.db",
|
||||||
|
suppress_errors=True,
|
||||||
|
)
|
||||||
|
|
||||||
|
assert updater.update(MagicMock(), ["EURUSD"]) is False
|
||||||
|
assert updater.last_update_monotonic is None
|
||||||
|
|||||||
@@ -487,7 +487,7 @@ wheels = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "mt5cli"
|
name = "mt5cli"
|
||||||
version = "0.5.2"
|
version = "0.6.0"
|
||||||
source = { editable = "." }
|
source = { editable = "." }
|
||||||
dependencies = [
|
dependencies = [
|
||||||
{ name = "click" },
|
{ name = "click" },
|
||||||
@@ -836,11 +836,11 @@ wheels = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "pygments"
|
name = "pygments"
|
||||||
version = "2.19.2"
|
version = "2.20.0"
|
||||||
source = { registry = "https://pypi.org/simple" }
|
source = { registry = "https://pypi.org/simple" }
|
||||||
sdist = { url = "https://files.pythonhosted.org/packages/b0/77/a5b8c569bf593b0140bde72ea885a803b82086995367bf2037de0159d924/pygments-2.19.2.tar.gz", hash = "sha256:636cb2477cec7f8952536970bc533bc43743542f70392ae026374600add5b887", size = 4968631, upload-time = "2025-06-21T13:39:12.283Z" }
|
sdist = { url = "https://files.pythonhosted.org/packages/c3/b2/bc9c9196916376152d655522fdcebac55e66de6603a76a02bca1b6414f6c/pygments-2.20.0.tar.gz", hash = "sha256:6757cd03768053ff99f3039c1a36d6c0aa0b263438fcab17520b30a303a82b5f", size = 4955991, upload-time = "2026-03-29T13:29:33.898Z" }
|
||||||
wheels = [
|
wheels = [
|
||||||
{ url = "https://files.pythonhosted.org/packages/c7/21/705964c7812476f378728bdf590ca4b771ec72385c533964653c68e86bdc/pygments-2.19.2-py3-none-any.whl", hash = "sha256:86540386c03d588bb81d44bc3928634ff26449851e99741617ecb9037ee5ec0b", size = 1225217, upload-time = "2025-06-21T13:39:07.939Z" },
|
{ url = "https://files.pythonhosted.org/packages/f4/7e/a72dd26f3b0f4f2bf1dd8923c85f7ceb43172af56d63c7383eb62b332364/pygments-2.20.0-py3-none-any.whl", hash = "sha256:81a9e26dd42fd28a23a2d169d86d7ac03b46e2f8b59ed4698fb4785f946d0176", size = 1231151, upload-time = "2026-03-29T13:29:30.038Z" },
|
||||||
]
|
]
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
|
|||||||
Reference in New Issue
Block a user