Compare commits
3 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| b2bb2ad0a0 | |||
| 756faf747b | |||
| c4232bf44d |
@@ -57,6 +57,7 @@ python -m mt5cli -o account.csv account-info
|
||||
| `rates-range` | Export rates for a date range |
|
||||
| `ticks-from` | Export ticks from a start date |
|
||||
| `ticks-range` | Export ticks for a date range |
|
||||
| `ticks-recent` | Export ticks from a recent trailing window |
|
||||
| `account-info` | Export account information |
|
||||
| `terminal-info` | Export terminal information |
|
||||
| `version` | Export MetaTrader 5 version information |
|
||||
@@ -64,6 +65,7 @@ python -m mt5cli -o account.csv account-info
|
||||
| `symbols` | Export symbol list |
|
||||
| `symbol-info` | Export symbol details |
|
||||
| `symbol-info-tick` | Export the last tick for a symbol |
|
||||
| `minimum-margins` | Export minimum-volume buy and sell margin requirements |
|
||||
| `market-book` | Export market depth (order book) |
|
||||
| `orders` | Export active orders |
|
||||
| `positions` | Export open positions |
|
||||
@@ -87,7 +89,49 @@ mt5cli -o history.db collect-history \
|
||||
--timeframe M1 --flags ALL --if-exists append --with-views
|
||||
```
|
||||
|
||||
History orders and deals are fetched per symbol and concatenated, so the symbol filter is applied consistently across all datasets. The `cash_events` view is derived from symbol-filtered `history_deals`, so account-level cash events with empty or non-matching symbols may be excluded. The `rates` table records the requested `timeframe` so appended runs at different timeframes remain distinguishable. The `positions_reconstructed` view aggregates trade deals by `position_id`, excludes positions without closing deals, and uses volume-weighted open/close prices; reversal deals (`DEAL_ENTRY_INOUT`) are reported via `volume_reversal` / `reversal_count` columns and do not contribute to the weighted prices.
|
||||
History orders and deals are fetched per symbol and concatenated, so the symbol filter is applied consistently across all datasets. The `cash_events` view is derived from symbol-filtered `history_deals`, so account-level cash events with empty or non-matching symbols may be excluded. The `rates` table records the requested `timeframe` so appended runs at different timeframes remain distinguishable. The `positions_reconstructed` view aggregates trade deals by `position_id`, excludes positions without closing-side entries, and uses volume-weighted open/close prices; reversal deals (`DEAL_ENTRY_INOUT`) are reported via `volume_reversal` / `reversal_count` columns.
|
||||
|
||||
### Incremental history SDK
|
||||
|
||||
For automated pipelines, use the importable incremental API instead of re-fetching fixed date ranges:
|
||||
|
||||
```python
|
||||
from pdmt5 import Mt5Config, Mt5DataClient
|
||||
from mt5cli import Dataset, update_history, update_history_with_config
|
||||
|
||||
# Reuse an already-connected pdmt5 client (does not open/close MT5)
|
||||
client = Mt5DataClient(config=Mt5Config(login=12345))
|
||||
client.initialize_and_login_mt5()
|
||||
try:
|
||||
update_history(
|
||||
client=client,
|
||||
output="history.db",
|
||||
symbols=["EURUSD", "GBPUSD"],
|
||||
datasets={Dataset.rates, Dataset.history_deals},
|
||||
timeframes=["M1", "H1"], # default: all fixed MT5 timeframes
|
||||
lookback_hours=24,
|
||||
create_rate_views=True,
|
||||
with_views=True,
|
||||
include_account_events=True,
|
||||
)
|
||||
finally:
|
||||
client.shutdown()
|
||||
|
||||
# Standalone wrapper that opens and closes MT5 for you
|
||||
update_history_with_config(
|
||||
output="history.db",
|
||||
symbols=["EURUSD"],
|
||||
config=Mt5Config(login=12345),
|
||||
)
|
||||
```
|
||||
|
||||
- **`collect-history`**: explicit date-range export into SQLite.
|
||||
- **`update_history`**: incremental append based on existing SQLite `MAX(time)` per symbol (and timeframe for rates); account-level deals use a separate cursor when `include_account_events=True`.
|
||||
- **`rates` table**: normalized storage with `symbol` and `timeframe` columns.
|
||||
- **Rate compatibility views**: mt5cli manages all `rate_*` views. Naming is `rate_<symbol>__<timeframe>` when a symbol has one timeframe, otherwise `rate_<symbol>__<granularity>_<timeframe>` (for example `rate_EURUSD__M1_1`). Stale `rate_*` views are dropped and recreated when rates change for offline tools such as mteor optimize.
|
||||
- **Rate view resolution**: use `mt5cli.history.resolve_rate_view_name()` / `resolve_rate_view_names()` to map symbols and granularities to existing SQLite compatibility views without creating databases.
|
||||
- **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.
|
||||
|
||||
## Requirements
|
||||
|
||||
|
||||
@@ -0,0 +1,166 @@
|
||||
# History Collection (SQLite)
|
||||
|
||||
::: mt5cli.history
|
||||
|
||||
## `collect-history` schema
|
||||
|
||||
The `collect-history` command (and the matching `collect_history` SDK function) writes
|
||||
selected MT5 datasets into one SQLite database. Each dataset becomes a table; column
|
||||
names and types mirror the pdmt5 DataFrame schema for that export, with two additions:
|
||||
|
||||
- `symbol` is prepended on every table.
|
||||
- `timeframe` is prepended on `rates` so appended runs at different bar sizes stay
|
||||
distinguishable.
|
||||
|
||||
SQLite does not declare foreign keys. Rows are linked logically by `symbol`, time
|
||||
windows, and (for deals) `position_id` / `order`. Duplicate rows are removed on
|
||||
append using dataset-specific keys (for example `ticket` on history tables, or
|
||||
`(symbol, timeframe, time)` on rates).
|
||||
|
||||
Optional views are created when `--with-views` is set and the `history-deals` dataset
|
||||
was written.
|
||||
|
||||
### Entity-relationship diagram
|
||||
|
||||
Sample layout for a full collection with `--with-views`:
|
||||
|
||||
```mermaid
|
||||
erDiagram
|
||||
rates {
|
||||
TEXT symbol "dedup key"
|
||||
INTEGER timeframe "dedup key"
|
||||
TEXT time "dedup key"
|
||||
REAL open
|
||||
REAL high
|
||||
REAL low
|
||||
REAL close
|
||||
INTEGER tick_volume
|
||||
INTEGER spread
|
||||
INTEGER real_volume
|
||||
}
|
||||
|
||||
ticks {
|
||||
TEXT symbol "dedup key"
|
||||
TEXT time "dedup key"
|
||||
INTEGER time_msc "dedup key (preferred)"
|
||||
REAL bid
|
||||
REAL ask
|
||||
REAL last
|
||||
INTEGER volume
|
||||
INTEGER flags
|
||||
REAL volume_real
|
||||
}
|
||||
|
||||
history_orders {
|
||||
INTEGER ticket "dedup key"
|
||||
TEXT symbol
|
||||
TEXT time
|
||||
INTEGER type
|
||||
INTEGER state
|
||||
REAL volume_initial
|
||||
REAL price_open
|
||||
REAL price_current
|
||||
INTEGER magic
|
||||
}
|
||||
|
||||
history_deals {
|
||||
INTEGER ticket "dedup key"
|
||||
INTEGER order
|
||||
INTEGER position_id "groups position view"
|
||||
TEXT symbol
|
||||
TEXT time
|
||||
INTEGER type "0/1 trade, else cash event"
|
||||
INTEGER entry "0 IN, 1 OUT, 2 INOUT, 3 OUT_BY"
|
||||
REAL volume
|
||||
REAL price
|
||||
REAL profit
|
||||
REAL commission
|
||||
REAL swap
|
||||
REAL fee
|
||||
}
|
||||
|
||||
cash_events {
|
||||
INTEGER ticket
|
||||
TEXT symbol
|
||||
TEXT time
|
||||
INTEGER type
|
||||
REAL profit
|
||||
}
|
||||
|
||||
positions_reconstructed {
|
||||
INTEGER position_id
|
||||
TEXT symbol
|
||||
TEXT open_time
|
||||
TEXT close_time
|
||||
INTEGER direction
|
||||
REAL volume_open
|
||||
REAL volume_close
|
||||
REAL volume_reversal
|
||||
REAL open_price
|
||||
REAL close_price
|
||||
REAL total_profit
|
||||
INTEGER reversal_count
|
||||
INTEGER deals_count
|
||||
}
|
||||
|
||||
rates ||--o{ history_deals : "symbol (logical)"
|
||||
ticks ||--o{ history_deals : "symbol (logical)"
|
||||
history_orders ||--o{ history_deals : "order ~ ticket (logical)"
|
||||
history_deals ||--|| cash_events : "VIEW: type NOT IN (0,1)"
|
||||
history_deals ||--o{ positions_reconstructed : "VIEW: GROUP BY position_id"
|
||||
```
|
||||
|
||||
### Tables and views
|
||||
|
||||
| Object | Kind | Source | Notes |
|
||||
| ------------------------- | ----- | -------------------- | ------------------------------------------------------------------------------------------- |
|
||||
| `rates` | table | `copy_rates_range` | Indexed on `(symbol, timeframe, time)` when columns exist. |
|
||||
| `ticks` | table | `copy_ticks_range` | Indexed on `(symbol, time)` when columns exist. |
|
||||
| `history_orders` | table | `history_orders_get` | Fetched per `--symbol`, then concatenated. |
|
||||
| `history_deals` | table | `history_deals_get` | Fetched per `--symbol`, then concatenated. Indexed on `(position_id, symbol)` when present. |
|
||||
| `cash_events` | view | `history_deals` | Non-trade deal types (deposits, balance ops, etc.). Requires `type` column. |
|
||||
| `positions_reconstructed` | view | `history_deals` | One row per closed `position_id`; volume-weighted prices and reversal stats. |
|
||||
|
||||
Column sets can vary with terminal and pdmt5 version. Views are skipped with a warning
|
||||
when required columns are missing.
|
||||
|
||||
### Incremental collection
|
||||
|
||||
The `update_history` SDK path uses the same base tables and optional
|
||||
`cash_events` / `positions_reconstructed` views. It additionally maintains
|
||||
`rate_<symbol>__<timeframe>` compatibility views when `create_rate_views=True`.
|
||||
|
||||
### Rate view resolution
|
||||
|
||||
Downstream tools can resolve mt5cli-managed compatibility view names from an
|
||||
existing SQLite history database without creating files or guessing legacy
|
||||
naming schemes:
|
||||
|
||||
```python
|
||||
from pathlib import Path
|
||||
|
||||
from mt5cli.history import resolve_rate_view_name, resolve_rate_view_names
|
||||
|
||||
# Single symbol and granularity
|
||||
view = resolve_rate_view_name(Path("history.db"), "EURUSD", "M1")
|
||||
|
||||
# Batch resolution in row-major order
|
||||
views = resolve_rate_view_names(
|
||||
Path("history.db"),
|
||||
["EURUSD", "GBPUSD"],
|
||||
["M1", "H1"],
|
||||
)
|
||||
```
|
||||
|
||||
Resolution rules:
|
||||
|
||||
- Returns `rate_<symbol>__<timeframe>` when a symbol stores one timeframe.
|
||||
- Returns `rate_<symbol>__<granularity>_<timeframe>` when multiple timeframes
|
||||
are stored for the same symbol.
|
||||
- When multiple naming candidates apply, prefers an existing managed
|
||||
`rate_*__*` view from the candidate list.
|
||||
- Falls back to single-timeframe naming when the database path is missing or
|
||||
`rates` metadata is unavailable.
|
||||
- Pass `require_existing=True` to raise `ValueError` instead of returning a
|
||||
best-guess name when the database or view is missing.
|
||||
- Accepts either a SQLite path or an open `sqlite3.Connection`.
|
||||
@@ -18,6 +18,10 @@ Utility module providing constants, enums, Click parameter types, and helper fun
|
||||
|
||||
Programmatic SDK for read-only MetaTrader 5 data collection. Returns pandas DataFrames and provides `collect_history` for SQLite bulk collection.
|
||||
|
||||
### [History Collection (SQLite)](history.md)
|
||||
|
||||
SQLite storage helpers for the `collect-history` command schema, incremental updates, deduplication, indexes, and optional views.
|
||||
|
||||
## Architecture Overview
|
||||
|
||||
The package follows a simple architecture built on top of pdmt5:
|
||||
@@ -61,12 +65,18 @@ from datetime import UTC, datetime
|
||||
from pathlib import Path
|
||||
|
||||
from mt5cli import (
|
||||
Dataset,
|
||||
IfExists,
|
||||
Mt5CliClient,
|
||||
collect_history,
|
||||
copy_rates_range,
|
||||
detect_format,
|
||||
export_dataframe,
|
||||
export_dataframe_to_sqlite,
|
||||
minimum_margins,
|
||||
recent_ticks,
|
||||
)
|
||||
from mt5cli.history import resolve_rate_view_name
|
||||
|
||||
# Fetch rates programmatically
|
||||
rates = copy_rates_range(
|
||||
@@ -82,6 +92,20 @@ fmt = detect_format(Path("output.parquet")) # Returns "parquet"
|
||||
# Export a DataFrame
|
||||
export_dataframe(rates, Path("output.csv"), "csv")
|
||||
|
||||
# Append to SQLite with deduplication
|
||||
export_dataframe_to_sqlite(
|
||||
rates,
|
||||
Path("history.db"),
|
||||
"rates",
|
||||
if_exists=IfExists.APPEND,
|
||||
deduplicate_on=("symbol", "timeframe", "time"),
|
||||
)
|
||||
|
||||
# Resolve rate compatibility views and fetch recent ticks
|
||||
view = resolve_rate_view_name(Path("history.db"), "EURUSD", "M1")
|
||||
ticks = recent_ticks("EURUSD", seconds=300)
|
||||
margins = minimum_margins("EURUSD")
|
||||
|
||||
# Collect history into SQLite
|
||||
collect_history(
|
||||
Path("history.db"),
|
||||
|
||||
+26
-6
@@ -22,13 +22,22 @@ pip install mt5cli
|
||||
|
||||
## Programmatic usage / SDK usage
|
||||
|
||||
mt5cli can be used as a small Python SDK for read-only MetaTrader 5 data collection. SDK functions return pandas DataFrames without writing files. Use `export_dataframe` when you need to persist results.
|
||||
mt5cli can be used as a small Python SDK for read-only MetaTrader 5 data collection. SDK functions return pandas DataFrames without writing files. Use `export_dataframe` or `export_dataframe_to_sqlite` when you need to persist results.
|
||||
|
||||
```python
|
||||
from datetime import UTC, datetime
|
||||
from pathlib import Path
|
||||
|
||||
from mt5cli import Mt5CliClient, collect_history, copy_rates_range, export_dataframe
|
||||
from mt5cli import (
|
||||
Mt5CliClient,
|
||||
collect_history,
|
||||
copy_rates_range,
|
||||
export_dataframe,
|
||||
export_dataframe_to_sqlite,
|
||||
minimum_margins,
|
||||
recent_ticks,
|
||||
)
|
||||
from mt5cli.history import resolve_rate_view_name
|
||||
|
||||
# One-off fetch with module-level helpers
|
||||
rates = copy_rates_range(
|
||||
@@ -39,6 +48,13 @@ rates = copy_rates_range(
|
||||
)
|
||||
export_dataframe(rates, Path("rates.csv"), "csv")
|
||||
|
||||
# Resolve SQLite rate compatibility views for downstream tools
|
||||
view = resolve_rate_view_name(Path("history.db"), "EURUSD", "M1")
|
||||
|
||||
# Recent tick window and minimum margin summary
|
||||
ticks = recent_ticks("EURUSD", seconds=300)
|
||||
margins = minimum_margins("EURUSD")
|
||||
|
||||
# Reuse one MT5 connection for multiple calls
|
||||
with Mt5CliClient(login=12345, password="secret", server="Broker-Demo") as client:
|
||||
account = client.account_info()
|
||||
@@ -92,10 +108,11 @@ mt5cli --login 12345 --password mypass --server MyBroker-Demo \
|
||||
|
||||
### Ticks
|
||||
|
||||
| Command | Description |
|
||||
| ------------- | ------------------------------ |
|
||||
| `ticks-from` | Export ticks from a start date |
|
||||
| `ticks-range` | Export ticks for a date range |
|
||||
| Command | Description |
|
||||
| -------------- | ----------------------------------- |
|
||||
| `ticks-from` | Export ticks from a start date |
|
||||
| `ticks-range` | Export ticks for a date range |
|
||||
| `ticks-recent` | Export ticks from a trailing window |
|
||||
|
||||
### Information
|
||||
|
||||
@@ -108,6 +125,7 @@ mt5cli --login 12345 --password mypass --server MyBroker-Demo \
|
||||
| `symbols` | Export symbol list |
|
||||
| `symbol-info` | Export symbol details |
|
||||
| `symbol-info-tick` | Export the last tick for a symbol |
|
||||
| `minimum-margins` | Export minimum-volume margin summary |
|
||||
| `market-book` | Export market depth (order book) |
|
||||
|
||||
### Trading
|
||||
@@ -152,6 +170,8 @@ mt5cli -o history.db collect-history \
|
||||
|
||||
History orders and deals are fetched per symbol and concatenated, so the symbol filter is applied consistently across all datasets. The `cash_events` view is derived from symbol-filtered `history_deals`, so account-level cash events with empty or non-matching symbols may be excluded. The `positions_reconstructed` view excludes positions with no closing deal, uses volume-weighted open/close prices, and reports reversal deals (`DEAL_ENTRY_INOUT`) via `volume_reversal` / `reversal_count`.
|
||||
|
||||
See the [History schema diagram](api/history.md#entity-relationship-diagram) for a sample ER layout of the resulting database.
|
||||
|
||||
## Global Options
|
||||
|
||||
| Option | Description |
|
||||
|
||||
@@ -24,6 +24,7 @@ theme:
|
||||
features:
|
||||
- content.code.annotate
|
||||
- content.code.copy
|
||||
- content.code.mermaid
|
||||
- navigation.indexes
|
||||
- navigation.sections
|
||||
- navigation.tabs
|
||||
@@ -57,6 +58,7 @@ nav:
|
||||
- Overview: api/index.md
|
||||
- CLI: api/cli.md
|
||||
- SDK: api/sdk.md
|
||||
- History Collection (SQLite): api/history.md
|
||||
- Utils: api/utils.md
|
||||
|
||||
markdown_extensions:
|
||||
|
||||
+18
-1
@@ -16,21 +16,33 @@ from .sdk import (
|
||||
history_orders,
|
||||
last_error,
|
||||
market_book,
|
||||
minimum_margins,
|
||||
orders,
|
||||
positions,
|
||||
recent_ticks,
|
||||
symbol_info,
|
||||
symbol_info_tick,
|
||||
symbols,
|
||||
terminal_info,
|
||||
update_history,
|
||||
update_history_with_config,
|
||||
)
|
||||
from .sdk import (
|
||||
version as mt5_version,
|
||||
)
|
||||
from .utils import detect_format, export_dataframe
|
||||
from .utils import (
|
||||
Dataset,
|
||||
IfExists,
|
||||
detect_format,
|
||||
export_dataframe,
|
||||
export_dataframe_to_sqlite,
|
||||
)
|
||||
|
||||
__version__ = version(__package__) if __package__ else None
|
||||
|
||||
__all__ = [
|
||||
"Dataset",
|
||||
"IfExists",
|
||||
"Mt5CliClient",
|
||||
"account_info",
|
||||
"build_config",
|
||||
@@ -42,15 +54,20 @@ __all__ = [
|
||||
"copy_ticks_range",
|
||||
"detect_format",
|
||||
"export_dataframe",
|
||||
"export_dataframe_to_sqlite",
|
||||
"history_deals",
|
||||
"history_orders",
|
||||
"last_error",
|
||||
"market_book",
|
||||
"minimum_margins",
|
||||
"mt5_version",
|
||||
"orders",
|
||||
"positions",
|
||||
"recent_ticks",
|
||||
"symbol_info",
|
||||
"symbol_info_tick",
|
||||
"symbols",
|
||||
"terminal_info",
|
||||
"update_history",
|
||||
"update_history_with_config",
|
||||
]
|
||||
|
||||
@@ -300,6 +300,44 @@ def ticks_range(
|
||||
)
|
||||
|
||||
|
||||
@app.command()
|
||||
def ticks_recent(
|
||||
ctx: typer.Context,
|
||||
symbol: Annotated[str, typer.Option(help="Symbol name.")],
|
||||
seconds: Annotated[
|
||||
float,
|
||||
typer.Option(help="Lookback window in seconds."),
|
||||
],
|
||||
date_to: Annotated[
|
||||
datetime | None,
|
||||
typer.Option(click_type=DATETIME_TYPE, help="Window end date."),
|
||||
] = None,
|
||||
count: Annotated[
|
||||
int,
|
||||
typer.Option(help="Maximum number of ticks to return."),
|
||||
] = 10000,
|
||||
flags: Annotated[
|
||||
int,
|
||||
typer.Option(
|
||||
click_type=TICK_FLAGS_TYPE,
|
||||
help="Tick flags (ALL, INFO, TRADE, or integer).",
|
||||
),
|
||||
] = 1,
|
||||
) -> None:
|
||||
"""Export ticks from a recent time window."""
|
||||
client = _sdk_client(ctx)
|
||||
_execute_export(
|
||||
ctx,
|
||||
lambda: client.recent_ticks(
|
||||
symbol,
|
||||
seconds,
|
||||
date_to=date_to,
|
||||
count=count,
|
||||
flags=flags,
|
||||
),
|
||||
)
|
||||
|
||||
|
||||
@app.command()
|
||||
def account_info(ctx: typer.Context) -> None:
|
||||
"""Export account information."""
|
||||
@@ -335,6 +373,16 @@ def symbol_info(
|
||||
_execute_export(ctx, lambda: client.symbol_info(symbol))
|
||||
|
||||
|
||||
@app.command()
|
||||
def minimum_margins(
|
||||
ctx: typer.Context,
|
||||
symbol: Annotated[str, typer.Option(help="Symbol name.")],
|
||||
) -> None:
|
||||
"""Export minimum-volume buy and sell margin requirements."""
|
||||
client = _sdk_client(ctx)
|
||||
_execute_export(ctx, lambda: client.minimum_margins(symbol))
|
||||
|
||||
|
||||
@app.command()
|
||||
def orders(
|
||||
ctx: typer.Context,
|
||||
|
||||
+1401
File diff suppressed because it is too large
Load Diff
+393
-330
@@ -5,12 +5,24 @@ from __future__ import annotations
|
||||
import logging
|
||||
import sqlite3
|
||||
from contextlib import contextmanager
|
||||
from datetime import datetime
|
||||
from pathlib import Path # noqa: TC003
|
||||
from dataclasses import dataclass
|
||||
from datetime import UTC, datetime, timedelta
|
||||
from pathlib import Path
|
||||
from typing import TYPE_CHECKING, Self, TypeVar
|
||||
|
||||
import pandas as pd
|
||||
from pdmt5 import Mt5Config, Mt5DataClient
|
||||
|
||||
from .history import (
|
||||
create_cash_events_view,
|
||||
create_history_indexes,
|
||||
create_positions_reconstructed_view,
|
||||
resolve_history_datasets,
|
||||
resolve_history_tick_flags,
|
||||
resolve_history_timeframes,
|
||||
write_collected_datasets,
|
||||
write_incremental_datasets,
|
||||
)
|
||||
from .utils import (
|
||||
Dataset,
|
||||
IfExists,
|
||||
@@ -20,9 +32,7 @@ from .utils import (
|
||||
)
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from collections.abc import Callable, Iterator
|
||||
|
||||
import pandas as pd
|
||||
from collections.abc import Callable, Iterator, Sequence
|
||||
|
||||
T = TypeVar("T")
|
||||
|
||||
@@ -42,28 +52,19 @@ __all__ = [
|
||||
"history_orders",
|
||||
"last_error",
|
||||
"market_book",
|
||||
"minimum_margins",
|
||||
"orders",
|
||||
"positions",
|
||||
"recent_ticks",
|
||||
"symbol_info",
|
||||
"symbol_info_tick",
|
||||
"symbols",
|
||||
"terminal_info",
|
||||
"update_history",
|
||||
"update_history_with_config",
|
||||
"version",
|
||||
]
|
||||
|
||||
_TRADE_DEAL_TYPES: tuple[int, int] = (0, 1)
|
||||
_TRADE_DEAL_TYPES_SQL = f"({', '.join(str(value) for value in _TRADE_DEAL_TYPES)})"
|
||||
_POSITIONS_VIEW_REQUIRED_COLUMNS: frozenset[str] = frozenset({
|
||||
"position_id",
|
||||
"symbol",
|
||||
"time",
|
||||
"type",
|
||||
"entry",
|
||||
"volume",
|
||||
"price",
|
||||
"profit",
|
||||
})
|
||||
|
||||
|
||||
def _coerce_timeframe(timeframe: int | str) -> int:
|
||||
if isinstance(timeframe, int):
|
||||
@@ -89,6 +90,89 @@ def _coerce_datetime(value: datetime | str | None) -> datetime | None:
|
||||
return parse_datetime(value)
|
||||
|
||||
|
||||
def _coerce_tick_time(value: object) -> datetime:
|
||||
if isinstance(value, datetime):
|
||||
return value
|
||||
if isinstance(value, str):
|
||||
return parse_datetime(value)
|
||||
if isinstance(value, (int, float)):
|
||||
return datetime.fromtimestamp(value, tz=UTC)
|
||||
msg = f"Unsupported tick time value: {value!r}"
|
||||
raise TypeError(msg)
|
||||
|
||||
|
||||
def _filter_ticks_to_end(frame: pd.DataFrame, end: datetime) -> pd.DataFrame:
|
||||
if frame.empty or "time" not in frame.columns:
|
||||
return frame
|
||||
times = pd.to_datetime(frame["time"], utc=True)
|
||||
return frame.loc[times <= end].reset_index(drop=True)
|
||||
|
||||
|
||||
def _fetch_recent_ticks(
|
||||
client: Mt5DataClient,
|
||||
symbol: str,
|
||||
seconds: float,
|
||||
date_to: datetime | None,
|
||||
count: int,
|
||||
flags: int,
|
||||
) -> pd.DataFrame:
|
||||
if date_to is not None:
|
||||
end = date_to
|
||||
else:
|
||||
tick = client.symbol_info_tick(symbol)
|
||||
end = _coerce_tick_time(tick.time)
|
||||
start = end - timedelta(seconds=seconds)
|
||||
if count > 0:
|
||||
from_frame = _filter_ticks_to_end(
|
||||
client.copy_ticks_from_as_df(
|
||||
symbol=symbol,
|
||||
date_from=start,
|
||||
count=count,
|
||||
flags=flags,
|
||||
),
|
||||
end,
|
||||
)
|
||||
if len(from_frame) < count:
|
||||
return from_frame
|
||||
frame = client.copy_ticks_range_as_df(
|
||||
symbol=symbol,
|
||||
date_from=start,
|
||||
date_to=end,
|
||||
flags=flags,
|
||||
)
|
||||
if count > 0 and len(frame) > count:
|
||||
return frame.tail(count).reset_index(drop=True)
|
||||
return frame
|
||||
|
||||
|
||||
def _fetch_minimum_margins(client: Mt5DataClient, symbol: str) -> pd.DataFrame:
|
||||
sym = client.symbol_info(symbol)
|
||||
account = client.account_info()
|
||||
tick = client.symbol_info_tick(symbol)
|
||||
volume_min = sym.volume_min
|
||||
buy_margin = client.order_calc_margin(
|
||||
client.mt5.ORDER_TYPE_BUY,
|
||||
symbol,
|
||||
volume_min,
|
||||
tick.ask,
|
||||
)
|
||||
sell_margin = client.order_calc_margin(
|
||||
client.mt5.ORDER_TYPE_SELL,
|
||||
symbol,
|
||||
volume_min,
|
||||
tick.bid,
|
||||
)
|
||||
return pd.DataFrame([
|
||||
{
|
||||
"symbol": symbol,
|
||||
"account_currency": account.currency,
|
||||
"volume_min": volume_min,
|
||||
"buy_margin": buy_margin,
|
||||
"sell_margin": sell_margin,
|
||||
}
|
||||
])
|
||||
|
||||
|
||||
def build_config(
|
||||
*,
|
||||
path: str | None = None,
|
||||
@@ -418,326 +502,271 @@ class Mt5CliClient:
|
||||
"""Return market depth for a symbol."""
|
||||
return self._fetch(lambda c: c.market_book_get_as_df(symbol=symbol))
|
||||
|
||||
def recent_ticks(
|
||||
self,
|
||||
symbol: str,
|
||||
seconds: float,
|
||||
*,
|
||||
date_to: datetime | str | None = None,
|
||||
count: int = 10000,
|
||||
flags: int | str = "ALL",
|
||||
) -> pd.DataFrame:
|
||||
"""Return ticks from a recent time window.
|
||||
|
||||
def _create_cash_events_view(
|
||||
conn: sqlite3.Connection,
|
||||
deals_columns: set[str],
|
||||
) -> bool:
|
||||
"""Create the cash_events SQLite view derived from history_deals.
|
||||
Args:
|
||||
symbol: Symbol name.
|
||||
seconds: Lookback window in seconds ending at ``date_to``.
|
||||
date_to: Window end time. When ``None``, uses the latest
|
||||
``symbol_info_tick().time`` rather than wall-clock now.
|
||||
count: Maximum ticks to return. Values ``<= 0`` return the full
|
||||
window without trimming. Positive values keep the most recent
|
||||
ticks; when the window is sparse, ``copy_ticks_from`` avoids
|
||||
fetching the entire range.
|
||||
flags: Tick flags as ``ALL``, ``INFO``, ``TRADE``, or an integer.
|
||||
|
||||
Returns:
|
||||
Tick DataFrame with MT5 tick columns such as ``time``, ``bid``,
|
||||
``ask``, ``last``, and ``volume``.
|
||||
"""
|
||||
tick_flags = _coerce_tick_flags(flags)
|
||||
end = _coerce_datetime(date_to)
|
||||
return self._fetch(
|
||||
lambda c: _fetch_recent_ticks(
|
||||
c,
|
||||
symbol,
|
||||
seconds,
|
||||
end,
|
||||
count,
|
||||
tick_flags,
|
||||
),
|
||||
)
|
||||
|
||||
def minimum_margins(self, symbol: str) -> pd.DataFrame:
|
||||
"""Return minimum-volume buy and sell margin requirements.
|
||||
|
||||
Args:
|
||||
symbol: Symbol name.
|
||||
|
||||
Returns:
|
||||
One-row DataFrame with columns ``symbol``, ``account_currency``,
|
||||
``volume_min``, ``buy_margin``, and ``sell_margin``.
|
||||
"""
|
||||
return self._fetch(lambda c: _fetch_minimum_margins(c, symbol))
|
||||
|
||||
|
||||
def _resolve_incremental_settings(
|
||||
selected_datasets: set[Dataset],
|
||||
timeframes: Sequence[int | str] | None,
|
||||
flags: int | str,
|
||||
) -> tuple[list[int], int]:
|
||||
"""Resolve dataset-specific incremental update settings.
|
||||
|
||||
Returns:
|
||||
True if the view was created, False if required columns are missing.
|
||||
Tuple of resolved rate timeframes and tick copy flags.
|
||||
|
||||
Raises:
|
||||
ValueError: If timeframe or tick flag values are invalid.
|
||||
"""
|
||||
if "type" not in deals_columns:
|
||||
logger.warning("Skipping cash_events view: history_deals.type is missing")
|
||||
return False
|
||||
conn.execute("DROP VIEW IF EXISTS cash_events")
|
||||
conn.execute(
|
||||
"CREATE VIEW cash_events AS" # noqa: S608
|
||||
f" SELECT * FROM history_deals WHERE type NOT IN {_TRADE_DEAL_TYPES_SQL}",
|
||||
)
|
||||
return True
|
||||
resolved_timeframes: list[int] = []
|
||||
if Dataset.rates in selected_datasets:
|
||||
try:
|
||||
resolved_timeframes = resolve_history_timeframes(timeframes)
|
||||
except ValueError as exc:
|
||||
msg = str(exc)
|
||||
raise ValueError(msg) from exc
|
||||
resolved_tick_flags = 0
|
||||
if Dataset.ticks in selected_datasets:
|
||||
try:
|
||||
resolved_tick_flags = resolve_history_tick_flags(flags)
|
||||
except ValueError as exc:
|
||||
msg = str(exc)
|
||||
raise ValueError(msg) from exc
|
||||
return resolved_timeframes, resolved_tick_flags
|
||||
|
||||
|
||||
def _create_positions_reconstructed_view(
|
||||
conn: sqlite3.Connection,
|
||||
deals_columns: set[str],
|
||||
) -> bool:
|
||||
"""Create the positions_reconstructed SQLite view derived from history_deals.
|
||||
@dataclass(frozen=True)
|
||||
class _UpdateHistoryRequest:
|
||||
selected: set[Dataset]
|
||||
end: datetime
|
||||
fallback_start: datetime
|
||||
resolved_timeframes: list[int]
|
||||
resolved_tick_flags: int
|
||||
output_path: Path
|
||||
|
||||
|
||||
def _resolve_update_history_request(
|
||||
*,
|
||||
output: Path | str,
|
||||
symbols: Sequence[str],
|
||||
datasets: set[Dataset] | None,
|
||||
timeframes: Sequence[int | str] | None,
|
||||
flags: int | str,
|
||||
lookback_hours: float,
|
||||
date_to: datetime | str | None,
|
||||
) -> _UpdateHistoryRequest | None:
|
||||
"""Validate and resolve incremental history update inputs.
|
||||
|
||||
Returns:
|
||||
True if the view was created, False if required columns are missing.
|
||||
Resolved request parameters, or None when no datasets are selected.
|
||||
|
||||
Raises:
|
||||
ValueError: If symbols are empty, lookback_hours is not positive, or
|
||||
timeframe/flag values are invalid.
|
||||
"""
|
||||
if not _POSITIONS_VIEW_REQUIRED_COLUMNS.issubset(deals_columns):
|
||||
missing = ", ".join(sorted(_POSITIONS_VIEW_REQUIRED_COLUMNS - deals_columns))
|
||||
logger.warning(
|
||||
"Skipping positions_reconstructed view: history_deals missing columns: %s",
|
||||
missing,
|
||||
)
|
||||
return False
|
||||
conn.execute("DROP VIEW IF EXISTS positions_reconstructed")
|
||||
conn.execute(
|
||||
"CREATE VIEW positions_reconstructed AS" # noqa: S608
|
||||
" SELECT"
|
||||
" position_id,"
|
||||
" symbol,"
|
||||
" MIN(CASE WHEN entry = 0 THEN time END) AS open_time,"
|
||||
" MAX(CASE WHEN entry IN (1, 2, 3) THEN time END) AS close_time,"
|
||||
" MIN(CASE WHEN entry = 0 THEN type END) AS direction,"
|
||||
" SUM(CASE WHEN entry = 0 THEN volume ELSE 0 END) AS volume_open,"
|
||||
" SUM(CASE WHEN entry IN (1, 3) THEN volume ELSE 0 END) AS volume_close,"
|
||||
" SUM(CASE WHEN entry = 2 THEN volume ELSE 0 END) AS volume_reversal,"
|
||||
" CASE"
|
||||
" WHEN SUM(CASE WHEN entry = 0 THEN volume ELSE 0 END) > 0"
|
||||
" THEN SUM(CASE WHEN entry = 0 THEN price * volume ELSE 0 END)"
|
||||
" / SUM(CASE WHEN entry = 0 THEN volume ELSE 0 END)"
|
||||
" END AS open_price,"
|
||||
" CASE"
|
||||
" WHEN SUM(CASE WHEN entry IN (1, 3) THEN volume ELSE 0 END) > 0"
|
||||
" THEN SUM(CASE WHEN entry IN (1, 3) THEN price * volume ELSE 0 END)"
|
||||
" / SUM(CASE WHEN entry IN (1, 3) THEN volume ELSE 0 END)"
|
||||
" END AS close_price,"
|
||||
" SUM(profit) AS total_profit,"
|
||||
" SUM(CASE WHEN entry = 2 THEN 1 ELSE 0 END) AS reversal_count,"
|
||||
" COUNT(*) AS deals_count"
|
||||
" FROM history_deals"
|
||||
f" WHERE type IN {_TRADE_DEAL_TYPES_SQL} AND position_id != 0"
|
||||
" GROUP BY position_id, symbol"
|
||||
" HAVING SUM(CASE WHEN entry IN (1, 3) THEN 1 ELSE 0 END) > 0",
|
||||
)
|
||||
return True
|
||||
if lookback_hours <= 0:
|
||||
msg = "lookback_hours must be positive."
|
||||
raise ValueError(msg)
|
||||
selected = resolve_history_datasets(datasets)
|
||||
if not selected:
|
||||
logger.info("Skipping SQLite history update: no datasets selected.")
|
||||
return None
|
||||
if not symbols:
|
||||
msg = "At least one symbol is required."
|
||||
raise ValueError(msg)
|
||||
|
||||
|
||||
def _write_frame_to_sqlite(
|
||||
conn: sqlite3.Connection,
|
||||
frame: pd.DataFrame,
|
||||
table_name: str,
|
||||
if_exists: IfExists,
|
||||
) -> bool:
|
||||
"""Write a non-empty-schema frame to SQLite.
|
||||
|
||||
Returns:
|
||||
True if a table was written, False if the frame had no columns.
|
||||
"""
|
||||
if len(frame.columns) == 0:
|
||||
logger.warning("Skipping %s: dataset returned no columns", table_name)
|
||||
return False
|
||||
frame.to_sql( # type: ignore[reportUnknownMemberType]
|
||||
table_name,
|
||||
conn,
|
||||
if_exists=if_exists.value,
|
||||
index=False,
|
||||
chunksize=50_000,
|
||||
method="multi",
|
||||
)
|
||||
return True
|
||||
|
||||
|
||||
def _create_collect_history_indexes(
|
||||
conn: sqlite3.Connection,
|
||||
written_columns: dict[Dataset, set[str]],
|
||||
) -> None:
|
||||
"""Create useful indexes for collected history tables when present."""
|
||||
if {"symbol", "time"}.issubset(written_columns.get(Dataset.rates, set())):
|
||||
conn.execute(
|
||||
"CREATE INDEX IF NOT EXISTS idx_rates_symbol_time ON rates(symbol, time)",
|
||||
)
|
||||
if {"symbol", "time"}.issubset(written_columns.get(Dataset.ticks, set())):
|
||||
conn.execute(
|
||||
"CREATE INDEX IF NOT EXISTS idx_ticks_symbol_time ON ticks(symbol, time)",
|
||||
)
|
||||
if {"position_id", "symbol"}.issubset(
|
||||
written_columns.get(Dataset.history_deals, set())
|
||||
):
|
||||
conn.execute(
|
||||
"CREATE INDEX IF NOT EXISTS idx_history_deals_position_symbol"
|
||||
" ON history_deals(position_id, symbol)",
|
||||
)
|
||||
|
||||
|
||||
def _record_written_columns(
|
||||
written_columns: dict[Dataset, set[str]],
|
||||
dataset: Dataset,
|
||||
frame: pd.DataFrame,
|
||||
) -> None:
|
||||
"""Remember columns for datasets written during streaming collection."""
|
||||
columns = set(frame.columns)
|
||||
if dataset in written_columns:
|
||||
written_columns[dataset].update(columns)
|
||||
if date_to is not None:
|
||||
resolved_end = _coerce_datetime(date_to)
|
||||
else:
|
||||
written_columns[dataset] = columns
|
||||
|
||||
|
||||
def _write_streamed_frame(
|
||||
conn: sqlite3.Connection,
|
||||
frame: pd.DataFrame,
|
||||
dataset: Dataset,
|
||||
table_exists: bool,
|
||||
if_exists: IfExists,
|
||||
written_columns: dict[Dataset, set[str]],
|
||||
) -> bool:
|
||||
"""Write one streamed dataset frame and track table state.
|
||||
|
||||
Returns:
|
||||
True if the dataset table exists after this write attempt.
|
||||
"""
|
||||
write_mode = IfExists.APPEND if table_exists else if_exists
|
||||
if _write_frame_to_sqlite(
|
||||
conn,
|
||||
frame,
|
||||
dataset.table_name,
|
||||
write_mode,
|
||||
):
|
||||
_record_written_columns(written_columns, dataset, frame)
|
||||
return True
|
||||
return table_exists
|
||||
|
||||
|
||||
def _write_rates_dataset(
|
||||
conn: sqlite3.Connection,
|
||||
client: Mt5DataClient,
|
||||
symbols: list[str],
|
||||
timeframe: int,
|
||||
date_from: datetime,
|
||||
date_to: datetime,
|
||||
if_exists: IfExists,
|
||||
written_columns: dict[Dataset, set[str]],
|
||||
) -> bool:
|
||||
"""Stream rates frames into SQLite.
|
||||
|
||||
Returns:
|
||||
True if the rates table was written.
|
||||
"""
|
||||
table_exists = False
|
||||
for sym in symbols:
|
||||
frame = client.copy_rates_range_as_df(
|
||||
symbol=sym,
|
||||
timeframe=timeframe,
|
||||
date_from=date_from,
|
||||
date_to=date_to,
|
||||
)
|
||||
frame.insert(0, "symbol", sym)
|
||||
frame.insert(1, "timeframe", timeframe)
|
||||
table_exists = _write_streamed_frame(
|
||||
conn,
|
||||
frame,
|
||||
Dataset.rates,
|
||||
table_exists,
|
||||
if_exists,
|
||||
written_columns,
|
||||
)
|
||||
return table_exists
|
||||
|
||||
|
||||
def _write_ticks_dataset(
|
||||
conn: sqlite3.Connection,
|
||||
client: Mt5DataClient,
|
||||
symbols: list[str],
|
||||
flags: int,
|
||||
date_from: datetime,
|
||||
date_to: datetime,
|
||||
if_exists: IfExists,
|
||||
written_columns: dict[Dataset, set[str]],
|
||||
) -> bool:
|
||||
"""Stream ticks frames into SQLite.
|
||||
|
||||
Returns:
|
||||
True if the ticks table was written.
|
||||
"""
|
||||
table_exists = False
|
||||
for sym in symbols:
|
||||
frame = client.copy_ticks_range_as_df(
|
||||
symbol=sym,
|
||||
date_from=date_from,
|
||||
date_to=date_to,
|
||||
flags=flags,
|
||||
)
|
||||
frame.insert(0, "symbol", sym)
|
||||
table_exists = _write_streamed_frame(
|
||||
conn,
|
||||
frame,
|
||||
Dataset.ticks,
|
||||
table_exists,
|
||||
if_exists,
|
||||
written_columns,
|
||||
)
|
||||
return table_exists
|
||||
|
||||
|
||||
def _write_history_dataset(
|
||||
conn: sqlite3.Connection,
|
||||
fetch: Callable[..., pd.DataFrame],
|
||||
dataset: Dataset,
|
||||
symbols: list[str],
|
||||
date_from: datetime,
|
||||
date_to: datetime,
|
||||
if_exists: IfExists,
|
||||
written_columns: dict[Dataset, set[str]],
|
||||
) -> bool:
|
||||
"""Stream a history dataset into SQLite with exact symbol filtering.
|
||||
|
||||
Returns:
|
||||
True if the history table was written.
|
||||
"""
|
||||
table_exists = False
|
||||
for sym in symbols:
|
||||
frame = fetch(date_from=date_from, date_to=date_to, symbol=sym)
|
||||
if "symbol" in frame.columns:
|
||||
frame = frame[frame["symbol"] == sym]
|
||||
table_exists = _write_streamed_frame(
|
||||
conn,
|
||||
frame,
|
||||
dataset,
|
||||
table_exists,
|
||||
if_exists,
|
||||
written_columns,
|
||||
)
|
||||
return table_exists
|
||||
|
||||
|
||||
def _write_collected_datasets(
|
||||
conn: sqlite3.Connection,
|
||||
client: Mt5DataClient,
|
||||
symbols: list[str],
|
||||
datasets: set[Dataset],
|
||||
timeframe: int,
|
||||
flags: int,
|
||||
date_from: datetime,
|
||||
date_to: datetime,
|
||||
if_exists: IfExists,
|
||||
) -> tuple[set[Dataset], dict[Dataset, set[str]]]:
|
||||
"""Collect selected datasets and stream each symbol frame into SQLite.
|
||||
|
||||
Returns:
|
||||
Written datasets and their columns.
|
||||
"""
|
||||
written_columns: dict[Dataset, set[str]] = {}
|
||||
written_tables: set[Dataset] = set()
|
||||
if Dataset.rates in datasets and _write_rates_dataset(
|
||||
conn,
|
||||
client,
|
||||
symbols,
|
||||
timeframe,
|
||||
date_from,
|
||||
date_to,
|
||||
if_exists,
|
||||
written_columns,
|
||||
):
|
||||
written_tables.add(Dataset.rates)
|
||||
if Dataset.ticks in datasets and _write_ticks_dataset(
|
||||
conn,
|
||||
client,
|
||||
symbols,
|
||||
resolved_end = datetime.now(UTC)
|
||||
end = resolved_end if resolved_end is not None else datetime.now(UTC)
|
||||
fallback_start = end - timedelta(hours=lookback_hours)
|
||||
resolved_timeframes, resolved_tick_flags = _resolve_incremental_settings(
|
||||
selected,
|
||||
timeframes,
|
||||
flags,
|
||||
date_from,
|
||||
date_to,
|
||||
if_exists,
|
||||
written_columns,
|
||||
):
|
||||
written_tables.add(Dataset.ticks)
|
||||
if Dataset.history_orders in datasets and _write_history_dataset(
|
||||
conn,
|
||||
client.history_orders_get_as_df,
|
||||
Dataset.history_orders,
|
||||
symbols,
|
||||
date_from,
|
||||
date_to,
|
||||
if_exists,
|
||||
written_columns,
|
||||
):
|
||||
written_tables.add(Dataset.history_orders)
|
||||
if Dataset.history_deals in datasets and _write_history_dataset(
|
||||
conn,
|
||||
client.history_deals_get_as_df,
|
||||
Dataset.history_deals,
|
||||
symbols,
|
||||
date_from,
|
||||
date_to,
|
||||
if_exists,
|
||||
written_columns,
|
||||
):
|
||||
written_tables.add(Dataset.history_deals)
|
||||
return written_tables, written_columns
|
||||
)
|
||||
return _UpdateHistoryRequest(
|
||||
selected=selected,
|
||||
end=end,
|
||||
fallback_start=fallback_start,
|
||||
resolved_timeframes=resolved_timeframes,
|
||||
resolved_tick_flags=resolved_tick_flags,
|
||||
output_path=Path(output),
|
||||
)
|
||||
|
||||
|
||||
def update_history( # noqa: PLR0913
|
||||
*,
|
||||
client: Mt5DataClient,
|
||||
output: Path | str,
|
||||
symbols: Sequence[str],
|
||||
datasets: set[Dataset] | None = None,
|
||||
timeframes: Sequence[int | str] | None = None,
|
||||
flags: int | str = "ALL",
|
||||
lookback_hours: float = 24.0,
|
||||
date_to: datetime | str | None = None,
|
||||
deduplicate: bool = True,
|
||||
create_rate_views: bool = True,
|
||||
with_views: bool = False,
|
||||
include_account_events: bool = True,
|
||||
) -> None:
|
||||
"""Incrementally append MT5 history into a SQLite database.
|
||||
|
||||
Uses an already-connected ``Mt5DataClient`` and does not create or close
|
||||
the MT5 connection. For first-time tables, data is fetched from
|
||||
``date_to - lookback_hours``. Subsequent runs resume from existing
|
||||
``MAX(time)`` per symbol (and timeframe for rates); when
|
||||
``include_account_events=True``, account-level deals use a separate cursor
|
||||
over ``type NOT IN (0, 1)`` / empty-symbol rows.
|
||||
|
||||
Args:
|
||||
client: Connected MT5 data client.
|
||||
output: SQLite database path.
|
||||
symbols: Symbols to update.
|
||||
datasets: Datasets to include (defaults to all).
|
||||
timeframes: Rate timeframes to update (defaults to all fixed MT5
|
||||
timeframes when None).
|
||||
flags: Tick copy flags as integer or name (e.g. ``ALL``).
|
||||
lookback_hours: First-run lookback when a table has no prior rows.
|
||||
date_to: Optional update end datetime. Defaults to now (UTC).
|
||||
deduplicate: Remove duplicate rows after append, keeping latest ROWID.
|
||||
create_rate_views: Create ``rate_<symbol>__<timeframe>`` views.
|
||||
with_views: Create ``cash_events`` and ``positions_reconstructed`` views.
|
||||
include_account_events: Include account-level cash events in
|
||||
``history_deals`` when True.
|
||||
"""
|
||||
request = _resolve_update_history_request(
|
||||
output=output,
|
||||
symbols=symbols,
|
||||
datasets=datasets,
|
||||
timeframes=timeframes,
|
||||
flags=flags,
|
||||
lookback_hours=lookback_hours,
|
||||
date_to=date_to,
|
||||
)
|
||||
if request is None:
|
||||
return
|
||||
logger.info(
|
||||
"Updating history in SQLite: symbols=%s, datasets=%s, path=%s",
|
||||
list(symbols),
|
||||
sorted(dataset.value for dataset in request.selected),
|
||||
request.output_path,
|
||||
)
|
||||
with sqlite3.connect(request.output_path) as conn:
|
||||
conn.execute("PRAGMA journal_mode=WAL")
|
||||
conn.execute("PRAGMA synchronous=NORMAL")
|
||||
write_incremental_datasets(
|
||||
conn,
|
||||
client,
|
||||
symbols,
|
||||
request.selected,
|
||||
request.resolved_timeframes,
|
||||
request.resolved_tick_flags,
|
||||
request.fallback_start,
|
||||
request.end,
|
||||
deduplicate=deduplicate,
|
||||
create_rate_views=create_rate_views,
|
||||
with_views=with_views,
|
||||
include_account_events=include_account_events,
|
||||
)
|
||||
|
||||
|
||||
def update_history_with_config( # noqa: PLR0913
|
||||
*,
|
||||
output: Path | str,
|
||||
symbols: Sequence[str],
|
||||
config: Mt5Config | None = None,
|
||||
datasets: set[Dataset] | None = None,
|
||||
timeframes: Sequence[int | str] | None = None,
|
||||
flags: int | str = "ALL",
|
||||
lookback_hours: float = 24.0,
|
||||
date_to: datetime | str | None = None,
|
||||
deduplicate: bool = True,
|
||||
create_rate_views: bool = True,
|
||||
with_views: bool = False,
|
||||
include_account_events: bool = True,
|
||||
) -> None:
|
||||
"""Incrementally append MT5 history, opening and closing the MT5 connection.
|
||||
|
||||
Convenience wrapper around :func:`update_history` for standalone use.
|
||||
"""
|
||||
request = _resolve_update_history_request(
|
||||
output=output,
|
||||
symbols=symbols,
|
||||
datasets=datasets,
|
||||
timeframes=timeframes,
|
||||
flags=flags,
|
||||
lookback_hours=lookback_hours,
|
||||
date_to=date_to,
|
||||
)
|
||||
if request is None:
|
||||
return
|
||||
mt5_config = config or build_config()
|
||||
with _connected_client(mt5_config) as client:
|
||||
update_history(
|
||||
client=client,
|
||||
output=output,
|
||||
symbols=symbols,
|
||||
datasets=datasets,
|
||||
timeframes=timeframes,
|
||||
flags=flags,
|
||||
lookback_hours=lookback_hours,
|
||||
date_to=date_to,
|
||||
deduplicate=deduplicate,
|
||||
create_rate_views=create_rate_views,
|
||||
with_views=with_views,
|
||||
include_account_events=include_account_events,
|
||||
)
|
||||
|
||||
|
||||
def collect_history(
|
||||
@@ -776,7 +805,7 @@ def collect_history(
|
||||
with _connected_client(mt5_config) as client, sqlite3.connect(output) as conn:
|
||||
conn.execute("PRAGMA journal_mode=WAL")
|
||||
conn.execute("PRAGMA synchronous=NORMAL")
|
||||
written_tables, written_columns = _write_collected_datasets(
|
||||
written_tables, written_columns = write_collected_datasets(
|
||||
conn,
|
||||
client,
|
||||
symbols,
|
||||
@@ -787,10 +816,10 @@ def collect_history(
|
||||
end,
|
||||
if_exists,
|
||||
)
|
||||
_create_collect_history_indexes(conn, written_columns)
|
||||
create_history_indexes(conn, written_columns)
|
||||
if with_views and Dataset.history_deals in written_tables:
|
||||
_create_cash_events_view(conn, written_columns[Dataset.history_deals])
|
||||
_create_positions_reconstructed_view(
|
||||
create_cash_events_view(conn, written_columns[Dataset.history_deals])
|
||||
create_positions_reconstructed_view(
|
||||
conn,
|
||||
written_columns[Dataset.history_deals],
|
||||
)
|
||||
@@ -1021,3 +1050,37 @@ def market_book(
|
||||
) -> pd.DataFrame:
|
||||
"""Return market depth for a symbol."""
|
||||
return _make_client(config=config).market_book(symbol)
|
||||
|
||||
|
||||
def recent_ticks(
|
||||
symbol: str,
|
||||
seconds: float,
|
||||
*,
|
||||
date_to: datetime | str | None = None,
|
||||
count: int = 10000,
|
||||
flags: int | str = "ALL",
|
||||
config: Mt5Config | None = None,
|
||||
) -> pd.DataFrame:
|
||||
"""Return ticks from a recent time window ending at ``date_to`` or now.
|
||||
|
||||
See ``Mt5CliClient.recent_ticks`` for parameter and return details.
|
||||
"""
|
||||
return _make_client(config=config).recent_ticks(
|
||||
symbol,
|
||||
seconds,
|
||||
date_to=date_to,
|
||||
count=count,
|
||||
flags=flags,
|
||||
)
|
||||
|
||||
|
||||
def minimum_margins(
|
||||
symbol: str,
|
||||
*,
|
||||
config: Mt5Config | None = None,
|
||||
) -> pd.DataFrame:
|
||||
"""Return minimum-volume buy and sell margin requirements.
|
||||
|
||||
See ``Mt5CliClient.minimum_margins`` for return details.
|
||||
"""
|
||||
return _make_client(config=config).minimum_margins(symbol)
|
||||
|
||||
+55
-10
@@ -2,16 +2,18 @@
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import importlib
|
||||
import json
|
||||
import sqlite3
|
||||
from datetime import UTC, datetime
|
||||
from enum import StrEnum
|
||||
from pathlib import Path
|
||||
from typing import TYPE_CHECKING, Any, TypeGuard, cast
|
||||
from typing import TYPE_CHECKING, Any, TypeGuard
|
||||
|
||||
import click
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from collections.abc import Sequence
|
||||
|
||||
import pandas as pd
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
@@ -260,6 +262,50 @@ def detect_format(
|
||||
raise ValueError(msg)
|
||||
|
||||
|
||||
def export_dataframe_to_sqlite(
|
||||
df: pd.DataFrame,
|
||||
output_path: Path,
|
||||
table_name: str = "data",
|
||||
*,
|
||||
if_exists: IfExists = IfExists.APPEND,
|
||||
index: bool = False,
|
||||
index_label: str | None = None,
|
||||
deduplicate_on: Sequence[str] | None = None,
|
||||
) -> None:
|
||||
"""Write a DataFrame to SQLite with configurable append and deduplication.
|
||||
|
||||
Args:
|
||||
df: DataFrame to export.
|
||||
output_path: SQLite database path.
|
||||
table_name: Target table name.
|
||||
if_exists: Conflict behavior when the table already exists.
|
||||
index: Whether to write the DataFrame index as a column.
|
||||
index_label: Column name for the index when ``index=True``.
|
||||
deduplicate_on: Optional key columns to deduplicate after writing,
|
||||
keeping the latest ``ROWID`` per key group. Deduplication scans the
|
||||
full table, so repeated appends cost O(table size); index the key
|
||||
columns when appending frequently.
|
||||
"""
|
||||
with sqlite3.connect(output_path) as conn:
|
||||
df.to_sql( # type: ignore[reportUnknownMemberType]
|
||||
table_name,
|
||||
conn,
|
||||
if_exists=if_exists.value,
|
||||
index=index,
|
||||
index_label=index_label,
|
||||
)
|
||||
if deduplicate_on:
|
||||
from .history import drop_duplicates_in_table # noqa: PLC0415
|
||||
|
||||
drop_duplicates_in_table(
|
||||
conn.cursor(),
|
||||
table_name,
|
||||
list(deduplicate_on),
|
||||
keep="last",
|
||||
)
|
||||
conn.commit()
|
||||
|
||||
|
||||
def export_dataframe(
|
||||
df: pd.DataFrame,
|
||||
output_path: Path,
|
||||
@@ -289,14 +335,13 @@ def export_dataframe(
|
||||
elif output_format == "parquet":
|
||||
df.to_parquet(output_path, index=False)
|
||||
elif output_format == "sqlite3":
|
||||
sqlite3 = cast("Any", importlib.import_module("sqlite3"))
|
||||
with sqlite3.connect(output_path) as conn:
|
||||
df.to_sql( # type: ignore[reportUnknownMemberType]
|
||||
table_name,
|
||||
conn,
|
||||
if_exists="replace",
|
||||
index=False,
|
||||
)
|
||||
export_dataframe_to_sqlite(
|
||||
df,
|
||||
output_path,
|
||||
table_name,
|
||||
if_exists=IfExists.REPLACE,
|
||||
index=False,
|
||||
)
|
||||
else:
|
||||
msg = f"Unsupported output format: {output_format}"
|
||||
raise ValueError(msg)
|
||||
|
||||
+2
-1
@@ -1,6 +1,6 @@
|
||||
[project]
|
||||
name = "mt5cli"
|
||||
version = "0.4.0"
|
||||
version = "0.4.3"
|
||||
description = "Command-line tool for MetaTrader 5"
|
||||
authors = [{name = "dceoy", email = "dceoy@users.noreply.github.com"}]
|
||||
maintainers = [{name = "dceoy", email = "dceoy@users.noreply.github.com"}]
|
||||
@@ -124,6 +124,7 @@ ignore = [
|
||||
]
|
||||
|
||||
[tool.ruff.lint.per-file-ignores]
|
||||
"mt5cli/history.py" = ["TC003"]
|
||||
"tests/**/*.py" = [
|
||||
"DOC201", # Missing return documentation
|
||||
"DOC501", # Raised exception missing from docstring
|
||||
|
||||
+63
-4
@@ -6,7 +6,7 @@ import json
|
||||
import logging
|
||||
import re
|
||||
import sqlite3
|
||||
from datetime import UTC, datetime
|
||||
from datetime import UTC, datetime, timedelta
|
||||
from typing import TYPE_CHECKING
|
||||
from unittest.mock import MagicMock
|
||||
|
||||
@@ -316,6 +316,65 @@ class TestCommands:
|
||||
flags=2,
|
||||
)
|
||||
|
||||
def test_ticks_recent(
|
||||
self,
|
||||
tmp_path: Path,
|
||||
mock_client: MagicMock,
|
||||
) -> None:
|
||||
"""Test ticks-recent command."""
|
||||
output = tmp_path / "out.csv"
|
||||
result = runner.invoke(
|
||||
app,
|
||||
[
|
||||
"-o",
|
||||
str(output),
|
||||
"ticks-recent",
|
||||
"--symbol",
|
||||
"EURUSD",
|
||||
"--seconds",
|
||||
"120",
|
||||
"--date-to",
|
||||
"2024-01-02",
|
||||
"--count",
|
||||
"500",
|
||||
"--flags",
|
||||
"ALL",
|
||||
],
|
||||
)
|
||||
assert result.exit_code == 0, result.output
|
||||
mock_client.copy_ticks_from_as_df.assert_called_once_with(
|
||||
symbol="EURUSD",
|
||||
date_from=datetime(2024, 1, 2, tzinfo=UTC) - timedelta(seconds=120),
|
||||
count=500,
|
||||
flags=1,
|
||||
)
|
||||
mock_client.copy_ticks_range_as_df.assert_not_called()
|
||||
|
||||
def test_minimum_margins(
|
||||
self,
|
||||
tmp_path: Path,
|
||||
mock_client: MagicMock,
|
||||
) -> None:
|
||||
"""Test minimum-margins command."""
|
||||
sym = MagicMock(volume_min=0.01)
|
||||
account = MagicMock(currency="USD")
|
||||
tick = MagicMock(ask=1.1010, bid=1.1000)
|
||||
mock_client.symbol_info.return_value = sym
|
||||
mock_client.account_info.return_value = account
|
||||
mock_client.symbol_info_tick.return_value = tick
|
||||
mock_client.order_calc_margin.side_effect = [12.5, 12.4]
|
||||
mock_client.mt5.ORDER_TYPE_BUY = 0
|
||||
mock_client.mt5.ORDER_TYPE_SELL = 1
|
||||
output = tmp_path / "out.csv"
|
||||
result = runner.invoke(
|
||||
app,
|
||||
["-o", str(output), "minimum-margins", "--symbol", "EURUSD"],
|
||||
)
|
||||
assert result.exit_code == 0, result.output
|
||||
mock_client.symbol_info.assert_called_once_with("EURUSD")
|
||||
mock_client.order_calc_margin.assert_any_call(0, "EURUSD", 0.01, 1.1010)
|
||||
mock_client.order_calc_margin.assert_any_call(1, "EURUSD", 0.01, 1.1000)
|
||||
|
||||
def test_orders(
|
||||
self,
|
||||
tmp_path: Path,
|
||||
@@ -1089,7 +1148,7 @@ class TestCollectHistory:
|
||||
assert all(row[0] not in {0, 1} for row in cash)
|
||||
# Position 100 (BUY 1@1.10 + BUY 3@1.20 then SELL 4@1.50) is closed.
|
||||
# Position 200 (BUY 2@2.00 then SELL 2@2.20) is closed.
|
||||
# Position 300 (open-only) and 400 (reversal-only) are excluded.
|
||||
# Position 400 (reversal-only with non-trade deal type) stays excluded.
|
||||
assert set(positions) == {100, 200, 500, 600}
|
||||
pos_100 = positions[100]
|
||||
tol = 1e-9
|
||||
@@ -1106,10 +1165,10 @@ class TestCollectHistory:
|
||||
assert abs(pos_500[5] - 1.05) < tol
|
||||
pos_600 = positions[600]
|
||||
assert abs(pos_600[1] - 3.0) < tol
|
||||
assert abs(pos_600[2] - 3.0) < tol
|
||||
assert abs(pos_600[2] - 4.0) < tol # reversal + close volumes
|
||||
assert abs(pos_600[3] - 1.0) < tol
|
||||
assert abs(pos_600[4] - 1.10) < tol
|
||||
assert abs(pos_600[5] - 1.40) < tol
|
||||
assert abs(pos_600[5] - 3.5475) < tol
|
||||
assert pos_600[6] == 1
|
||||
|
||||
def test_collect_history_filters_history_symbols_exactly(
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
+517
-1
@@ -4,7 +4,7 @@ from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import sqlite3
|
||||
from datetime import UTC, datetime
|
||||
from datetime import UTC, datetime, timedelta
|
||||
from typing import TYPE_CHECKING
|
||||
from unittest.mock import MagicMock
|
||||
|
||||
@@ -16,6 +16,7 @@ if TYPE_CHECKING:
|
||||
from pathlib import Path
|
||||
|
||||
from mt5cli import sdk
|
||||
from mt5cli.history import DEFAULT_HISTORY_TIMEFRAMES
|
||||
from mt5cli.sdk import (
|
||||
Mt5CliClient,
|
||||
account_info,
|
||||
@@ -30,12 +31,16 @@ from mt5cli.sdk import (
|
||||
history_orders,
|
||||
last_error,
|
||||
market_book,
|
||||
minimum_margins,
|
||||
orders,
|
||||
positions,
|
||||
recent_ticks,
|
||||
symbol_info,
|
||||
symbol_info_tick,
|
||||
symbols,
|
||||
terminal_info,
|
||||
update_history,
|
||||
update_history_with_config,
|
||||
version,
|
||||
)
|
||||
from mt5cli.utils import Dataset
|
||||
@@ -464,3 +469,514 @@ class TestCollectHistory:
|
||||
}
|
||||
assert "cash_events" not in views
|
||||
assert "positions_reconstructed" not in views
|
||||
|
||||
|
||||
class TestUpdateHistory:
|
||||
"""Tests for update_history SDK functions."""
|
||||
|
||||
@pytest.fixture
|
||||
def connected_client(self) -> MagicMock:
|
||||
"""Create a connected mock client without MT5 lifecycle patching."""
|
||||
return MagicMock()
|
||||
|
||||
def test_update_history_appends_incrementally(
|
||||
self,
|
||||
connected_client: MagicMock,
|
||||
mocker: MockerFixture,
|
||||
tmp_path: Path,
|
||||
) -> None:
|
||||
"""Test sequential SQLite history updates use existing max timestamps."""
|
||||
date_to = datetime(2024, 1, 2, tzinfo=UTC)
|
||||
first_expected_start = datetime(2024, 1, 1, tzinfo=UTC)
|
||||
second_expected_start = datetime(2024, 1, 1, 12, tzinfo=UTC)
|
||||
rate_starts: list[datetime] = []
|
||||
deal_starts: list[datetime] = []
|
||||
|
||||
def make_rates(**kwargs: object) -> pd.DataFrame:
|
||||
assert kwargs["symbol"] == "EURUSD"
|
||||
assert kwargs["timeframe"] == 1
|
||||
assert kwargs["date_to"] == date_to
|
||||
rate_starts.append(kwargs["date_from"]) # type: ignore[arg-type]
|
||||
return pd.DataFrame({
|
||||
"time": ["2024-01-01T12:00:00+00:00"],
|
||||
"open": [1.0 + len(rate_starts) / 10],
|
||||
})
|
||||
|
||||
def make_deals(**kwargs: object) -> pd.DataFrame:
|
||||
assert kwargs["date_to"] == date_to
|
||||
deal_starts.append(kwargs["date_from"]) # type: ignore[arg-type]
|
||||
return pd.DataFrame({
|
||||
"ticket": [10],
|
||||
"position_id": [100],
|
||||
"symbol": ["EURUSD"],
|
||||
"time": ["2024-01-01T12:00:00+00:00"],
|
||||
"type": [0],
|
||||
"entry": [0],
|
||||
"volume": [1.0],
|
||||
"price": [1.1],
|
||||
"profit": [0.0],
|
||||
})
|
||||
|
||||
connected_client.copy_rates_range_as_df.side_effect = make_rates
|
||||
connected_client.history_deals_get_as_df.side_effect = make_deals
|
||||
mocker.patch("mt5cli.sdk.Mt5DataClient")
|
||||
output = tmp_path / "incremental-history.db"
|
||||
|
||||
for _ in range(2):
|
||||
update_history(
|
||||
client=connected_client,
|
||||
output=output,
|
||||
symbols=["EURUSD"],
|
||||
datasets={Dataset.rates, Dataset.history_deals},
|
||||
timeframes=["M1"],
|
||||
lookback_hours=24,
|
||||
date_to=date_to,
|
||||
with_views=True,
|
||||
)
|
||||
|
||||
assert rate_starts == [first_expected_start, second_expected_start]
|
||||
assert deal_starts == [first_expected_start, first_expected_start]
|
||||
connected_client.initialize_and_login_mt5.assert_not_called()
|
||||
connected_client.shutdown.assert_not_called()
|
||||
with sqlite3.connect(output) as conn:
|
||||
assert conn.execute("SELECT COUNT(*) FROM rates").fetchone() == (1,)
|
||||
assert conn.execute("SELECT open FROM rates").fetchone() == (1.2,)
|
||||
assert conn.execute(
|
||||
"SELECT COUNT(*) FROM history_deals",
|
||||
).fetchone() == (1,)
|
||||
assert conn.execute(
|
||||
"SELECT name FROM sqlite_master WHERE name = 'cash_events'",
|
||||
).fetchone() == ("cash_events",)
|
||||
|
||||
def test_update_history_rejects_invalid_inputs(
|
||||
self,
|
||||
connected_client: MagicMock,
|
||||
tmp_path: Path,
|
||||
) -> None:
|
||||
"""Test validation errors for incremental history updates."""
|
||||
output = tmp_path / "invalid-update.db"
|
||||
with pytest.raises(ValueError, match="At least one symbol"):
|
||||
update_history(
|
||||
client=connected_client,
|
||||
output=output,
|
||||
symbols=[],
|
||||
)
|
||||
with pytest.raises(ValueError, match="lookback_hours must be positive"):
|
||||
update_history(
|
||||
client=connected_client,
|
||||
output=output,
|
||||
symbols=["EURUSD"],
|
||||
lookback_hours=0,
|
||||
)
|
||||
with pytest.raises(ValueError, match="Invalid timeframe"):
|
||||
update_history(
|
||||
client=connected_client,
|
||||
output=output,
|
||||
symbols=["EURUSD"],
|
||||
datasets={Dataset.rates},
|
||||
timeframes=["BAD"],
|
||||
)
|
||||
with pytest.raises(ValueError, match="Invalid tick flags"):
|
||||
update_history(
|
||||
client=connected_client,
|
||||
output=output,
|
||||
symbols=["EURUSD"],
|
||||
datasets={Dataset.ticks},
|
||||
flags="BAD",
|
||||
)
|
||||
|
||||
def test_update_history_noops_for_empty_datasets(
|
||||
self,
|
||||
connected_client: MagicMock,
|
||||
mocker: MockerFixture,
|
||||
tmp_path: Path,
|
||||
) -> None:
|
||||
"""Test empty dataset selection skips MT5 and SQLite writes."""
|
||||
writer = mocker.patch("mt5cli.sdk.write_incremental_datasets")
|
||||
connect = mocker.patch("mt5cli.sdk.sqlite3.connect")
|
||||
update_history(
|
||||
client=connected_client,
|
||||
output=tmp_path / "empty-datasets.db",
|
||||
symbols=["EURUSD"],
|
||||
datasets=set(),
|
||||
)
|
||||
writer.assert_not_called()
|
||||
connect.assert_not_called()
|
||||
|
||||
def test_update_history_uses_all_default_timeframes(
|
||||
self,
|
||||
connected_client: MagicMock,
|
||||
mocker: MockerFixture,
|
||||
tmp_path: Path,
|
||||
) -> None:
|
||||
"""Test that timeframes=None writes rates for all default MT5 timeframes."""
|
||||
timeframes_written: list[int] = []
|
||||
|
||||
def capture(
|
||||
*args: object,
|
||||
**_kwargs: object,
|
||||
) -> tuple[set[Dataset], dict[Dataset, set[str]]]:
|
||||
timeframes_written.extend(args[4]) # type: ignore[arg-type]
|
||||
return set(), {}
|
||||
|
||||
mocker.patch("mt5cli.sdk.write_incremental_datasets", side_effect=capture)
|
||||
update_history(
|
||||
client=connected_client,
|
||||
output=tmp_path / "default-timeframes.db",
|
||||
symbols=["EURUSD"],
|
||||
datasets={Dataset.rates},
|
||||
timeframes=None,
|
||||
lookback_hours=1,
|
||||
date_to=datetime(2024, 1, 1, tzinfo=UTC),
|
||||
)
|
||||
assert len(timeframes_written) == len(DEFAULT_HISTORY_TIMEFRAMES)
|
||||
|
||||
def test_update_history_uses_specified_timeframes(
|
||||
self,
|
||||
connected_client: MagicMock,
|
||||
mocker: MockerFixture,
|
||||
tmp_path: Path,
|
||||
) -> None:
|
||||
"""Test explicit timeframes limit rate updates."""
|
||||
timeframes_written: list[int] = []
|
||||
|
||||
def capture(
|
||||
*args: object,
|
||||
**_kwargs: object,
|
||||
) -> tuple[set[Dataset], dict[Dataset, set[str]]]:
|
||||
timeframes_written.extend(args[4]) # type: ignore[arg-type]
|
||||
return set(), {}
|
||||
|
||||
mocker.patch("mt5cli.sdk.write_incremental_datasets", side_effect=capture)
|
||||
update_history(
|
||||
client=connected_client,
|
||||
output=tmp_path / "specific-timeframes.db",
|
||||
symbols=["EURUSD"],
|
||||
datasets={Dataset.rates},
|
||||
timeframes=["M1", "H1"],
|
||||
lookback_hours=1,
|
||||
date_to=datetime(2024, 1, 1, tzinfo=UTC),
|
||||
)
|
||||
assert timeframes_written == [1, 16385]
|
||||
|
||||
def test_update_history_updates_ticks_and_orders(
|
||||
self,
|
||||
connected_client: MagicMock,
|
||||
tmp_path: Path,
|
||||
) -> None:
|
||||
"""Test incremental update writes selected ticks and orders datasets."""
|
||||
date_to = datetime(2024, 1, 2, tzinfo=UTC)
|
||||
expected_start = datetime(2024, 1, 1, tzinfo=UTC)
|
||||
|
||||
def make_ticks(**kwargs: object) -> pd.DataFrame:
|
||||
assert kwargs["symbol"] == "EURUSD"
|
||||
assert kwargs["date_from"] == expected_start
|
||||
assert kwargs["date_to"] == date_to
|
||||
assert kwargs["flags"] == 1
|
||||
return pd.DataFrame({
|
||||
"time": ["2024-01-01T12:00:00+00:00"],
|
||||
"time_msc": [1_704_110_400_000],
|
||||
"bid": [1.1],
|
||||
})
|
||||
|
||||
def make_orders(**kwargs: object) -> pd.DataFrame:
|
||||
assert kwargs["symbol"] == "EURUSD"
|
||||
assert kwargs["date_from"] == expected_start
|
||||
assert kwargs["date_to"] == date_to
|
||||
return pd.DataFrame({
|
||||
"ticket": [1],
|
||||
"symbol": ["EURUSD"],
|
||||
"time": ["2024-01-01T12:00:00+00:00"],
|
||||
"type": [0],
|
||||
})
|
||||
|
||||
connected_client.copy_ticks_range_as_df.side_effect = make_ticks
|
||||
connected_client.history_orders_get_as_df.side_effect = make_orders
|
||||
output = tmp_path / "ticks-orders.db"
|
||||
update_history(
|
||||
client=connected_client,
|
||||
output=output,
|
||||
symbols=["EURUSD"],
|
||||
datasets={Dataset.ticks, Dataset.history_orders},
|
||||
lookback_hours=24,
|
||||
date_to=date_to,
|
||||
)
|
||||
with sqlite3.connect(output) as conn:
|
||||
assert conn.execute("SELECT COUNT(*) FROM ticks").fetchone() == (1,)
|
||||
assert conn.execute(
|
||||
"SELECT COUNT(*) FROM history_orders",
|
||||
).fetchone() == (1,)
|
||||
|
||||
def test_update_history_with_config_opens_and_closes_connection(
|
||||
self,
|
||||
mocker: MockerFixture,
|
||||
tmp_path: Path,
|
||||
) -> None:
|
||||
"""Test update_history_with_config manages MT5 connection lifecycle."""
|
||||
mock_client = MagicMock()
|
||||
mocker.patch("mt5cli.sdk.Mt5DataClient", return_value=mock_client)
|
||||
updater = mocker.patch("mt5cli.sdk.update_history")
|
||||
update_history_with_config(
|
||||
output=tmp_path / "config-wrapper.db",
|
||||
symbols=["EURUSD"],
|
||||
datasets={Dataset.history_deals},
|
||||
timeframes=["M1"],
|
||||
flags="ALL",
|
||||
lookback_hours=1,
|
||||
date_to=datetime(2024, 1, 1, tzinfo=UTC),
|
||||
deduplicate=False,
|
||||
create_rate_views=False,
|
||||
with_views=True,
|
||||
include_account_events=False,
|
||||
)
|
||||
mock_client.initialize_and_login_mt5.assert_called_once()
|
||||
mock_client.shutdown.assert_called_once()
|
||||
updater.assert_called_once()
|
||||
assert updater.call_args.kwargs == {
|
||||
"client": mock_client,
|
||||
"output": tmp_path / "config-wrapper.db",
|
||||
"symbols": ["EURUSD"],
|
||||
"datasets": {Dataset.history_deals},
|
||||
"timeframes": ["M1"],
|
||||
"flags": "ALL",
|
||||
"lookback_hours": 1,
|
||||
"date_to": datetime(2024, 1, 1, tzinfo=UTC),
|
||||
"deduplicate": False,
|
||||
"create_rate_views": False,
|
||||
"with_views": True,
|
||||
"include_account_events": False,
|
||||
}
|
||||
|
||||
def test_update_history_with_config_validates_before_connecting(
|
||||
self,
|
||||
mocker: MockerFixture,
|
||||
tmp_path: Path,
|
||||
) -> None:
|
||||
"""Test invalid inputs fail before MT5 is initialized."""
|
||||
mock_client = MagicMock()
|
||||
mocker.patch("mt5cli.sdk.Mt5DataClient", return_value=mock_client)
|
||||
with pytest.raises(ValueError, match="lookback_hours must be positive"):
|
||||
update_history_with_config(
|
||||
output=tmp_path / "invalid-config.db",
|
||||
symbols=["EURUSD"],
|
||||
lookback_hours=0,
|
||||
)
|
||||
mock_client.initialize_and_login_mt5.assert_not_called()
|
||||
mock_client.shutdown.assert_not_called()
|
||||
|
||||
def test_update_history_with_config_noops_for_empty_datasets(
|
||||
self,
|
||||
mocker: MockerFixture,
|
||||
tmp_path: Path,
|
||||
) -> None:
|
||||
"""Test empty dataset selection skips MT5 initialization."""
|
||||
mock_client = MagicMock()
|
||||
mocker.patch("mt5cli.sdk.Mt5DataClient", return_value=mock_client)
|
||||
updater = mocker.patch("mt5cli.sdk.update_history")
|
||||
update_history_with_config(
|
||||
output=tmp_path / "empty-config.db",
|
||||
symbols=["EURUSD"],
|
||||
datasets=set(),
|
||||
)
|
||||
mock_client.initialize_and_login_mt5.assert_not_called()
|
||||
mock_client.shutdown.assert_not_called()
|
||||
updater.assert_not_called()
|
||||
|
||||
def test_update_history_defaults_date_to_now(
|
||||
self,
|
||||
connected_client: MagicMock,
|
||||
mocker: MockerFixture,
|
||||
tmp_path: Path,
|
||||
) -> None:
|
||||
"""Test update_history uses current UTC time when date_to is omitted."""
|
||||
captured: dict[str, datetime] = {}
|
||||
|
||||
def capture(
|
||||
*args: object,
|
||||
**_kwargs: object,
|
||||
) -> tuple[set[Dataset], dict[Dataset, set[str]]]:
|
||||
captured["end"] = args[7] # type: ignore[assignment]
|
||||
return set(), {}
|
||||
|
||||
mocker.patch("mt5cli.sdk.write_incremental_datasets", side_effect=capture)
|
||||
before = datetime.now(UTC)
|
||||
update_history(
|
||||
client=connected_client,
|
||||
output=tmp_path / "now-default.db",
|
||||
symbols=["EURUSD"],
|
||||
datasets={Dataset.rates},
|
||||
timeframes=["M1"],
|
||||
lookback_hours=12,
|
||||
)
|
||||
after = datetime.now(UTC)
|
||||
assert before <= captured["end"] <= after
|
||||
|
||||
|
||||
class TestRecentTicks:
|
||||
"""Tests for recent_ticks helper."""
|
||||
|
||||
def test_recent_ticks_uses_explicit_date_to_window(
|
||||
self,
|
||||
mocker: MockerFixture,
|
||||
) -> None:
|
||||
"""Test recent_ticks fetches the requested trailing window."""
|
||||
client = MagicMock()
|
||||
end = datetime(2024, 1, 2, 12, 0, 0, tzinfo=UTC)
|
||||
client.copy_ticks_from_as_df.return_value = pd.DataFrame({
|
||||
"time": [end],
|
||||
"bid": [1.0],
|
||||
})
|
||||
mocker.patch("mt5cli.sdk.Mt5DataClient", return_value=client)
|
||||
result = recent_ticks(
|
||||
"EURUSD",
|
||||
60,
|
||||
date_to=end,
|
||||
count=100,
|
||||
flags="INFO",
|
||||
config=build_config(login=123),
|
||||
)
|
||||
assert isinstance(result, pd.DataFrame)
|
||||
client.copy_ticks_from_as_df.assert_called_once_with(
|
||||
symbol="EURUSD",
|
||||
date_from=end - timedelta(seconds=60),
|
||||
count=100,
|
||||
flags=2,
|
||||
)
|
||||
client.copy_ticks_range_as_df.assert_not_called()
|
||||
|
||||
def test_recent_ticks_uses_latest_tick_when_date_to_omitted(
|
||||
self,
|
||||
mocker: MockerFixture,
|
||||
) -> None:
|
||||
"""Test recent_ticks anchors the window on the latest tick time."""
|
||||
client = MagicMock()
|
||||
tick = MagicMock()
|
||||
tick.time = datetime(2024, 1, 2, 12, 0, 0, tzinfo=UTC)
|
||||
client.symbol_info_tick.return_value = tick
|
||||
client.copy_ticks_from_as_df.return_value = pd.DataFrame({
|
||||
"time": [1, 2],
|
||||
"bid": [1.0, 1.1],
|
||||
})
|
||||
client.copy_ticks_range_as_df.return_value = pd.DataFrame({
|
||||
"time": [1, 2, 3],
|
||||
"bid": [1.0, 1.1, 1.2],
|
||||
})
|
||||
mocker.patch("mt5cli.sdk.Mt5DataClient", return_value=client)
|
||||
result = Mt5CliClient().recent_ticks("EURUSD", 30, count=2, flags="ALL")
|
||||
assert len(result) == 2
|
||||
client.symbol_info_tick.assert_called_once_with("EURUSD")
|
||||
client.copy_ticks_from_as_df.assert_called_once()
|
||||
_, kwargs = client.copy_ticks_range_as_df.call_args
|
||||
assert kwargs["symbol"] == "EURUSD"
|
||||
assert kwargs["date_to"] == tick.time
|
||||
assert kwargs["date_from"] == tick.time - timedelta(seconds=30)
|
||||
assert kwargs["flags"] == 1
|
||||
|
||||
def test_recent_ticks_rejects_unsupported_tick_time(
|
||||
self,
|
||||
mocker: MockerFixture,
|
||||
) -> None:
|
||||
"""Test recent_ticks raises when the latest tick time is unsupported."""
|
||||
client = MagicMock()
|
||||
tick = MagicMock()
|
||||
tick.time = object()
|
||||
client.symbol_info_tick.return_value = tick
|
||||
mocker.patch("mt5cli.sdk.Mt5DataClient", return_value=client)
|
||||
with pytest.raises(TypeError, match="Unsupported tick time value"):
|
||||
Mt5CliClient().recent_ticks("EURUSD", 30)
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
"tick_time",
|
||||
[
|
||||
"2024-01-02T12:00:00+00:00",
|
||||
1704196800,
|
||||
],
|
||||
)
|
||||
def test_recent_ticks_coerces_string_and_unix_tick_times(
|
||||
self,
|
||||
mocker: MockerFixture,
|
||||
tick_time: str | int,
|
||||
) -> None:
|
||||
"""Test recent_ticks accepts string and unix tick timestamps."""
|
||||
client = MagicMock()
|
||||
tick = MagicMock()
|
||||
tick.time = tick_time
|
||||
client.symbol_info_tick.return_value = tick
|
||||
expected_end = (
|
||||
datetime(2024, 1, 2, 12, 0, 0, tzinfo=UTC)
|
||||
if isinstance(tick_time, str)
|
||||
else datetime.fromtimestamp(tick_time, tz=UTC)
|
||||
)
|
||||
client.copy_ticks_from_as_df.return_value = pd.DataFrame({
|
||||
"time": [expected_end],
|
||||
})
|
||||
mocker.patch("mt5cli.sdk.Mt5DataClient", return_value=client)
|
||||
Mt5CliClient().recent_ticks("EURUSD", 30)
|
||||
_, kwargs = client.copy_ticks_from_as_df.call_args
|
||||
assert kwargs["date_from"] == expected_end - timedelta(seconds=30)
|
||||
|
||||
def test_recent_ticks_returns_full_frame_when_count_not_positive(
|
||||
self,
|
||||
mocker: MockerFixture,
|
||||
) -> None:
|
||||
"""Test non-positive count returns the full range without trimming."""
|
||||
client = MagicMock()
|
||||
end = datetime(2024, 1, 2, 12, 0, 0, tzinfo=UTC)
|
||||
client.copy_ticks_range_as_df.return_value = pd.DataFrame({
|
||||
"time": [1, 2, 3],
|
||||
"bid": [1.0, 1.1, 1.2],
|
||||
})
|
||||
mocker.patch("mt5cli.sdk.Mt5DataClient", return_value=client)
|
||||
result = recent_ticks(
|
||||
"EURUSD",
|
||||
60,
|
||||
date_to=end,
|
||||
count=0,
|
||||
config=build_config(login=123),
|
||||
)
|
||||
assert len(result) == 3
|
||||
client.copy_ticks_from_as_df.assert_not_called()
|
||||
client.copy_ticks_range_as_df.assert_called_once_with(
|
||||
symbol="EURUSD",
|
||||
date_from=end - timedelta(seconds=60),
|
||||
date_to=end,
|
||||
flags=1,
|
||||
)
|
||||
|
||||
|
||||
class TestMinimumMargins:
|
||||
"""Tests for minimum_margins helper."""
|
||||
|
||||
def test_minimum_margins_shape(
|
||||
self,
|
||||
mocker: MockerFixture,
|
||||
) -> None:
|
||||
"""Test minimum_margins returns the expected summary columns."""
|
||||
client = MagicMock()
|
||||
sym = MagicMock(volume_min=0.01)
|
||||
account = MagicMock(currency="USD")
|
||||
tick = MagicMock(ask=1.1010, bid=1.1000)
|
||||
client.symbol_info.return_value = sym
|
||||
client.account_info.return_value = account
|
||||
client.symbol_info_tick.return_value = tick
|
||||
client.order_calc_margin.side_effect = [12.5, 12.4]
|
||||
client.mt5.ORDER_TYPE_BUY = 0
|
||||
client.mt5.ORDER_TYPE_SELL = 1
|
||||
mocker.patch("mt5cli.sdk.Mt5DataClient", return_value=client)
|
||||
|
||||
result = minimum_margins("EURUSD", config=build_config(login=123))
|
||||
|
||||
pd.testing.assert_frame_equal(
|
||||
result,
|
||||
pd.DataFrame([
|
||||
{
|
||||
"symbol": "EURUSD",
|
||||
"account_currency": "USD",
|
||||
"volume_min": 0.01,
|
||||
"buy_margin": 12.5,
|
||||
"sell_margin": 12.4,
|
||||
}
|
||||
]),
|
||||
)
|
||||
client.order_calc_margin.assert_any_call(0, "EURUSD", 0.01, 1.1010)
|
||||
client.order_calc_margin.assert_any_call(1, "EURUSD", 0.01, 1.1000)
|
||||
|
||||
@@ -21,8 +21,10 @@ from mt5cli.utils import (
|
||||
TIMEFRAME_MAP,
|
||||
TIMEFRAME_TYPE,
|
||||
Dataset,
|
||||
IfExists,
|
||||
detect_format,
|
||||
export_dataframe,
|
||||
export_dataframe_to_sqlite,
|
||||
parse_datetime,
|
||||
parse_request,
|
||||
parse_tick_flags,
|
||||
@@ -130,6 +132,112 @@ class TestExportDataframe:
|
||||
export_dataframe(sample_df, tmp_path / "out.txt", "xml")
|
||||
|
||||
|
||||
class TestExportDataframeToSqlite:
|
||||
"""Tests for export_dataframe_to_sqlite."""
|
||||
|
||||
def test_append_preserves_existing_rows(self, tmp_path: Path) -> None:
|
||||
"""Test append mode keeps prior rows in the SQLite table."""
|
||||
output = tmp_path / "append.db"
|
||||
first = pd.DataFrame({"id": [1], "value": ["a"]})
|
||||
second = pd.DataFrame({"id": [2], "value": ["b"]})
|
||||
export_dataframe_to_sqlite(first, output, "items", if_exists=IfExists.REPLACE)
|
||||
export_dataframe_to_sqlite(second, output, "items", if_exists=IfExists.APPEND)
|
||||
with sqlite3.connect(output) as conn:
|
||||
result = pd.read_sql( # type: ignore[reportUnknownMemberType]
|
||||
"SELECT id, value FROM items ORDER BY id",
|
||||
conn,
|
||||
)
|
||||
pd.testing.assert_frame_equal(
|
||||
result,
|
||||
pd.DataFrame({"id": [1, 2], "value": ["a", "b"]}),
|
||||
)
|
||||
|
||||
def test_deduplicate_keeps_latest_row(self, tmp_path: Path) -> None:
|
||||
"""Test deduplication keeps the latest ROWID for key columns."""
|
||||
output = tmp_path / "dedup.db"
|
||||
first = pd.DataFrame({
|
||||
"symbol": ["EURUSD", "EURUSD"],
|
||||
"time": ["2024-01-01", "2024-01-01"],
|
||||
"bid": [1.0, 1.1],
|
||||
})
|
||||
second = pd.DataFrame({
|
||||
"symbol": ["EURUSD"],
|
||||
"time": ["2024-01-01"],
|
||||
"bid": [1.2],
|
||||
})
|
||||
export_dataframe_to_sqlite(
|
||||
first,
|
||||
output,
|
||||
"ticks",
|
||||
if_exists=IfExists.REPLACE,
|
||||
deduplicate_on=("symbol", "time"),
|
||||
)
|
||||
export_dataframe_to_sqlite(
|
||||
second,
|
||||
output,
|
||||
"ticks",
|
||||
if_exists=IfExists.APPEND,
|
||||
deduplicate_on=("symbol", "time"),
|
||||
)
|
||||
with sqlite3.connect(output) as conn:
|
||||
result = pd.read_sql( # type: ignore[reportUnknownMemberType]
|
||||
"SELECT symbol, time, bid FROM ticks",
|
||||
conn,
|
||||
)
|
||||
pd.testing.assert_frame_equal(
|
||||
result.reset_index(drop=True),
|
||||
pd.DataFrame({
|
||||
"symbol": ["EURUSD"],
|
||||
"time": ["2024-01-01"],
|
||||
"bid": [1.2],
|
||||
}),
|
||||
)
|
||||
|
||||
def test_default_if_exists_appends_without_dropping_rows(
|
||||
self,
|
||||
tmp_path: Path,
|
||||
) -> None:
|
||||
"""Test the default append mode keeps prior rows."""
|
||||
output = tmp_path / "default-append.db"
|
||||
first = pd.DataFrame({"id": [1], "value": ["a"]})
|
||||
second = pd.DataFrame({"id": [2], "value": ["b"]})
|
||||
export_dataframe_to_sqlite(first, output, "items")
|
||||
export_dataframe_to_sqlite(second, output, "items")
|
||||
with sqlite3.connect(output) as conn:
|
||||
result = pd.read_sql( # type: ignore[reportUnknownMemberType]
|
||||
"SELECT id, value FROM items ORDER BY id",
|
||||
conn,
|
||||
)
|
||||
pd.testing.assert_frame_equal(
|
||||
result,
|
||||
pd.DataFrame({"id": [1, 2], "value": ["a", "b"]}),
|
||||
)
|
||||
|
||||
def test_writes_index_with_label(self, tmp_path: Path) -> None:
|
||||
"""Test optional index export with a custom label."""
|
||||
output = tmp_path / "index.db"
|
||||
frame = pd.DataFrame(
|
||||
{"value": [1.0]}, index=pd.Index(["EURUSD"], name="symbol")
|
||||
)
|
||||
export_dataframe_to_sqlite(
|
||||
frame,
|
||||
output,
|
||||
"margins",
|
||||
if_exists=IfExists.REPLACE,
|
||||
index=True,
|
||||
index_label="symbol",
|
||||
)
|
||||
with sqlite3.connect(output) as conn:
|
||||
result = pd.read_sql( # type: ignore[reportUnknownMemberType]
|
||||
"SELECT symbol, value FROM margins",
|
||||
conn,
|
||||
)
|
||||
pd.testing.assert_frame_equal(
|
||||
result,
|
||||
pd.DataFrame({"symbol": ["EURUSD"], "value": [1.0]}),
|
||||
)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Parse helpers
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
Reference in New Issue
Block a user